from __future__ import annotations import argparse import csv import hashlib import importlib.metadata import json import math import os import platform import re import shutil import subprocess import sys import time from dataclasses import dataclass from pathlib import Path from typing import Any, Iterable import cv2 import librosa import matplotlib matplotlib.use("Agg") import matplotlib.pyplot as plt import numpy as np import torch from openpyxl import load_workbook from scipy.fft import dct, irfft, rfft from scipy.spatial.transform import Rotation from sklearn.linear_model import LogisticRegression, Ridge from sklearn.metrics import accuracy_score, f1_score, mean_absolute_error from sklearn.model_selection import GroupKFold from transformers import AutoModel, AutoModelForCTC, AutoTokenizer ROOT = Path(__file__).resolve().parents[1] MATH_DIR = Path(__file__).resolve().parent from ..data_paths import ATTACHMENT1 DATA_DIR = ATTACHMENT1 LABEL_FILE = DATA_DIR / "label-100.xlsx" OUTPUT_DIR = ROOT / "output" / "q1" / "model_comparison" CACHE_DIR = MATH_DIR / "cache" / "native" OPENFACE_RUNNER = MATH_DIR / "run_openface.ps1" OPENFACE_PACKAGE = MATH_DIR / "cache" / "openface" / "OpenFace_2.2.0_win_x64" TEXT_MODEL_ID = "google-bert/bert-base-uncased" SPEECH_MODEL_ID = "facebook/wav2vec2-base-960h" SAMPLE_RATE = 16_000 FRAME_LENGTH = 400 FRAME_STEP = 160 N_FFT = 512 MEL_COUNT = 40 GRID_STEP_S = 0.1 CTC_FRAME_STEP_S = 320 / SAMPLE_RATE CTC_RECEPTIVE_SAMPLES = 400 SEED = 20260923 N_FOLDS = 5 AU_CODES = (1, 2, 4, 5, 6, 7, 9, 10, 12, 14, 15, 17, 20, 23, 25, 26, 45) VISION_NAMES = ( *[f"AU{code:02d}_r" for code in AU_CODES], "pose_Rx", "pose_Ry", "pose_Rz", "pose_Tx", "pose_Ty", "pose_Tz", "gaze_0_x", "gaze_0_y", "gaze_0_z", "gaze_1_x", "gaze_1_y", "gaze_1_z", "geometry_eye_open_left", "geometry_eye_open_right", "geometry_mouth_open", "geometry_mouth_width", "geometry_brow_eye_left", "geometry_brow_eye_right", ) AUDIO_NAMES = ( *[f"logmel_{index:02d}" for index in range(40)], *[f"mfcc_{index:02d}" for index in range(13)], *[f"delta_mfcc_{index:02d}" for index in range(13)], "log_energy", "log_f0_hz", "voiced_probability", "spectral_centroid_hz", "spectral_bandwidth_hz", "spectral_flux", "zero_crossing_rate", "hnr_db", ) @dataclass class Sample: sample_id: str video_id: str clip_id: str duration_s: float sentiment: float polarity: int text_words: list[str] text_features: np.ndarray text_valid: np.ndarray hard_intervals: np.ndarray hard_valid: np.ndarray hard_quality: np.ndarray audio_times: np.ndarray audio_features: np.ndarray audio_observed: np.ndarray audio_quality: np.ndarray audio_quality_available: np.ndarray vision_times: np.ndarray vision_frame_indices: np.ndarray vision_features: np.ndarray vision_observed: np.ndarray vision_quality: np.ndarray vision_quality_available: np.ndarray ctc_times: np.ndarray ctc_occupancy: np.ndarray speech_features: np.ndarray boundary_summary: list[dict[str, float]] video_sha256: str @dataclass class View: features: dict[str, np.ndarray] observed: dict[str, np.ndarray] coverage: dict[str, np.ndarray] edges: np.ndarray def _sha256(path: Path) -> str: digest = hashlib.sha256() with path.open("rb") as stream: for chunk in iter(lambda: stream.read(1024 * 1024), b""): digest.update(chunk) return digest.hexdigest() def _version(package: str) -> str: try: return importlib.metadata.version(package) except importlib.metadata.PackageNotFoundError: return "not-installed" def _safe_name(value: str) -> str: return re.sub(r"[^A-Za-z0-9_.-]+", "_", value) def _str_id(value: Any) -> str: if isinstance(value, float) and value.is_integer(): return str(int(value)) return str(value).strip() def _read_labels(path: Path) -> list[dict[str, Any]]: workbook = load_workbook(path, read_only=True, data_only=True) sheet = workbook.active rows = sheet.iter_rows(values_only=True) header = [str(value).strip() if value is not None else "" for value in next(rows)] index = {name: position for position, name in enumerate(header)} required = {"video_id", "clip_id", "text", "label", "annotation"} missing = required - set(index) if missing: raise ValueError(f"label workbook lacks columns: {sorted(missing)}") records: list[dict[str, Any]] = [] seen: set[str] = set() for row_number, row in enumerate(rows, start=2): video_id = _str_id(row[index["video_id"]]) clip_id = _str_id(row[index["clip_id"]]) sample_id = f"{video_id}/{clip_id}" if sample_id in seen: raise ValueError(f"duplicate source key in label workbook: {sample_id}") seen.add(sample_id) label = float(row[index["label"]]) annotation = str(row[index["annotation"]]).strip().lower() polarity = 0 if label < 0 else 1 if label == 0 else 2 expected = {"negative": 0, "neutral": 1, "positive": 2}.get(annotation) records.append({ "row_number": row_number, "sample_id": sample_id, "video_id": video_id, "clip_id": clip_id, "text": str(row[index["text"]] or ""), "sentiment": label, "polarity": polarity, "annotation": annotation, "label_consistent": expected == polarity, }) workbook.close() if len(records) != 100: raise ValueError(f"expected 100 label rows, received {len(records)}") return records def _decode_audio(path: Path) -> np.ndarray: result = subprocess.run( [ "ffmpeg", "-nostdin", "-v", "error", "-i", str(path), "-vn", "-ac", "1", "-ar", str(SAMPLE_RATE), "-f", "f32le", "pipe:1", ], check=True, capture_output=True, ) audio = np.frombuffer(result.stdout, dtype=" float: result = subprocess.run( [ "ffprobe", "-v", "error", "-show_entries", "format=duration", "-of", "default=noprint_wrappers=1:nokey=1", str(path), ], check=True, capture_output=True, text=True, ) value = float(result.stdout.strip()) if not math.isfinite(value) or value <= 0: raise ValueError(f"invalid duration for {path}: {value}") return value def _mel_filterbank(sample_rate: int, n_fft: int, count: int) -> np.ndarray: def to_mel(freq: np.ndarray | float) -> np.ndarray | float: return 2595.0 * np.log10(1.0 + np.asarray(freq) / 700.0) def from_mel(mel: np.ndarray) -> np.ndarray: return 700.0 * (10.0 ** (mel / 2595.0) - 1.0) mel_points = np.linspace(to_mel(0.0), to_mel(sample_rate / 2), count + 2) hz_points = from_mel(mel_points) bins_hz = np.fft.rfftfreq(n_fft, 1 / sample_rate) bank = np.zeros((count, len(bins_hz)), dtype=np.float32) for band in range(count): left, center, right = hz_points[band:band + 3] bank[band] = np.maximum( 0.0, np.minimum((bins_hz - left) / max(center - left, 1e-9), (right - bins_hz) / max(right - center, 1e-9)), ) return bank def _audio_features(audio: np.ndarray, duration_s: float) -> tuple[np.ndarray, np.ndarray, np.ndarray]: starts = np.arange(0, max(len(audio), 1), FRAME_STEP, dtype=np.int64) frames = np.zeros((len(starts), FRAME_LENGTH), dtype=np.float32) for index, start in enumerate(starts): chunk = audio[start:start + FRAME_LENGTH] frames[index, :len(chunk)] = chunk window = np.hanning(FRAME_LENGTH).astype(np.float32) windowed = frames * window[None, :] spectrum = np.abs(rfft(windowed, n=N_FFT, axis=1)).astype(np.float32) power = spectrum ** 2 mel = power @ _mel_filterbank(SAMPLE_RATE, N_FFT, MEL_COUNT).T log_mel = np.log(np.maximum(mel, 1e-8)).astype(np.float32) mfcc = dct(log_mel, type=2, norm="ortho", axis=1)[:, :13].astype(np.float32) delta = np.zeros_like(mfcc) if len(mfcc) > 1: for frame in range(len(mfcc)): left, right = max(0, frame - 2), min(len(mfcc), frame + 3) offsets = np.arange(left, right, dtype=np.float32) - frame denominator = float(np.square(offsets).sum()) if denominator > 0: delta[frame] = (offsets[:, None] * mfcc[left:right]).sum(axis=0) / denominator frequencies = np.fft.rfftfreq(N_FFT, 1 / SAMPLE_RATE).astype(np.float32) magnitude_sum = np.maximum(spectrum.sum(axis=1), 1e-8) centroid = (spectrum * frequencies[None, :]).sum(axis=1) / magnitude_sum bandwidth = np.sqrt( (spectrum * (frequencies[None, :] - centroid[:, None]) ** 2).sum(axis=1) / magnitude_sum ) normalized_spectrum = spectrum / magnitude_sum[:, None] flux = np.zeros(len(frames), dtype=np.float32) if len(frames) > 1: flux[1:] = np.square(normalized_spectrum[1:] - normalized_spectrum[:-1]).sum(axis=1) zcr = (np.diff(np.signbit(frames), axis=1) != 0).mean(axis=1).astype(np.float32) energy = np.mean(np.square(frames), axis=1) log_f0 = np.zeros(len(frames), dtype=np.float32) voiced_probability = np.zeros(len(frames), dtype=np.float32) hnr = np.zeros(len(frames), dtype=np.float32) pitch_observed = np.zeros(len(frames), dtype=np.bool_) voicing_observed = np.zeros(len(frames), dtype=np.bool_) hnr_observed = np.zeros(len(frames), dtype=np.bool_) if len(audio) >= FRAME_LENGTH: f0, voiced_flag, voiced_prob = librosa.pyin( audio.astype(np.float32), fmin=80.0, fmax=400.0, sr=SAMPLE_RATE, frame_length=FRAME_LENGTH, hop_length=FRAME_STEP, center=False, fill_na=np.nan, ) rows = min(len(frames), len(f0), len(voiced_prob)) for frame_index in range(rows): if np.isfinite(voiced_prob[frame_index]): voiced_probability[frame_index] = float(np.clip(voiced_prob[frame_index], 0.0, 1.0)) voicing_observed[frame_index] = True if not np.isfinite(f0[frame_index]) or not bool(voiced_flag[frame_index]): continue log_f0[frame_index] = math.log(float(f0[frame_index])) pitch_observed[frame_index] = True frame = windowed[frame_index] autocorr = irfft(np.abs(rfft(frame, n=1024)) ** 2, n=1024)[:FRAME_LENGTH] if autocorr[0] <= 1e-10: continue lag = int(round(SAMPLE_RATE / float(f0[frame_index]))) if lag <= 0 or lag >= len(autocorr): continue periodicity = float(np.clip(autocorr[lag] / autocorr[0], 1e-5, 1.0 - 1e-5)) hnr[frame_index] = 10.0 * math.log10(periodicity / (1.0 - periodicity)) hnr_observed[frame_index] = True prosody = np.column_stack(( np.log(np.maximum(energy, 1e-8)), log_f0, voiced_probability, centroid, bandwidth, flux, zcr, hnr, )).astype(np.float32) features = np.column_stack((log_mel, mfcc, delta, prosody)).astype(np.float32) if features.shape[1] != 74: raise AssertionError(f"internal audio dimension error: {features.shape}") observed = np.isfinite(features) observed[:, 67] = pitch_observed observed[:, 68] = voicing_observed observed[:, 73] = hnr_observed if len(frames) < 2: observed[:, 53:66] = False features = np.nan_to_num(features, nan=0.0, posinf=0.0, neginf=0.0) times = np.minimum((starts + FRAME_LENGTH / 2) / SAMPLE_RATE, duration_s).astype(np.float32) return times, features, observed def _text_features(words: list[str], tokenizer: Any, model: Any, device: torch.device) -> tuple[np.ndarray, np.ndarray]: output = np.zeros((len(words), int(model.config.hidden_size)), dtype=np.float32) valid = np.zeros(len(words), dtype=np.bool_) if not words: return output, valid token_ids: list[int] = [] word_ids: list[int] = [] for word_index, word in enumerate(words): word_tokens = tokenizer(word, add_special_tokens=False)["input_ids"] token_ids.extend(int(value) for value in word_tokens) word_ids.extend([word_index] * len(word_tokens)) token_sums = np.zeros((len(token_ids), output.shape[1]), dtype=np.float64) token_weights = np.zeros(len(token_ids), dtype=np.float64) window_size, stride = 510, 384 offsets = range(0, max(len(token_ids), 1), stride) with torch.inference_mode(): for offset in offsets: end = min(len(token_ids), offset + window_size) if offset >= end: break content = token_ids[offset:end] special_content = [int(tokenizer.cls_token_id), *content, int(tokenizer.sep_token_id)] encoded = { "input_ids": torch.tensor([special_content], dtype=torch.long, device=device), "attention_mask": torch.ones((1, len(special_content)), dtype=torch.long, device=device), } hidden_states = model(**encoded, output_hidden_states=True).hidden_states hidden = torch.stack(tuple(hidden_states[-4:]), dim=0).mean(dim=0)[0, 1:1 + len(content)] values = hidden.float().cpu().numpy() for local_index, value in enumerate(values): global_index = offset + local_index edge_weight = float(min(local_index + 1, len(content) - local_index)) token_sums[global_index] += edge_weight * value token_weights[global_index] += edge_weight accum = np.zeros_like(output, dtype=np.float64) counts = np.zeros(len(words), dtype=np.float64) for token_index, word_index in enumerate(word_ids): if word_index is None or not (0 <= word_index < len(words)) or token_weights[token_index] <= 0: continue accum[word_index] += token_sums[token_index] / token_weights[token_index] counts[word_index] += 1.0 valid = counts > 0 output[valid] = (accum[valid] / counts[valid, None]).astype(np.float32) return output, valid def _ctc_targets(words: list[str], tokenizer: Any) -> tuple[list[int], list[list[int]]]: vocab = tokenizer.get_vocab() delimiter_id = int(tokenizer.convert_tokens_to_ids(tokenizer.word_delimiter_token or "|")) unknown_id = int(tokenizer.unk_token_id) targets: list[int] = [] per_word: list[list[int]] = [[] for _ in words] for word_index, raw_word in enumerate(words): if word_index: targets.append(delimiter_id) normalized = re.sub(r"[^a-z']", "", raw_word.lower().replace("’", "'")) for char in normalized: per_word[word_index].append(len(targets)) targets.append(int(vocab.get(char, unknown_id))) return targets, per_word def _ctc_path_and_posterior(log_probs: np.ndarray, targets: list[int], blank_id: int) -> tuple[np.ndarray | None, np.ndarray | None, float]: if not targets or log_probs.ndim != 2: return None, None, -math.inf target = np.asarray(targets, dtype=np.int64) states = np.full(2 * len(target) + 1, blank_id, dtype=np.int64) states[1::2] = target frames, state_count = log_probs.shape[0], len(states) if frames == 0 or frames < len(target): return None, None, -math.inf alpha = np.full((frames, state_count), -np.inf, dtype=np.float64) alpha[0, 0] = float(log_probs[0, blank_id]) alpha[0, 1] = float(log_probs[0, states[1]]) skip_ok = np.zeros(state_count, dtype=np.bool_) if state_count > 2: indices = np.arange(state_count) skip_ok[2:] = (states[2:] != blank_id) & (states[2:] != states[:-2]) for frame in range(1, frames): prev = alpha[frame - 1] incoming = prev.copy() incoming[1:] = np.logaddexp(incoming[1:], prev[:-1]) valid_skip = np.flatnonzero(skip_ok) if len(valid_skip): incoming[valid_skip] = np.logaddexp(incoming[valid_skip], prev[valid_skip - 2]) alpha[frame] = incoming + log_probs[frame, states] log_z = float(np.logaddexp(alpha[-1, -1], alpha[-1, -2])) if not math.isfinite(log_z): return None, None, -math.inf beta = np.full((frames, state_count), -np.inf, dtype=np.float64) beta[-1, -1] = 0.0 beta[-1, -2] = 0.0 for frame in range(frames - 2, -1, -1): nxt = beta[frame + 1] + log_probs[frame + 1, states] # A CTC state may stay active for any number of frames. Start with the # same-state transition, then add transitions to the next states. values = nxt.copy() values[:-1] = np.logaddexp(values[:-1], nxt[1:]) if state_count > 2: dest = np.flatnonzero(skip_ok) values[dest - 2] = np.logaddexp(values[dest - 2], nxt[dest]) beta[frame] = values posterior = np.exp(np.clip(alpha + beta - log_z, -745.0, 0.0)).astype(np.float32) back = np.zeros((frames, state_count), dtype=np.uint8) previous = np.full(state_count, -np.inf, dtype=np.float64) previous[0] = float(log_probs[0, blank_id]) previous[1] = float(log_probs[0, states[1]]) for frame in range(1, frames): stay = previous one = np.full(state_count, -np.inf, dtype=np.float64) one[1:] = previous[:-1] two = np.full(state_count, -np.inf, dtype=np.float64) two[skip_ok] = previous[np.flatnonzero(skip_ok) - 2] candidates = np.stack((stay, one, two), axis=0) choice = candidates.argmax(axis=0).astype(np.uint8) previous = candidates[choice, np.arange(state_count)] + log_probs[frame, states] back[frame] = choice state = state_count - 1 if previous[-1] >= previous[-2] else state_count - 2 path = np.empty(frames, dtype=np.int32) path[-1] = state for frame in range(frames - 1, 0, -1): state -= int(back[frame, state]) path[frame - 1] = state return path, posterior, log_z def _entry_distribution(log_probs: np.ndarray, states: np.ndarray, alpha: np.ndarray, beta: np.ndarray, log_z: float, state: int) -> np.ndarray: frames = len(log_probs) mass = np.zeros(frames, dtype=np.float64) if state == 1: mass[0] = math.exp(float(log_probs[0, states[state]] + beta[0, state] - log_z)) predecessors = [state - 1] if state >= 2 and states[state] != states[state - 2]: predecessors.append(state - 2) for frame in range(1, frames): terms = [alpha[frame - 1, prev] for prev in predecessors if prev >= 0] if terms: log_mass = float(np.logaddexp.reduce(terms) + log_probs[frame, states[state]] + beta[frame, state] - log_z) mass[frame] = math.exp(max(-745.0, min(log_mass, 0.0))) return mass def _exit_distribution(log_probs: np.ndarray, states: np.ndarray, alpha: np.ndarray, beta: np.ndarray, log_z: float, state: int) -> np.ndarray: frames = len(log_probs) mass = np.zeros(frames + 1, dtype=np.float64) successors = [state + 1] if state + 2 < len(states) and states[state] != states[state + 2]: successors.append(state + 2) for frame in range(1, frames): terms = [ alpha[frame - 1, state] + log_probs[frame, states[nxt]] + beta[frame, nxt] for nxt in successors if nxt < len(states) ] if terms: mass[frame] = math.exp(max(-745.0, min(float(np.logaddexp.reduce(terms) - log_z), 0.0))) if state == len(states) - 2: mass[frames] = math.exp(max(-745.0, min(float(alpha[-1, state] - log_z), 0.0))) return mass def _weighted_quantile(values: np.ndarray, weights: np.ndarray, quantile: float) -> float: if not len(values) or float(weights.sum()) <= 0: return math.nan order = np.argsort(values) values, weights = values[order], weights[order] cumulative = np.cumsum(weights) / weights.sum() return float(values[min(int(np.searchsorted(cumulative, quantile, side="left")), len(values) - 1)]) def _ctc_boundaries( log_probs: np.ndarray, targets: list[int], word_targets: list[list[int]], blank_id: int, duration_s: float, center_s: float, ) -> tuple[np.ndarray, list[dict[str, float]], np.ndarray]: path, posterior, log_z = _ctc_path_and_posterior(log_probs, targets, blank_id) occupancy = np.zeros((len(log_probs), len(word_targets)), dtype=np.float32) summaries: list[dict[str, float]] = [ {"start_mean_s": math.nan, "end_mean_s": math.nan, "start_p05_s": math.nan, "start_p95_s": math.nan, "end_p05_s": math.nan, "end_p95_s": math.nan, "start_width90_s": math.nan, "end_width90_s": math.nan} for _ in word_targets ] if path is None or posterior is None: return np.empty((0, len(word_targets)), dtype=np.float32), summaries, path if path is not None else np.empty((0,), dtype=np.int32) state_ids = np.full(2 * len(targets) + 1, blank_id, dtype=np.int64) state_ids[1::2] = np.asarray(targets, dtype=np.int64) # Recompute the dynamic-programming tables for boundary transition marginals. frames, count = len(log_probs), len(state_ids) alpha = np.full((frames, count), -np.inf, dtype=np.float64) alpha[0, 0] = log_probs[0, blank_id] alpha[0, 1] = log_probs[0, state_ids[1]] skip = np.zeros(count, dtype=np.bool_) if count > 2: skip[2:] = (state_ids[2:] != blank_id) & (state_ids[2:] != state_ids[:-2]) for frame in range(1, frames): incoming = alpha[frame - 1].copy() incoming[1:] = np.logaddexp(incoming[1:], alpha[frame - 1, :-1]) skip_states = np.flatnonzero(skip) if len(skip_states): incoming[skip_states] = np.logaddexp(incoming[skip_states], alpha[frame - 1, skip_states - 2]) alpha[frame] = incoming + log_probs[frame, state_ids] beta = np.full((frames, count), -np.inf, dtype=np.float64) beta[-1, -1] = beta[-1, -2] = 0.0 for frame in range(frames - 2, -1, -1): nxt = beta[frame + 1] + log_probs[frame + 1, state_ids] current = nxt.copy() current[:-1] = np.logaddexp(current[:-1], nxt[1:]) skip_states = np.flatnonzero(skip) if len(skip_states): current[skip_states - 2] = np.logaddexp(current[skip_states - 2], nxt[skip_states]) beta[frame] = current frame_step = CTC_FRAME_STEP_S first_edge = max(0.0, center_s - frame_step / 2) frame_edges = np.clip(first_edge + np.arange(frames + 1, dtype=np.float64) * frame_step, 0.0, duration_s) for word_index, token_indices in enumerate(word_targets): if not token_indices: continue char_states = np.asarray([2 * target_index + 1 for target_index in token_indices], dtype=np.int64) occupancy[:, word_index] = posterior[:, char_states].sum(axis=1) start_mass = _entry_distribution(log_probs, state_ids, alpha, beta, log_z, int(char_states[0])) end_mass = _exit_distribution(log_probs, state_ids, alpha, beta, log_z, int(char_states[-1])) start_total = float(start_mass.sum()) end_total = float(end_mass.sum()) if not math.isclose(start_total, 1.0, rel_tol=2e-4, abs_tol=2e-4): raise FloatingPointError(f"CTC start-boundary marginal sums to {start_total:.6g}, expected 1") if not math.isclose(end_total, 1.0, rel_tol=2e-4, abs_tol=2e-4): raise FloatingPointError(f"CTC end-boundary marginal sums to {end_total:.6g}, expected 1") start_times = frame_edges[:-1] end_times = frame_edges start_p05 = _weighted_quantile(start_times, start_mass, 0.05) start_p95 = _weighted_quantile(start_times, start_mass, 0.95) end_p05 = _weighted_quantile(end_times, end_mass, 0.05) end_p95 = _weighted_quantile(end_times, end_mass, 0.95) summaries[word_index] = { "start_mean_s": float(np.sum(start_mass * start_times) / start_total), "end_mean_s": float(np.sum(end_mass * end_times) / end_total), "start_p05_s": start_p05, "start_p95_s": start_p95, "end_p05_s": end_p05, "end_p95_s": end_p95, "start_width90_s": start_p95 - start_p05, "end_width90_s": end_p95 - end_p05, } return occupancy, summaries, path def _posterior_hard_intervals( log_probs: np.ndarray, path: np.ndarray | None, targets: list[int], word_targets: list[list[int]], duration_s: float, frame_step: float, center_s: float, ) -> tuple[np.ndarray, np.ndarray, np.ndarray]: intervals = np.zeros((len(word_targets), 2), dtype=np.float32) valid = np.zeros(len(word_targets), dtype=np.bool_) quality = np.zeros(len(word_targets), dtype=np.float32) if path is None or not len(path): return intervals, valid, quality for word_index, target_indices in enumerate(word_targets): if not target_indices: continue char_states = np.asarray([2 * index + 1 for index in target_indices], dtype=np.int32) frame_indices = np.flatnonzero(np.isin(path, char_states)) if not len(frame_indices): continue first, last = int(frame_indices[0]), int(frame_indices[-1]) start = max(0.0, first * frame_step + center_s - frame_step / 2) end = min(duration_s, (last + 1) * frame_step + center_s - frame_step / 2) if end <= start: continue intervals[word_index] = (start, end) valid[word_index] = True scores = [] for target_index in target_indices: selected = np.flatnonzero(path == 2 * target_index + 1) if len(selected): scores.extend(log_probs[selected, targets[target_index]].tolist()) quality[word_index] = float(np.exp(np.mean(scores))) if scores else 0.0 return intervals, valid, quality def _openface_float(row: dict[str, str], name: str) -> float: try: value = float(row[name]) except (KeyError, TypeError, ValueError): return math.nan return value if math.isfinite(value) else math.nan def _openface_geometry(row: dict[str, str]) -> tuple[np.ndarray, np.ndarray]: points = np.full((68, 2), np.nan, dtype=np.float64) for index in range(68): points[index, 0] = _openface_float(row, f"x_{index}") points[index, 1] = _openface_float(row, f"y_{index}") values = np.zeros(6, dtype=np.float32) valid = np.zeros(6, dtype=bool) if not np.isfinite(points[36:48]).all(): return values, valid left_eye = points[36:42].mean(axis=0) right_eye = points[42:48].mean(axis=0) interocular = float(np.linalg.norm(left_eye - right_eye)) if interocular <= 1e-6: return values, valid def distance(left: int, right: int) -> float: return float(np.linalg.norm(points[left] - points[right])) pairs = ((37, 41), (38, 40), (43, 47), (44, 46), (61, 67), (62, 66), (63, 65), (48, 54)) needed = {index for pair in pairs for index in pair} | set(range(17, 27)) complete = {index: np.isfinite(points[index]).all() for index in needed} if complete[37] and complete[41] and complete[38] and complete[40]: values[0] = (distance(37, 41) + distance(38, 40)) / (2.0 * interocular) valid[0] = True if complete[43] and complete[47] and complete[44] and complete[46]: values[1] = (distance(43, 47) + distance(44, 46)) / (2.0 * interocular) valid[1] = True mouth_pairs = ((61, 67), (62, 66), (63, 65)) if all(complete[a] and complete[b] for a, b in mouth_pairs): values[2] = sum(distance(a, b) for a, b in mouth_pairs) / (3.0 * interocular) valid[2] = True if complete[48] and complete[54]: values[3] = distance(48, 54) / interocular valid[3] = True if all(complete[index] for index in range(17, 22)): values[4] = float(np.linalg.norm(points[17:22].mean(axis=0) - left_eye) / interocular) valid[4] = True if all(complete[index] for index in range(22, 27)): values[5] = float(np.linalg.norm(points[22:27].mean(axis=0) - right_eye) / interocular) valid[5] = True valid &= np.isfinite(values) values[~valid] = 0.0 return values, valid def _run_openface(path: Path, output_name: str) -> Path: try: relative_path = path.resolve().relative_to(ROOT).as_posix() except ValueError: relative_path = str(path.resolve()) output_dir = MATH_DIR / "cache" / "openface_outputs" output_dir.mkdir(parents=True, exist_ok=True) csv_path = output_dir / f"{output_name}.csv" if csv_path.is_file() and csv_path.stat().st_mtime >= path.stat().st_mtime: return csv_path powershell = shutil.which("powershell.exe") or shutil.which("pwsh.exe") if powershell is None or not OPENFACE_RUNNER.is_file(): raise FileNotFoundError( "OpenFace 2.2.0 is required for compliant vision features. Install its Windows package under " f"{OPENFACE_PACKAGE} and run through {OPENFACE_RUNNER}; no MediaPipe proxy fallback is allowed." ) completed = subprocess.run( [powershell, "-NoProfile", "-ExecutionPolicy", "Bypass", "-File", str(OPENFACE_RUNNER), "-VideoRelativePath", relative_path, "-OutputName", output_name], check=False, capture_output=True, text=True, errors="replace", ) if completed.returncode != 0 or not csv_path.is_file(): raise RuntimeError( f"OpenFace failed for {path}: exit={completed.returncode}; " f"stdout={completed.stdout[-2000:]}; stderr={completed.stderr[-2000:]}" ) return csv_path def _vision_features(path: Path, duration_s: float, output_name: str) -> tuple[np.ndarray, np.ndarray, np.ndarray, np.ndarray, np.ndarray, np.ndarray]: """Read 35-D OpenFace features and keep raw PTS, frame IDs and confidence.""" csv_path = _run_openface(path, output_name) timestamps: list[float] = [] frame_indices: list[int] = [] values: list[np.ndarray] = [] masks: list[np.ndarray] = [] qualities: list[float] = [] quality_available: list[bool] = [] with csv_path.open("r", encoding="utf-8-sig", newline="") as stream: reader = csv.DictReader(stream, skipinitialspace=True) for raw_row in reader: row = {str(key).strip(): str(value).strip() for key, value in raw_row.items() if key is not None} timestamp = _openface_float(row, "timestamp") frame = _openface_float(row, "frame") confidence = _openface_float(row, "confidence") success = _openface_float(row, "success") if not math.isfinite(timestamp) or not math.isfinite(frame): continue timestamp = float(np.clip(timestamp, 0.0, duration_s)) vector = np.zeros(35, dtype=np.float32) observed = np.zeros(35, dtype=bool) good_face = success >= 1.0 and math.isfinite(confidence) and confidence >= 0.8 if good_face: for index, code in enumerate(AU_CODES): value = _openface_float(row, f"AU{code:02d}_r") if math.isfinite(value): vector[index] = value observed[index] = True for offset, name in enumerate(("pose_Rx", "pose_Ry", "pose_Rz", "pose_Tx", "pose_Ty", "pose_Tz"), start=17): value = _openface_float(row, name) if math.isfinite(value): vector[offset] = value observed[offset] = True gaze_names = ("gaze_0_x", "gaze_0_y", "gaze_0_z", "gaze_1_x", "gaze_1_y", "gaze_1_z") for offset, name in enumerate(gaze_names, start=23): value = _openface_float(row, name) if math.isfinite(value): vector[offset] = value observed[offset] = True geometry, geometry_mask = _openface_geometry(row) vector[29:35] = geometry observed[29:35] = geometry_mask timestamps.append(timestamp) frame_indices.append(int(frame)) values.append(vector) masks.append(observed) qualities.append(float(np.clip(confidence, 0.0, 1.0)) if math.isfinite(confidence) and success >= 1.0 else 1.0) quality_available.append(bool(math.isfinite(confidence) and success >= 1.0)) if not timestamps: return (np.empty(0, np.float32), np.empty(0, np.int32), np.empty((0, 35), np.float32), np.empty((0, 35), bool), np.empty(0, np.float32), np.empty(0, bool)) order = np.argsort(np.asarray(timestamps), kind="stable") return ( np.asarray(timestamps, np.float32)[order], np.asarray(frame_indices, np.int32)[order], np.stack(values)[order], np.stack(masks)[order], np.asarray(qualities, np.float32)[order], np.asarray(quality_available, bool)[order], ) def _load_models(device: torch.device) -> tuple[Any, Any, Any, Any, dict[str, Any]]: text_tokenizer = AutoTokenizer.from_pretrained(TEXT_MODEL_ID, use_fast=True) text_model = AutoModel.from_pretrained(TEXT_MODEL_ID, output_hidden_states=True).to(device).eval() speech_tokenizer = AutoTokenizer.from_pretrained(SPEECH_MODEL_ID) speech_model = AutoModelForCTC.from_pretrained(SPEECH_MODEL_ID).to(device).eval() config = speech_model.config stride_samples = int(np.prod(config.conv_stride)) receptive = 1 jump = 1 for kernel, stride in zip(config.conv_kernel, config.conv_stride): receptive += (int(kernel) - 1) * jump jump *= int(stride) info = { "text_id": TEXT_MODEL_ID, "text_revision": getattr(text_model.config, "_commit_hash", None), "speech_id": SPEECH_MODEL_ID, "speech_revision": getattr(config, "_commit_hash", None), "speech_stride_samples": stride_samples, "speech_receptive_field_samples": receptive, } return text_tokenizer, text_model, speech_tokenizer, speech_model, info def _extract_native_sample( record: dict[str, Any], models: tuple[Any, Any, Any, Any, dict[str, Any]], device: torch.device, video_sha256: str, ) -> Sample: text_tokenizer, text_model, ctc_tokenizer, speech_model, model_info = models video_path = DATA_DIR / record["video_id"] / f"{record['clip_id']}.mp4" duration_s = _duration(video_path) words = record["text"].split() text_vectors, text_valid = _text_features(words, text_tokenizer, text_model, device) waveform = _decode_audio(video_path) audio_times, audio_vectors, audio_observed = _audio_features(waveform, duration_s) output_name = f"{_safe_name(record['video_id'])}__{_safe_name(record['clip_id'])}" (vision_times, vision_frame_indices, vision_vectors, vision_observed, vision_quality, vision_quality_available) = _vision_features(video_path, duration_s, output_name) targets, word_targets = _ctc_targets(words, ctc_tokenizer) blank_id = int(ctc_tokenizer.pad_token_id) if targets: with torch.inference_mode(): input_values = torch.from_numpy(waveform).to(device).unsqueeze(0) output = speech_model(input_values=input_values, output_hidden_states=True) logits = output.logits[0].float() log_probs = torch.log_softmax(logits, dim=-1).cpu().numpy().astype(np.float32) deep_states = torch.stack(tuple(output.hidden_states[-4:]), dim=0).mean(dim=0)[0].float().cpu().numpy() else: log_probs = np.empty((0, len(ctc_tokenizer)), dtype=np.float32) deep_states = np.empty((0, int(speech_model.config.hidden_size)), dtype=np.float32) stride_samples = int(model_info["speech_stride_samples"]) frame_step = stride_samples / SAMPLE_RATE center_s = float(model_info["speech_receptive_field_samples"]) / (2 * SAMPLE_RATE) if len(log_probs) and targets: occupancy, boundary_summary, path = _ctc_boundaries( log_probs, targets, word_targets, blank_id, duration_s, center_s ) hard_intervals, hard_valid, hard_quality = _posterior_hard_intervals( log_probs, path, targets, word_targets, duration_s, frame_step, center_s ) else: occupancy = np.zeros((len(log_probs), len(words)), dtype=np.float32) boundary_summary = [ {"start_mean_s": math.nan, "end_mean_s": math.nan, "start_p05_s": math.nan, "start_p95_s": math.nan, "end_p05_s": math.nan, "end_p95_s": math.nan, "start_width90_s": math.nan, "end_width90_s": math.nan} for _ in words ] hard_intervals = np.zeros((len(words), 2), dtype=np.float32) hard_valid = np.zeros(len(words), dtype=np.bool_) hard_quality = np.zeros(len(words), dtype=np.float32) ctc_frame_times = (np.arange(len(log_probs), dtype=np.float32) * frame_step + center_s).clip(0.0, duration_s) return Sample( sample_id=record["sample_id"], video_id=record["video_id"], clip_id=record["clip_id"], duration_s=duration_s, sentiment=record["sentiment"], polarity=record["polarity"], text_words=words, text_features=text_vectors.astype(np.float16), text_valid=text_valid, hard_intervals=hard_intervals, hard_valid=hard_valid, hard_quality=hard_quality, audio_times=audio_times, audio_features=audio_vectors.astype(np.float16), audio_observed=audio_observed, audio_quality=np.clip((len(waveform) - np.arange(len(audio_times)) * FRAME_STEP) / FRAME_LENGTH, 0.0, 1.0).astype(np.float32), audio_quality_available=np.ones(len(audio_times), bool), vision_times=vision_times, vision_frame_indices=vision_frame_indices, vision_features=vision_vectors.astype(np.float16), vision_observed=vision_observed, vision_quality=vision_quality, vision_quality_available=vision_quality_available, ctc_times=ctc_frame_times, ctc_occupancy=occupancy.astype(np.float16), speech_features=deep_states.astype(np.float16), boundary_summary=boundary_summary, video_sha256=video_sha256, ) def _save_cache(path: Path, sample: Sample, cache_schema: str) -> None: path.parent.mkdir(parents=True, exist_ok=True) temporary = path.with_suffix(".npz.tmp") with temporary.open("wb") as stream: np.savez_compressed( stream, cache_schema=np.asarray(cache_schema), sample_id=np.asarray(sample.sample_id), video_id=np.asarray(sample.video_id), clip_id=np.asarray(sample.clip_id), duration_s=np.asarray(sample.duration_s, dtype=np.float32), sentiment=np.asarray(sample.sentiment, dtype=np.float32), polarity=np.asarray(sample.polarity, dtype=np.int8), text_words=np.asarray(sample.text_words, dtype="U160"), text_features=sample.text_features.astype(np.float16), text_valid=sample.text_valid, hard_intervals=sample.hard_intervals.astype(np.float32), hard_valid=sample.hard_valid, hard_quality=sample.hard_quality.astype(np.float32), audio_times=sample.audio_times.astype(np.float32), audio_features=sample.audio_features.astype(np.float16), audio_observed=sample.audio_observed, audio_quality=sample.audio_quality.astype(np.float32), audio_quality_available=sample.audio_quality_available, vision_times=sample.vision_times.astype(np.float32), vision_frame_indices=sample.vision_frame_indices.astype(np.int32), vision_features=sample.vision_features.astype(np.float16), vision_observed=sample.vision_observed, vision_quality=sample.vision_quality.astype(np.float32), vision_quality_available=sample.vision_quality_available, ctc_times=sample.ctc_times.astype(np.float32), ctc_occupancy=sample.ctc_occupancy.astype(np.float16), speech_features=sample.speech_features.astype(np.float16), boundary_summary=json.dumps(sample.boundary_summary, ensure_ascii=False), video_sha256=np.asarray(sample.video_sha256), ) temporary.replace(path) def _load_cache(path: Path, record: dict[str, Any], source_hash: str, cache_schema: str) -> Sample: with np.load(path, allow_pickle=False) as data: if str(data["cache_schema"].item()) != cache_schema: raise ValueError("feature cache schema mismatch") if str(data["sample_id"].item()) != record["sample_id"]: raise ValueError("feature cache sample key mismatch") if str(data["video_sha256"].item()) != source_hash: raise ValueError("feature cache source hash mismatch") return Sample( sample_id=record["sample_id"], video_id=record["video_id"], clip_id=record["clip_id"], duration_s=float(data["duration_s"]), sentiment=float(data["sentiment"]), polarity=int(data["polarity"]), text_words=data["text_words"].astype(str).tolist(), text_features=np.asarray(data["text_features"], dtype=np.float16), text_valid=np.asarray(data["text_valid"], dtype=np.bool_), hard_intervals=np.asarray(data["hard_intervals"], dtype=np.float32), hard_valid=np.asarray(data["hard_valid"], dtype=np.bool_), hard_quality=np.asarray(data["hard_quality"], dtype=np.float32), audio_times=np.asarray(data["audio_times"], dtype=np.float32), audio_features=np.asarray(data["audio_features"], dtype=np.float16), audio_observed=np.asarray(data["audio_observed"], dtype=np.bool_), audio_quality=np.asarray(data["audio_quality"], dtype=np.float32), audio_quality_available=np.asarray(data["audio_quality_available"], dtype=np.bool_), vision_times=np.asarray(data["vision_times"], dtype=np.float32), vision_frame_indices=np.asarray(data["vision_frame_indices"], dtype=np.int32), vision_features=np.asarray(data["vision_features"], dtype=np.float16), vision_observed=np.asarray(data["vision_observed"], dtype=np.bool_), vision_quality=np.asarray(data["vision_quality"], dtype=np.float32), vision_quality_available=np.asarray(data["vision_quality_available"], dtype=np.bool_), ctc_times=np.asarray(data["ctc_times"], dtype=np.float32), ctc_occupancy=np.asarray(data["ctc_occupancy"], dtype=np.float16), speech_features=np.asarray(data["speech_features"], dtype=np.float16), boundary_summary=json.loads(str(data["boundary_summary"].item())), video_sha256=source_hash, ) def _voronoi_intervals(times: np.ndarray, duration_s: float) -> np.ndarray: if not len(times): return np.empty((0, 2), dtype=np.float32) centers = np.clip(np.asarray(times, dtype=np.float64), 0.0, duration_s) if np.any(np.diff(centers) < -1e-7): order = np.argsort(centers) centers = centers[order] else: order = np.arange(len(centers)) # Duplicate PTS values are coalesced to a tiny increasing interval for assignment. for index in range(1, len(centers)): if centers[index] <= centers[index - 1]: centers[index] = min(duration_s, centers[index - 1] + 1e-6) edges = np.empty(len(centers) + 1, dtype=np.float64) edges[0], edges[-1] = 0.0, duration_s if len(centers) > 1: edges[1:-1] = (centers[:-1] + centers[1:]) / 2 edges = np.maximum.accumulate(np.clip(edges, 0.0, duration_s)) intervals = np.column_stack((edges[:-1], edges[1:])).astype(np.float32) if len(order) != len(intervals) or not np.array_equal(order, np.arange(len(order))): restored = np.zeros_like(intervals) restored[order] = intervals return restored return intervals def _grid_edges(duration_s: float) -> np.ndarray: count = max(1, int(math.ceil(duration_s / GRID_STEP_S))) edges = np.arange(count + 1, dtype=np.float64) * GRID_STEP_S edges[-1] = duration_s return edges def _project_rows( values: np.ndarray, observed: np.ndarray, source_intervals: np.ndarray, qualities: np.ndarray, target_edges: np.ndarray, ) -> tuple[np.ndarray, np.ndarray, np.ndarray]: dim = values.shape[1] if values.ndim == 2 else 0 count = len(target_edges) - 1 output = np.zeros((count, dim), dtype=np.float32) mask = np.zeros((count, dim), dtype=np.bool_) coverage = np.zeros((count, dim), dtype=np.float32) if not dim or not len(source_intervals): return output, mask, coverage intervals = np.asarray(source_intervals, dtype=np.float64) qualities = np.maximum(np.asarray(qualities, dtype=np.float64), 0.0) values64 = np.asarray(values, dtype=np.float64) for target_index, (left, right) in enumerate(zip(target_edges[:-1], target_edges[1:])): if right <= left: continue physical_overlap = np.maximum( 0.0, np.minimum(intervals[:, 1], right) - np.maximum(intervals[:, 0], left), ) weighted_overlap = physical_overlap * qualities candidate = np.flatnonzero(physical_overlap > 0) if not len(candidate): continue local_mask = observed[candidate] local_weight = weighted_overlap[candidate, None] * local_mask denominator = local_weight.sum(axis=0) physical_coverage = (physical_overlap[candidate, None] * local_mask).sum(axis=0) good = denominator > 0 if good.any(): output[target_index, good] = ( (values64[candidate] * local_weight).sum(axis=0)[good] / denominator[good] ).astype(np.float32) mask[target_index, good] = True coverage[target_index, good] = np.minimum(1.0, physical_coverage[good] / (right - left)).astype(np.float32) return output, mask, coverage def _so3_weighted_mean(euler_xyz: np.ndarray, weights: np.ndarray) -> np.ndarray | None: positive = weights > 0 if not positive.any(): return None rotations = Rotation.from_euler("xyz", np.asarray(euler_xyz[positive], dtype=np.float64)) local_weights = np.asarray(weights[positive], dtype=np.float64) local_weights /= local_weights.sum() matrices = rotations.as_matrix() mean_matrix = np.einsum("n,nij->ij", local_weights, matrices) u, _, vh = np.linalg.svd(mean_matrix) mean_matrix = u @ vh if np.linalg.det(mean_matrix) < 0: u[:, -1] *= -1 mean_matrix = u @ vh mean_rotation = Rotation.from_matrix(mean_matrix) for _ in range(30): residual = (mean_rotation.inv() * rotations).as_rotvec() delta = np.einsum("n,nd->d", local_weights, residual) if np.linalg.norm(delta) < 1e-8: break mean_rotation = mean_rotation * Rotation.from_rotvec(delta) return mean_rotation.as_euler("xyz").astype(np.float32) def _project_vision_geometry(sample: Sample, target_edges: np.ndarray, geometry_aware: bool) -> tuple[np.ndarray, np.ndarray, np.ndarray]: intervals = _voronoi_intervals(sample.vision_times, sample.duration_s) quality = sample.vision_quality values, masks, coverage = _project_rows( sample.vision_features.astype(np.float32), sample.vision_observed, intervals, quality, target_edges ) if not geometry_aware or not len(intervals): return values, masks, coverage raw = sample.vision_features.astype(np.float32) for target_index, (left, right) in enumerate(zip(target_edges[:-1], target_edges[1:])): overlap = np.maximum(0.0, np.minimum(intervals[:, 1], right) - np.maximum(intervals[:, 0], left)) if not np.any(overlap > 0): continue weighted_overlap = overlap * quality pose_rows = (weighted_overlap > 0) & sample.vision_observed[:, 17:20].all(axis=1) if pose_rows.any(): vector = _so3_weighted_mean(raw[pose_rows, 17:20], weighted_overlap[pose_rows]) if vector is not None: values[target_index, 17:20] = vector masks[target_index, 17:20] = True for start in (23, 26): valid_rows = (weighted_overlap > 0) & sample.vision_observed[:, start:start + 3].all(axis=1) if valid_rows.any(): weights = weighted_overlap[valid_rows].astype(np.float64) vector = np.average(raw[valid_rows, start:start + 3], axis=0, weights=weights) norm = float(np.linalg.norm(vector)) if norm > 1e-6: values[target_index, start:start + 3] = vector / norm masks[target_index, start:start + 3] = True else: values[target_index, start:start + 3] = 0.0 masks[target_index, start:start + 3] = False return values, masks, coverage def _project_posterior_text(sample: Sample, target_edges: np.ndarray) -> tuple[np.ndarray, np.ndarray, np.ndarray]: words, dim = sample.text_features.shape output = np.zeros((len(target_edges) - 1, dim), dtype=np.float32) mask = np.zeros((len(target_edges) - 1, dim), dtype=np.bool_) coverage = np.zeros((len(target_edges) - 1, dim), dtype=np.float32) if not words or not len(sample.ctc_times): return output, mask, coverage frame_step = CTC_FRAME_STEP_S first_edge = max(0.0, float(sample.ctc_times[0]) - frame_step / 2) frame_edges = np.clip( first_edge + np.arange(len(sample.ctc_times) + 1, dtype=np.float64) * frame_step, 0.0, _audio_end_time(sample), ) occupancy = sample.ctc_occupancy.astype(np.float64) for target_index, (left, right) in enumerate(zip(target_edges[:-1], target_edges[1:])): overlap = np.maximum(0.0, np.minimum(frame_edges[1:], right) - np.maximum(frame_edges[:-1], left)) if not np.any(overlap > 0): continue alpha = (occupancy * overlap[:, None]).sum(axis=0) alpha *= sample.text_valid total = float(alpha.sum()) if total <= 0: continue output[target_index] = (alpha @ sample.text_features.astype(np.float64) / total).astype(np.float32) mask[target_index] = True coverage[target_index] = min(1.0, total / max(right - left, 1e-8)) return output, mask, coverage def _normalized_hard_quality(sample: Sample) -> np.ndarray: text_quality = np.ones_like(sample.hard_quality, dtype=np.float32) valid_quality = sample.hard_valid & (sample.hard_quality > 0) if valid_quality.any(): text_quality[valid_quality] = np.clip(sample.hard_quality[valid_quality], 0.0, 1.0) return text_quality def _audio_end_time(sample: Sample) -> float: """Recover decoded sample support from the final partial analysis window.""" valid = np.flatnonzero(np.asarray(sample.audio_quality) > 0) if not len(valid): return 0.0 index = int(valid[-1]) quality = float(np.clip(sample.audio_quality[index], 0.0, 1.0)) sample_count = index * FRAME_STEP + int(round(quality * FRAME_LENGTH)) return float(min(sample.duration_s, sample_count / SAMPLE_RATE)) def _make_views(sample: Sample) -> tuple[View, View, View]: edges = _grid_edges(sample.duration_s) audio_end_s = _audio_end_time(sample) hard_intervals = np.clip(sample.hard_intervals.astype(np.float32), 0.0, audio_end_s) hard_valid = sample.hard_valid & (hard_intervals[:, 1] > hard_intervals[:, 0]) hard_text_observed = np.broadcast_to(sample.text_valid[:, None], sample.text_features.shape).copy() text_quality = _normalized_hard_quality(sample) text_quality[~hard_valid] = 0.0 text_hard = _project_rows( sample.text_features.astype(np.float32), hard_text_observed, hard_intervals, text_quality, edges, ) audio = _project_rows( sample.audio_features.astype(np.float32), sample.audio_observed, _voronoi_intervals(sample.audio_times, audio_end_s), sample.audio_quality, edges, ) vision_b0 = _project_vision_geometry(sample, edges, geometry_aware=False) vision_b1 = _project_vision_geometry(sample, edges, geometry_aware=True) speech = _project_rows( sample.speech_features.astype(np.float32), np.ones(sample.speech_features.shape, dtype=np.bool_), _voronoi_intervals(sample.ctc_times, audio_end_s), np.ones(len(sample.ctc_times), dtype=np.float32), edges, ) post_text = _project_posterior_text(sample, edges) view_b0 = View( features={"text": text_hard[0], "audio": audio[0], "vision": vision_b0[0], "speech": speech[0]}, observed={"text": text_hard[1], "audio": audio[1], "vision": vision_b0[1], "speech": speech[1]}, coverage={"text": text_hard[2], "audio": audio[2], "vision": vision_b0[2], "speech": speech[2]}, edges=edges, ) view_b1 = View( features={"text": text_hard[0], "audio": audio[0], "vision": vision_b1[0], "speech": speech[0]}, observed={"text": text_hard[1], "audio": audio[1], "vision": vision_b1[1], "speech": speech[1]}, coverage={"text": text_hard[2], "audio": audio[2], "vision": vision_b1[2], "speech": speech[2]}, edges=edges, ) # B2 changes only text's time projection; media timestamps and their masks stay fixed. view_b2 = View( features={**view_b1.features, "text": post_text[0]}, observed={**view_b1.observed, "text": post_text[1]}, coverage={**view_b1.coverage, "text": post_text[2]}, edges=edges, ) return view_b0, view_b1, view_b2 def _fit_feature_scaler(samples: list[Sample], views: list[View], indices: np.ndarray, mode: str) -> dict[str, tuple[np.ndarray, np.ndarray]]: modalities = ("text", "audio", "vision", "speech") scalers: dict[str, tuple[np.ndarray, np.ndarray]] = {} for modality in modalities: dims = views[int(indices[0])].features[modality].shape[1] columns: list[np.ndarray] = [] for dim in range(dims): values = [ views[int(index)].features[modality][:, dim][views[int(index)].observed[modality][:, dim]] for index in indices ] nonempty = [array for array in values if len(array)] columns.append(np.concatenate(nonempty).astype(np.float64) if nonempty else np.empty(0, dtype=np.float64)) center = np.zeros(dims, dtype=np.float64) scale = np.ones(dims, dtype=np.float64) for dim, values in enumerate(columns): if not len(values): continue if mode == "robust": center[dim] = np.median(values) mad = 1.4826 * np.median(np.abs(values - center[dim])) std = np.std(values) scale[dim] = mad if mad > 1e-8 else std if std > 1e-8 else 1.0 else: center[dim] = np.mean(values) std = np.std(values) scale[dim] = std if std > 1e-8 else 1.0 scalers[modality] = (center.astype(np.float32), scale.astype(np.float32)) return scalers def _scaled_view(view: View, scalers: dict[str, tuple[np.ndarray, np.ndarray]], clip: bool) -> View: features: dict[str, np.ndarray] = {} for modality, value in view.features.items(): center, scale = scalers[modality] transformed = (value - center[None, :]) / scale[None, :] if clip: transformed = np.clip(transformed, -8.0, 8.0) transformed = np.where(view.observed[modality], transformed, 0.0) features[modality] = transformed.astype(np.float32) return View(features=features, observed=view.observed, coverage=view.coverage, edges=view.edges) def _path_signature(times: np.ndarray, values: np.ndarray) -> np.ndarray: # Time-augmented piecewise-linear path; return first level and antisymmetric area. dimension = values.shape[1] + 1 if len(times) < 3: return np.zeros(dimension + dimension * (dimension - 1) // 2, dtype=np.float32) local_time = (times - times[0]) / max(times[-1] - times[0], 1e-8) path = np.column_stack((local_time, values)).astype(np.float64) increments = np.diff(path, axis=0) first = increments.sum(axis=0) areas: list[float] = [] for left in range(dimension): for right in range(left + 1, dimension): prefix_left = np.cumsum(increments[:, left]) - increments[:, left] prefix_right = np.cumsum(increments[:, right]) - increments[:, right] areas.append(float(0.5 * np.sum(prefix_left * increments[:, right] - prefix_right * increments[:, left]))) return np.concatenate((first, np.asarray(areas, dtype=np.float64))).astype(np.float32) def _signature_for_segment( times: np.ndarray, values: np.ndarray, valid: np.ndarray, left: float, right: float, ) -> tuple[np.ndarray, float]: dim = values.shape[1] output_dim = (dim + 1) + (dim + 1) * dim // 2 selected = (times >= left) & (times < right) & valid indices = np.flatnonzero(selected) if len(indices) < 3: return np.zeros(output_dim, dtype=np.float32), 0.0 # Split at every missing source position; never connect across a missing interval. chunks = np.split(indices, np.flatnonzero(np.diff(indices) > 1) + 1) signatures: list[np.ndarray] = [] weights: list[float] = [] for chunk in chunks: if len(chunk) < 3: continue duration = float(max(times[chunk[-1]] - times[chunk[0]], 0.0)) if duration <= 0: continue signatures.append(_path_signature(times[chunk], values[chunk])) weights.append(duration) if not weights: return np.zeros(output_dim, dtype=np.float32), 0.0 weight_array = np.asarray(weights, dtype=np.float64) signature = np.average(np.stack(signatures), axis=0, weights=weight_array) valid_duration = min(float(np.sum(weight_array)), max(right - left, 1e-8)) return signature.astype(np.float32), float(np.clip(valid_duration / (right - left), 0.0, 1.0)) def _dynamic_signature(sample: Sample, scaled_view: View) -> np.ndarray: duration = sample.duration_s segments = np.linspace(0.0, duration, 6) centers = (scaled_view.edges[:-1] + scaled_view.edges[1:]) / 2 audio = scaled_view.features["audio"] audio_mask = scaled_view.observed["audio"] vision = scaled_view.features["vision"] vision_mask = scaled_view.observed["vision"] # log-F0, log-energy; OpenFace AU12 intensity and normalized mouth opening. audio_values = audio[:, [67, 66]] audio_valid = audio_mask[:, [67, 66]].all(axis=1) vision_values = vision[:, [10, 31]] vision_valid = vision_mask[:, [10, 31]].all(axis=1) result: list[float] = [] for index in range(5): left, right = float(segments[index]), float(segments[index + 1]) signature, coverage = _signature_for_segment(centers, audio_values, audio_valid, left, right) result.extend(signature.tolist()) result.append(coverage) signature, coverage = _signature_for_segment(centers, vision_values, vision_valid, left, right) result.extend(signature.tolist()) result.append(coverage) return np.asarray(result, dtype=np.float32) def _pool_progress(view: View, include_speech: bool, extra: np.ndarray | None = None) -> np.ndarray: duration = float(view.edges[-1]) progress_edges = np.linspace(0.0, duration, 6) centers = (view.edges[:-1] + view.edges[1:]) / 2 cell_lengths = np.diff(view.edges) names = ["text", "audio", "vision"] + (["speech"] if include_speech else []) pooled: list[np.ndarray] = [] for name in names: values, masks, coverage = view.features[name], view.observed[name], view.coverage[name] rows = [] for left, right in zip(progress_edges[:-1], progress_edges[1:]): select = (centers >= left) & (centers < right) weights = cell_lengths[select, None] * coverage[select] * masks[select] denom = weights.sum(axis=0) good = denom > 0 out = np.zeros(values.shape[1], dtype=np.float32) if good.any(): out[good] = ((values[select] * weights).sum(axis=0)[good] / denom[good]).astype(np.float32) rows.append(out) pooled.append(np.stack(rows).reshape(-1)) if extra is not None: pooled.append(extra) return np.concatenate(pooled).astype(np.float32) def _view_coverage(view: View, modality: str) -> float: mask = view.observed[modality] cov = view.coverage[modality] if not mask.size: return 0.0 return float(np.mean(cov[mask])) if mask.any() else 0.0 def _metrics(y_class: np.ndarray, p_class: np.ndarray, y_value: np.ndarray, p_value: np.ndarray) -> dict[str, float]: pearson = float(np.corrcoef(y_value, p_value)[0, 1]) if len(y_value) > 1 and np.std(p_value) > 1e-12 else math.nan return { "accuracy": float(accuracy_score(y_class, p_class)), "macro_f1": float(f1_score(y_class, p_class, labels=[0, 1, 2], average="macro", zero_division=0)), "mae": float(mean_absolute_error(y_value, p_value)), "pearson": pearson, } def _write_csv(path: Path, rows: list[dict[str, Any]], fields: list[str]) -> None: path.parent.mkdir(parents=True, exist_ok=True) with path.open("w", encoding="utf-8-sig", newline="") as stream: writer = csv.DictWriter(stream, fieldnames=fields, extrasaction="ignore") writer.writeheader() writer.writerows(rows) def _group_bootstrap_deltas( samples: list[Sample], predictions: dict[str, dict[int, tuple[int, float]]], repeats: int, ) -> list[dict[str, Any]]: rng = np.random.default_rng(SEED) groups: dict[str, list[int]] = {} for index, sample in enumerate(samples): groups.setdefault(sample.video_id, []).append(index) group_names = sorted(groups) y_class = np.asarray([sample.polarity for sample in samples], dtype=np.int64) y_value = np.asarray([sample.sentiment for sample in samples], dtype=np.float64) records: list[dict[str, Any]] = [] reference = "B1" for method in ("B0", "B2", "B3", "B4"): point_ref = np.asarray([predictions[reference][i][0] for i in range(len(samples))], dtype=np.int64) point_method = np.asarray([predictions[method][i][0] for i in range(len(samples))], dtype=np.int64) value_ref = np.asarray([predictions[reference][i][1] for i in range(len(samples))], dtype=np.float64) value_method = np.asarray([predictions[method][i][1] for i in range(len(samples))], dtype=np.float64) for metric in ("accuracy", "macro_f1", "mae", "pearson"): deltas: list[float] = [] for _ in range(repeats): chosen = rng.choice(group_names, size=len(group_names), replace=True) indices = np.asarray([index for name in chosen for index in groups[str(name)]], dtype=np.int64) a = _metrics(y_class[indices], point_ref[indices], y_value[indices], value_ref[indices])[metric] b = _metrics(y_class[indices], point_method[indices], y_value[indices], value_method[indices])[metric] if math.isfinite(a) and math.isfinite(b): deltas.append(b - a) ref_metric = _metrics(y_class, point_ref, y_value, value_ref)[metric] method_metric = _metrics(y_class, point_method, y_value, value_method)[metric] records.append({ "comparison": f"{method}-B1", "metric": metric, "reference_oof": ref_metric, "method_oof": method_metric, "delta_oof": method_metric - ref_metric if math.isfinite(method_metric) and math.isfinite(ref_metric) else math.nan, "group_bootstrap_ci95_low": float(np.quantile(deltas, 0.025)) if deltas else math.nan, "group_bootstrap_ci95_high": float(np.quantile(deltas, 0.975)) if deltas else math.nan, "bootstrap_repeats": len(deltas), }) return records def _evaluate(samples: list[Sample], views: list[tuple[View, View, View]], bootstrap_repeats: int) -> tuple[list[dict[str, Any]], list[dict[str, Any]], list[dict[str, Any]], list[dict[str, Any]], list[dict[str, Any]]]: y_class = np.asarray([sample.polarity for sample in samples], dtype=np.int64) y_value = np.asarray([sample.sentiment for sample in samples], dtype=np.float64) groups = np.asarray([sample.video_id for sample in samples]) splitter = GroupKFold(n_splits=N_FOLDS) folds = list(splitter.split(np.zeros(len(samples)), y_class, groups)) methods = ("B0", "B1", "B2", "B3", "B4") class_oof = {method: np.full(len(samples), -1, dtype=np.int64) for method in methods} value_oof = {method: np.full(len(samples), np.nan, dtype=np.float64) for method in methods} predictions: dict[str, dict[int, tuple[int, float]]] = {method: {} for method in methods} fold_rows: list[dict[str, Any]] = [] split_rows: list[dict[str, Any]] = [] for fold_index, (train_indices, test_indices) in enumerate(folds, start=1): train_groups = set(groups[train_indices].tolist()) test_groups = set(groups[test_indices].tolist()) if train_groups & test_groups: raise AssertionError("video_id leakage across grouped fold") for index in test_indices: split_rows.append({ "sample_id": samples[index].sample_id, "video_id": samples[index].video_id, "fold": fold_index, "split": "valid_oof", }) variant_views = { "B0": [item[0] for item in views], "B1": [item[1] for item in views], "B2": [item[2] for item in views], "B3": [item[1] for item in views], "B4": [item[1] for item in views], } for method in methods: source_views = variant_views[method] scale_mode = "standard" if method == "B0" else "robust" scalers = _fit_feature_scaler(samples, source_views, train_indices, scale_mode) scaled = [_scaled_view(view, scalers, clip=(scale_mode == "robust")) for view in source_views] add_speech = method == "B4" x = np.stack([ _pool_progress( scaled[int(index)], include_speech=add_speech, extra=_dynamic_signature(samples[int(index)], scaled[int(index)]) if method == "B3" else None, ) for index in train_indices ]) x_valid = np.stack([ _pool_progress( scaled[int(index)], include_speech=add_speech, extra=_dynamic_signature(samples[int(index)], scaled[int(index)]) if method == "B3" else None, ) for index in test_indices ]) classifier = LogisticRegression(C=0.05, max_iter=2500, solver="lbfgs", random_state=SEED) classifier.fit(x, y_class[train_indices]) predicted_class = classifier.predict(x_valid) regressor = Ridge(alpha=25.0, solver="lsqr") regressor.fit(x, y_value[train_indices]) predicted_value = np.clip(regressor.predict(x_valid), -3.0, 3.0) class_oof[method][test_indices] = predicted_class value_oof[method][test_indices] = predicted_value for index, pred_class, pred_value in zip(test_indices, predicted_class, predicted_value): predictions[method][int(index)] = (int(pred_class), float(pred_value)) fold_metric = _metrics(y_class[test_indices], predicted_class, y_value[test_indices], predicted_value) fold_rows.append({ "method": method, "fold": fold_index, "train_samples": int(len(train_indices)), "valid_samples": int(len(test_indices)), "train_video_groups": len(train_groups), "valid_video_groups": len(test_groups), "accuracy": fold_metric["accuracy"], "macro_f1": fold_metric["macro_f1"], "mae": fold_metric["mae"], "pearson": fold_metric["pearson"], "feature_dimension": int(x.shape[1]), }) print(f"[CV] fold {fold_index}/{N_FOLDS}: train={len(train_indices)} clips/{len(train_groups)} videos; valid={len(test_indices)} clips/{len(test_groups)} videos", flush=True) summary_rows: list[dict[str, Any]] = [] for method in methods: fold_metrics = [row for row in fold_rows if row["method"] == method] overall = _metrics(y_class, class_oof[method], y_value, value_oof[method]) row: dict[str, Any] = { "method": method, "sample_count": len(samples), "video_group_count": len(set(groups.tolist())), "oof_accuracy": overall["accuracy"], "oof_macro_f1": overall["macro_f1"], "oof_mae": overall["mae"], "oof_pearson": overall["pearson"], } for metric in ("accuracy", "macro_f1", "mae", "pearson"): values = np.asarray([item[metric] for item in fold_metrics if math.isfinite(item[metric])], dtype=np.float64) row[f"fold_{metric}_mean"] = float(values.mean()) if len(values) else math.nan row[f"fold_{metric}_sd"] = float(values.std(ddof=1)) if len(values) > 1 else math.nan summary_rows.append(row) prediction_rows = [] for index, sample in enumerate(samples): row: dict[str, Any] = { "sample_id": sample.sample_id, "video_id": sample.video_id, "clip_id": sample.clip_id, "true_polarity": sample.polarity, "true_polarity_name": ("negative", "neutral", "positive")[sample.polarity], "true_sentiment": sample.sentiment, } for method in methods: row[f"{method}_predicted_polarity"] = int(class_oof[method][index]) row[f"{method}_predicted_sentiment"] = float(value_oof[method][index]) prediction_rows.append(row) delta_rows = _group_bootstrap_deltas(samples, predictions, bootstrap_repeats) return summary_rows, fold_rows, prediction_rows, split_rows, delta_rows def _sample_output_rows(samples: list[Sample], views: list[tuple[View, View, View]]) -> tuple[list[dict[str, Any]], list[dict[str, Any]]]: sample_rows: list[dict[str, Any]] = [] modality_rows: list[dict[str, Any]] = [] for sample, (view0, view1, view2) in zip(samples, views): finite_boundaries = [item for item in sample.boundary_summary if math.isfinite(item["start_width90_s"]) and math.isfinite(item["end_width90_s"])] hard_text_coverage = _view_coverage(view0, "text") post_text_coverage = _view_coverage(view2, "text") sample_rows.append({ "sample_id": sample.sample_id, "video_id": sample.video_id, "clip_id": sample.clip_id, "source_video_sha256": sample.video_sha256, "duration_s": sample.duration_s, "word_count": len(sample.text_words), "hard_aligned_word_count": int(sample.hard_valid.sum()), "unlocated_word_count": int((~sample.hard_valid).sum()), "posterior_usable_word_count": len(finite_boundaries), "mean_start_interval_width90_s": float(np.mean([item["start_width90_s"] for item in finite_boundaries])) if finite_boundaries else math.nan, "mean_end_interval_width90_s": float(np.mean([item["end_width90_s"] for item in finite_boundaries])) if finite_boundaries else math.nan, "text_hard_grid_coverage": hard_text_coverage, "text_posterior_grid_occupancy": post_text_coverage, "audio_grid_coverage": _view_coverage(view1, "audio"), "vision_grid_coverage": _view_coverage(view1, "vision"), "audio_native_rows": len(sample.audio_times), "vision_native_rows": len(sample.vision_times), "status": "ok", }) for modality, native_length, native_dim, view, key in ( ("text", len(sample.text_words), 768, view1, "text"), ("audio", len(sample.audio_times), 74, view1, "audio"), ("vision", len(sample.vision_times), 35, view1, "vision"), ): coverage = view.coverage[key] observed = view.observed[key] observed_duration = float((coverage * np.diff(view.edges)[:, None] * observed).sum(axis=0).mean()) if observed.size else 0.0 modality_rows.append({ "sample_id": sample.sample_id, "video_id": sample.video_id, "clip_id": sample.clip_id, "modality": modality, "source_duration_s": sample.duration_s, "observed_duration_s_mean_dimension": observed_duration, "native_length": native_length, "native_dimension": native_dim, "main_grid_length": len(view.edges) - 1, "grid_step_s": GRID_STEP_S, "grid_dimension": native_dim, "mean_grid_coverage": _view_coverage(view, key), "status": "ok" if native_length else "no_native_rows", "source_video_sha256": sample.video_sha256, }) return sample_rows, modality_rows def _create_plot(rows: list[dict[str, Any]], output_path: Path) -> None: labels = [row["method"] for row in rows] figure, axes = plt.subplots(1, 3, figsize=(14, 4.5), constrained_layout=True) metrics = (("oof_macro_f1", "Macro F1 (higher is better)"), ("oof_mae", "MAE (lower is better)"), ("oof_pearson", "Pearson r (higher is better)")) colors = ("#718096", "#2463A6", "#6A8E3A", "#C77D27", "#8357A5") for axis, (column, title) in zip(axes, metrics): values = [float(row[column]) for row in rows] bars = axis.bar(labels, values, color=colors) axis.set_title(title) axis.grid(axis="y", alpha=0.25) for bar, value in zip(bars, values): axis.text(bar.get_x() + bar.get_width() / 2, bar.get_height(), f"{value:.3f}", ha="center", va="bottom", fontsize=8) figure.suptitle("B0-B4 emotion probe on grouped out-of-fold predictions") output_path.parent.mkdir(parents=True, exist_ok=True) figure.savefig(output_path, dpi=170) plt.close(figure) def _create_typical_alignment_figure( samples: list[Sample], views: list[tuple[View, View, View]], output_path: Path, ) -> str: eligible = [ index for index, sample in enumerate(samples) if sample.text_words and sample.hard_valid.all() and sample.vision_observed.any() ] if not eligible: return "" median_duration = float(np.median([samples[index].duration_s for index in eligible])) chosen_index = min(eligible, key=lambda index: abs(samples[index].duration_s - median_duration)) sample = samples[chosen_index] view_b0, _, view_b2 = views[chosen_index] middle = len(sample.text_words) // 2 left_word = max(0, middle - 5) right_word = min(len(sample.text_words), left_word + 10) left_word = max(0, right_word - 10) word_indices = np.arange(left_word, right_word) figure = plt.figure(figsize=(15, 11), constrained_layout=True) grid = figure.add_gridspec(4, 1, height_ratios=(2.6, 1.1, 1.4, 2.1)) alignment_axis = figure.add_subplot(grid[0, 0]) coverage_axis = figure.add_subplot(grid[1, 0]) feature_axis = figure.add_subplot(grid[2, 0]) image_grid = grid[3, 0].subgridspec(1, 5, wspace=0.04) image_axes = [figure.add_subplot(image_grid[0, i]) for i in range(5)] for y_position, word_index in enumerate(word_indices): hard_start, hard_end = sample.hard_intervals[word_index] alignment_axis.plot((hard_start, hard_end), (y_position, y_position), color="#2563a6", lw=5, solid_capstyle="butt") boundary = sample.boundary_summary[word_index] if math.isfinite(boundary["start_p05_s"]) and math.isfinite(boundary["start_p95_s"]): center = boundary["start_mean_s"] alignment_axis.errorbar( center, y_position + 0.19, xerr=np.asarray([[max(0.0, center - boundary["start_p05_s"]),], [max(0.0, boundary["start_p95_s"] - center),]]), fmt="o", color="#d17a00", markersize=3, capsize=2, lw=1, ) if math.isfinite(boundary["end_p05_s"]) and math.isfinite(boundary["end_p95_s"]): center = boundary["end_mean_s"] alignment_axis.errorbar( center, y_position - 0.19, xerr=np.asarray([[max(0.0, center - boundary["end_p05_s"]),], [max(0.0, boundary["end_p95_s"] - center),]]), fmt="s", color="#638b32", markersize=3, capsize=2, lw=1, ) alignment_axis.set_yticks(np.arange(len(word_indices)), labels=[sample.text_words[index] for index in word_indices], fontsize=8) alignment_axis.set_ylim(-0.65, len(word_indices) - 0.35) start_values = [float(sample.hard_intervals[index, 0]) for index in word_indices] end_values = [float(sample.hard_intervals[index, 1]) for index in word_indices] alignment_axis.set_xlim(max(0.0, min(start_values) - 0.25), min(sample.duration_s, max(end_values) + 0.25)) alignment_axis.set_xlabel("seconds from clip start") alignment_axis.set_title("Hard CTC word interval (blue); 90% internal start/end boundary intervals (orange/green)") alignment_axis.grid(axis="x", alpha=0.25) times = (view_b0.edges[:-1] + view_b0.edges[1:]) / 2 coverage_axis.plot(times, view_b0.coverage["text"].mean(axis=1), label="hard word coverage", color="#2563a6") coverage_axis.plot(times, view_b2.coverage["text"].mean(axis=1), label="posterior word occupancy", color="#d17a00") coverage_axis.plot(times, view_b0.coverage["audio"].mean(axis=1), label="audio coverage", color="#638b32") coverage_axis.plot(times, view_b0.coverage["vision"].mean(axis=1), label="face feature coverage", color="#8357a5") coverage_axis.set_ylim(0, 1.05) coverage_axis.set_xlim(0, sample.duration_s) coverage_axis.set_ylabel("coverage / occupancy") coverage_axis.legend(ncol=4, fontsize=8, loc="upper right") coverage_axis.grid(alpha=0.2) audio = sample.audio_features.astype(np.float32) audio_time = sample.audio_times feature_axis.plot(audio_time, audio[:, 66], color="#2563a6", lw=1, label="log energy") pitch_valid = sample.audio_observed[:, 67] if pitch_valid.any(): pitch_axis = feature_axis.twinx() pitch_axis.plot(audio_time[pitch_valid], np.exp(audio[pitch_valid, 67]), color="#d17a00", lw=0.8, alpha=0.8, label="F0 (Hz)") pitch_axis.set_ylabel("F0 (Hz)") vision_valid = sample.vision_observed[:, 31] if vision_valid.any(): feature_axis.scatter(sample.vision_times[vision_valid], sample.vision_features[vision_valid, 31].astype(np.float32), s=9, color="#638b32", label="mouth aperture") feature_axis.set_xlim(0, sample.duration_s) feature_axis.set_xlabel("seconds from clip start") feature_axis.set_ylabel("log energy / mouth aperture") feature_axis.set_title("Example native acoustic and facial feature trajectories") feature_axis.grid(alpha=0.2) handles, labels = feature_axis.get_legend_handles_labels() if len(handles): feature_axis.legend(handles, labels, loc="upper right", fontsize=8) video_path = DATA_DIR / sample.video_id / f"{sample.clip_id}.mp4" capture = cv2.VideoCapture(str(video_path)) frame_times = np.linspace(0.0, max(sample.duration_s - 0.05, 0.0), len(image_axes)) for axis, timestamp in zip(image_axes, frame_times): capture.set(cv2.CAP_PROP_POS_MSEC, float(timestamp) * 1000.0) ok, bgr = capture.read() if ok: axis.imshow(cv2.cvtColor(bgr, cv2.COLOR_BGR2RGB)) axis.set_title(f"{timestamp:.2f}s", fontsize=8) else: axis.text(0.5, 0.5, "frame unavailable", ha="center", va="center") axis.axis("off") capture.release() figure.suptitle(f"Typical sample alignment audit: {sample.sample_id} ({sample.duration_s:.2f}s)", fontsize=14) output_path.parent.mkdir(parents=True, exist_ok=True) figure.savefig(output_path, dpi=170) plt.close(figure) return sample.sample_id def _write_report( output_dir: Path, summary: list[dict[str, Any]], deltas: list[dict[str, Any]], sample_rows: list[dict[str, Any]], modality_rows: list[dict[str, Any]], manifest: dict[str, Any], ) -> None: best_f1 = max(summary, key=lambda row: row["oof_macro_f1"]) best_mae = min(summary, key=lambda row: row["oof_mae"]) b1 = next(row for row in summary if row["method"] == "B1") delta = {(row["comparison"], row["metric"]): row for row in deltas} b4_f1 = delta[("B4-B1", "macro_f1")] b0_mae = delta[("B0-B1", "mae")] b0_pearson = delta[("B0-B1", "pearson")] b2_f1 = delta[("B2-B1", "macro_f1")] b2_mae = delta[("B2-B1", "mae")] b3_f1 = delta[("B3-B1", "macro_f1")] mean_start_width = float(np.nanmean([row["mean_start_interval_width90_s"] for row in sample_rows])) mean_end_width = float(np.nanmean([row["mean_end_interval_width90_s"] for row in sample_rows])) mean_hard_coverage = float(np.mean([row["text_hard_grid_coverage"] for row in sample_rows])) mean_posterior_occupancy = float(np.mean([row["text_posterior_grid_occupancy"] for row in sample_rows])) mean_vision_coverage = float(np.mean([row["vision_grid_coverage"] for row in sample_rows])) complete_hard_count = sum(row["unlocated_word_count"] == 0 for row in sample_rows) consistent = manifest["inputs"]["label_polarity_consistency_count"] lines = [ "# 问题一 B0-B4 模型对比结果", "", "本报告按 `math/E题V2.pdf` 第 4.9.3 节实施五组单因素对照。全部 100 条附件一视频均使用同一套冻结特征抽取器、0.1 秒主时间网格、观测掩码与按 `video_id` 分组的五折划分。情感标签仅用于折内的轻量预测探针,不进入时间对齐、CTC 后验或动态签名计算。", "", "## 实际特征和比较版本", "", "| 版本 | 实施内容 |", "| --- | --- |", "| B0 | 固定转写 CTC Viterbi 词区间与原生音频/视频时间归属区间的质量加权投影;训练折均值/标准差缩放。 |", "| B1 | 在 B0 上改用训练折中位数与 MAD 的稳健尺度(MAD 为零时回退到训练折标准差),并对头部旋转用 SO(3) 均值、视线用归一化向量均值。 |", "| B2 | 在 B1 上只将文本词向量的硬边界投影替换为固定转写 CTC 状态图的前向–后向占据概率投影;媒体时间戳仍按物理时间。 |", "| B3 | 在 B1 上仅附加音频 log-F0/能量与 OpenFace AU12/嘴部开合率的时间增广一、二阶路径签名;缺测点之间不连线。 |", "| B4 | 在 B1 上仅附加冻结 Wav2Vec2-base-960h 最后四层的时间投影表示。 |", "", "附件一现有标签表和原始视频被直接使用。文本为 768 维冻结 BERT-base-uncased 末四层均值;音频为 74 维(log-Mel 40、MFCC 13、ΔMFCC 13、韵律/谱统计 8);视觉为 35 维 OpenFace(17 个 AU 强度、6 维头姿、6 维双眼视线、6 维无量纲几何),置信度低于 0.8 的面部量按缺失处理。基频/有声概率使用 librosa pYIN(80–400 Hz、400 样本帧、160 样本步长),谐噪比由预测基频处的自相关周期性另行计算;pYIN 对 80 Hz 与 25 ms 窗发出周期数不足提示,因此该基频下界结果需谨慎解释。", "", "## 结果", "", "分类按连续标签的严格符号构造 Negative/Neutral/Positive,强度预测限幅到 [-3, 3]。每折训练折单独拟合标准化器、逻辑回归(C=0.05)与 Ridge(alpha=25);固定五折由 GroupKFold 按原始 `video_id` 分组。总体 OOF 指标按 100 条留组预测计算。", "", "| 方法 | OOF Accuracy | OOF Macro-F1 | OOF MAE | OOF Pearson | Macro-F1 折均值±SD |", "| --- | ---: | ---: | ---: | ---: | ---: |", ] for row in summary: lines.append( f"| {row['method']} | {row['oof_accuracy']:.3f} | {row['oof_macro_f1']:.3f} | {row['oof_mae']:.3f} | {row['oof_pearson']:.3f} | {row['fold_macro_f1_mean']:.3f} ± {row['fold_macro_f1_sd']:.3f} |" ) lines += [ "", f"按当前 OOF 探针,Macro-F1 最高的是 {best_f1['method']}({best_f1['oof_macro_f1']:.3f}),MAE 最低的是 {best_mae['method']}({best_mae['oof_mae']:.3f})。相对 B1,B4 的 Macro-F1 差为 {b4_f1['delta_oof']:+.3f},视频组 Bootstrap 95% 区间 [{b4_f1['group_bootstrap_ci95_low']:+.3f}, {b4_f1['group_bootstrap_ci95_high']:+.3f}],区间跨过零。B0 的 MAE 差为 {b0_mae['delta_oof']:+.3f},区间 [{b0_mae['group_bootstrap_ci95_low']:+.3f}, {b0_mae['group_bootstrap_ci95_high']:+.3f}],Pearson 差为 {b0_pearson['delta_oof']:+.3f}。B2 的 Macro-F1 差为 {b2_f1['delta_oof']:+.3f},区间 [{b2_f1['group_bootstrap_ci95_low']:+.3f}, {b2_f1['group_bootstrap_ci95_high']:+.3f}],MAE 差为 {b2_mae['delta_oof']:+.3f};B3 的 Macro-F1 差为 {b3_f1['delta_oof']:+.3f},区间 [{b3_f1['group_bootstrap_ci95_low']:+.3f}, {b3_f1['group_bootstrap_ci95_high']:+.3f}]。组 Bootstrap 只反映当前 37 个来源视频上的抽样不确定性,未作多重比较校正。", "", "## 覆盖与对齐核验", "", f"样本数:{len(sample_rows)};标签表连续标签与极性一致 {consistent}/100;CTC 硬对齐完整的样本 {complete_hard_count}/100。B0 硬区间平均文本覆盖率为 {mean_hard_coverage:.3f},B2 平均词占据后验质量为 {mean_posterior_occupancy:.3f}(这是期望占据率,不是覆盖率);视觉检测覆盖率(按 0.1 秒网格平均)为 {mean_vision_coverage:.3f}。后验内部 90% 边界区间平均宽度为起点 {mean_start_width:.3f} 秒、终点 {mean_end_width:.3f} 秒。没有人工边界子集,因此这些宽度只是模型内部不确定性摘要,不是经验校准率或边界误差。", "", "B0–B4 的特征来源、有效掩码、连续覆盖率、词边界后验宽度、折内预测及视频哈希均分别保存在配套 CSV/JSON 中。没有人工词边界时不报告 IoU/MATE,也不把不同方法的预测探针分数解释为对齐真值。", "", "## 文件", "", "- `comparison_summary.csv`:五种方法总体 OOF 与折均值指标。", "- `fold_metrics.csv`:每折分类/回归指标与特征维数。", "- `oof_predictions.csv`:每条样本的留组预测。", "- `group_bootstrap_deltas.csv`:相对 B1 的视频组配对 Bootstrap 95% 区间。", "- `sample_alignment_summary.csv`:100 条样本、硬/概率文本覆盖和后验边界宽度。", "- `modality_summary.csv`:300 行样本-模态明细。", "- `word_alignment_posterior.csv`:逐词硬边界、实际相对质量权重及后验边界区间。", "- `split_assignments.csv`:逐样本折号。", "- `comparison.png`:核心 OOF 指标图。", "- `typical_alignment_example.png`:中位时长样本的词边界、覆盖率、声学/视觉轨迹与视频帧核验图。", "- `run_manifest.json`:环境、模型 revision、参数、哈希与运行时间。", "", "## 复现", "", "在仓库根目录执行:", "", "```bash", "cd math", "uv sync", "uv run python compare_models.py", "```", "", "特征缓存写在 `q1/cache/native/`,OpenFace 原始 CSV 写在 `q1/cache/openface_outputs/`,正式结果写在 `output/q1/model_comparison/`。删除缓存后会从附件一重新提取。运行前需准备 BERT、Wav2Vec2 权重及 OpenFace 2.2.0 Windows 包。", ] (output_dir / "report.md").write_text("\n".join(lines) + "\n", encoding="utf-8") def _write_boundary_detail(output_dir: Path, samples: list[Sample]) -> None: rows: list[dict[str, Any]] = [] for sample in samples: quality_weights = _normalized_hard_quality(sample) for word_index, (word, interval, valid, quality, posterior) in enumerate(zip( sample.text_words, sample.hard_intervals, sample.hard_valid, sample.hard_quality, sample.boundary_summary, )): rows.append({ "sample_id": sample.sample_id, "video_id": sample.video_id, "clip_id": sample.clip_id, "word_index": word_index, "word": word, "hard_start_s": float(interval[0]) if valid else math.nan, "hard_end_s": float(interval[1]) if valid else math.nan, "hard_valid": bool(valid), "ctc_quality_score_uncalibrated": float(quality), "relative_quality_weight": float(quality_weights[word_index]) if valid else math.nan, **posterior, }) _write_csv(output_dir / "word_alignment_posterior.csv", rows, [ "sample_id", "video_id", "clip_id", "word_index", "word", "hard_start_s", "hard_end_s", "hard_valid", "ctc_quality_score_uncalibrated", "relative_quality_weight", "start_mean_s", "end_mean_s", "start_p05_s", "start_p95_s", "end_p05_s", "end_p95_s", "start_width90_s", "end_width90_s", ]) def run(args: argparse.Namespace) -> None: output_dir = Path(args.output_dir).resolve() cache_dir = Path(args.cache_dir).resolve() input_dir = Path(args.data_dir).resolve() label_file = input_dir / "label-100.xlsx" records = _read_labels(label_file) source_hashes: dict[str, str] = {} for record in records: video_path = input_dir / record["video_id"] / f"{record['clip_id']}.mp4" if not video_path.is_file(): raise FileNotFoundError(f"missing source video for {record['sample_id']}: {video_path}") source_hashes[record["sample_id"]] = _sha256(video_path) output_dir.mkdir(parents=True, exist_ok=True) cache_dir.mkdir(parents=True, exist_ok=True) cache_schema = "q1-b0b4-v4-openface" samples: list[Sample] = [] missing_records: list[dict[str, Any]] = [] cache_paths = { record["sample_id"]: cache_dir / f"{_safe_name(record['video_id'])}__{_safe_name(record['clip_id'])}.npz" for record in records } if not args.force_extract: for record in records: path = cache_paths[record["sample_id"]] if path.is_file(): try: samples.append(_load_cache(path, record, source_hashes[record["sample_id"]], cache_schema)) except (ValueError, KeyError, OSError): missing_records.append(record) else: missing_records.append(record) else: missing_records = records.copy() # Restore source-row order even when resuming a partial cache. cached_by_id = {sample.sample_id: sample for sample in samples} if missing_records: device = torch.device(args.device if args.device != "auto" else ("cuda" if torch.cuda.is_available() else "cpu")) print(f"Extracting {len(missing_records)} uncached clips on {device}.", flush=True) models = _load_models(device) for index, record in enumerate(missing_records, start=1): started = time.perf_counter() sample = _extract_native_sample( record, models, device, source_hashes[record["sample_id"]] ) _save_cache(cache_paths[record["sample_id"]], sample, cache_schema) cached_by_id[sample.sample_id] = sample print(f"[features] {index}/{len(missing_records)} {sample.sample_id}: {time.perf_counter()-started:.2f}s", flush=True) model_info = models[-1] else: manifest_path = output_dir / "run_manifest.json" prior = json.loads(manifest_path.read_text(encoding="utf-8")) if manifest_path.is_file() else {} model_info = prior.get("models", {"text_id": TEXT_MODEL_ID, "speech_id": SPEECH_MODEL_ID, "revision_status": "features resumed from cache"}) samples = [cached_by_id[record["sample_id"]] for record in records] views = [_make_views(sample) for sample in samples] label_consistent_count = sum(bool(record["label_consistent"]) for record in records) sample_rows, modality_rows = _sample_output_rows(samples, views) summary, folds, predictions, splits, deltas = _evaluate(samples, views, args.bootstrap_repeats) _write_csv(output_dir / "comparison_summary.csv", summary, [ "method", "sample_count", "video_group_count", "oof_accuracy", "oof_macro_f1", "oof_mae", "oof_pearson", "fold_accuracy_mean", "fold_accuracy_sd", "fold_macro_f1_mean", "fold_macro_f1_sd", "fold_mae_mean", "fold_mae_sd", "fold_pearson_mean", "fold_pearson_sd", ]) _write_csv(output_dir / "fold_metrics.csv", folds, [ "method", "fold", "train_samples", "valid_samples", "train_video_groups", "valid_video_groups", "accuracy", "macro_f1", "mae", "pearson", "feature_dimension", ]) _write_csv(output_dir / "oof_predictions.csv", predictions, [ "sample_id", "video_id", "clip_id", "true_polarity", "true_polarity_name", "true_sentiment", *[f"{method}_predicted_polarity" for method in ("B0", "B1", "B2", "B3", "B4")], *[f"{method}_predicted_sentiment" for method in ("B0", "B1", "B2", "B3", "B4")], ]) _write_csv(output_dir / "group_bootstrap_deltas.csv", deltas, [ "comparison", "metric", "reference_oof", "method_oof", "delta_oof", "group_bootstrap_ci95_low", "group_bootstrap_ci95_high", "bootstrap_repeats", ]) _write_csv(output_dir / "split_assignments.csv", sorted(splits, key=lambda row: row["sample_id"]), ["sample_id", "video_id", "fold", "split"]) _write_csv(output_dir / "sample_alignment_summary.csv", sample_rows, list(sample_rows[0])) _write_csv(output_dir / "modality_summary.csv", modality_rows, list(modality_rows[0])) _write_boundary_detail(output_dir, samples) _create_plot(summary, output_dir / "comparison.png") typical_sample_id = _create_typical_alignment_figure(samples, views, output_dir / "typical_alignment_example.png") manifest = { "created_at_local": time.strftime("%Y-%m-%d %H:%M:%S %z"), "python": sys.version, "platform": platform.platform(), "uv_version": subprocess.run([shutil.which("uv"), "--version"], capture_output=True, text=True).stdout.strip() if shutil.which("uv") else "uv run parent; executable not on child PATH", "packages": {name: _version(name) for name in ("torch", "transformers", "numpy", "scipy", "scikit-learn", "opencv-python-headless", "librosa", "openpyxl", "matplotlib")}, "device": args.device if args.device != "auto" else ("cuda" if torch.cuda.is_available() else "cpu"), "cuda_available": bool(torch.cuda.is_available()), "gpu_name": torch.cuda.get_device_name(0) if torch.cuda.is_available() and args.device != "cpu" else None, "models": model_info, "openface": {"version": "2.2.0", "executable": "FeatureExtraction", "landmark_model": "main_clnf_general.txt", "confidence_threshold": 0.8}, "inputs": { "label_file": str(label_file), "label_file_sha256": _sha256(label_file), "sample_count": len(samples), "video_sha256_by_sample": source_hashes, "label_polarity_consistency_count": label_consistent_count, }, "experiment": { "methods": { "B0": "quality-weighted hard CTC word projection + physical source-time pooling + train-fold standard scaling", "B1": "B0 + train-fold median/MAD scaling + SO(3) rotation mean + normalized gaze vector mean", "B2": "B1 + fixed-transcript CTC forward-backward word occupancy projection for text", "B3": "B1 + second-order time-augmented log signatures of selected audio/face channels", "B4": "B1 + frozen Wav2Vec2 last-four-layer speech representation", }, "text_dimension": 768, "audio_dimension": 74, "vision_dimension": 35, "audio_features": list(AUDIO_NAMES), "vision_features": list(VISION_NAMES), "vision_feature_note": "OpenFace 2.2.0 AU intensities, pose, gaze, and six geometry values defined in E题V2.pdf section 4.2.3", "grid_step_s": GRID_STEP_S, "audio_window_samples": FRAME_LENGTH, "audio_step_samples": FRAME_STEP, "audio_fft": N_FFT, "audio_mel_bands": MEL_COUNT, "pitch_estimator": {"name": "librosa.pyin", "version": _version("librosa"), "f0_range_hz": [80.0, 400.0], "frame_length_samples": FRAME_LENGTH, "hop_length_samples": FRAME_STEP, "center": False}, "hnr_estimator": "autocorrelation periodicity at pYIN-detected F0; dB", "vision_sample_rate_hz": "native OpenFace frame timestamps; no fixed-rate assumption", "grouping": "GroupKFold by video_id", "fold_count": N_FOLDS, "probe": {"classifier": "LogisticRegression", "C": 0.05, "regressor": "Ridge", "alpha": 25.0, "intensity_clip": [-3.0, 3.0]}, "bootstrap": {"unit": "video_id", "repeats": args.bootstrap_repeats, "seed": SEED}, "human_boundary_reference_count": 0, "boundary_interval_level": 0.90, "hard_text_quality_weight": "raw fixed-transcript CTC word path score in [0,1]; unavailable scores use q*=1 and are flagged unavailable", "cache_schema": cache_schema, }, "output_dir": str(output_dir), "typical_sample_id": typical_sample_id, "elapsed_seconds": time.time() - args.started_at, } (output_dir / "run_manifest.json").write_text(json.dumps(manifest, ensure_ascii=False, indent=2) + "\n", encoding="utf-8") _write_report(output_dir, summary, deltas, sample_rows, modality_rows, manifest) print(f"Saved B0-B4 comparison to {output_dir}", flush=True) def _parse_args() -> argparse.Namespace: parser = argparse.ArgumentParser(description="Run Q1 B0-B4 model comparison on Attachment 1.") parser.add_argument("--data-dir", type=Path, default=DATA_DIR) parser.add_argument("--output-dir", type=Path, default=OUTPUT_DIR) parser.add_argument("--cache-dir", type=Path, default=CACHE_DIR) parser.add_argument("--device", choices=("auto", "cuda", "cpu"), default="auto") parser.add_argument("--bootstrap-repeats", type=int, default=2000) parser.add_argument("--force-extract", action="store_true", help="recompute native features from the 100 raw videos") args = parser.parse_args() if args.bootstrap_repeats < 100: parser.error("--bootstrap-repeats must be at least 100") args.started_at = time.time() return args if __name__ == "__main__": run(_parse_args())