Files

405 lines
30 KiB
Python

from __future__ import annotations
import csv
import hashlib
import json
import subprocess
import sys
from dataclasses import replace
from pathlib import Path
from typing import Any
import numpy as np
from scipy import sparse
ROOT = Path(__file__).resolve().parents[2]
Q1 = Path(__file__).resolve().parent
sys.path.insert(0, str(Q1))
import compare_models as cm # noqa: E402
CACHE = Q1 / "cache" / "native"
RUN = Q1 / "source_cache_manifest.json"
LABELS = cm.LABEL_FILE
OUT = Q1 / "features_v2"
MODALITIES = ("text", "audio", "vision")
CONFIG = {
"schema": "q1-v2-source-first-2026-09",
"source_cache_schema": "q1-b0b4-v4-openface",
"common_step_s": 0.1,
"multi_context_windows_s": [0.1, 0.3, 0.7],
"relative_progress_bins": 50,
"text_encoder": "google-bert/bert-base-uncased; frozen last-four-layer mean per whitespace word; window=510, stride=384",
"alignment": "fixed-transcript CTC monotone state graph; hard Viterbi path plus stored forward-backward occupancy; scores uncalibrated",
"audio": "74-D: 40 log-Mel + 13 MFCC + 13 two-sided local-linear delta MFCC + [log-energy, log-F0, librosa.pyin voiced probability, spectral centroid, bandwidth, flux, zero-crossing rate, HNR dB]; 16 kHz, 400-sample window, 160-sample hop, FFT 512; HNR autocorrelation uses n=1024 at detected F0; preserve per-field validity; librosa.pyin warns cycle support is short at fmin=80 Hz for a 25 ms frame",
"vision": "OpenFace 2.2.0 native-frame output: 17 AU intensities, Pose6, Gaze6, Geometry6; confidence threshold 0.8; actual frame PTS and indices retained",
"quality": "B0 uses raw CTC word-path score in [0,1]; unavailable quality falls back to q*=1 and is flagged; audio/vision retain explicit unknown or detector-confidence states",
"content_probe": "disabled; physical time/query maps remain separate from content similarity",
"time": "left-closed/right-open seconds relative to video stream origin; source indices retained",
"branch_states": {"sec": "primary B0 view and materialized", "B1_vision_geometry": "geometry-aware 0.1 s view materialized", "word": "materialized", "phase50": "materialized", "multi": "derived from native rows on demand by q1_io", "posterior": "native CTC posterior stored; sec projection available through q1_io", "dynamics": "materialized if continuous support is sufficient", "b4_speech": "source-cache-backed on demand through q1_io to stay within the 50 MiB feature-package cap", "query_word": "physical H_time matrices materialized", "content_probe": "disabled"},
}
def digest(path: Path) -> str:
h = hashlib.sha256()
with path.open("rb") as f:
for block in iter(lambda: f.read(1024 * 1024), b""):
h.update(block)
return h.hexdigest()
def finite_boundary_summary(rows: list[dict[str, float]], offset: float, end_s: float) -> tuple[np.ndarray, str]:
values=[]
json_rows=[]
for row in rows:
vals=[row.get("start_p05_s",np.nan)+offset,row.get("start_p95_s",np.nan)+offset,row.get("end_p05_s",np.nan)+offset,row.get("end_p95_s",np.nan)+offset]
vals=[float(np.clip(v,0.0,end_s)) if np.isfinite(v) else np.nan for v in vals]
values.append([v if np.isfinite(v) else -1.0 for v in vals])
absolute_fields={"start_mean_s","end_mean_s","start_p05_s","start_p95_s","end_p05_s","end_p95_s"}
json_rows.append({k:(float(np.clip(v+offset,0.0,end_s)) if np.isfinite(v) and k in absolute_fields else (float(v) if np.isfinite(v) else None)) for k,v in row.items()})
return np.asarray(values,np.float32),json.dumps(json_rows,ensure_ascii=False)
def audio_sample_count_from_quality(sample: cm.Sample) -> int:
"""The final partial 25 ms frame records its exact decoded sample fraction."""
valid=np.flatnonzero(np.asarray(sample.audio_quality)>0)
if not len(valid):
return 0
index=int(valid[-1])
fraction=float(np.clip(sample.audio_quality[index],0.0,1.0))
return index*cm.FRAME_STEP+int(round(fraction*cm.FRAME_LENGTH))
def effective_text_quality(raw: np.ndarray, valid: np.ndarray) -> tuple[np.ndarray, np.ndarray]:
"""Use bounded raw CTC scores; q*=1 means no extra discount when unavailable."""
raw = np.asarray(raw, np.float32)
available = np.asarray(valid, bool) & (raw > 0) & np.isfinite(raw)
effective = np.ones(len(raw), np.float32)
effective[available] = np.clip(raw[available], 0.0, 1.0)
return effective, available
def intervals_from_times(times: np.ndarray, duration: float, step: float, gap_factor: float = 1.5) -> np.ndarray:
times = np.asarray(times, dtype=np.float64)
if not len(times):
return np.empty((0, 2), np.float32)
starts = np.empty(len(times), dtype=np.float64)
ends = np.empty(len(times), dtype=np.float64)
breaks = np.diff(times) > gap_factor * step
for i, t in enumerate(times):
starts[i] = (times[i - 1] + t) / 2 if i and not breaks[i - 1] else max(0.0, t - step / 2)
ends[i] = (t + times[i + 1]) / 2 if i + 1 < len(times) and not breaks[i] else min(duration, t + step / 2)
if starts[0] <= step:
starts[0] = 0.0
if ends[-1] >= duration - step:
ends[-1] = duration
starts = np.clip(starts, 0.0, duration)
ends = np.clip(ends, 0.0, duration)
return np.column_stack((starts, np.maximum(starts, ends))).astype(np.float32)
def probe_media(path: Path) -> dict[str, Any]:
command = ["ffprobe", "-v", "error", "-show_entries", "format=start_time,duration:stream=codec_type,start_time,time_base,sample_rate,avg_frame_rate", "-of", "json", str(path)]
try:
obj = json.loads(subprocess.run(command, check=True, capture_output=True, text=True).stdout)
streams = obj.get("streams", [])
fmt = obj.get("format", {})
vs = next((s for s in streams if s.get("codec_type") == "video"), {})
aus = next((s for s in streams if s.get("codec_type") == "audio"), {})
v0 = float(vs.get("start_time", fmt.get("start_time", 0)) or 0)
a0 = float(aus.get("start_time", fmt.get("start_time", 0)) or 0)
return {"format_start_s": float(fmt.get("start_time", 0) or 0), "format_duration_s": float(fmt.get("duration", 0) or 0), "video_start_s": v0, "audio_start_s": a0, "audio_offset_from_video_s": a0 - v0, "video_time_base": vs.get("time_base"), "audio_time_base": aus.get("time_base"), "audio_sample_rate": aus.get("sample_rate"), "avg_frame_rate": vs.get("avg_frame_rate"), "status": "ok"}
except Exception as exc:
return {"format_start_s": 0.0, "format_duration_s": 0.0, "video_start_s": 0.0, "audio_start_s": 0.0, "audio_offset_from_video_s": 0.0, "status": f"ffprobe_error:{type(exc).__name__}"}
def aggregate(values: np.ndarray, observed: np.ndarray, source_intervals: np.ndarray, quality: np.ndarray, edges: np.ndarray, quality_available: np.ndarray | None = None) -> dict[str, np.ndarray]:
values = np.asarray(values, np.float64)
observed = np.asarray(observed, bool)
source_intervals = np.asarray(source_intervals, np.float64)
quality = np.maximum(0.0, np.asarray(quality, np.float64))
quality_available = np.zeros(len(values), dtype=bool) if quality_available is None else np.asarray(quality_available, dtype=bool)
n, d = len(edges) - 1, values.shape[1]
mean = np.zeros((n, d), np.float32)
var = np.zeros((n, d), np.float32)
count = np.zeros((n, d), np.uint8)
coverage = np.zeros((n, d), np.float32)
qbar = np.ones((n, d), np.float32)
qavail = np.zeros((n, d), np.float32)
mask = np.zeros((n, d), bool)
for i, (left, right) in enumerate(zip(edges[:-1], edges[1:])):
if right <= left or not len(source_intervals):
continue
overlap = np.maximum(0.0, np.minimum(source_intervals[:, 1], right) - np.maximum(source_intervals[:, 0], left))
phys = overlap[:, None] * observed
w = phys * quality[:, None]
wsum = w.sum(axis=0)
good = wsum > 0
if good.any():
sx = (values * w).sum(axis=0)
sx2 = (np.square(values) * w).sum(axis=0)
mean[i, good] = (sx[good] / wsum[good]).astype(np.float32)
var[i, good] = np.maximum(0.0, sx2[good] / wsum[good] - mean[i, good].astype(np.float64) ** 2).astype(np.float32)
mask[i, good] = True
coverage[i] = np.minimum(1.0, phys.sum(axis=0) / (right - left)).astype(np.float32)
count[i] = np.minimum(255, ((overlap[:, None] > 0) & observed).sum(axis=0)).astype(np.uint8)
base = phys.sum(axis=0)
valid = base > 0
qbar[i, valid] = ((phys * quality[:, None]).sum(axis=0)[valid] / base[valid]).astype(np.float32)
qavail[i, valid] = ((phys * quality_available[:, None]).sum(axis=0)[valid] / base[valid]).astype(np.float32)
# Audio HNR/prosodic rows can have genuinely large within-window variance;
# retain variance in float32 rather than silently overflowing float16.
return {"x": mean.astype(np.float16), "var": var.astype(np.float32), "count": count, "coverage_u8": np.rint(coverage * 255).astype(np.uint8), "mask": mask, "qbar": qbar.astype(np.float16), "quality_available_fraction": qavail.astype(np.float16)}
def sparse_overlap(source_intervals: np.ndarray, target_intervals: np.ndarray, quality: np.ndarray | None = None) -> sparse.csr_matrix:
rows, cols, vals = [], [], []
q = np.ones(len(source_intervals), np.float32) if quality is None else np.asarray(quality, np.float32)
for i, (left, right) in enumerate(target_intervals):
overlap = np.maximum(0.0, np.minimum(source_intervals[:, 1], right) - np.maximum(source_intervals[:, 0], left)) * q
js = np.flatnonzero(overlap > 0)
rows.extend([i] * len(js)); cols.extend(js.tolist()); vals.extend(overlap[js].tolist())
return sparse.csr_matrix((np.asarray(vals, np.float32), (rows, cols)), shape=(len(target_intervals), len(source_intervals)))
def sparse_word_map(word_intervals: np.ndarray, source_intervals: np.ndarray, chi: np.ndarray) -> sparse.csr_matrix:
rows, cols, vals = [], [], []
for k, (left, right) in enumerate(word_intervals):
if right <= left:
continue
overlap = np.maximum(0.0, np.minimum(source_intervals[:, 1], right) - np.maximum(source_intervals[:, 0], left)) * chi
js = np.flatnonzero(overlap > 0)
total = float(overlap[js].sum())
if total > 0:
rows.extend([k] * len(js)); cols.extend(js.tolist()); vals.extend((overlap[js] / total).tolist())
return sparse.csr_matrix((np.asarray(vals, np.float32), (rows, cols)), shape=(len(word_intervals), len(source_intervals)))
def pack_csr(arrays: dict[str, Any], prefix: str, matrix: sparse.csr_matrix) -> None:
arrays[prefix + "_data"] = matrix.data.astype(np.float16)
arrays[prefix + "_indices"] = matrix.indices.astype(np.int32)
arrays[prefix + "_indptr"] = matrix.indptr.astype(np.int32)
arrays[prefix + "_shape"] = np.asarray(matrix.shape, np.int32)
def text_provenance(record: dict[str, Any], words: list[str], tokenizer: Any) -> tuple[np.ndarray, np.ndarray, np.ndarray]:
text = record["text"]
spans = np.full((len(words), 2), -1, np.int32)
subword = np.zeros((len(words), 2), np.int32)
all_ids: list[int] = []
cursor = 0
for i, word in enumerate(words):
start = text.find(word, cursor)
if start >= 0:
spans[i] = (start, start + len(word)); cursor = start + len(word)
ids = tokenizer(word, add_special_tokens=False)["input_ids"] if tokenizer is not None else []
subword[i] = (len(all_ids), len(all_ids) + len(ids))
all_ids.extend(int(v) for v in ids)
return spans, subword, np.asarray(all_ids, np.int32)
def make_dynamics(sample: cm.Sample, audio_intervals: np.ndarray, vision_intervals: np.ndarray, edges: np.ndarray) -> tuple[np.ndarray, np.ndarray, np.ndarray]:
# Time-augmented level-2 log-signature on continuous 0.7 s contexts; this is auxiliary only.
names = [("audio", 67, 66), ("vision", 10, 31)]
output = np.zeros((len(edges)-1, 2, 6), np.float32)
valid = np.zeros((len(edges)-1, 2), bool)
bounds = np.zeros((len(edges)-1, 2, 2), np.float32)
for t,(l,r) in enumerate(zip(edges[:-1],edges[1:])):
center=(l+r)/2
a=max(0.0,center-0.35); b=min(sample.duration_s,center+0.35)
bounds[t,:,:]=[a,b]
for branch,(mod,j,k) in enumerate(names):
if mod=="audio":
times=sample.audio_times; x=sample.audio_features.astype(np.float32); mask=sample.audio_observed
ints=audio_intervals; required=mask[:,[j,k]].all(axis=1)
else:
times=sample.vision_times; x=sample.vision_features.astype(np.float32); mask=sample.vision_observed
ints=vision_intervals; required=mask[:,[j,k]].all(axis=1)
ix=np.flatnonzero((ints[:,0]>=a-1e-7)&(ints[:,1]<=b+1e-7)&required)
if len(ix)<3: continue
# Break the path at missing or unusually long source gaps.
expected=0.01 if mod=="audio" else max(float(np.median(np.diff(times))),1e-6)
if np.any(np.diff(times[ix])>1.5*expected): continue
vals=np.column_stack(((times[ix]-a)/max(b-a,1e-6),x[ix,j],x[ix,k])).astype(np.float64)
dz=np.diff(vals,axis=0)
first=dz.sum(axis=0)
area=np.zeros(3,np.float64); pairs=((0,1),(0,2),(1,2))
for p,(u,v) in enumerate(pairs):
area[p]=0.5*sum(dz[a0,u]*dz[b0,v]-dz[a0,v]*dz[b0,u] for a0 in range(len(dz)) for b0 in range(a0+1,len(dz)))
output[t,branch]=np.concatenate((first,area)).astype(np.float32)
valid[t,branch]=True
return output.astype(np.float16),valid,bounds
def geometry_aware_vision(sample: cm.Sample, intervals: np.ndarray, bounds: np.ndarray, base: dict[str, np.ndarray]) -> tuple[np.ndarray, np.ndarray]:
"""B1 SO(3) mean for head pose and unit-vector means for gaze channels."""
values = base["x"].astype(np.float32).copy()
masks = base["mask"].astype(bool).copy()
raw = sample.vision_features.astype(np.float32)
observed = sample.vision_observed
for target_index, (left, right) in enumerate(bounds):
overlap = np.maximum(0.0, np.minimum(intervals[:, 1], right) - np.maximum(intervals[:, 0], left))
weighted_overlap=overlap*sample.vision_quality
pose_rows = (weighted_overlap > 0) & observed[:, 17:20].all(axis=1)
if pose_rows.any():
mean_rot = cm._so3_weighted_mean(raw[pose_rows, 17:20], weighted_overlap[pose_rows])
if mean_rot is not None:
values[target_index, 17:20] = mean_rot
masks[target_index, 17:20] = True
for start in (23, 26):
rows = (weighted_overlap > 0) & observed[:, start:start + 3].all(axis=1)
if not rows.any():
continue
mean = np.average(raw[rows, start:start + 3], axis=0, weights=weighted_overlap[rows])
norm = float(np.linalg.norm(mean))
if norm > 1e-6:
values[target_index, start:start + 3] = mean / 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.astype(np.float16), masks
def main() -> None:
if not CACHE.is_dir() or not RUN.is_file():
raise FileNotFoundError("Existing Q1 native cache and run manifest are required")
OUT.mkdir(parents=True, exist_ok=True)
run = json.loads(RUN.read_text(encoding="utf-8"))
records = cm._read_labels(LABELS)
if len(records) != 100:
raise ValueError(f"expected 100 attachment-1 rows, got {len(records)}")
tokenizer = None
try:
tokenizer = cm.AutoTokenizer.from_pretrained(cm.TEXT_MODEL_ID, use_fast=True, local_files_only=True)
except Exception:
pass
config_text = json.dumps(CONFIG, ensure_ascii=False, sort_keys=True, separators=(",", ":"))
config_hash = hashlib.sha256(config_text.encode("utf-8")).hexdigest()
grid_rows: list[dict[str, Any]] = []
sample_rows: list[dict[str, Any]] = []
manifest_lines: list[dict[str, Any]] = []
for record in records:
sid=record["sample_id"]
source_hash=run["inputs"]["video_sha256_by_sample"][sid]
cache=CACHE/f"{cm._safe_name(record['video_id'])}__{cm._safe_name(record['clip_id'])}.npz"
sample=cm._load_cache(cache,record,source_hash,"q1-b0b4-v4-openface")
video_path=cm.DATA_DIR/record["video_id"]/f"{record['clip_id']}.mp4"
media=probe_media(video_path)
audio_offset=float(media.get("audio_offset_from_video_s",0.0))
sample_count=audio_sample_count_from_quality(sample)
audio_end_s=float(np.clip(audio_offset+sample_count/cm.SAMPLE_RATE,0.0,sample.duration_s))
audio_times=sample.audio_times.astype(np.float32)+audio_offset
ctc_times=sample.ctc_times.astype(np.float32)+audio_offset
hard=sample.hard_intervals.astype(np.float32).copy()+np.float32(audio_offset)
hard=np.clip(hard,0.0,audio_end_s)
hard_valid=sample.hard_valid & (hard[:,1]>hard[:,0])
text_intervals=np.where(hard_valid[:,None],hard,0.0)
audio_starts=np.arange(len(audio_times),dtype=np.int64)*cm.FRAME_STEP
audio_ranges=np.column_stack((audio_starts,np.minimum(audio_starts+cm.FRAME_LENGTH,sample_count))).astype(np.int32)
audio_intervals=intervals_from_times(audio_times,audio_end_s,cm.FRAME_STEP/cm.SAMPLE_RATE)
audio_intervals=np.clip(audio_intervals,0.0,sample.duration_s)
video_step=float(np.median(np.diff(sample.vision_times))) if len(sample.vision_times)>1 else 1.0/30.0
vision_intervals=intervals_from_times(sample.vision_times,sample.duration_s,video_step)
frame_indices=sample.vision_frame_indices.astype(np.int32)
words=sample.text_words
char_spans,subword_spans,subword_ids=text_provenance(record,words,tokenizer)
duration=sample.duration_s
sec_edges=cm._grid_edges(duration)
sec_bounds=np.column_stack((sec_edges[:-1],sec_edges[1:])).astype(np.float32)
phase_edges=np.linspace(0.0,duration,51,dtype=np.float64)
phase_bounds=np.column_stack((phase_edges[:-1],phase_edges[1:])).astype(np.float32)
word_bounds=text_intervals.copy()
# Keep the media timestamps in the video-origin timebase for dynamics too.
sample.audio_times = audio_times
text_quality, text_quality_available = effective_text_quality(sample.hard_quality, hard_valid)
source={
"text":(sample.text_features.astype(np.float32),np.broadcast_to(sample.text_valid[:,None],sample.text_features.shape).copy(),text_intervals,text_quality,text_quality_available),
"audio":(sample.audio_features.astype(np.float32),sample.audio_observed,audio_intervals,sample.audio_quality,sample.audio_quality_available),
"vision":(sample.vision_features.astype(np.float32),sample.vision_observed,vision_intervals,sample.vision_quality,sample.vision_quality_available),
}
sec={name:aggregate(source[name][0],source[name][1],source[name][2],source[name][3],sec_edges,source[name][4]) for name in MODALITIES}
phase={name:aggregate(source[name][0],source[name][1],source[name][2],source[name][3],phase_edges,source[name][4]) for name in MODALITIES}
vision_b1_x,vision_b1_mask=geometry_aware_vision(sample,vision_intervals,sec_bounds,sec["vision"])
# Word-level physical correspondence and query matrices.
word_views={}
for name in MODALITIES:
result={k:[] for k in ("x","var","count","coverage_u8","mask","qbar","quality_available_fraction")}
for aa,bb in word_bounds:
v=aggregate(source[name][0],source[name][1],source[name][2],source[name][3],np.asarray([aa,bb],np.float64),source[name][4])
for k in result: result[k].append(v[k][0])
word_views[name]={k:np.asarray(v) for k,v in result.items()}
# Sparse target-to-native overlap maps for the primary physical view.
maps={name:sparse_overlap(source[name][2],sec_bounds) for name in MODALITIES}
# Query channels are fixed for display only: acoustic log-energy and OpenFace AU12.
h_audio=sparse_word_map(word_bounds,audio_intervals,sample.audio_observed[:,66].astype(np.float32))
h_vision=sparse_word_map(word_bounds,vision_intervals,sample.vision_observed[:,10].astype(np.float32))
dyn,dyn_valid,dyn_bounds=make_dynamics(sample,audio_intervals,vision_intervals,sec_edges)
delta_centers=np.arange(len(audio_times),dtype=np.int64)
delta_left=np.maximum(0,delta_centers-2)*cm.FRAME_STEP
delta_right=np.minimum(sample_count,np.minimum(len(audio_times)-1,delta_centers+2)*cm.FRAME_STEP+cm.FRAME_LENGTH)
delta_ranges=np.column_stack((delta_left,delta_right)).astype(np.int32)
boundary_p05_p95,boundary_summary_json=finite_boundary_summary(sample.boundary_summary,audio_offset,audio_end_s)
cache_file=OUT/f"{cm._safe_name(record['video_id'])}__{cm._safe_name(record['clip_id'])}.npz"
payload:dict[str,Any]={
"meta_json":np.asarray(json.dumps({"sample_id":sid,"video_id":record["video_id"],"clip_id":record["clip_id"],"duration_s":duration,"sentiment":sample.sentiment,"polarity":sample.polarity,"annotation":record["annotation"],"source_video_sha256":source_hash,"media":media,"config_hash":config_hash,"vision_backend":"OpenFace 2.2.0 FeatureExtraction","quality_available":{"text":bool(text_quality_available.any()),"audio":bool(sample.audio_quality_available.any()),"vision":bool(sample.vision_quality_available.any())}},ensure_ascii=False)),
"native_text_words":np.asarray(words,dtype=str),"native_text_features":sample.text_features.astype(np.float16),"native_text_observed":sample.text_valid.astype(bool),"native_text_intervals":text_intervals.astype(np.float32),"native_text_hard_quality_uncalibrated":sample.hard_quality.astype(np.float16),"native_text_quality_effective":text_quality.astype(np.float16),"native_text_quality_available":text_quality_available,"native_text_boundary_p05_p95":boundary_p05_p95,"native_text_char_spans":char_spans,"native_text_subword_spans":subword_spans,"native_text_subword_ids":subword_ids,
"native_audio_times":audio_times,"native_audio_features":sample.audio_features.astype(np.float16),"native_audio_mask":sample.audio_observed.astype(bool),"native_audio_quality":sample.audio_quality.astype(np.float16),"native_audio_quality_available":sample.audio_quality_available.astype(bool),"native_audio_source_range_samples":audio_ranges,"native_audio_delta_calc_range_samples":delta_ranges,"native_audio_intervals":audio_intervals,
"native_vision_times":sample.vision_times.astype(np.float32),"native_vision_features":sample.vision_features.astype(np.float16),"native_vision_mask":sample.vision_observed.astype(bool),"native_vision_quality":sample.vision_quality.astype(np.float16),"native_vision_quality_available":sample.vision_quality_available.astype(bool),"native_vision_frame_index":frame_indices,"native_vision_intervals":vision_intervals,
"native_ctc_times":ctc_times,"native_ctc_occupancy":sample.ctc_occupancy.astype(np.float16),"native_ctc_boundary_summary_json":np.asarray(boundary_summary_json),
"views_sec_time_bounds_s":sec_bounds,"views_phase50_time_bounds_s":phase_bounds,"views_word_time_bounds_s":word_bounds,
"views_sec_vision_b1_x":vision_b1_x,"views_sec_vision_b1_mask":vision_b1_mask,
"views_dynamics_x":dyn,"views_dynamics_mask":dyn_valid,"views_dynamics_time_bounds_s":dyn_bounds,
}
for name in MODALITIES:
for k,v in sec[name].items(): payload[f"views_sec_{name}_{k}"]=v
for k,v in phase[name].items(): payload[f"views_phase50_{name}_{k}"]=v
for k,v in word_views[name].items(): payload[f"views_word_{name}_{k}"]=v
pack_csr(payload,f"alignment_map_{name}",maps[name])
pack_csr(payload,"query_word_audio_H_time",h_audio)
pack_csr(payload,"query_word_vision_H_time",h_vision)
np.savez_compressed(cache_file,**payload)
feature_hash=digest(cache_file)
observed_summary={}
for name in MODALITIES:
vals,mask,ints,q,q_available=source[name]
support=np.zeros(mask.shape[1],np.float64)
for r,(a,b) in enumerate(ints):
support += max(0.0,float(b-a))*mask[r]*(q[r]>0)
coverage=np.clip(support/max(duration,1e-9),0,1)
observed_summary[name]={"observed_duration_mean_s":float(np.mean(coverage)*duration),"mean_coverage":float(np.mean(coverage)),"row_count":int(len(vals)),"dimension":int(vals.shape[1]),"quality_available":bool(q_available.any()),"quality_available_fraction":float(q_available.mean()) if len(q_available) else 0.0}
gran={"text":"word-level, variable duration; fixed-transcript CTC source interval","audio":"native 10 ms feature step; 25 ms windows; context ranges retained separately","vision":"native OpenFace frame timestamps and frame indices"}[name]
unobserved_rows=int((~mask.any(axis=1)).sum()) if len(mask) else 0
if name=="text":
missing_reason=f"{int((~hard_valid).sum())} word intervals are not reliable under the fixed-transcript CTC path" if (~hard_valid).any() else "none"
elif name=="audio":
undefined_pitch=int((~mask[:,67]).sum()) if mask.shape[1]>67 else 0
missing_reason=f"F0/HNR naturally undefined or unavailable in {undefined_pitch} rows; support masks retained"
else:
missing_reason=f"OpenFace confidence/success gate rejected {unobserved_rows} native frames" if unobserved_rows else "none"
grid_rows.append({"sample_id":sid,"video_id":record["video_id"],"clip_id":record["clip_id"],"modality":name,"source_duration_s":duration,"observed_duration_s":observed_summary[name]["observed_duration_mean_s"],"native_length":len(vals),"native_dim":vals.shape[1],"aligned_length":len(sec_bounds),"aligned_dim":vals.shape[1],"alignment_type":"fixed 0.1 s physical grid; direct native interval projection","granularity":gran,"mean_coverage":observed_summary[name]["mean_coverage"],"missing_reason":missing_reason,"quality_available":observed_summary[name]["quality_available"],"quality_available_fraction":observed_summary[name]["quality_available_fraction"],"quality_path":f"{cache_file.name}::native_{name}_*quality*","query_map_path":f"{cache_file.name}::query_word_{name}_H_time" if name in ("audio","vision") else "not_applicable","probe_status":"content probe disabled" if name in ("audio","vision") else "not_applicable","feature_path":str(cache_file.relative_to(Q1)),"source_hash":source_hash,"status":"success" if mask.any() else "failed_or_unavailable","config_hash":config_hash})
sample_status="success" if all(source[m][1].any() for m in MODALITIES) else "partial"
sample_rows.append({"sample_id":sid,"video_id":record["video_id"],"clip_id":record["clip_id"],"source_duration_s":duration,"source_video_sha256":source_hash,"feature_path":str(cache_file.relative_to(Q1)),"feature_sha256":feature_hash,"sec_length":len(sec_bounds),"word_count":len(words),"hard_aligned_word_count":int(hard_valid.sum()),"unlocated_word_count":int((~hard_valid).sum()),"status":sample_status,"config_hash":config_hash})
manifest_lines.append({"sample_id":sid,"video_id":record["video_id"],"clip_id":record["clip_id"],"feature_path":str(cache_file.relative_to(Q1)),"source_video_sha256":source_hash,"feature_sha256":feature_hash,"duration_s":duration,"status":sample_rows[-1]["status"],"config_hash":config_hash})
print(f"{sid}: {len(words)} words, {len(sec_bounds)} sec steps, {cache_file.stat().st_size/1024:.1f} KiB")
def write_csv(path:Path,rows:list[dict[str,Any]])->None:
with path.open("w",encoding="utf-8-sig",newline="") as f:
w=csv.DictWriter(f,fieldnames=list(rows[0]));w.writeheader();w.writerows(rows)
write_csv(OUT/"sample_summary.csv",sample_rows)
write_csv(OUT/"modality_summary.csv",grid_rows)
(OUT/"manifest_q1.jsonl").write_text("\n".join(json.dumps(x,ensure_ascii=False) for x in manifest_lines)+"\n",encoding="utf-8")
configuration={"config":CONFIG,"config_hash":config_hash,"source_run_manifest":RUN.relative_to(Q1).as_posix(),"source_cache_schema":"q1-b0b4-v4-openface","source_label_sha256":digest(LABELS),"python":run.get("python"),"packages":run.get("packages",{}),"models":run.get("models",{}),"samples":len(records),"modalities":300,"total_bytes":sum(p.stat().st_size for p in OUT.iterdir() if p.is_file())}
(OUT/"feature_manifest.json").write_text(json.dumps(configuration,ensure_ascii=False,indent=2),encoding="utf-8")
median_sample=sorted(sample_rows,key=lambda x:x["source_duration_s"])[len(sample_rows)//2]
with np.load(Q1/median_sample["feature_path"],allow_pickle=False) as d:
k=int(min(5,len(d["native_text_words"])-1)); sid=median_sample["sample_id"]
card={"sample_id":sid,"source_video_sha256":median_sample["source_video_sha256"],"duration_s":median_sample["source_duration_s"],"word":str(d["native_text_words"][k]),"word_char_span":d["native_text_char_spans"][k].tolist(),"hard_interval_s":d["native_text_intervals"][k].tolist(),"boundary_interval_90_s":d["native_text_boundary_p05_p95"][k].tolist(),"audio_source_rows":sparse.csr_matrix((d["query_word_audio_H_time_data"].astype(np.float32),d["query_word_audio_H_time_indices"],d["query_word_audio_H_time_indptr"]),shape=tuple(d["query_word_audio_H_time_shape"])).getrow(k).indices.tolist(),"vision_source_rows":sparse.csr_matrix((d["query_word_vision_H_time_data"].astype(np.float32),d["query_word_vision_H_time_indices"],d["query_word_vision_H_time_indptr"]),shape=tuple(d["query_word_vision_H_time_shape"])).getrow(k).indices.tolist(),"manual_boundary_audit":"not available; no empirical boundary error or calibration claimed","vision_backend":"OpenFace 2.2.0"}
(OUT/"typical_alignment_card.json").write_text(json.dumps(card,ensure_ascii=False,indent=2),encoding="utf-8")
configuration["total_bytes"]=sum(p.stat().st_size for p in OUT.rglob("*") if p.is_file())
(OUT/"feature_manifest.json").write_text(json.dumps(configuration,ensure_ascii=False,indent=2),encoding="utf-8")
print(f"Wrote {len(sample_rows)} samples and {len(grid_rows)} modality rows; total {configuration['total_bytes']/1024**2:.2f} MiB")
if __name__=="__main__":
main()