405 lines
30 KiB
Python
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()
|