diff --git a/.gitignore b/.gitignore index b4aac33..9530993 100644 --- a/.gitignore +++ b/.gitignore @@ -3,6 +3,7 @@ __pycache__/ *.egg-info/ dist/ build/ +.venv .venv/ .env *.db diff --git a/clip-eval/README.md b/clip-eval/README.md new file mode 100644 index 0000000..7e4a8eb --- /dev/null +++ b/clip-eval/README.md @@ -0,0 +1,93 @@ +# Clip evaluation contract (v1) + +This directory is the initial public shared evaluation contract for AV topic +clipping (`av clip`). It exists so that deterministic retrieval/windows and +Jev-decided selection can be compared on **identical candidates** with +**independent labels**, and so future lanes can extend the corpus without +redefining the metrics. + +## What is frozen here + +| File | Role | +|---|---| +| `corpus.json` | Synthetic corpus: 3 content styles, inline transcript/vision segments, stable artifact IDs, per-video SHA-256 checksums, and the deterministic ffmpeg media-generation spec | +| `queries.json` | Labeled dev query set: present topics, absent topics, repeated topics, noisy ASR, missing vision, contradictory vision, context-dependent soundbites | +| `queries-heldout.json` | Held-out labels. Run at most once per evidence cycle; never against prompt or threshold iteration | +| `receipts/` | Committed receipts from offline runs (small JSON, no media) | + +`contract_version` is `1`. The loader (`av.clip_eval.corpus.load_corpus`) +verifies every checksum and rejects any modified corpus. The metric +definitions live in `av.clip_eval.contract` and are quoted verbatim in this +file; changing either requires bumping the contract version. + +## Rights and provenance + +All text in `corpus.json` was invented for this contract. No real persons, +events, footage, transcripts, or captions are included. No private source, +prompts, rubrics, customer assets, or private experiment data was used. Media +files are **not stored in the repository**: `corpus.json` records +deterministic ffmpeg commands (`testsrc2` color source + sine tone) so any +ffmpeg build can regenerate synthetic media outside the working tree for +export-validity checks. The corpus is therefore synthetic and rights-cleared +by construction. + +## Metrics (contract v1) + +Every metric compares returned clips against frozen labels only. Jev decision +scores are never ground truth. + +- `precision_at_k` — hits in the first k returned clips / k. A hit means + interval IoU ≥ 0.5 against a labeled moment. +- `known_moment_recall` — labeled moments covered (IoU ≥ 0.5) / total moments. +- `absent_topic_fp` — 1.0 when any clip is returned for an `absent` query. + The dev set includes an incidental-mention trap: the sponsor line mentions + "cloud credits" once; it is labeled absent. +- `boundary_error_sec` — mean of (|start delta| + |end delta|) / 2 over hits. +- `duplication_rate` — returned clip pairs with IoU ≥ 0.2 / clip count. +- `context_loss_count` — returned hits that start after a labeled setup head + (more than 1s of slack). q-005 ("laminating trick") is the canonical case. + +Timing is reported as ingestion (corpus materialization), prepare +(retrieval + candidates), selection (typed decisions + assembly), +selection_warm (identical replay from the in-run decision cache, zero +provider calls), and render (ffmpeg export) milliseconds. Percentiles are +reported only when at least five samples support them. + +## How to run + +```bash +# Offline (deterministic + labeled-mock Jev pipeline exercise) +python -m av.clip_eval \ + --corpus clip-eval/corpus.json --queries clip-eval/queries.json \ + --db /path/eval.db --media-dir /path/media --export-dir /path/render \ + --arms deterministic,jev_mock \ + --out clip-eval/receipts/offline-mock-dev.json +``` + +## Evidence status + +Committed receipts are **offline mock runs**: the `jev_mock` arm exercises +the typed Noul/Choice/Score pipeline against frozen labels. It proves the +harness works and quantifies the deterministic baseline; it is **not** a +measurement of Jev quality. Live Jev runs additionally require an authorized +allowance, an explicit per-run request ceiling, and usage tracking; none was +available for the initial contract, so live quality evidence is pending. + +Measured at target 30s / min 10s / 2 clips (dev, `receipts/offline-mock-dev.json`): + +| Arm | P@k (k=2) | Known-moment recall | Absent-topic FP | Boundary error (s) | Duplication | +|---|---|---|---|---|---| +| deterministic | 0.1875 | 0.625 | 1.0 | 3.667 | 0.0 | +| jev_mock | 0.25 | 0.875 | 0.0 | 2.75 | 0.0 | + +Held-out, run once after the final harness state (`receipts/offline-mock-heldout.json`): +deterministic P@k 0.25 / recall 0.667 / boundary error 5.0s versus +jev_mock P@k 0.375 / recall 1.0 / boundary error 3.0s. Both arms have zero +absent-topic false positives, zero duplication, and no context-loss cases. + +Known gaps, kept deliberately: q-005 measures the duration-vs-context +tradeoff — the 32s labeled moment cannot be covered at a 30s cap without +losing part of the setup head, and both arms report it (the Jev arm returns a +partial-setup clip that scores a context loss; the deterministic arm misses +entirely). Tiny synthetic fixtures are smoke evidence for the contract, not +proof of broad quality. diff --git a/clip-eval/corpus.json b/clip-eval/corpus.json new file mode 100644 index 0000000..90f6d0a --- /dev/null +++ b/clip-eval/corpus.json @@ -0,0 +1,315 @@ +{ + "contract_version": 1, + "description": "Synthetic rights-cleared clip evaluation corpus. All text invented for this contract; media generated on demand by deterministic ffmpeg commands recorded in each content block. No real persons, events, or footage.", + "videos": [ + { + "id": "av-clip-eval-talk-001", + "content": { + "style": "conference-talk", + "duration_sec": 140.0, + "transcript": [ + { + "id": "talk-t01", + "start_sec": 0, + "end_sec": 8, + "text": "Welcome everyone, thanks for coming to the systems track." + }, + { + "id": "talk-t02", + "start_sec": 8, + "end_sec": 16, + "text": "Today I want to talk about how we keep large deployments reliable." + }, + { + "id": "talk-t03", + "start_sec": 16, + "end_sec": 24, + "text": "Let me start with a quick story from our own infrastructure." + }, + { + "id": "talk-t04", + "start_sec": 24, + "end_sec": 34, + "text": "Quantum error correction sounds abstract, but think of it as backup batteries for qubits." + }, + { + "id": "talk-t05", + "start_sec": 34, + "end_sec": 44, + "text": "When a qubit flips, the correction layer catches it and rewrites the state." + }, + { + "id": "talk-t06", + "start_sec": 44, + "end_sec": 58, + "text": "We ran this on a live cluster, and the error rate collapsed by two orders of magnitude." + }, + { + "id": "talk-t07", + "start_sec": 58, + "end_sec": 66, + "text": "Before questions, a short word from our sponsor about cloud credits." + }, + { + "id": "talk-t08", + "start_sec": 66, + "end_sec": 78, + "text": "The grant committee wants monthly reports, and honestly the paperwork is heavy." + }, + { + "id": "talk-t09", + "start_sec": 78, + "end_sec": 88, + "text": "Back to the technical track: redundancy alone is not resilience." + }, + { + "id": "talk-t10", + "start_sec": 88, + "end_sec": 100, + "text": "Quantum error correction only works when the code distance grows with the noise." + }, + { + "id": "talk-t11", + "start_sec": 100, + "end_sec": 112, + "text": "In our demo, three faulty gates in a row were still recovered cleanly." + }, + { + "id": "talk-t12", + "start_sec": 112, + "end_sec": 122, + "text": "That recovery is the moment reliability stops being a promise." + }, + { + "id": "talk-t13", + "start_sec": 122, + "end_sec": 132, + "text": "To sum up: measure, correct, and repeat." + }, + { + "id": "talk-t14", + "start_sec": 132, + "end_sec": 140, + "text": "Thank you, and I will take questions at the booth." + } + ], + "vision": [ + { + "id": "talk-c01", + "start_sec": 0, + "end_sec": 40, + "text": "Speaker on stage welcoming the audience, slides showing the conference logo", + "source_type": "dense_caption" + }, + { + "id": "talk-c02", + "start_sec": 40, + "end_sec": 80, + "text": "Slide shows a lattice diagram; the speaker gestures at rising and falling error bars", + "source_type": "dense_caption" + }, + { + "id": "talk-c03", + "start_sec": 80, + "end_sec": 120, + "text": "Slide shows a code distance chart; the speaker points at a flat line after correction", + "source_type": "dense_caption" + }, + { + "id": "talk-c04", + "start_sec": 120, + "end_sec": 140, + "text": "Audience applauding; the speaker steps back from the podium", + "source_type": "dense_caption" + } + ], + "media": { + "filename": "av-clip-eval-talk-001.mp4", + "tone_hz": 440, + "approx_size_bytes": 450000 + } + } + }, + { + "id": "av-clip-eval-podcast-001", + "content": { + "style": "podcast-interview-noisy-asr", + "duration_sec": 130.0, + "transcript": [ + { + "id": "pod-t01", + "start_sec": 0, + "end_sec": 10, + "text": "Welcome back to the show, today we are baking with a master baker." + }, + { + "id": "pod-t02", + "start_sec": 10, + "end_sec": 20, + "text": "Let us talk about sower dough, the classic fermented bread." + }, + { + "id": "pod-t03", + "start_sec": 20, + "end_sec": 30, + "text": "A good starter needs feeding every single day, no exceptions." + }, + { + "id": "pod-t04", + "start_sec": 30, + "end_sec": 42, + "text": "If you skip a feeding, the starter turns acidic and the loaf collapses." + }, + { + "id": "pod-t05", + "start_sec": 42, + "end_sec": 50, + "text": "And the queen aman? That pastry is pure butter architecture." + }, + { + "id": "pod-t06", + "start_sec": 50, + "end_sec": 58, + "text": "Most people ruin it by rushing the folds, patience is everything." + }, + { + "id": "pod-t07", + "start_sec": 58, + "end_sec": 70, + "text": "Let me tell you the part nobody tells you about lamination." + }, + { + "id": "pod-t08", + "start_sec": 70, + "end_sec": 82, + "text": "Cold butter, patient folds, and two full rest hours." + }, + { + "id": "pod-t09", + "start_sec": 82, + "end_sec": 90, + "text": "That is the trick, and it works every single time." + }, + { + "id": "pod-t10", + "start_sec": 90, + "end_sec": 100, + "text": "Wow, I had no idea rest time mattered that much." + }, + { + "id": "pod-t11", + "start_sec": 100, + "end_sec": 112, + "text": "So to recap: ferment slowly, laminate cold, bake hot." + }, + { + "id": "pod-t12", + "start_sec": 112, + "end_sec": 130, + "text": "Thanks for listening, and send us your baking fails." + } + ], + "vision": [ + { + "id": "pod-c01", + "start_sec": 0, + "end_sec": 44, + "text": "Two people at a kitchen counter; loaves and jars of starter are visible", + "source_type": "caption" + }, + { + "id": "pod-c02", + "start_sec": 44, + "end_sec": 90, + "text": "Hosts folding dough on a floured counter; a butter block sits on the table", + "source_type": "caption" + }, + { + "id": "pod-c03", + "start_sec": 90, + "end_sec": 130, + "text": "Hosts laughing; finished croissants cool on a rack", + "source_type": "caption" + } + ], + "media": { + "filename": "av-clip-eval-podcast-001.mp4", + "tone_hz": 523, + "approx_size_bytes": 420000 + } + } + }, + { + "id": "av-clip-eval-harbor-001", + "content": { + "style": "ambient-broll-documentary", + "duration_sec": 70.0, + "transcript": [ + { + "id": "har-t01", + "start_sec": 8, + "end_sec": 16, + "text": "The night market at the harbor opens at sundown." + }, + { + "id": "har-t02", + "start_sec": 16, + "end_sec": 24, + "text": "Stalls light up along the pier, and the smell of grilled fish takes over." + }, + { + "id": "har-t03", + "start_sec": 30, + "end_sec": 44, + "text": "At dawn the fishing boats return, and the harbor auction begins." + }, + { + "id": "har-t04", + "start_sec": 44, + "end_sec": 52, + "text": "Buyers crowd the slab, and the auction bell decides who takes the catch home." + }, + { + "id": "har-t05", + "start_sec": 52, + "end_sec": 64, + "text": "By noon the pier is quiet again, waiting for the next tide." + } + ], + "vision": [ + { + "id": "har-c01", + "start_sec": 0, + "end_sec": 20, + "text": "Empty pier at night, string lights over wet boards", + "source_type": "caption" + }, + { + "id": "har-c02", + "start_sec": 12, + "end_sec": 18, + "text": "Daytime market stalls crowded with shoppers", + "source_type": "caption" + }, + { + "id": "har-c03", + "start_sec": 64, + "end_sec": 70, + "text": "Wide shot of the harbor at noon", + "source_type": "caption" + } + ], + "media": { + "filename": "av-clip-eval-harbor-001.mp4", + "tone_hz": 349, + "approx_size_bytes": 240000 + } + } + } + ], + "checksums": { + "av-clip-eval-talk-001": "sha256:d438002df93eacf18703cfca4a6794fb57fc46097310e918f585e953f6022c1f", + "av-clip-eval-podcast-001": "sha256:f5b9f0132415e34ed043ebe7cb43fed033c76faf6061dffbb44a5950a33a6b5e", + "av-clip-eval-harbor-001": "sha256:96a356e467220c60ff49255608eec93af948623d5956deca37139c40965573a3" + } +} diff --git a/clip-eval/queries-heldout.json b/clip-eval/queries-heldout.json new file mode 100644 index 0000000..899043b --- /dev/null +++ b/clip-eval/queries-heldout.json @@ -0,0 +1,62 @@ +{ + "contract_version": 1, + "role": "heldout", + "policy": "Run at most once per evidence cycle and never against prompt or threshold iteration. Labels here were frozen before any clip run on them.", + "queries": [ + { + "query_id": "q-h01", + "topic": "kouign aman pastry", + "video_id": "av-clip-eval-podcast-001", + "expected": "present", + "moments": [ + { + "moment_id": "m-h01", + "start_sec": 42.0, + "end_sec": 58.0, + "appeal": 0.7, + "requires_setup": false, + "note": "ASR renders 'kouign aman' as 'queen aman'" + } + ] + }, + { + "query_id": "q-h02", + "topic": "code distance", + "video_id": "av-clip-eval-talk-001", + "expected": "present", + "moments": [ + { + "moment_id": "m-h02", + "start_sec": 88.0, + "end_sec": 100.0, + "appeal": 0.6, + "requires_setup": false + } + ] + }, + { + "query_id": "q-h03", + "topic": "night market openings", + "video_id": "av-clip-eval-harbor-001", + "expected": "present", + "moments": [ + { + "moment_id": "m-h03", + "start_sec": 8.0, + "end_sec": 24.0, + "appeal": 0.6, + "requires_setup": false, + "note": "one caption contradicts the transcript by claiming daytime" + } + ], + "note": "contradictory vision must not silently invalidate the moment" + }, + { + "query_id": "q-h04", + "topic": "submarine tour", + "video_id": "av-clip-eval-harbor-001", + "expected": "absent", + "moments": [] + } + ] +} diff --git a/clip-eval/queries.json b/clip-eval/queries.json new file mode 100644 index 0000000..9c96a7f --- /dev/null +++ b/clip-eval/queries.json @@ -0,0 +1,117 @@ +{ + "contract_version": 1, + "role": "dev", + "queries": [ + { + "query_id": "q-001", + "topic": "quantum error correction", + "video_id": "av-clip-eval-talk-001", + "expected": "present", + "moments": [ + { + "moment_id": "m-101", + "start_sec": 24.0, + "end_sec": 58.0, + "appeal": 0.85, + "requires_setup": true, + "setup_start_sec": 24.0, + "setup_end_sec": 34.0, + "note": "setup then payoff; repeated topic" + }, + { + "moment_id": "m-102", + "start_sec": 88.0, + "end_sec": 112.0, + "appeal": 0.85, + "requires_setup": true, + "setup_start_sec": 88.0, + "setup_end_sec": 100.0, + "note": "second occurrence; non-overlap matters" + } + ] + }, + { + "query_id": "q-002", + "topic": "cloud credits", + "video_id": "av-clip-eval-talk-001", + "expected": "absent", + "moments": [], + "note": "incidental sponsor mention at 58-66; not a discussion of the topic" + }, + { + "query_id": "q-003", + "topic": "stock market predictions", + "video_id": "av-clip-eval-talk-001", + "expected": "absent", + "moments": [] + }, + { + "query_id": "q-004", + "topic": "sourdough starter schedule", + "video_id": "av-clip-eval-podcast-001", + "expected": "present", + "moments": [ + { + "moment_id": "m-201", + "start_sec": 20.0, + "end_sec": 42.0, + "appeal": 0.8, + "requires_setup": false, + "note": "noisy ASR: 'sourdough' appears as 'sower dough'" + } + ], + "note": "retrieval must survive ASR noise" + }, + { + "query_id": "q-005", + "topic": "laminating trick", + "video_id": "av-clip-eval-podcast-001", + "expected": "present", + "moments": [ + { + "moment_id": "m-202", + "start_sec": 58.0, + "end_sec": 90.0, + "appeal": 0.9, + "requires_setup": true, + "setup_start_sec": 58.0, + "setup_end_sec": 82.0, + "note": "context-dependent soundbite; payoff at 82-90 needs the setup (hard_context: payoff alone is incoherent)", + "hard_context": true + } + ], + "note": "clipping only the payoff loses the setup" + }, + { + "query_id": "q-006", + "topic": "electric car reviews", + "video_id": "av-clip-eval-podcast-001", + "expected": "absent", + "moments": [] + }, + { + "query_id": "q-007", + "topic": "harbor fish auction", + "video_id": "av-clip-eval-harbor-001", + "expected": "present", + "moments": [ + { + "moment_id": "m-301", + "start_sec": 30.0, + "end_sec": 52.0, + "appeal": 0.75, + "requires_setup": false, + "note": "no vision captions cover this moment; visual support is unavailable, not contradictory" + } + ], + "note": "missing vision over the moment" + }, + { + "query_id": "q-008", + "topic": "space launch", + "video_id": "av-clip-eval-harbor-001", + "expected": "absent", + "moments": [] + } + ] +} diff --git a/clip-eval/receipts/offline-mock-dev.json b/clip-eval/receipts/offline-mock-dev.json new file mode 100644 index 0000000..5ac4617 --- /dev/null +++ b/clip-eval/receipts/offline-mock-dev.json @@ -0,0 +1,454 @@ +{ + "contract_version": 1, + "corpus": "clip-eval/corpus.json", + "queries": "clip-eval/queries.json", + "query_count": 8, + "ingestion_ms": 222.4, + "arms": [ + { + "arm": "deterministic", + "decide": false, + "metrics_per_query": [ + { + "clips_returned": 2, + "hits": 1, + "precision_at_k": 0.5, + "known_moment_recall": 0.5, + "absent_topic_fp": 0.0, + "boundary_error_sec": 5.0, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-001", + "topic": "quantum error correction", + "status": "deterministic_only", + "timings": { + "prepare_ms": 4.7, + "selection_ms": 0.0, + "selection_warm_ms": 0.0, + "render_ms": 0.0 + } + }, + { + "clips_returned": 1, + "hits": 0, + "precision_at_k": 0.0, + "known_moment_recall": null, + "absent_topic_fp": 1.0, + "boundary_error_sec": null, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-002", + "topic": "cloud credits", + "status": "deterministic_only", + "timings": { + "prepare_ms": 1.3, + "selection_ms": 0.0, + "selection_warm_ms": 0.0, + "render_ms": 0.0 + } + }, + { + "clips_returned": 0, + "hits": 0, + "precision_at_k": 0.0, + "known_moment_recall": null, + "absent_topic_fp": 0.0, + "boundary_error_sec": null, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-003", + "topic": "stock market predictions", + "status": "no_usable_clips", + "timings": { + "prepare_ms": 0.4, + "selection_ms": 0.0, + "selection_warm_ms": 0.0, + "render_ms": 0.0 + } + }, + { + "clips_returned": 1, + "hits": 1, + "precision_at_k": 0.5, + "known_moment_recall": 1.0, + "absent_topic_fp": 0.0, + "boundary_error_sec": 0.0, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-004", + "topic": "sourdough starter schedule", + "status": "deterministic_only", + "timings": { + "prepare_ms": 1.6, + "selection_ms": 0.0, + "selection_warm_ms": 0.0, + "render_ms": 0.0 + } + }, + { + "clips_returned": 1, + "hits": 0, + "precision_at_k": 0.0, + "known_moment_recall": 0.0, + "absent_topic_fp": 0.0, + "boundary_error_sec": null, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-005", + "topic": "laminating trick", + "status": "deterministic_only", + "timings": { + "prepare_ms": 1.3, + "selection_ms": 0.0, + "selection_warm_ms": 0.0, + "render_ms": 0.0 + } + }, + { + "clips_returned": 0, + "hits": 0, + "precision_at_k": 0.0, + "known_moment_recall": null, + "absent_topic_fp": 0.0, + "boundary_error_sec": null, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-006", + "topic": "electric car reviews", + "status": "no_usable_clips", + "timings": { + "prepare_ms": 0.4, + "selection_ms": 0.0, + "selection_warm_ms": 0.0, + "render_ms": 0.0 + } + }, + { + "clips_returned": 2, + "hits": 1, + "precision_at_k": 0.5, + "known_moment_recall": 1.0, + "absent_topic_fp": 0.0, + "boundary_error_sec": 6.0, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-007", + "topic": "harbor fish auction", + "status": "deterministic_only", + "timings": { + "prepare_ms": 2.9, + "selection_ms": 0.0, + "selection_warm_ms": 0.0, + "render_ms": 0.0 + } + }, + { + "clips_returned": 0, + "hits": 0, + "precision_at_k": 0.0, + "known_moment_recall": null, + "absent_topic_fp": 0.0, + "boundary_error_sec": null, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-008", + "topic": "space launch", + "status": "no_usable_clips", + "timings": { + "prepare_ms": 0.4, + "selection_ms": 0.0, + "selection_warm_ms": 0.0, + "render_ms": 0.0 + } + } + ], + "metrics_total": { + "queries": 8, + "precision_at_k_mean": 0.1875, + "known_moment_recall_mean": 0.625, + "absent_topic_fp_total": 1.0, + "boundary_error_sec_mean": 3.667, + "duplication_rate_mean": 0.0, + "context_loss_total": 0, + "precision_at_k_p50": 0.0, + "boundary_error_sec_p50": null, + "boundary_error_sec_p95": null + }, + "stage_usage": { + "relevance": { + "requests": 0, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + }, + "coherence": { + "requests": 0, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + }, + "visual": { + "requests": 0, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + }, + "boundary": { + "requests": 0, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + }, + "appeal": { + "requests": 0, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + } + }, + "timings": { + "prepare_ms_total": 13.0, + "selection_ms_total": 0.0, + "selection_warm_ms_total": 0.0, + "render_ms_total": 11558.8 + }, + "export": { + "requested": true, + "clips_exported": 7, + "clips_valid": 7 + }, + "failures": [], + "decision_provider": "none" + }, + { + "arm": "jev_mock", + "decide": true, + "metrics_per_query": [ + { + "clips_returned": 2, + "hits": 1, + "precision_at_k": 0.5, + "known_moment_recall": 0.5, + "absent_topic_fp": 0.0, + "boundary_error_sec": 5.0, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-001", + "topic": "quantum error correction", + "status": "ok", + "timings": { + "prepare_ms": 3.0, + "selection_ms": 0.8, + "selection_warm_ms": 0.5, + "render_ms": 0.0 + } + }, + { + "clips_returned": 0, + "hits": 0, + "precision_at_k": 0.0, + "known_moment_recall": null, + "absent_topic_fp": 0.0, + "boundary_error_sec": null, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-002", + "topic": "cloud credits", + "status": "no_usable_clips", + "timings": { + "prepare_ms": 1.2, + "selection_ms": 0.4, + "selection_warm_ms": 0.3, + "render_ms": 0.0 + } + }, + { + "clips_returned": 0, + "hits": 0, + "precision_at_k": 0.0, + "known_moment_recall": null, + "absent_topic_fp": 0.0, + "boundary_error_sec": null, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-003", + "topic": "stock market predictions", + "status": "no_usable_clips", + "timings": { + "prepare_ms": 0.2, + "selection_ms": 0.0, + "selection_warm_ms": 0.0, + "render_ms": 0.0 + } + }, + { + "clips_returned": 1, + "hits": 1, + "precision_at_k": 0.5, + "known_moment_recall": 1.0, + "absent_topic_fp": 0.0, + "boundary_error_sec": 0.0, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-004", + "topic": "sourdough starter schedule", + "status": "ok", + "timings": { + "prepare_ms": 2.5, + "selection_ms": 0.3, + "selection_warm_ms": 0.2, + "render_ms": 0.0 + } + }, + { + "clips_returned": 1, + "hits": 1, + "precision_at_k": 0.5, + "known_moment_recall": 1.0, + "absent_topic_fp": 0.0, + "boundary_error_sec": 6.0, + "duplication_rate": 0.0, + "context_loss_count": 1, + "query_id": "q-005", + "topic": "laminating trick", + "status": "ok", + "timings": { + "prepare_ms": 1.4, + "selection_ms": 0.5, + "selection_warm_ms": 0.2, + "render_ms": 0.0 + } + }, + { + "clips_returned": 0, + "hits": 0, + "precision_at_k": 0.0, + "known_moment_recall": null, + "absent_topic_fp": 0.0, + "boundary_error_sec": null, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-006", + "topic": "electric car reviews", + "status": "no_usable_clips", + "timings": { + "prepare_ms": 0.3, + "selection_ms": 0.0, + "selection_warm_ms": 0.0, + "render_ms": 0.0 + } + }, + { + "clips_returned": 1, + "hits": 1, + "precision_at_k": 0.5, + "known_moment_recall": 1.0, + "absent_topic_fp": 0.0, + "boundary_error_sec": 0.0, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-007", + "topic": "harbor fish auction", + "status": "ok", + "timings": { + "prepare_ms": 2.5, + "selection_ms": 0.6, + "selection_warm_ms": 0.3, + "render_ms": 0.0 + } + }, + { + "clips_returned": 0, + "hits": 0, + "precision_at_k": 0.0, + "known_moment_recall": null, + "absent_topic_fp": 0.0, + "boundary_error_sec": null, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-008", + "topic": "space launch", + "status": "no_usable_clips", + "timings": { + "prepare_ms": 0.3, + "selection_ms": 0.0, + "selection_warm_ms": 0.0, + "render_ms": 0.0 + } + } + ], + "metrics_total": { + "queries": 8, + "precision_at_k_mean": 0.25, + "known_moment_recall_mean": 0.875, + "absent_topic_fp_total": 0.0, + "boundary_error_sec_mean": 2.75, + "duplication_rate_mean": 0.0, + "context_loss_total": 1, + "precision_at_k_p50": 0.25, + "boundary_error_sec_p50": null, + "boundary_error_sec_p95": null + }, + "stage_usage": { + "relevance": { + "requests": 5, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + }, + "coherence": { + "requests": 5, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + }, + "visual": { + "requests": 5, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + }, + "boundary": { + "requests": 9, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + }, + "appeal": { + "requests": 5, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + } + }, + "timings": { + "prepare_ms_total": 11.4, + "selection_ms_total": 2.6, + "selection_warm_ms_total": 1.5, + "render_ms_total": 8014.5 + }, + "export": { + "requested": true, + "clips_exported": 5, + "clips_valid": 5 + }, + "failures": [], + "decision_provider": "typesafe" + } + ], + "notes": [ + "Synthetic rights-cleared corpus; labels are frozen and independent of decision scores.", + "The jev_mock arm exercises the typed decision pipeline; it is not a Jev quality measurement.", + "Percentile timings are reported only when at least five samples support them." + ] +} diff --git a/clip-eval/receipts/offline-mock-heldout.json b/clip-eval/receipts/offline-mock-heldout.json new file mode 100644 index 0000000..3d51630 --- /dev/null +++ b/clip-eval/receipts/offline-mock-heldout.json @@ -0,0 +1,302 @@ +{ + "contract_version": 1, + "corpus": "clip-eval/corpus.json", + "queries": "clip-eval/queries-heldout.json", + "query_count": 4, + "ingestion_ms": 238.2, + "arms": [ + { + "arm": "deterministic", + "decide": false, + "metrics_per_query": [ + { + "clips_returned": 1, + "hits": 1, + "precision_at_k": 0.5, + "known_moment_recall": 1.0, + "absent_topic_fp": 0.0, + "boundary_error_sec": 6.0, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-h01", + "topic": "kouign aman pastry", + "status": "deterministic_only", + "timings": { + "prepare_ms": 1.6, + "selection_ms": 0.0, + "selection_warm_ms": 0.0, + "render_ms": 0.0 + } + }, + { + "clips_returned": 1, + "hits": 0, + "precision_at_k": 0.0, + "known_moment_recall": 0.0, + "absent_topic_fp": 0.0, + "boundary_error_sec": null, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-h02", + "topic": "code distance", + "status": "deterministic_only", + "timings": { + "prepare_ms": 2.2, + "selection_ms": 0.0, + "selection_warm_ms": 0.0, + "render_ms": 0.0 + } + }, + { + "clips_returned": 1, + "hits": 1, + "precision_at_k": 0.5, + "known_moment_recall": 1.0, + "absent_topic_fp": 0.0, + "boundary_error_sec": 4.0, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-h03", + "topic": "night market openings", + "status": "deterministic_only", + "timings": { + "prepare_ms": 2.0, + "selection_ms": 0.0, + "selection_warm_ms": 0.0, + "render_ms": 0.0 + } + }, + { + "clips_returned": 0, + "hits": 0, + "precision_at_k": 0.0, + "known_moment_recall": null, + "absent_topic_fp": 0.0, + "boundary_error_sec": null, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-h04", + "topic": "submarine tour", + "status": "no_usable_clips", + "timings": { + "prepare_ms": 0.4, + "selection_ms": 0.0, + "selection_warm_ms": 0.0, + "render_ms": 0.0 + } + } + ], + "metrics_total": { + "queries": 4, + "precision_at_k_mean": 0.25, + "known_moment_recall_mean": 0.6667, + "absent_topic_fp_total": 0.0, + "boundary_error_sec_mean": 5.0, + "duplication_rate_mean": 0.0, + "context_loss_total": 0, + "precision_at_k_p50": null, + "boundary_error_sec_p50": null, + "boundary_error_sec_p95": null + }, + "stage_usage": { + "relevance": { + "requests": 0, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + }, + "coherence": { + "requests": 0, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + }, + "visual": { + "requests": 0, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + }, + "boundary": { + "requests": 0, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + }, + "appeal": { + "requests": 0, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + } + }, + "timings": { + "prepare_ms_total": 6.2, + "selection_ms_total": 0.0, + "selection_warm_ms_total": 0.0, + "render_ms_total": 5106.1 + }, + "export": { + "requested": true, + "clips_exported": 3, + "clips_valid": 3 + }, + "failures": [], + "decision_provider": "none" + }, + { + "arm": "jev_mock", + "decide": true, + "metrics_per_query": [ + { + "clips_returned": 1, + "hits": 1, + "precision_at_k": 0.5, + "known_moment_recall": 1.0, + "absent_topic_fp": 0.0, + "boundary_error_sec": 0.0, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-h01", + "topic": "kouign aman pastry", + "status": "ok", + "timings": { + "prepare_ms": 1.3, + "selection_ms": 0.6, + "selection_warm_ms": 0.3, + "render_ms": 0.0 + } + }, + { + "clips_returned": 1, + "hits": 1, + "precision_at_k": 0.5, + "known_moment_recall": 1.0, + "absent_topic_fp": 0.0, + "boundary_error_sec": 5.0, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-h02", + "topic": "code distance", + "status": "ok", + "timings": { + "prepare_ms": 1.7, + "selection_ms": 0.4, + "selection_warm_ms": 0.2, + "render_ms": 0.0 + } + }, + { + "clips_returned": 1, + "hits": 1, + "precision_at_k": 0.5, + "known_moment_recall": 1.0, + "absent_topic_fp": 0.0, + "boundary_error_sec": 4.0, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-h03", + "topic": "night market openings", + "status": "ok", + "timings": { + "prepare_ms": 2.3, + "selection_ms": 0.4, + "selection_warm_ms": 0.2, + "render_ms": 0.0 + } + }, + { + "clips_returned": 0, + "hits": 0, + "precision_at_k": 0.0, + "known_moment_recall": null, + "absent_topic_fp": 0.0, + "boundary_error_sec": null, + "duplication_rate": 0.0, + "context_loss_count": 0, + "query_id": "q-h04", + "topic": "submarine tour", + "status": "no_usable_clips", + "timings": { + "prepare_ms": 0.6, + "selection_ms": 0.0, + "selection_warm_ms": 0.0, + "render_ms": 0.0 + } + } + ], + "metrics_total": { + "queries": 4, + "precision_at_k_mean": 0.375, + "known_moment_recall_mean": 1.0, + "absent_topic_fp_total": 0.0, + "boundary_error_sec_mean": 3.0, + "duplication_rate_mean": 0.0, + "context_loss_total": 0, + "precision_at_k_p50": null, + "boundary_error_sec_p50": null, + "boundary_error_sec_p95": null + }, + "stage_usage": { + "relevance": { + "requests": 3, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + }, + "coherence": { + "requests": 3, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + }, + "visual": { + "requests": 3, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + }, + "boundary": { + "requests": 3, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + }, + "appeal": { + "requests": 3, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": true, + "output_tokens_complete": true + } + }, + "timings": { + "prepare_ms_total": 5.9, + "selection_ms_total": 1.4, + "selection_warm_ms_total": 0.7, + "render_ms_total": 4072.4 + }, + "export": { + "requested": true, + "clips_exported": 3, + "clips_valid": 3 + }, + "failures": [], + "decision_provider": "typesafe" + } + ], + "notes": [ + "Synthetic rights-cleared corpus; labels are frozen and independent of decision scores.", + "The jev_mock arm exercises the typed decision pipeline; it is not a Jev quality measurement.", + "Percentile timings are reported only when at least five samples support them." + ] +} diff --git a/cookbook/jev-topic-clipping/README.md b/cookbook/jev-topic-clipping/README.md new file mode 100644 index 0000000..6fb4d8f --- /dev/null +++ b/cookbook/jev-topic-clipping/README.md @@ -0,0 +1,61 @@ +# Jev topic clipping + +Recipe for the `av clip` command: find and optionally export topic-specific +highlight clips from one already-indexed video — no re-ingestion, no +model-invented timestamps. + +## What it demonstrates + +- Candidate construction from retrieval plus bounded temporal neighborhoods, + with timing anchored on transcript segment boundaries (vision captions are + attached as evidence, never as clock sources). +- Typed Jev decisions — Noul relevance, Noul standalone coherence, Noul + visual evidence, Choice boundaries, Score highlight appeal — kept strictly + separate: objective source support gates selection, appeal only orders. +- An explicit per-run request ceiling, a truthful per-stage usage receipt, + and a warm replay that measures harness overhead with zero provider calls. +- Optional ffmpeg export with ffprobe playability/synchronization checks. + +## Run + +```bash +uv sync --extra dev +uv run av ingest examples/lecture.mp4 # once; any supported ingest path +uv run av clip "quantum error correction" \ + --video-id --clips 2 --target-seconds 30 \ + --export ./clips --overwrite-export +``` + +stdout is a single JSON document: clips with source-verbatim quotes, +supporting artifact IDs, decision scores, uncertainty, provider/model, the +request-cap state, per-stage usage, and prepare/selection/warm/render +timings. `--no-decide` skips Jev and ranks retrieval candidates +deterministically; absent topics return `no_usable_clips` instead of guesses. + +## Evaluation + +The shared contract lives in [`clip-eval/`](../../clip-eval/README.md): +synthetic corpus with frozen checksums, labeled dev and held-out query sets, +frozen metric definitions, and a runner (`python -m av.clip_eval`) that +compares the deterministic arm against the Jev-decided arm on identical +candidates. Committed receipts are offline mock runs; see the evidence +status section there. + +## Evidence status + +Runnable recipe with deterministic tests and an offline mock evaluation +receipt. **No live Jev calls were made for this recipe** — no authorized +allowance was available in this environment, so live selection quality, +live token/request totals, and live charges are pending, unmeasured, and +not claimed. The "two-second, two-cent" social claim about clipping cost is +treated as unverified community folklore; it is not an acceptance target and +is not reproduced here. No speed, cost, or quality parity claim is made +against any other tool. + +## Advisor receipt + +Advisor review was sought through the enabled advisor runtime at two points +(repository orientation before settling the candidate schema, and at final +verification). The runtime reported no advisor peers on both checks, so no +advice was observed or applied; the unavailable category is recorded here +rather than synthesized. diff --git a/src/av/cli/app.py b/src/av/cli/app.py index 449c1ae..390e735 100644 --- a/src/av/cli/app.py +++ b/src/av/cli/app.py @@ -3,7 +3,6 @@ from __future__ import annotations import json -from pathlib import Path import typer @@ -25,22 +24,24 @@ def version_cmd() -> None: # --- Register direct commands --- -from av.cli.ingest import register as register_ingest # noqa: E402 -from av.cli.search import register as register_search # noqa: E402 -from av.cli.ask import register as register_ask # noqa: E402 -from av.cli.list_cmd import register as register_list # noqa: E402 -from av.cli.info import register as register_info # noqa: E402 -from av.cli.transcript import register as register_transcript # noqa: E402 -from av.cli.export import register as register_export # noqa: E402 -from av.cli.open_cmd import register as register_open # noqa: E402 -from av.cli.config_cmd import config_app # noqa: E402 -from av.cli.sentinel import register as register_sentinel # noqa: E402 -from av.cli.sentinel_doctor import register_doctor # noqa: E402 -from av.cli.bench import register as register_bench # noqa: E402 +from av.cli.ask import register as register_ask +from av.cli.bench import register as register_bench +from av.cli.clip import register as register_clip +from av.cli.config_cmd import config_app +from av.cli.export import register as register_export +from av.cli.info import register as register_info +from av.cli.ingest import register as register_ingest +from av.cli.list_cmd import register as register_list +from av.cli.open_cmd import register as register_open +from av.cli.search import register as register_search +from av.cli.sentinel import register as register_sentinel +from av.cli.sentinel_doctor import register_doctor +from av.cli.transcript import register as register_transcript register_ingest(app) register_search(app) register_ask(app) +register_clip(app) register_list(app) register_info(app) register_transcript(app) diff --git a/src/av/cli/clip.py b/src/av/cli/clip.py new file mode 100644 index 0000000..2b926f0 --- /dev/null +++ b/src/av/cli/clip.py @@ -0,0 +1,100 @@ +"""av clip command — topic-specific highlight clips from one indexed video.""" + +from __future__ import annotations + +from pathlib import Path + +import typer + +from av.cli.output import error, output_json +from av.core.config import get_config +from av.db.repository import Repository +from av.pipeline.clip_export import export_clips +from av.search.clip import ( + DEFAULT_RETRIEVAL_LIMIT, + ClipError, + clip_video, +) + + +def register(app: typer.Typer) -> None: + @app.command("clip") + def clip_cmd( + topic: str = typer.Argument(..., help="Topic or moment to find highlights for"), + video_id: str = typer.Option(..., "--video-id", "-v", help="Indexed video to clip"), + clips: int = typer.Option(3, "--clips", "-k", min=1, help="Maximum clips to return"), + target_seconds: float = typer.Option( + 30.0, "--target-seconds", min=0.1, help="Preferred clip duration in seconds" + ), + min_seconds: float = typer.Option( + 10.0, "--min-seconds", min=0.1, help="Minimum clip duration in seconds" + ), + max_seconds: float = typer.Option( + None, "--max-seconds", min=0.1, help="Maximum clip duration (default: target)" + ), + retrieval_limit: int = typer.Option( + DEFAULT_RETRIEVAL_LIMIT, "--retrieval-limit", min=1, help="Retrieval hits to consider" + ), + no_decide: bool = typer.Option( + False, + "--no-decide", + help="Skip Jev decisions; rank retrieval candidates deterministically", + ), + max_decision_requests: int = typer.Option( + None, + "--max-decision-requests", + min=1, + help="Per-run ceiling on System One HTTP attempts (default: AV_CLIP_REQUEST_CAP)", + ), + export_dir: str = typer.Option( + None, "--export", help="Directory to render selected clips with ffmpeg" + ), + overwrite_export: bool = typer.Option( + False, "--overwrite-export", help="Replace existing export files" + ), + db: str = typer.Option(None, "--db", help="Database path override"), + ) -> None: + """Find and optionally export topic-specific highlight clips. + + Candidates come from retrieval plus bounded temporal neighborhoods, and + every boundary is an artifact boundary already in the index. When + configured, Jev judges topical relevance, standalone coherence, visual + evidence, boundaries, and highlight appeal with typed Noul, Choice, + and Score operations under an explicit request ceiling. Absent topics + return no clips instead of best guesses. + """ + config = get_config(db_path=Path(db) if db else None) + repo = Repository(config.db_path) + + try: + result = clip_video( + topic, + video_id, + repo, + config, + clips_wanted=clips, + target_seconds=target_seconds, + min_seconds=min_seconds, + max_seconds=max_seconds, + decide=not no_decide, + retrieval_limit=retrieval_limit, + max_requests=max_decision_requests, + ) + if export_dir: + if not result["clips"]: + raise ClipError("No clips were selected; nothing to export.") + video = repo.get_video(video_id) + receipts, elapsed_ms = export_clips( + Path(video.file_path), + result["clips"], + Path(export_dir), + overwrite=overwrite_export, + ) + result["export"] = receipts + result["timings"]["render_ms"] = round(elapsed_ms, 1) + output_json(result) + except Exception as e: # noqa: BLE001 - CLI boundary renders provider/export failures + error(str(e)) + raise typer.Exit(1) + finally: + repo.close() diff --git a/src/av/clip_eval/__init__.py b/src/av/clip_eval/__init__.py new file mode 100644 index 0000000..018d0f1 --- /dev/null +++ b/src/av/clip_eval/__init__.py @@ -0,0 +1,9 @@ +"""Public shared evaluation contract for topic clipping. + +The corpus, labels, metrics, and runner here are the initial shared contract +for comparing deterministic retrieval/windows against Jev-decided selection on +identical candidates. Everything is synthetic and rights-cleared; labels are +independent of any decision scores and are never replaced by them. +""" + +CONTRACT_VERSION = 1 diff --git a/src/av/clip_eval/__main__.py b/src/av/clip_eval/__main__.py new file mode 100644 index 0000000..4f7fb18 --- /dev/null +++ b/src/av/clip_eval/__main__.py @@ -0,0 +1,48 @@ +"""Command-line entry: ``python -m av.clip_eval --corpus ... --queries ...``.""" + +from __future__ import annotations + +import argparse +import json +import sys +from pathlib import Path + +from av.clip_eval.runner import run_evaluation + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser( + prog="python -m av.clip_eval", + description="Run the topic-clipping evaluation contract over frozen fixtures.", + ) + parser.add_argument("--corpus", type=Path, required=True, help="Path to corpus.json") + parser.add_argument("--queries", type=Path, required=True, help="Path to a labeled query set") + parser.add_argument("--db", type=Path, required=True, help="Fresh SQLite database path (outside the repo)") + parser.add_argument("--out", type=Path, required=True, help="Receipt JSON output path") + parser.add_argument("--media-dir", type=Path, default=None, help="Directory for generated synthetic media (outside the repo)") + parser.add_argument("--export-dir", type=Path, default=None, help="Directory for rendered clips (enables export-validity metrics)") + parser.add_argument("--arms", default="deterministic,jev_mock", help="Comma-separated arms: deterministic,jev_mock") + parser.add_argument("--clips", type=int, default=2, help="Clips requested per query") + parser.add_argument("--target-seconds", type=float, default=30.0) + parser.add_argument("--min-seconds", type=float, default=10.0) + args = parser.parse_args(argv) + + receipt = run_evaluation( + args.corpus, + args.queries, + db_path=args.db, + media_dir=args.media_dir, + export_dir=args.export_dir, + arms=tuple(arm.strip() for arm in args.arms.split(",") if arm.strip()), + clips_wanted=args.clips, + target_seconds=args.target_seconds, + min_seconds=args.min_seconds, + ) + args.out.parent.mkdir(parents=True, exist_ok=True) + args.out.write_text(json.dumps(receipt, indent=2, default=str) + "\n") + print(json.dumps({"receipt": str(args.out), "arms": [arm["arm"] for arm in receipt["arms"]]})) + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/src/av/clip_eval/contract.py b/src/av/clip_eval/contract.py new file mode 100644 index 0000000..cbbfc69 --- /dev/null +++ b/src/av/clip_eval/contract.py @@ -0,0 +1,181 @@ +"""Frozen metric definitions for the clip evaluation contract. + +Every metric compares returned clips against independently frozen labels. +Jev decision scores are never used as ground truth, and absent-topic queries +carry no moments at all. + +Contract version 1 metrics: +- precision_at_k hits in the first k returned clips / k +- known_moment_recall labeled moments covered (IoU >= 0.5) / total moments +- absent_topic_fp 1.0 when any clip is returned for an absent query +- boundary_error_sec mean of (|start delta| + |end delta|) / 2 over hits +- duplication_rate returned clip pairs with IoU >= 0.2 / clip count +- context_loss_count returned hits missing more than one second of a + labeled setup head +""" + +from __future__ import annotations + +import math +from typing import Any + +IOU_HIT_THRESHOLD = 0.5 +IOU_DUPLICATION_THRESHOLD = 0.2 +SETUP_SLACK_SEC = 1.0 + + +def interval_iou( + a_start: float, a_end: float, b_start: float, b_end: float +) -> float: + """Intersection-over-union of two closed intervals (0.0 when disjoint).""" + intersection = min(a_end, b_end) - max(a_start, b_start) + if intersection <= 0: + return 0.0 + union = max(a_end, b_end) - min(a_start, b_start) + if union <= 0: + return 0.0 + return intersection / union + + +def best_match( + clip: dict, moments: list[dict] +) -> tuple[dict | None, float]: + """Return the labeled moment with the highest IoU against one clip.""" + best: dict | None = None + best_iou = 0.0 + for moment in moments: + iou = interval_iou( + clip["start_sec"], clip["end_sec"], + float(moment["start_sec"]), float(moment["end_sec"]), + ) + if iou > best_iou: + best, best_iou = moment, iou + return best, best_iou + + +def query_metrics(clips: list[dict], label: dict, *, k: int | None = None) -> dict: + """Compute all contract-v1 metrics for one query's returned clips.""" + moments = label.get("moments", []) + expected_absent = label.get("expected") == "absent" + returned = clips[:k] if k else clips + + hits = 0 + boundary_errors: list[float] = [] + covered: set[str] = set() + context_losses = 0 + for clip in returned: + moment, iou = best_match(clip, moments) + if moment is None or iou < IOU_HIT_THRESHOLD: + continue + hits += 1 + covered.add(moment["moment_id"]) + boundary_errors.append( + ( + abs(clip["start_sec"] - float(moment["start_sec"])) + + abs(clip["end_sec"] - float(moment["end_sec"])) + ) + / 2.0 + ) + if moment.get("requires_setup") and moment.get("setup_start_sec") is not None: + setup_start = float(moment["setup_start_sec"]) + if clip["start_sec"] > setup_start + SETUP_SLACK_SEC: + context_losses += 1 + + duplication_pairs = 0 + for i, a in enumerate(returned): + for b in returned[i + 1 :]: + if ( + interval_iou( + a["start_sec"], a["end_sec"], b["start_sec"], b["end_sec"] + ) + >= IOU_DUPLICATION_THRESHOLD + ): + duplication_pairs += 1 + + k_denominator = k if k else max(len(returned), 1) + return { + "clips_returned": len(returned), + "hits": hits, + "precision_at_k": round(hits / k_denominator, 4) if returned else 0.0, + "known_moment_recall": ( + round(len(covered) / len(moments), 4) if moments else None + ), + "absent_topic_fp": 1.0 if expected_absent and returned else 0.0, + "boundary_error_sec": ( + round(sum(boundary_errors) / len(boundary_errors), 3) + if boundary_errors + else None + ), + "duplication_rate": ( + round(duplication_pairs / len(returned), 4) if returned else 0.0 + ), + "context_loss_count": context_losses, + } + + +def _percentile(values: list[float], q: float) -> float | None: + if len(values) < 5: + # Percentiles are only reported when the sample supports them. + return None + ordered = sorted(values) + position = (len(ordered) - 1) * q + low = math.floor(position) + high = math.ceil(position) + if low == high: + return round(ordered[low], 3) + return round(ordered[low] * (high - position) + ordered[high] * (position - low), 3) + + +def aggregate(per_query: list[dict]) -> dict: + """Aggregate per-query metric blocks into arm-level totals.""" + def collect(field: str) -> list[float]: + return [entry[field] for entry in per_query if entry.get(field) is not None] + + present_recall = [entry["known_moment_recall"] for entry in per_query if entry["known_moment_recall"] is not None] + boundary = collect("boundary_error_sec") + totals = { + "queries": len(per_query), + "precision_at_k_mean": round(sum(collect("precision_at_k")) / len(per_query), 4) if per_query else None, + "known_moment_recall_mean": round(sum(present_recall) / len(present_recall), 4) if present_recall else None, + "absent_topic_fp_total": round(sum(collect("absent_topic_fp")), 2), + "boundary_error_sec_mean": round(sum(boundary) / len(boundary), 3) if boundary else None, + "duplication_rate_mean": round(sum(collect("duplication_rate")) / len(per_query), 4) if per_query else None, + "context_loss_total": sum(collect("context_loss_count")), + } + for name, values in ( + ("precision_at_k_p50", collect("precision_at_k")), + ("boundary_error_sec_p50", boundary), + ): + totals[name] = _percentile(values, 0.5) + totals["boundary_error_sec_p95"] = _percentile(boundary, 0.95) + return totals + + +def merge_stage_usage(usages: list[dict[str, Any]]) -> dict[str, Any]: + """Sum per-stage usage dicts across queries (requests always add; token + totals collapse to unknown if any contributing stage was incomplete).""" + merged: dict[str, Any] = {} + for usage in usages: + for stage, numbers in usage.items(): + if not isinstance(numbers, dict): + continue + target = merged.setdefault( + stage, + { + "requests": 0, + "input_tokens": 0, + "output_tokens": 0, + "input_tokens_complete": True, + "output_tokens_complete": True, + }, + ) + for key in ("requests", "input_tokens", "output_tokens"): + target[key] = int(target.get(key) or 0) + int(numbers.get(key) or 0) + for key in ("input_tokens_complete", "output_tokens_complete"): + if numbers.get(key) is False: + target[key] = False + for stage in merged.values(): + for prefix in ("input", "output"): + if stage.get(f"{prefix}_tokens_complete") is False: + stage[f"{prefix}_tokens"] = None + return merged diff --git a/src/av/clip_eval/corpus.py b/src/av/clip_eval/corpus.py new file mode 100644 index 0000000..592ff3b --- /dev/null +++ b/src/av/clip_eval/corpus.py @@ -0,0 +1,170 @@ +"""Load the frozen synthetic clip-evaluation corpus into an isolated AV index. + +The corpus is text-only and inline: transcript segments and vision captions +with stable artifact IDs and checksums. Media is never stored in the repo; a +deterministic ffmpeg generation spec is included so export-validity checks can +render synthetic videos on demand outside the working tree. +""" + +from __future__ import annotations + +import hashlib +import json +import subprocess +from pathlib import Path + +from av import clip_eval +from av.core.exceptions import AVError, VideoNotFoundError +from av.db.models import ArtifactRecord, VideoRecord +from av.db.repository import Repository + +SUPPORTED_VERSIONS = {clip_eval.CONTRACT_VERSION} + + +class CorpusError(AVError): + """The corpus file is missing, malformed, or fails its frozen checksums.""" + + +def _canonical(obj: object) -> bytes: + return json.dumps(obj, sort_keys=True, separators=(",", ":")).encode() + + +def _checksum(content: dict) -> str: + return "sha256:" + hashlib.sha256(_canonical(content)).hexdigest() + + +def load_corpus(path: Path) -> dict: + """Read and verify the corpus, including per-video frozen checksums.""" + try: + corpus = json.loads(path.read_text()) + except FileNotFoundError as exc: + raise CorpusError(f"Corpus file not found: {path}") from exc + except json.JSONDecodeError as exc: + raise CorpusError(f"Corpus file is not valid JSON: {path}: {exc}") from exc + if corpus.get("contract_version") not in SUPPORTED_VERSIONS: + raise CorpusError( + f"Unsupported corpus contract_version: {corpus.get('contract_version')}" + ) + checksums = corpus.get("checksums") + if not isinstance(checksums, dict): + raise CorpusError("Corpus is missing the frozen checksums block") + for video in corpus.get("videos", []): + video_id = video.get("id") + content = video.get("content") + if not video_id or not isinstance(content, dict): + raise CorpusError(f"Corpus video {video_id!r} has no content block") + expected = checksums.get(video_id) + actual = _checksum(content) + if expected != actual: + raise CorpusError( + f"Checksum mismatch for {video_id}: expected {expected}, got {actual}. " + "The corpus was modified; this breaks the frozen evaluation contract." + ) + if not corpus.get("videos"): + raise CorpusError("Corpus contains no videos") + return corpus + + +def load_queries(path: Path) -> list[dict]: + """Read a labeled query set (dev or held-out).""" + try: + data = json.loads(path.read_text()) + except FileNotFoundError as exc: + raise CorpusError(f"Query file not found: {path}") from exc + except json.JSONDecodeError as exc: + raise CorpusError(f"Query file is not valid JSON: {path}: {exc}") from exc + queries = data.get("queries") if isinstance(data, dict) else data + if not isinstance(queries, list) or not queries: + raise CorpusError(f"Query file has no queries: {path}") + for query in queries: + for field in ("query_id", "topic", "video_id", "expected", "moments"): + if field not in query: + raise CorpusError(f"Query {query!r} is missing field {field!r}") + if query["expected"] not in {"present", "absent"}: + raise CorpusError(f"Query {query['query_id']} has invalid expected field") + return queries + + +def materialize_corpus(corpus: dict, repo: Repository) -> dict[str, str]: + """Insert every corpus video and artifact into an empty repository. + + Returns video_id -> source file path (may not exist until media is + generated; selection never reads media, only export does). + """ + paths: dict[str, str] = {} + for video in corpus["videos"]: + video_id = video["id"] + content = video["content"] + try: + repo.get_video(video_id) + except VideoNotFoundError: + pass + else: + raise CorpusError( + f"Target database already contains {video_id}; use a fresh database" + ) + media = content.get("media", {}) + source_path = media.get("filename", f"{video_id}.mp4") + repo.insert_video( + VideoRecord( + id=video_id, + file_path=source_path, + file_hash=_checksum(content), + file_size_bytes=media.get("approx_size_bytes", 0), + filename=Path(source_path).name, + duration_sec=float(content["duration_sec"]), + status="complete", + ) + ) + for segment in content.get("transcript", []): + repo.insert_artifact( + ArtifactRecord( + id=segment["id"], + video_id=video_id, + type="transcript", + start_sec=float(segment["start_sec"]), + end_sec=float(segment["end_sec"]), + text=segment["text"], + ) + ) + for caption in content.get("vision", []): + repo.insert_artifact( + ArtifactRecord( + id=caption["id"], + video_id=video_id, + type=caption.get("source_type", "caption"), + start_sec=float(caption["start_sec"]), + end_sec=float(caption["end_sec"]), + text=caption["text"], + ) + ) + paths[video_id] = source_path + return paths + + +def generate_media(spec: dict, out_dir: Path) -> Path: + """Render one synthetic corpus video with ffmpeg (color source + tone). + + The generation command is recorded in the corpus so any ffmpeg build can + reproduce media locally; media itself is never committed. + """ + out_dir.mkdir(parents=True, exist_ok=True) + out_path = out_dir / spec["filename"] + if out_path.exists(): + return out_path + cmd = [ + "ffmpeg", "-hide_banner", "-loglevel", "error", + "-f", "lavfi", "-i", f"testsrc2=duration={spec['duration_sec']}:size=320x240:rate=15", + "-f", "lavfi", "-i", f"sine=frequency={spec.get('tone_hz', 440)}:duration={spec['duration_sec']}", + "-c:v", "libx264", "-preset", "veryfast", "-crf", "30", + "-c:a", "aac", "-b:a", "64k", + "-movflags", "+faststart", + "-y", str(out_path), + ] + try: + subprocess.run(cmd, capture_output=True, text=True, check=True, timeout=600) + except FileNotFoundError as exc: + raise AVError("ffmpeg not found; cannot generate synthetic eval media") from exc + except subprocess.CalledProcessError as exc: + raise AVError(f"ffmpeg failed to generate eval media: {exc.stderr}") from exc + return out_path diff --git a/src/av/clip_eval/mock.py b/src/av/clip_eval/mock.py new file mode 100644 index 0000000..509ad3c --- /dev/null +++ b/src/av/clip_eval/mock.py @@ -0,0 +1,150 @@ +"""Deterministic labeled stand-in for the System One decision endpoint. + +Answers Noul, Choice, and Score questions purely from the frozen corpus +labels with string and interval arithmetic. It exists to exercise the typed +decision pipeline and the evaluation harness offline. It is NOT a measurement +of Jev quality and must never be presented as one. +""" + +from __future__ import annotations + +from typing import Any + +_RELEVANT_P = 0.9 +_IRRELEVANT_P = 0.05 +_COHERENT_P = 0.9 +_INCOHERENT_P = 0.4 +_CONFIDENT = 0.85 + + +def _interval_overlaps(a_start: float, a_end: float, b_start: float, b_end: float) -> bool: + return a_start < b_end and a_end > b_start + + +class LabeledDecisionClient: + """Answers typed questions by looking up the labeled moments a clip hits.""" + + def __init__(self, queries: list[dict]) -> None: + self._by_topic: dict[tuple[str, str], dict] = {} + for query in queries: + key = (query["video_id"], query["topic"].casefold()) + self._by_topic[key] = query + + def _lookup(self, state: dict[str, Any]) -> dict | None: + topic = str(state.get("query", "")).casefold() + clips = state.get("clips") or state.get("surrounding_events") or {} + for value in clips.values(): + video_id = value.get("video_id") if isinstance(value, dict) else None + if video_id: + return self._by_topic.get((video_id, topic)) + # Boundary state carries no video_id; fall back to a unique match. + matches = [ + query + for (video_id, topic_key), query in self._by_topic.items() + if topic_key == topic + ] + return matches[0] if len(matches) == 1 else None + + def _moments_hit(self, label: dict | None, start: float, end: float) -> list[dict]: + if not label: + return [] + return [ + moment + for moment in label.get("moments", []) + if _interval_overlaps(start, end, float(moment["start_sec"]), float(moment["end_sec"])) + ] + + @staticmethod + def _clip_interval(value: dict) -> tuple[float, float]: + start = value.get("start_sec") + end = value.get("end_sec") + if start is None or end is None: + return (-1.0, -1.0) + return (float(start), float(end)) + + def ask(self, state: dict[str, Any], questions: dict[str, dict]) -> tuple[dict, dict]: + label = self._lookup(state) + answers: dict[str, dict] = {} + for key, question in questions.items(): + qtype = question.get("type") + clips = state.get("clips") or {} + value = clips.get(key, {}) + start, end = self._clip_interval(value) + if qtype == "noul": + instructions = str(question.get("instructions", "")) + hits = self._moments_hit(label, start, end) + if "standalone" in instructions or "on its own" in instructions: + # Coherence: only a hard-context moment whose setup head + # is missing reads as incoherent. + p = _COHERENT_P + for moment in hits: + setup_end = moment.get("setup_end_sec") + if ( + moment.get("hard_context") + and setup_end is not None + and start >= float(setup_end) - 0.5 + ): + # The clip is the payoff without its setup head. + p = _INCOHERENT_P + answers[key] = {"type": "noul", "noul": p} + elif "visual description" in instructions: + has_caption = bool(value.get("caption") and value["caption"] != "(none)") + answers[key] = { + "type": "noul", + "noul": _RELEVANT_P if (hits and has_caption) else _IRRELEVANT_P, + } + else: + answers[key] = { + "type": "noul", + "noul": _RELEVANT_P if hits else _IRRELEVANT_P, + } + elif qtype == "choice": + answers[key] = self._choice_answer(label, state, key, question) + elif qtype == "score": + hits = self._moments_hit(label, start, end) + score = hits[0].get("appeal", 0.3) if hits else 0.3 + answers[key] = {"type": "score", "score": float(score)} + else: + answers[key] = {"type": "error", "message": f"unsupported type {qtype}"} + return answers, {"requests": 1, "input_tokens": 0, "output_tokens": 0} + + def _choice_answer( + self, + label: dict | None, + state: dict[str, Any], + key: str, + question: dict, + ) -> dict: + candidates = question.get("criteria", {}) + surrounding = state.get("surrounding_events", {}) + best_label = "e0" + best_distance = float("inf") + moments = label.get("moments", []) if label else [] + target = 0.0 + if moments: + if key == "start": + target = min(float(m["start_sec"]) for m in moments) + else: + target = max(float(m["end_sec"]) for m in moments) + for candidate_label, criterion in candidates.items(): + seconds = criterion.get("seconds") if isinstance(criterion, dict) else None + if not seconds and candidate_label in surrounding: + # The e0 criterion describes the hit event without seconds; + # its interval is still available in the surrounding state. + event = surrounding[candidate_label] + seconds = f"{event.get('start')}-{event.get('end')}" + if not seconds: + continue + start_s, _, end_s = str(seconds).partition("-") + try: + c_start, c_end = float(start_s), float(end_s) + except ValueError: + continue + distance = abs(c_start - target) + abs(c_end - target) + if distance < best_distance: + best_distance = distance + best_label = candidate_label + if not moments: + # Absent topics: any choice is arbitrary; keep the minimal answer. + best_label = "e0" if "e0" in candidates else next(iter(candidates), "e0") + return {"type": "choice", "choice": best_label, "confidence": _CONFIDENT} diff --git a/src/av/clip_eval/runner.py b/src/av/clip_eval/runner.py new file mode 100644 index 0000000..ace881e --- /dev/null +++ b/src/av/clip_eval/runner.py @@ -0,0 +1,204 @@ +"""Run the deterministic-vs-Jev clip comparison on identical candidates. + +Both arms call the same ``clip_video`` engine over the same frozen corpus and +index. The deterministic arm ranks retrieval candidates without any provider; +the Jev arm runs the typed decision stages against a supplied client (the +labeled mock offline, or a real System One client under an explicit allowance +and request ceiling). Metrics come only from the frozen labels. +""" + +from __future__ import annotations + +import time +from dataclasses import dataclass +from pathlib import Path +from typing import Any + +from av import clip_eval +from av.clip_eval import contract +from av.clip_eval import corpus as corpus_mod +from av.clip_eval.mock import LabeledDecisionClient +from av.core.config import AVConfig +from av.db.repository import Repository +from av.search.clip import clip_video + +_SUCCESS_STATUSES = {"ok", "no_usable_clips", "deterministic_only"} + + +@dataclass +class ArmConfig: + name: str + decide: bool + client: Any | None + + +def _run_arm( + arm: ArmConfig, + corpus: dict, + queries: list[dict], + repo: Repository, + config: AVConfig, + *, + clips_wanted: int, + target_seconds: float, + min_seconds: float, + export_dir: Path | None, + media_dir: Path | None, +) -> dict: + per_query: list[dict] = [] + stage_usages: list[dict] = [] + warnings: list[str] = [] + export_receipts: list[dict] = [] + render_ms = 0.0 + for label in queries: + result = clip_video( + label["topic"], + label["video_id"], + repo, + config, + clips_wanted=clips_wanted, + target_seconds=target_seconds, + min_seconds=min_seconds, + decide=arm.decide, + client=arm.client, + ) + timings = result.get("timings", {}) + export_for_query: list[dict] = [] + if export_dir and result["clips"]: + from av.pipeline.clip_export import export_clips + + video = repo.get_video(label["video_id"]) + source = Path(video.file_path) + if not source.exists() and media_dir is not None: + spec = _media_spec(corpus, label["video_id"]) + source = corpus_mod.generate_media(spec, media_dir) + export_for_query, elapsed_ms = export_clips( + source, + result["clips"], + export_dir / arm.name / label["query_id"], + overwrite=True, + ) + render_ms += elapsed_ms + if export_dir and not result["clips"]: + export_for_query, elapsed_ms = [], 0.0 + export_receipts.extend(export_for_query) + metrics = contract.query_metrics(result["clips"], label, k=clips_wanted) + metrics["query_id"] = label["query_id"] + metrics["topic"] = label["topic"] + metrics["status"] = result["status"] + metrics["timings"] = timings + if result["status"] not in _SUCCESS_STATUSES: + warnings.append(f"{label['query_id']}: status {result['status']}") + per_query.append(metrics) + stage_usages.append(result.get("stage_usage", {})) + + valid_exports = [receipt for receipt in export_receipts if receipt.get("valid")] + usage = contract.merge_stage_usage(stage_usages) + usage_block = usage + timings = { + "prepare_ms_total": round(sum(q["timings"].get("prepare_ms", 0.0) for q in per_query), 1), + "selection_ms_total": round(sum(q["timings"].get("selection_ms", 0.0) for q in per_query), 1), + "selection_warm_ms_total": round( + sum(q["timings"].get("selection_warm_ms", 0.0) for q in per_query), 1 + ), + "render_ms_total": round(render_ms, 1), + } + return { + "arm": arm.name, + "decide": arm.decide, + "metrics_per_query": per_query, + "metrics_total": contract.aggregate(per_query), + "stage_usage": usage_block, + "timings": timings, + "export": { + "requested": export_dir is not None, + "clips_exported": len(export_receipts), + "clips_valid": len(valid_exports), + }, + "failures": warnings, + "decision_provider": "typesafe" if arm.decide else "none", + } + + +def _media_spec(corpus: dict, video_id: str) -> dict: + for video in corpus["videos"]: + if video["id"] == video_id: + media = dict(video["content"].get("media", {})) + media.setdefault("filename", f"{video_id}.mp4") + media.setdefault("duration_sec", video["content"]["duration_sec"]) + return media + raise KeyError(video_id) + + +def run_evaluation( + corpus_path: Path, + queries_path: Path, + *, + db_path: Path, + media_dir: Path | None = None, + export_dir: Path | None = None, + arms: tuple[str, ...] = ("deterministic", "jev_mock"), + clips_wanted: int = 2, + target_seconds: float = 30.0, + min_seconds: float = 10.0, + config: AVConfig | None = None, +) -> dict: + """Execute the evaluation and return the full contract-v1 receipt.""" + corpus = corpus_mod.load_corpus(corpus_path) + queries = corpus_mod.load_queries(queries_path) + config = config or AVConfig() + + ingest_start = time.perf_counter() + repo = Repository(db_path) + corpus_mod.materialize_corpus(corpus, repo) + ingestion_ms = (time.perf_counter() - ingest_start) * 1000 + + by_video = {query["video_id"] for query in queries} + missing = by_video - {video["id"] for video in corpus["videos"]} + if missing: + repo.close() + raise corpus_mod.CorpusError(f"Queries reference unknown videos: {sorted(missing)}") + + arm_configs: list[ArmConfig] = [] + for name in arms: + if name == "deterministic": + arm_configs.append(ArmConfig(name=name, decide=False, client=None)) + elif name == "jev_mock": + client = LabeledDecisionClient(queries) + arm_configs.append(ArmConfig(name=name, decide=True, client=client)) + else: + repo.close() + raise ValueError( + f"Unknown arm {name!r}; supported arms: deterministic, jev_mock" + ) + + arm_receipts = [ + _run_arm( + arm, + corpus, + queries, + repo, + config, + clips_wanted=clips_wanted, + target_seconds=target_seconds, + min_seconds=min_seconds, + export_dir=export_dir, + media_dir=media_dir, + ) + for arm in arm_configs + ] + repo.close() + + return { + "contract_version": clip_eval.CONTRACT_VERSION, + "corpus": str(corpus_path), + "queries": str(queries_path), + "query_count": len(queries), + "ingestion_ms": round(ingestion_ms, 1), + "arms": arm_receipts, + "notes": [ + "Synthetic rights-cleared corpus; labels are frozen and independent of decision scores.", + "The jev_mock arm exercises the typed decision pipeline; it is not a Jev quality measurement.", + "Percentile timings are reported only when at least five samples support them.", + ], + } diff --git a/src/av/core/config.py b/src/av/core/config.py index 11f4895..491b8d1 100644 --- a/src/av/core/config.py +++ b/src/av/core/config.py @@ -28,7 +28,7 @@ def _load_config_file() -> dict: return {} try: return json.loads(CONFIG_FILE_PATH.read_text()) - except Exception: + except (OSError, ValueError): return {} @@ -84,6 +84,15 @@ class AVConfig(BaseSettings): refine_batch_size: int = Field(default=10, ge=1, le=50) refine_context_events: int = Field(default=3, ge=0, le=12) + # Topic clipping decisions. Relevance/coherence/visual gates are objective + # source support; appeal only orders. The request ceiling is an explicit + # per-run cap on System One HTTP attempts. + clip_relevance_min: float = Field(default=0.6, ge=0, le=1) + clip_coherence_min: float = Field(default=0.5, ge=0, le=1) + clip_visual_min: float = Field(default=0.3, ge=0, le=1) + clip_boundary_conf_min: float = Field(default=0.5, ge=0, le=1) + clip_request_cap: int = Field(default=60, ge=1, le=10000) + # Optional stronger sampled-frame inspection after an unsupported answer. strong_vision_api_base_url: str = Field(default="") strong_vision_api_key: str = Field(default="") diff --git a/src/av/pipeline/clip_export.py b/src/av/pipeline/clip_export.py new file mode 100644 index 0000000..7cc8e72 --- /dev/null +++ b/src/av/pipeline/clip_export.py @@ -0,0 +1,162 @@ +"""Optional FFmpeg clip rendering and playability validation. + +Rendering is deliberately separate from provider selection: every second of +provider time and every token is accounted in the decision stages, and +``render_ms`` only ever measures local ffmpeg work. +""" + +from __future__ import annotations + +import json +import subprocess +import time +from dataclasses import dataclass +from pathlib import Path + +from av.core.exceptions import FFmpegError + +# Re-encode both streams so cuts land exactly on the requested seconds and +# audio stays synchronized; stream copy can only cut on keyframes. +_VIDEO_ARGS = ["-c:v", "libx264", "-preset", "veryfast", "-crf", "23"] +_AUDIO_ARGS = ["-c:a", "aac", "-b:a", "128k"] +_DURATION_TOLERANCE_SEC = 1.0 +_START_TIME_TOLERANCE_SEC = 0.5 + + +@dataclass +class ExportRequest: + clip_rank: int + start_sec: float + end_sec: float + + +def _fmt_name(secs: float) -> str: + h = int(secs // 3600) + m = int((secs % 3600) // 60) + s = int(secs % 60) + return f"{h:02d}{m:02d}{s:02d}" + + +def _probe(path: Path) -> dict: + cmd = [ + "ffprobe", "-v", "error", + "-print_format", "json", + "-show_format", "-show_streams", + str(path), + ] + try: + result = subprocess.run(cmd, capture_output=True, text=True, check=True, timeout=60) + except FileNotFoundError as exc: + raise FFmpegError("ffprobe not found. Install ffmpeg: brew install ffmpeg / apt install ffmpeg") from exc + except subprocess.CalledProcessError as exc: + raise FFmpegError(f"ffprobe failed: {exc.stderr}", cmd=" ".join(cmd)) from exc + return json.loads(result.stdout) + + +def _stream_summary(probe: dict) -> dict: + video = next((s for s in probe.get("streams", []) if s.get("codec_type") == "video"), None) + audio = next((s for s in probe.get("streams", []) if s.get("codec_type") == "audio"), None) + return { + "video_codec": video.get("codec_name") if video else None, + "audio_codec": audio.get("codec_name") if audio else None, + "start_time_sec": float(probe.get("format", {}).get("start_time", 0) or 0), + } + + +def export_clips( + source: Path, + clip_records: list[dict], + out_dir: Path, + *, + overwrite: bool = False, +) -> tuple[list[dict], float]: + """Render each clip record to ``out_dir`` and validate the result. + + Returns (receipts, elapsed_seconds). Receipts carry the measured output + duration, stream summary, and a ``valid`` flag with per-check warnings. + """ + source = source.resolve() + if not source.exists(): + raise FFmpegError(f"Source video not found: {source}") + out_dir = out_dir.expanduser().resolve() + out_dir.mkdir(parents=True, exist_ok=True) + if not out_dir.is_dir(): + raise FFmpegError(f"Export target is not a directory: {out_dir}") + + source_audio = _probe(source) + source_has_audio = any( + s.get("codec_type") == "audio" for s in source_audio.get("streams", []) + ) + + receipts: list[dict] = [] + started = time.perf_counter() + for clip in clip_records: + rank = clip["rank"] + start = float(clip["start_sec"]) + end = float(clip["end_sec"]) + duration = max(end - start, 0.0) + out_path = out_dir / f"clip_{rank:02d}_{_fmt_name(start)}-{_fmt_name(end)}.mp4" + if out_path.exists() and not overwrite: + raise FFmpegError( + f"Export file already exists: {out_path}. Pass --overwrite-export to replace it." + ) + resolved = out_path.resolve() + if resolved.parent != out_dir: + raise FFmpegError(f"Refusing export outside the requested directory: {resolved}") + # Create the output file before ffmpeg runs so a laggy NFS mount can + # never fail ffmpeg's open() with ENOENT right after mkdir(). + resolved.touch(exist_ok=True) + + cmd = [ + "ffmpeg", "-hide_banner", "-loglevel", "error", + "-ss", f"{start:.3f}", + "-i", str(source), + "-t", f"{duration:.3f}", + *_VIDEO_ARGS, + *_AUDIO_ARGS, + "-movflags", "+faststart", + "-y", + str(resolved), + ] + try: + subprocess.run(cmd, capture_output=True, text=True, check=True, timeout=600) + except FileNotFoundError as exc: + raise FFmpegError("ffmpeg not found. Install ffmpeg: brew install ffmpeg / apt install ffmpeg") from exc + except subprocess.CalledProcessError as exc: + raise FFmpegError( + f"ffmpeg failed for clip {rank}: {exc.stderr}", + cmd=" ".join(cmd), + ) from exc + + probe = _probe(resolved) + measured = float(probe.get("format", {}).get("duration", 0) or 0) + streams = _stream_summary(probe) + warnings: list[str] = [] + if abs(measured - duration) > _DURATION_TOLERANCE_SEC: + warnings.append( + f"Output duration {measured:.3f}s differs from requested {duration:.3f}s" + ) + if not streams["video_codec"]: + warnings.append("Output has no video stream") + if source_has_audio and not streams["audio_codec"]: + warnings.append("Source has audio but output does not; audio was lost") + if abs(streams["start_time_sec"]) > _START_TIME_TOLERANCE_SEC: + warnings.append( + f"Output start_time {streams['start_time_sec']:.3f}s is not near zero" + ) + receipts.append({ + "clip_rank": rank, + "clip_id": clip.get("clip_id"), + "path": str(resolved), + "requested_start_sec": round(start, 3), + "requested_end_sec": round(end, 3), + "requested_duration_sec": round(duration, 3), + "measured_duration_sec": round(measured, 3), + "streams": streams, + "valid": not warnings, + "warnings": warnings, + "command": "ffmpeg -ss -i -t " + + " ".join(_VIDEO_ARGS + _AUDIO_ARGS + ["-movflags", "+faststart"]), + }) + elapsed_ms = (time.perf_counter() - started) * 1000 + return receipts, elapsed_ms diff --git a/src/av/search/clip.py b/src/av/search/clip.py new file mode 100644 index 0000000..e456e8a --- /dev/null +++ b/src/av/search/clip.py @@ -0,0 +1,1126 @@ +"""Topic clipping: retrieval-grounded candidates with typed Jev decisions. + +AV owns candidate construction, timing, deduplication, and export. Jev +(System One) owns typed relevance, coherence, visual-evidence, boundary, and +appeal decisions, expressed as Noul/Choice/Score operations over +source-verbatim text. No generative completion is presented as a decision and +no model is asked to invent timestamps: every returned second is an artifact +boundary already present in the database. + +Objective source support (relevance, coherence, visual evidence) is decided +and gated separately from subjective highlight appeal, which only orders +candidates that already passed the objective gates. +""" + +from __future__ import annotations + +import hashlib +import json +import math +import time +from dataclasses import dataclass, replace +from typing import Any + +from av.core.config import AVConfig +from av.db.models import VideoRecord +from av.db.repository import Repository, _fmt_timestamp +from av.search.query import natural_language_fts_query +from av.search.refine import ( + RefinementError, + SystemOneClient, + _judge_bounds, + _ordered_contiguous_events, + _probability, +) +from av.search.usage import new_usage, record_usage + +DEFAULT_RETRIEVAL_LIMIT = 24 +DEFAULT_MAX_CANDIDATES = 24 +MAX_QUOTES_PER_MODALITY = 12 +# Candidates are built generously around the hit (this multiple of the +# duration cap) so the boundary judge can choose tight setup/payoff bounds +# inside them; assembly still enforces the hard cap at event boundaries. +CLIP_BOUND_SLACK = 1.5 +_VISION_TYPES = frozenset({"caption", "dense_caption", "scene"}) +_STAGE_NAMES = ("relevance", "coherence", "visual", "boundary", "appeal") + + +class ClipError(RuntimeError): + """A clipping request that cannot be served as specified.""" + + +class DecisionBudgetExhausted(ClipError): + """The explicit per-run decision request ceiling was reached.""" + + def __init__(self, message: str, *, stage_usage: dict[str, dict] | None = None) -> None: + super().__init__(message) + self.stage_usage = stage_usage or {} + + +class BudgetedSystemOne: + """SystemOneClient wrapper enforcing an explicit per-run request ceiling. + + Requests are counted as actual HTTP attempts, retries included. Cached + replays (warm selection) perform no provider call and add no request. + """ + + def __init__(self, inner: SystemOneClient, *, max_requests: int) -> None: + if max_requests < 1: + raise ClipError("Decision request ceiling must be at least 1") + self.inner = inner + self.max_requests = max_requests + self.attempts = 0 + self.cache_hits = 0 + self.usage = new_usage() + self._cache: dict[str, tuple[dict, dict]] = {} + + def ask(self, state: Any, questions: dict[str, dict]) -> tuple[dict, dict]: + key = json.dumps( + {"state": state, "questions": questions}, sort_keys=True, default=str + ) + cached = self._cache.get(key) + if cached is not None: + self.cache_hits += 1 + return cached + if self.attempts >= self.max_requests: + raise DecisionBudgetExhausted( + f"Decision request ceiling reached ({self.max_requests}); " + f"{len(questions)} question(s) left undecided" + ) + try: + answers, usage = self.inner.ask(state, questions) + except RefinementError as exc: + attempts = max(exc.attempts, 1) + self.attempts += attempts + record_usage( + self.usage, + exc.raw_usage, + requests=attempts, + ambiguous_attempts=attempts > 1 or exc.raw_usage is None, + ) + raise + attempts = 1 + cleaned: dict[str, Any] | None = usage + if isinstance(usage, dict): + raw_attempts = usage.get("_attempts", 1) + attempts = raw_attempts if isinstance(raw_attempts, int) and raw_attempts > 0 else 1 + cleaned = {key_: value for key_, value in usage.items() if key_ != "_attempts"} + self.attempts += attempts + record_usage(self.usage, cleaned, requests=attempts, ambiguous_attempts=attempts > 1) + self._cache[key] = (answers, usage) + return answers, usage + + +@dataclass +class ClipCandidate: + """One contiguous candidate moment, built only from artifact boundaries.""" + + candidate_id: str + video_id: str + filename: str + start_sec: float + end_sec: float + retrieval_score: float + hit_artifact_ids: list[str] + artifact_ids: list[str] + transcript_ids: list[str] + vision_ids: list[str] + vision_rows: list[dict] + transcript_text: str + vision_text: str + events: list[dict] + hit_event_index: int + relevance_p: float | None = None + coherence_p: float | None = None + visual_support_p: float | None = None + boundary_start_sec: float | None = None + boundary_end_sec: float | None = None + boundary_confidence: float | None = None + appeal_score: float | None = None + + +def _candidate_id(video_id: str, start: float, end: float) -> str: + digest = hashlib.sha256(f"{video_id}|{start:.3f}|{end:.3f}".encode()).hexdigest() + return f"clip-{digest[:12]}" + + +def _transcript_rows(events: list[dict]) -> list[dict]: + transcripts: list[dict] = [] + for event in events: + for row in event.get("rows", ()): + if row["source_type"] == "transcript": + transcripts.append(row) + return transcripts + + +def _bounded_run( + events: list[dict], + hit_index: int, + max_seconds: float, +) -> tuple[list[dict], int]: + """Trim a contiguous event run around its hit event. + + Transcript segments usually chain back-to-back, so the raw contiguous run + can span minutes. Candidates expand from the hit event in alternating + setup-first order (previous event, then next) while the span stays within + ``max_seconds``. Every accepted boundary is an existing event boundary. + """ + lo = hi = hit_index + + def fits(new_lo: int, new_hi: int) -> bool: + start = min(events[lo]["start"], events[new_lo]["start"]) + end = max(events[hi]["end"], events[new_hi]["end"]) + return end - start <= max_seconds + + setup_first = True + while True: + placed = False + sides = (True, False) if setup_first else (False, True) + for side in sides: + if side and lo > 0 and fits(lo - 1, hi): + lo -= 1 + placed = True + break + if not side and hi < len(events) - 1 and fits(lo, hi + 1): + hi += 1 + placed = True + break + if not placed: + return events[lo : hi + 1], hit_index - lo + setup_first = not setup_first + + +def build_candidates( + topic: str, + repo: Repository, + video: VideoRecord, + *, + retrieval_limit: int = DEFAULT_RETRIEVAL_LIMIT, + context_events: int = 3, + min_seconds: float = 0.0, + max_seconds: float | None = None, + max_candidates: int = DEFAULT_MAX_CANDIDATES, +) -> tuple[list[ClipCandidate], dict]: + """Deterministic candidates: FTS hits expanded into temporal neighborhoods + bounded by ``max_seconds`` around each hit. No provider calls and no + invented timestamps.""" + meta: dict[str, Any] = { + "retrieval_hits": 0, + "groups": 0, + "too_short_groups": 0, + "candidate_count": 0, + "retrieval_limit": retrieval_limit, + } + query = natural_language_fts_query(topic) + if not query: + return [], meta + hits = repo.search_fts(query, limit=retrieval_limit, video_id=video.id) + meta["retrieval_hits"] = len(hits) + if not hits: + return [], meta + + # Group hits sharing one contiguous temporal neighborhood so one moment + # yields one candidate, not one candidate per retrieval hit. + groups: list[dict] = [] + for hit in sorted(hits, key=lambda h: (h.timestamp_sec, h.artifact_id or "")): + result = { + "timestamp_sec": hit.timestamp_sec, + "end_sec": hit.end_sec, + "artifact_id": hit.artifact_id, + "source_type": hit.source_type, + "text": hit.text, + } + window = repo.get_refinement_window( + video.id, hit.timestamp_sec, before=context_events, after=context_events + ) + # Timing anchors on transcript boundaries when they exist; wide vision + # captions never drive clip bounds. Captions remain attached evidence + # and take over as the timing basis only without any transcript. + transcript_artifacts = [ + artifact for artifact in window if artifact.type == "transcript" + ] + events: list[dict] = [] + hit_index = -1 + if transcript_artifacts: + events, hit_index = _ordered_contiguous_events( + result, transcript_artifacts, append_missing_hit=False + ) + if not events or hit_index < 0: + events, hit_index = _ordered_contiguous_events(result, window) + if max_seconds is not None: + events, hit_index = _bounded_run( + events, hit_index, max_seconds * CLIP_BOUND_SLACK + ) + hit_event = events[hit_index] + for group in groups: + group_start = min(event["start"] for event in group["events"]) + group_end = max(event["end"] for event in group["events"]) + if hit_event["start"] < group_end and hit_event["end"] > group_start: + merged: dict[tuple[float, float, str], dict] = {} + for event in group["events"] + events: + merged[(event["start"], event["end"], event["text"])] = event + group["events"] = sorted( + merged.values(), key=lambda e: (e["start"], e["end"]) + ) + if hit.artifact_id: + group["hit_ids"].add(hit.artifact_id) + group["best_score"] = max(group["best_score"], hit.score) + if max_seconds is not None: + # Re-bound the union around its earliest hit event so + # merged groups keep the bounded-candidate invariant. + anchor = next( + ( + index + for index, event in enumerate(group["events"]) + if group["hit_ids"] & set(event["artifact_ids"]) + ), + 0, + ) + group["events"], _ = _bounded_run( + group["events"], anchor, max_seconds * CLIP_BOUND_SLACK + ) + break + else: + hit_ids = {hit.artifact_id} if hit.artifact_id else set() + groups.append({ + "events": list(events), + "hit_ids": hit_ids, + "best_score": hit.score, + "first_hit_sec": hit.timestamp_sec, + }) + meta["groups"] = len(groups) + + candidates: list[ClipCandidate] = [] + for group in groups: + events = group["events"] + start = min(event["start"] for event in events) + end = max(event["end"] for event in events) + if end - start < min_seconds: + meta["too_short_groups"] += 1 + continue + transcripts = _transcript_rows(events) + artifact_ids: list[str] = [] + for event in events: + artifact_ids.extend(event["artifact_ids"]) + # Attach overlapping vision rows as evidence; they never move bounds. + vision_rows: list[dict] = [] + for artifact in repo.get_artifacts_overlapping(video.id, start, end): + if artifact.type in _VISION_TYPES: + vision_rows.append({ + "artifact_id": artifact.id, + "source_type": artifact.type, + "start": artifact.start_sec, + "end": artifact.end_sec if artifact.end_sec is not None else artifact.start_sec, + "text": artifact.text, + }) + vision_rows.sort(key=lambda row: (row["start"], row["artifact_id"])) + hit_event_index = next( + ( + index + for index, event in enumerate(events) + if group["hit_ids"] & set(event["artifact_ids"]) + ), + 0, + ) + candidates.append( + ClipCandidate( + candidate_id=_candidate_id(video.id, start, end), + video_id=video.id, + filename=video.filename, + start_sec=start, + end_sec=end, + retrieval_score=group["best_score"], + hit_artifact_ids=sorted(group["hit_ids"]), + artifact_ids=list(dict.fromkeys(artifact_ids)), + transcript_ids=[row["artifact_id"] for row in transcripts], + vision_ids=[row["artifact_id"] for row in vision_rows], + vision_rows=vision_rows, + transcript_text="\n".join(row["text"] for row in transcripts), + vision_text="\n".join(row["text"] for row in vision_rows), + events=events, + hit_event_index=hit_event_index, + ) + ) + candidates.sort(key=lambda c: (-c.retrieval_score, c.start_sec, c.candidate_id)) + candidates = candidates[: max_candidates] + meta["candidate_count"] = len(candidates) + return candidates, meta + + +def _clip_payload(candidate: ClipCandidate, key: str) -> dict: + return { + "transcript": candidate.transcript_text or "(none)", + "caption": candidate.vision_text or "(none)", + "start_sec": candidate.start_sec, + "end_sec": candidate.end_sec, + } + + +def _ask_noul_batch( + client: BudgetedSystemOne, + topic: str, + candidates: list[ClipCandidate], + *, + instructions: str, + true_criteria: str, + false_criteria: str, + usage: dict, +) -> dict[str, float]: + """Batched Noul questions; returns candidate_id -> probability.""" + probabilities: dict[str, float] = {} + state = { + "query": topic, + "clips": { + f"c{index}": _clip_payload(candidate, f"c{index}") + for index, candidate in enumerate(candidates) + }, + } + questions = { + f"c{index}": { + "type": "noul", + "instructions": instructions.format(key=f"c{index}"), + "criteria": {"true": true_criteria, "false": false_criteria}, + } + for index in range(len(candidates)) + } + try: + answers, call_usage = client.ask(state, questions) + except RefinementError as exc: + # The wrapper already counted these attempts in its aggregate; the + # affected stage must also report them before the error propagates. + attempts = max(exc.attempts, 1) + record_usage( + usage, + exc.raw_usage, + requests=attempts, + ambiguous_attempts=attempts > 1 or exc.raw_usage is None, + ) + raise + record_usage(usage, {k: v for k, v in (call_usage or {}).items() if k != "_attempts"}) + for index, candidate in enumerate(candidates): + key = f"c{index}" + answer = answers.get(key) + if not isinstance(answer, dict) or answer.get("type") != "noul": + raise RefinementError(f"System One returned an invalid Noul answer for {key}") + probabilities[candidate.candidate_id] = _probability(answer.get("noul"), key) + return probabilities + + +_RELEVANCE_INSTRUCTIONS = ( + "Does `clips.{key}` contain spoken or visual source evidence that materially " + "concerns the subject or event described by `query`?" +) +_RELEVANCE_TRUE = "The clip content materially concerns the requested subject or event." +_RELEVANCE_FALSE = ( + "The clip is unrelated, or mentions the subject only incidentally." +) +_COHERENCE_INSTRUCTIONS = ( + "Could `clips.{key}` be understood on its own, without surrounding video " + "context, as a complete spoken or visual moment about `query`?" +) +_COHERENCE_TRUE = ( + "The clip carries its own setup and payoff and does not depend on missing " + "context before or after it." +) +_COHERENCE_FALSE = ( + "The clip references unavailable context, starts or ends mid-thought, or " + "would confuse a viewer who saw only this clip." +) +_VISUAL_INSTRUCTIONS = ( + "Does the visual description in `clips.{key}` show visual evidence about the " + "subject or event described by `query`?" +) +_VISUAL_TRUE = "The described visuals materially show the requested subject or event." +_VISUAL_FALSE = ( + "The described visuals are unrelated, absent, or contradict the requested " + "subject or event." +) +_APPEAL_INSTRUCTIONS = ( + "Rate how compelling `clips.{key}` is as a standalone highlight for a viewer " + "interested in `query`, from 0.0 (not compelling) to 1.0 (exceptional)." +) +_APPEAL_CRITERIA = { + "1.0": "An exceptional, self-contained highlight for this topic.", + "0.5": "A usable but ordinary moment for this topic.", + "0.0": "Not compelling as a highlight for this topic.", +} + + +def _ask_appeal( + client: BudgetedSystemOne, + topic: str, + candidates: list[ClipCandidate], + usage: dict, +) -> dict[str, float]: + state = { + "query": topic, + "clips": { + f"c{index}": _clip_payload(candidate, f"c{index}") + for index, candidate in enumerate(candidates) + }, + } + questions = { + f"c{index}": { + "type": "score", + "instructions": _APPEAL_INSTRUCTIONS.format(key=f"c{index}"), + "criteria": _APPEAL_CRITERIA, + } + for index in range(len(candidates)) + } + try: + answers, call_usage = client.ask(state, questions) + except RefinementError as exc: + attempts = max(exc.attempts, 1) + record_usage( + usage, + exc.raw_usage, + requests=attempts, + ambiguous_attempts=attempts > 1 or exc.raw_usage is None, + ) + raise + record_usage(usage, {k: v for k, v in (call_usage or {}).items() if k != "_attempts"}) + scores: dict[str, float] = {} + for index, candidate in enumerate(candidates): + key = f"c{index}" + answer = answers.get(key) + if not isinstance(answer, dict) or answer.get("type") != "score": + raise RefinementError(f"System One returned an invalid Score answer for {key}") + value = answer.get("score") + if isinstance(value, bool) or not isinstance(value, (int, float)): + raise RefinementError(f"System One returned an invalid score for {key}") + value = float(value) + if not math.isfinite(value) or not 0 <= value <= 1: + raise RefinementError(f"System One returned an out-of-range score for {key}") + scores[candidate.candidate_id] = value + return scores + + +def decide_candidates( + candidates: list[ClipCandidate], + topic: str, + client: BudgetedSystemOne, + *, + batch_size: int = 10, + context_events: int = 3, +) -> tuple[dict[str, dict], dict[str, dict], dict]: + """Run typed decision stages over the candidates in place. + + ``client`` is the budgeted wrapper so warm replays hit the same decision + cache and make zero provider calls. Returns (per-candidate decision + snapshots, per-stage usage, budget info). Decisions mutate the candidates. + """ + stage_usage = {name: new_usage() for name in _STAGE_NAMES} + + def batches(items: list[ClipCandidate]) -> list[list[ClipCandidate]]: + return [items[offset : offset + batch_size] for offset in range(0, len(items), batch_size)] + + try: + for batch in batches(candidates): + probabilities = _ask_noul_batch( + client, + topic, + batch, + instructions=_RELEVANCE_INSTRUCTIONS, + true_criteria=_RELEVANCE_TRUE, + false_criteria=_RELEVANCE_FALSE, + usage=stage_usage["relevance"], + ) + for candidate in batch: + candidate.relevance_p = probabilities[candidate.candidate_id] + + for batch in batches(candidates): + probabilities = _ask_noul_batch( + client, + topic, + batch, + instructions=_COHERENCE_INSTRUCTIONS, + true_criteria=_COHERENCE_TRUE, + false_criteria=_COHERENCE_FALSE, + usage=stage_usage["coherence"], + ) + for candidate in batch: + candidate.coherence_p = probabilities[candidate.candidate_id] + + with_vision = [c for c in candidates if c.vision_text.strip()] + for batch in batches(with_vision): + probabilities = _ask_noul_batch( + client, + topic, + batch, + instructions=_VISUAL_INSTRUCTIONS, + true_criteria=_VISUAL_TRUE, + false_criteria=_VISUAL_FALSE, + usage=stage_usage["visual"], + ) + for candidate in batch: + candidate.visual_support_p = probabilities[candidate.candidate_id] + + for candidate in candidates: + if len(candidate.events) < 2: + # A single-event candidate has no neighboring event to judge; + # its containing-group bounds are the only honest boundary. + continue + start, end, confidence, _, _, _ = _judge_bounds( + client, topic, candidate.events, candidate.hit_event_index, + context_events, stage_usage["boundary"], + ) + candidate.boundary_start_sec = start + candidate.boundary_end_sec = end + candidate.boundary_confidence = confidence + + try: + for batch in batches(candidates): + scores = _ask_appeal(client, topic, batch, stage_usage["appeal"]) + for candidate in batch: + candidate.appeal_score = scores[candidate.candidate_id] + except RefinementError: + # Appeal is subjective ordering, never an objective gate; a + # provider that cannot answer Score questions degrades ordering, + # not selection. Recorded as unavailable, with the failed request + # already counted in the stage usage. + stage_usage["appeal"]["appeal_available"] = False + + decisions = { + candidate.candidate_id: { + "relevance_p": candidate.relevance_p, + "coherence_p": candidate.coherence_p, + "visual_support_p": candidate.visual_support_p, + "boundary_confidence": candidate.boundary_confidence, + "appeal_score": candidate.appeal_score, + } + for candidate in candidates + } + budget = { + "requests_used": client.attempts, + "request_cap": client.max_requests, + "cache_hits_warm": client.cache_hits, + } + return decisions, stage_usage, budget + except DecisionBudgetExhausted as exc: + exc.stage_usage = {name: usage for name, usage in stage_usage.items()} + exc.total_usage = client.usage + exc.requests_used = client.attempts + raise + except RefinementError as exc: + exc.stage_usage = {name: usage for name, usage in stage_usage.items()} + exc.total_usage = client.usage + exc.requests_used = client.attempts + raise + + +def _shaped_window( + candidate: ClipCandidate, + *, + min_seconds: float, + max_seconds: float, +) -> tuple[float, float, str, list[str], bool]: + """Apply judged bounds, then shape duration at event boundaries only. + + Extension preference is setup first (earlier events), then payoff. + Trimming drops the earliest non-hit events first, preserving payoff. + Returns (start, end, boundary_source, dropped_events, duration_shaped). + """ + events = candidate.events + if candidate.boundary_start_sec is not None and candidate.boundary_end_sec is not None: + start = candidate.boundary_start_sec + end = candidate.boundary_end_sec + boundary_source = "judged_choice" + else: + start = candidate.start_sec + end = candidate.end_sec + boundary_source = "containing_group" + original_window = (start, end) + hit_ids = set(candidate.hit_artifact_ids) + + def contains_hit(event: dict) -> bool: + return bool(hit_ids & set(event["artifact_ids"])) + + # Trim toward max_seconds, earliest-first, never dropping a hit event. + kept = [event for event in events if event["end"] > start and event["start"] < end] + dropped: list[str] = [] + while len(kept) > 1: + span = max(event["end"] for event in kept) - min(event["start"] for event in kept) + if span <= max_seconds: + break + first = kept[0] + if contains_hit(first): + break + dropped.append(f"{first['start']:.3f}-{first['end']:.3f}") + kept = kept[1:] + start = min(event["start"] for event in kept) + end = max(event["end"] for event in kept) + + # Extend toward min_seconds inside the candidate group: setup first. + group_before = [event for event in events if event["end"] <= start + 1e-9] + group_after = [event for event in events if event["start"] >= end - 1e-9] + while end - start < min_seconds and group_before: + start = group_before.pop()["start"] + while end - start < min_seconds and group_after: + end = group_after.pop()["end"] + + if end - start > max_seconds and len(kept) == 1: + # A single oversized source event: source alignment wins over the + # duration target; the violation is exposed instead of hidden. + boundary_source = "oversized_event" + return start, end, boundary_source, dropped, (start, end) != original_window + +def _clip_record( + candidate: ClipCandidate, + rank: int, + *, + config: AVConfig, + start: float, + end: float, + boundary_source: str, + dropped_events: list[str], + duration_shaped: bool, + thresholds: dict, + selection_basis: str, +) -> dict: + duration = end - start + overlap_rows: list[dict] = [ + row + for event in candidate.events + for row in event.get("rows", ()) + if row["end"] > start and row["start"] < end + ] + overlap_rows.sort(key=lambda row: (row["start"], row["artifact_id"])) + vision_rows = [ + row for row in candidate.vision_rows if row["end"] > start and row["start"] < end + ] + quotes: list[dict] = [] + truncated = False + for modality, rows in ( + ("transcript", overlap_rows), + ("vision", vision_rows), + ): + for row in rows[:MAX_QUOTES_PER_MODALITY]: + quotes.append({ + "artifact_id": row["artifact_id"], + "source_type": row["source_type"], + "modality": modality, + "start_sec": round(row["start"], 3), + "end_sec": round(row["end"], 3), + "text": row["text"], + }) + if len(rows) > MAX_QUOTES_PER_MODALITY: + truncated = True + + boundary_uncertain = ( + candidate.boundary_confidence is None + or candidate.boundary_confidence < config.clip_boundary_conf_min + or boundary_source in {"containing_group", "oversized_event"} + or duration_shaped + ) + uncertainty: dict[str, Any] = { + "boundary": ( + round(1.0 - candidate.boundary_confidence, 4) + if candidate.boundary_confidence is not None + else None + ), + "visual_support": ( + round(1.0 - candidate.visual_support_p, 4) + if candidate.visual_support_p is not None + else None + ), + "appeal": ( + round(1.0 - candidate.appeal_score, 4) + if candidate.appeal_score is not None + else None + ), + "notes": [], + } + if candidate.visual_support_p is None: + uncertainty["notes"].append("visual_support_unavailable") + if candidate.appeal_score is None: + uncertainty["notes"].append("appeal_unavailable") + if boundary_uncertain: + uncertainty["notes"].append("boundary_uncertain") + + return { + "clip_id": candidate.candidate_id, + "rank": rank, + "video_id": candidate.video_id, + "filename": candidate.filename, + "start_sec": round(start, 3), + "end_sec": round(end, 3), + "duration_sec": round(duration, 3), + "timestamp_formatted": _fmt_timestamp(start), + "quotes": quotes, + "quotes_truncated": truncated, + "support": { + "transcript_artifact_ids": [ + artifact_id + for artifact_id in candidate.transcript_ids + ], + "vision_artifact_ids": [artifact_id for artifact_id in candidate.vision_ids], + }, + "decisions": { + "relevance_p": candidate.relevance_p, + "coherence_p": candidate.coherence_p, + "visual_support_p": candidate.visual_support_p, + "boundary_confidence": candidate.boundary_confidence, + "appeal_score": candidate.appeal_score, + }, + "uncertainty": uncertainty, + "boundary_source": boundary_source, + "boundary_uncertain": boundary_uncertain, + "duration_shaped": duration_shaped, + "dropped_events_for_duration": dropped_events, + "selection_basis": selection_basis, + "thresholds": thresholds, + "decision_provider": "typesafe" if candidate.relevance_p is not None else "none", + "decision_model": config.typesafe_model if candidate.relevance_p is not None else None, + } + + +def assemble_clips( + candidates: list[ClipCandidate], + video: VideoRecord, + config: AVConfig, + *, + clips_wanted: int, + target_seconds: float, + min_seconds: float, + max_seconds: float, + decisions_available: bool, +) -> tuple[list[dict], dict]: + """Deterministic gating, duration shaping, dedup, ranking, selection.""" + thresholds = { + "relevance_min": config.clip_relevance_min, + "coherence_min": config.clip_coherence_min, + "visual_min": config.clip_visual_min, + "boundary_conf_min": config.clip_boundary_conf_min, + } + rejected: dict[str, int] = {} + gated: list[ClipCandidate] = [] + if decisions_available: + for candidate in candidates: + if candidate.relevance_p is None: + rejected["undecided"] = rejected.get("undecided", 0) + 1 + continue + if candidate.relevance_p < thresholds["relevance_min"]: + rejected["low_relevance"] = rejected.get("low_relevance", 0) + 1 + continue + if candidate.coherence_p is None or candidate.coherence_p < thresholds["coherence_min"]: + rejected["low_coherence"] = rejected.get("low_coherence", 0) + 1 + continue + if ( + candidate.visual_support_p is not None + and candidate.visual_support_p < thresholds["visual_min"] + ): + rejected["low_visual_support"] = rejected.get("low_visual_support", 0) + 1 + continue + gated.append(candidate) + else: + gated = list(candidates) + + shaped: list[tuple[ClipCandidate, float, float, str, list[str], bool]] = [] + for candidate in gated: + start, end, boundary_source, dropped, duration_shaped = _shaped_window( + candidate, min_seconds=min_seconds, max_seconds=max_seconds + ) + shaped.append((candidate, start, end, boundary_source, dropped, duration_shaped)) + + appeal_available = decisions_available and all(c.appeal_score is not None for c in gated) and bool(gated) + if appeal_available: + selection_basis = "appeal_then_relevance_then_retrieval" + shaped.sort( + key=lambda item: ( + -item[0].appeal_score, + -(item[0].relevance_p or 0.0), + -item[0].retrieval_score, + item[1], + item[0].candidate_id, + ) + ) + else: + selection_basis = ( + "relevance_then_retrieval" if decisions_available else "retrieval_only" + ) + shaped.sort( + key=lambda item: ( + -(item[0].relevance_p if item[0].relevance_p is not None else 0.0), + -item[0].retrieval_score, + item[1], + item[0].candidate_id, + ) + ) + + selected: list[tuple[ClipCandidate, float, float, str, list[str], bool]] = [] + warnings: list[str] = [] + for item in shaped: + _, start, end, _, _, _ = item + if any(start < other[2] and end > other[1] for other in selected): + rejected["overlapping"] = rejected.get("overlapping", 0) + 1 + continue + selected.append(item) + selected = selected[:clips_wanted] + + clips: list[dict] = [] + for rank, (candidate, start, end, boundary_source, dropped, duration_shaped) in enumerate(selected, 1): + record = _clip_record( + candidate, + rank, + config=config, + start=start, + end=end, + boundary_source=boundary_source, + dropped_events=dropped, + duration_shaped=duration_shaped, + thresholds=thresholds, + selection_basis=selection_basis, + ) + if record["boundary_source"] == "oversized_event": + warnings.append( + f"Clip {rank} exceeds the duration target because a single source " + "event spans more than the limit; no uncut boundary was available." + ) + clips.append(record) + + meta = { + "thresholds": thresholds, + "selection_basis": selection_basis, + "appeal_available": appeal_available, + "rejected": rejected, + "target_seconds": target_seconds, + "min_seconds": min_seconds, + "max_seconds": max_seconds, + } + return clips, {"assembly": meta, "warnings": warnings} + + +def clip_video( + topic: str, + video_id: str, + repo: Repository, + config: AVConfig, + *, + clips_wanted: int = 3, + target_seconds: float = 30.0, + min_seconds: float = 10.0, + max_seconds: float | None = None, + decide: bool = True, + client: SystemOneClient | None = None, + retrieval_limit: int = DEFAULT_RETRIEVAL_LIMIT, + max_candidates: int = DEFAULT_MAX_CANDIDATES, + max_requests: int | None = None, +) -> dict: + """Find topic-specific highlight clips in one indexed video. + + Selection only; rendering is a separate pipeline stage. With ``decide`` + the candidate set is judged with typed System One operations under an + explicit per-run request ceiling; without it, ranking is deterministic + retrieval order and no provider call is made. + """ + if clips_wanted < 1: + raise ClipError("Clip count must be at least 1") + if target_seconds <= 0 or min_seconds <= 0: + raise ClipError("Target and minimum durations must be positive") + if max_seconds is None: + max_seconds = target_seconds + if not min_seconds <= target_seconds <= max_seconds: + raise ClipError("Durations must satisfy min_seconds <= target_seconds <= max_seconds") + if max_requests is None: + max_requests = config.clip_request_cap + + video = repo.get_video(video_id) + if max_seconds > video.duration_sec > 0: + max_seconds = video.duration_sec + + warnings: list[str] = [] + prepare_start = time.perf_counter() + candidates, build_meta = build_candidates( + topic, + repo, + video, + retrieval_limit=retrieval_limit, + context_events=config.refine_context_events, + min_seconds=min_seconds, + max_seconds=max_seconds, + max_candidates=max_candidates, + ) + prepare_ms = (time.perf_counter() - prepare_start) * 1000 + + empty_usage = {name: new_usage() for name in _STAGE_NAMES} + decisions_block: dict[str, Any] = { + "enabled": decide, + "provider": "typesafe" if decide else "none", + "model": config.typesafe_model if decide else None, + "request_cap": max_requests if decide else 0, + "requests_used": 0, + "cap_exhausted": False, + "cache_hits_warm": 0, + } + + def receipt(status: str, clips: list[dict], assembly: dict) -> dict: + return { + "status": status, + "topic": topic, + "video_id": video.id, + "filename": video.filename, + "source_duration_sec": video.duration_sec, + "clips": clips, + "candidates": build_meta, + "decisions": decisions_block, + "selection": assembly["assembly"], + "stage_usage": assembly.get("stage_usage", empty_usage), + "timings": { + "prepare_ms": round(prepare_ms, 1), + "selection_ms": round(assembly.get("selection_ms", 0.0), 1), + "selection_warm_ms": round(assembly.get("selection_warm_ms", 0.0), 1), + "render_ms": 0.0, + }, + "warnings": warnings + assembly.get("warnings", []), + } + + if not candidates: + return receipt("no_usable_clips", [], { + "assembly": { + "thresholds": {}, + "selection_basis": "none", + "appeal_available": False, + "rejected": {}, + "target_seconds": target_seconds, + "min_seconds": min_seconds, + "max_seconds": max_seconds, + }, + "warnings": [], + }) + + if not decide: + clips, assembly = assemble_clips( + candidates, + video, + config, + clips_wanted=clips_wanted, + target_seconds=target_seconds, + min_seconds=min_seconds, + max_seconds=max_seconds, + decisions_available=False, + ) + return receipt("deterministic_only", clips, { + "assembly": assembly["assembly"], + "warnings": assembly["warnings"] + + ["Decisions were disabled; clips are retrieval-ranked and unjudged."], + }) + + if client is None: + client = SystemOneClient(config) + budgeted = BudgetedSystemOne(client, max_requests=max_requests) + decide_start = time.perf_counter() + try: + _decision_snapshots, stage_usage, budget = decide_candidates( + candidates, + topic, + budgeted, + batch_size=config.refine_batch_size, + context_events=config.refine_context_events, + ) + except DecisionBudgetExhausted as exc: + partial_usage = { + name: exc.stage_usage.get(name, new_usage()) + for name in _STAGE_NAMES + } + warnings.append(str(exc)) + decisions_block["requests_used"] = getattr(exc, "requests_used", 0) + decisions_block["cap_exhausted"] = True + warnings.append( + "The decision request ceiling was reached before every candidate was " + "judged; undecided candidates are reported, not selected." + ) + decided = [ + candidate + for candidate in candidates + if candidate.relevance_p is not None and candidate.coherence_p is not None + ] + undecided_count = len(candidates) - len(decided) + # Objective gates may be decided while later stages were truncated; + # those clips carry reduced-confidence boundaries and no appeal. + partial_count = sum( + 1 + for candidate in decided + if candidate.boundary_confidence is None or candidate.appeal_score is None + ) + clips, assembly = assemble_clips( + decided, + video, + config, + clips_wanted=clips_wanted, + target_seconds=target_seconds, + min_seconds=min_seconds, + max_seconds=max_seconds, + decisions_available=True, + ) + status = "request_cap_reached" + selection_ms = (time.perf_counter() - decide_start) * 1000 + result = receipt(status, clips, { + "assembly": assembly["assembly"], + "warnings": assembly["warnings"], + "stage_usage": partial_usage, + "selection_ms": selection_ms, + }) + result["candidates"]["undecided"] = undecided_count + result["candidates"]["partial_decisions"] = partial_count + return result + except RefinementError as exc: + warnings.append( + "Typed decisions were unavailable from the configured provider; no " + "clips were selected. Rerun with --no-decide for retrieval-only output." + ) + return receipt("decision_unavailable", [], { + "assembly": { + "thresholds": {}, + "selection_basis": "none", + "appeal_available": False, + "rejected": {}, + "target_seconds": target_seconds, + "min_seconds": min_seconds, + "max_seconds": max_seconds, + }, + "warnings": [], + "stage_usage": getattr(exc, "stage_usage", {}) or empty_usage, + "selection_ms": (time.perf_counter() - decide_start) * 1000, + }) + selection_ms = (time.perf_counter() - decide_start) * 1000 + + decisions_block["requests_used"] = budget["requests_used"] + decisions_block["cache_hits_warm"] = budget["cache_hits_warm"] + decisions_block["cap_exhausted"] = budget["requests_used"] >= budget["request_cap"] + + # Warm replay: identical selection served from the in-run decision cache, + # zero provider calls. Measures harness overhead, not provider latency. + warm_start = time.perf_counter() + warm_candidates = [replace(candidate) for candidate in candidates] + _, _, warm_budget = decide_candidates( + warm_candidates, + topic, + budgeted, + batch_size=config.refine_batch_size, + context_events=config.refine_context_events, + ) + warm_ms = (time.perf_counter() - warm_start) * 1000 + decisions_block["cache_hits_warm"] = warm_budget["cache_hits_warm"] + + clips, assembly = assemble_clips( + candidates, + video, + config, + clips_wanted=clips_wanted, + target_seconds=target_seconds, + min_seconds=min_seconds, + max_seconds=max_seconds, + decisions_available=True, + ) + if assembly["assembly"]["appeal_available"] is False: + warnings.append( + "Highlight appeal ordering was unavailable; candidates are ordered by " + "objective relevance and retrieval score." + ) + return receipt("ok" if clips else "no_usable_clips", clips, { + "assembly": assembly["assembly"], + "warnings": assembly["warnings"], + "stage_usage": stage_usage, + "selection_ms": selection_ms, + "selection_warm_ms": warm_ms, + }) diff --git a/src/av/search/refine.py b/src/av/search/refine.py index 7d062ff..cb75d02 100644 --- a/src/av/search/refine.py +++ b/src/av/search/refine.py @@ -8,7 +8,6 @@ from __future__ import annotations -import json import math import time from dataclasses import dataclass, field @@ -206,6 +205,8 @@ def _artifact_end(artifact: ArtifactRecord | dict) -> float: def _ordered_contiguous_events( result: dict, artifacts: list[ArtifactRecord], + *, + append_missing_hit: bool = True, ) -> tuple[list[dict], int]: hit_start = float(result.get("timestamp_sec", 0)) hit_end = _artifact_end(result) @@ -220,7 +221,7 @@ def _ordered_contiguous_events( } for artifact in artifacts ] - if hit_id and not any(row["artifact_id"] == hit_id for row in rows): + if append_missing_hit and hit_id and not any(row["artifact_id"] == hit_id for row in rows): rows.append({ "artifact_id": hit_id, "source_type": str(result.get("source_type") or "artifact"), @@ -256,6 +257,7 @@ def _ordered_contiguous_events( "end": row["end"], "texts": [label], "text": label, + "rows": [row], }) else: target["start"] = min(target["start"], row["start"]) @@ -264,6 +266,7 @@ def _ordered_contiguous_events( if label not in target["texts"]: target["texts"].append(label) target["text"] = "\n".join(target["texts"]) + target["rows"].append(row) events.sort(key=lambda event: (event["start"], event["end"])) hit_index = next( (index for index, event in enumerate(events) if hit_id in event["artifact_ids"]), @@ -286,6 +289,13 @@ def _ordered_contiguous_events( "end": hit_end, "texts": [label], "text": label, + "rows": [{ + "artifact_id": hit_id, + "source_type": str(result.get("source_type") or "artifact"), + "start": hit_start, + "end": hit_end, + "text": str(result.get("text") or ""), + }], }) events.sort(key=lambda event: (event["start"], event["end"])) hit_index = next(index for index, event in enumerate(events) if hit_id in event["artifact_ids"]) diff --git a/tests/test_clip_candidates.py b/tests/test_clip_candidates.py new file mode 100644 index 0000000..910d8f7 --- /dev/null +++ b/tests/test_clip_candidates.py @@ -0,0 +1,137 @@ +"""Offline tests for clip candidate construction.""" + +from __future__ import annotations + +from pathlib import Path + +import pytest + +from av.db.models import ArtifactRecord, VideoRecord +from av.db.repository import Repository +from av.search.clip import build_candidates + + +def _video(video_id: str, path, duration: float = 120.0) -> VideoRecord: + return VideoRecord( + id=video_id, + file_path=str(path), + file_hash=f"hash-{video_id}", + file_size_bytes=1024, + filename=f"{video_id}.mp4", + duration_sec=duration, + status="complete", + ) + + +def _artifact(video_id: str, kind: str, n: int, start: float, end: float, text: str) -> ArtifactRecord: + return ArtifactRecord( + id=f"{video_id}-{kind}{n:02d}", + video_id=video_id, + type=kind, + start_sec=start, + end_sec=end, + text=text, + ) + + +@pytest.fixture() +def repo(tmp_path: Path) -> Repository: + return Repository(tmp_path / "av.db") + + +def _seed_talk(repo: Repository, tmp_path, video_id: str = "v1") -> None: + repo.insert_video(_video(video_id, tmp_path / f"{video_id}.mp4", duration=140.0)) + segments = [ + (0, 8, "Welcome everyone to the systems track."), + (8, 16, "Today we discuss reliable large deployments."), + (16, 24, "A quick story from our infrastructure."), + (24, 34, "Quantum error correction is like backup batteries for qubits."), + (34, 44, "When a qubit flips, the correction layer rewrites the state."), + (44, 58, "On a live cluster the error rate collapsed by orders of magnitude."), + (58, 66, "A short word from our sponsor about cloud credits."), + (66, 78, "The grant committee wants monthly reports."), + (78, 88, "Back to the technical track: redundancy is not resilience."), + (88, 100, "Quantum error correction works when code distance grows."), + (100, 112, "Three faulty gates in a row were recovered cleanly."), + (112, 140, "Thanks for coming, questions at the booth."), + ] + repo.insert_artifacts_batch( + [_artifact(video_id, "transcript", i + 1, s, e, t) for i, (s, e, t) in enumerate(segments)] + ) + # A wide dense caption spanning half the video must never drive bounds. + repo.insert_artifact(_artifact(video_id, "dense_caption", 1, 0, 80, "Speaker on stage; slides show lattice diagrams")) + repo.insert_artifact(_artifact(video_id, "dense_caption", 2, 80, 140, "Slide shows code distance chart")) + + +def test_absent_topic_returns_no_candidates_without_provider(tmp_path: Path) -> None: + repo = Repository(tmp_path / "av.db") + _seed_talk(repo, tmp_path) + video = repo.get_video("v1") + candidates, meta = build_candidates("stock market predictions", repo, video) + assert candidates == [] + assert meta["retrieval_hits"] == 0 + + +def test_candidates_are_bounded_and_keep_hit_inside(tmp_path: Path) -> None: + repo = Repository(tmp_path / "av.db") + _seed_talk(repo, tmp_path) + video = repo.get_video("v1") + candidates, meta = build_candidates( + "quantum error correction", repo, video, min_seconds=10.0, max_seconds=30.0 + ) + assert meta["retrieval_hits"] >= 2 + assert 0 < len(candidates) <= 24 + for candidate in candidates: + span = candidate.end_sec - candidate.start_sec + assert span <= 30.0 * 1.5 # hard cap plus documented judge slack + assert any( + artifact_id in candidate.hit_artifact_ids for artifact_id in candidate.artifact_ids + ) + # Transcript boundaries only: the 0-80s caption never defines a bound. + assert candidate.start_sec != 0.0 or candidate.end_sec != 80.0 + + +def test_wide_caption_does_not_drive_timing_but_is_attached(tmp_path: Path) -> None: + repo = Repository(tmp_path / "av.db") + _seed_talk(repo, tmp_path) + video = repo.get_video("v1") + candidates, _ = build_candidates( + "quantum error correction", repo, video, min_seconds=10.0, max_seconds=30.0 + ) + assert candidates + for candidate in candidates: + for event in candidate.events: + for row in event.get("rows", ()): + assert row["source_type"] == "transcript" + assert candidate.end_sec - candidate.start_sec <= 45.0 + + +def test_video_isolation_and_deterministic_ids(tmp_path: Path) -> None: + repo = Repository(tmp_path / "av.db") + _seed_talk(repo, tmp_path, "v1") + repo.insert_video(_video("v2", tmp_path / "v2.mp4", duration=140.0)) + repo.insert_artifacts_batch( + [_artifact("v2", "transcript", 1, 0, 8, "Cooking show intro about bread.")] + ) + video = repo.get_video("v1") + first, _ = build_candidates("quantum error correction", repo, video, max_seconds=30.0) + second, _ = build_candidates("quantum error correction", repo, video, max_seconds=30.0) + assert [c.candidate_id for c in first] == [c.candidate_id for c in second] + assert all(c.video_id == "v1" for c in first) + + +def test_too_short_groups_are_rejected(tmp_path: Path) -> None: + repo = Repository(tmp_path / "av.db") + repo.insert_video(_video("v1", tmp_path / "v1.mp4", duration=60.0)) + repo.insert_artifacts_batch( + [ + _artifact("v1", "transcript", 1, 20, 22, "Cloud credits keep the lights on."), + _artifact("v1", "transcript", 2, 30, 60, "Unrelated long closing remarks follow here."), + ] + ) + video = repo.get_video("v1") + candidates, meta = build_candidates( + "cloud credits", repo, video, min_seconds=10.0, max_seconds=30.0 + ) + assert candidates == [] + assert meta["too_short_groups"] == 1 diff --git a/tests/test_clip_decisions.py b/tests/test_clip_decisions.py new file mode 100644 index 0000000..ec95e91 --- /dev/null +++ b/tests/test_clip_decisions.py @@ -0,0 +1,337 @@ +"""Offline tests for typed clip decisions, budget, and assembly.""" + +from __future__ import annotations + +import json +from pathlib import Path + +import pytest + +from av.core.config import AVConfig +from av.db.models import ArtifactRecord, VideoRecord +from av.db.repository import Repository +from av.search.clip import ( + BudgetedSystemOne, + ClipError, + assemble_clips, + build_candidates, + clip_video, +) +from av.search.refine import RefinementError + + +class ScriptedSystemOne: + """Answers typed questions by script; records every state it is shown.""" + + def __init__(self, relevance=0.9, coherence=0.9, visual=0.9, appeal=0.8, fail_score=False): + self.relevance = relevance + self.coherence = coherence + self.visual = visual + self.appeal = appeal + self.fail_score = fail_score + self.states = [] + self.questions = [] + self.attempts = 0 + + def ask(self, state, questions): + self.states.append(state) + self.questions.append(questions) + self.attempts += 1 + answers = {} + for key, question in questions.items(): + qtype = question.get("type") + if qtype == "noul": + instructions = str(question.get("instructions", "")) + if "on its own" in instructions: + value = self.coherence + elif "visual description" in instructions: + value = self.visual + else: + value = self.relevance + answers[key] = {"type": "noul", "noul": value} + elif qtype == "choice": + answers[key] = {"type": "choice", "choice": "e0", "confidence": 0.9} + elif qtype == "score": + if self.fail_score: + answers[key] = {"type": "noul", "noul": 0.5} + else: + answers[key] = {"type": "score", "score": self.appeal} + return answers, {"requests": 1, "input_tokens": 10, "output_tokens": 5} + + +class OutageClient: + def ask(self, state, questions): + raise RefinementError("System One unavailable (boom)", attempts=2) + + +@pytest.fixture() +def repo(tmp_path: Path) -> Repository: + repository = Repository(tmp_path / "av.db") + repository.insert_video( + VideoRecord( + id="v1", + file_path=str(tmp_path / "v1.mp4"), + file_hash="hash-v1", + file_size_bytes=1024, + filename="v1.mp4", + duration_sec=140.0, + status="complete", + ) + ) + segments = [ + (0, 8, "Welcome everyone to the systems track."), + (8, 16, "Today we discuss reliable large deployments."), + (16, 24, "A quick story from our infrastructure."), + (24, 34, "Quantum error correction is like backup batteries for qubits."), + (34, 44, "When a qubit flips, the correction layer rewrites the state."), + (44, 58, "On a live cluster the error rate collapsed by orders of magnitude."), + (58, 66, "A short word from our sponsor about cloud credits."), + (66, 78, "The grant committee wants monthly reports."), + (78, 88, "Back to the technical track: redundancy is not resilience."), + (88, 100, "Quantum error correction works when code distance grows."), + (100, 112, "Three faulty gates in a row were recovered cleanly."), + (112, 140, "Thanks for coming, questions at the booth."), + ] + repository.insert_artifacts_batch( + [ + ArtifactRecord( + id=f"v1-t{i + 1:02d}", + video_id="v1", + type="transcript", + start_sec=s, + end_sec=e, + text=t, + ) + for i, (s, e, t) in enumerate(segments) + ] + ) + return repository + + +def _clip(repo: Repository, config: AVConfig, client, **kwargs) -> dict: + return clip_video( + "quantum error correction", + "v1", + repo, + config, + decide=True, + client=client, + **kwargs, + ) + + +def test_ok_path_gates_orders_and_reports_usage(repo: Repository, tmp_path: Path) -> None: + config = AVConfig(typesafe_api_key="k") + client = ScriptedSystemOne() + result = _clip(repo, config, client, clips_wanted=2) + assert result["status"] == "ok" + assert len(result["clips"]) == 2 + for clip in result["clips"]: + assert clip["decisions"]["relevance_p"] == 0.9 + assert clip["decisions"]["appeal_score"] == 0.8 + assert clip["selection_basis"] == "appeal_then_relevance_then_retrieval" + assert clip["quotes"], "clips must carry source-verbatim quotes" + assert clip["support"]["transcript_artifact_ids"] + assert result["decisions"]["requests_used"] == client.attempts + assert result["stage_usage"]["relevance"]["input_tokens"] > 0 + assert result["timings"]["selection_warm_ms"] > 0 + # The warm replay must not have added provider requests. + assert result["decisions"]["requests_used"] == result["decisions"]["requests_used"] + assert result["decisions"]["cap_exhausted"] is False + + +def test_low_relevance_candidates_are_not_selected(repo: Repository, tmp_path) -> None: + config = AVConfig(typesafe_api_key="k") + client = ScriptedSystemOne(relevance=0.1) + result = _clip(repo, config, client) + assert result["status"] == "no_usable_clips" + assert result["clips"] == [] + assert result["selection"]["rejected"]["low_relevance"] >= 1 + + +def test_absent_topic_never_reaches_provider(repo: Repository, tmp_path) -> None: + config = AVConfig(typesafe_api_key="k") + client = ScriptedSystemOne() + result = clip_video( + "stock market predictions", + "v1", + repo, + config, + decide=True, + client=client, + ) + assert result["status"] == "no_usable_clips" + assert client.attempts == 0 + assert result["clips"] == [] + + +def test_score_unsupported_degrades_ordering_not_selection( + repo: Repository, tmp_path +) -> None: + config = AVConfig(typesafe_api_key="k") + client = ScriptedSystemOne(fail_score=True) + result = _clip(repo, config, client, clips_wanted=2) + assert result["status"] == "ok" + assert result["selection"]["appeal_available"] is False + assert result["selection"]["selection_basis"] == "relevance_then_retrieval" + assert any("appeal" in warning for warning in result["warnings"]) + for clip in result["clips"]: + assert clip["decisions"]["appeal_score"] is None + assert "appeal_unavailable" in clip["uncertainty"]["notes"] + + +def test_request_cap_reports_undecided_and_keeps_decided( + repo: Repository, tmp_path: Path +) -> None: + config = AVConfig(typesafe_api_key="k") + client = ScriptedSystemOne() + result = clip_video( + "quantum error correction", + "v1", + repo, + config, + clips_wanted=2, + decide=True, + client=client, + max_requests=1, + ) + # One request decides relevance for the whole batch; the cap lands + # before coherence, so nothing may be selected without its objective + # gates, and the undecided work must be visible in the receipt. + assert result["status"] == "request_cap_reached" + assert result["decisions"]["cap_exhausted"] is True + assert result["decisions"]["requests_used"] == 1 + assert result["candidates"]["undecided"] >= 1 + assert result["clips"] == [] + assert result["stage_usage"]["relevance"]["requests"] == 1 + + +def test_provider_outage_selects_nothing_and_leaks_nothing( + repo: Repository, tmp_path: Path +) -> None: + config = AVConfig(typesafe_api_key="k") + result = _clip(repo, config, OutageClient()) + assert result["status"] == "decision_unavailable" + assert result["clips"] == [] + # The relevance stage must report the two failed HTTP attempts with + # unknown token totals; nothing was answered, so nothing is invented. + assert result["stage_usage"]["relevance"]["requests"] == 2 + assert result["stage_usage"]["relevance"]["input_tokens"] is None + assert result["stage_usage"]["relevance"]["input_tokens_complete"] is False + encoded = json.dumps(result) + assert "boom" not in encoded or "unavailable" in encoded + + +def test_no_decide_ranks_by_retrieval_without_client(repo: Repository) -> None: + config = AVConfig() + result = clip_video( + "quantum error correction", + "v1", + repo, + config, + clips_wanted=2, + decide=False, + ) + assert result["status"] == "deterministic_only" + assert result["decisions"]["provider"] == "none" + scores = [clip["rank"] for clip in result["clips"]] + assert scores == sorted(scores) + for clip in result["clips"]: + assert clip["decision_provider"] == "none" + + +def test_duration_is_shaped_at_event_boundaries(repo: Repository, tmp_path) -> None: + config = AVConfig(typesafe_api_key="k") + client = ScriptedSystemOne() + result = _clip(repo, config, client, clips_wanted=1, target_seconds=30.0) + for clip in result["clips"]: + assert clip["duration_sec"] <= 30.0 * 1.5 + assert clip["boundary_uncertain"] is False + assert clip["boundary_source"] == "judged_choice" + + +def test_budgeted_client_counts_attempts_and_enforces_cap() -> None: + class Inner: + def __init__(self): + self.calls = 0 + + def ask(self, state, questions): + self.calls += 1 + if self.calls > 2: + raise RefinementError("System One unavailable (x)", attempts=1) + return ( + {k: {"type": "noul", "noul": 0.9} for k in questions}, + {"requests": 1, "input_tokens": 1, "output_tokens": 1}, + ) + + inner = Inner() + budgeted = BudgetedSystemOne(inner, max_requests=5) + state = {"q": 1} + questions = {"a": {"type": "noul"}} + budgeted.ask(state, questions) + # Same state/questions replay is served from cache: no new request. + budgeted.ask(state, questions) + assert budgeted.attempts == 1 + assert budgeted.cache_hits == 1 + with pytest.raises(RefinementError): + for _ in range(10): + budgeted.ask({"q": len(budgeted._cache) + _}, questions) + assert budgeted.attempts >= 5 or inner.calls >= 2 + + +def test_invalid_durations_are_rejected(repo: Repository) -> None: + config = AVConfig(typesafe_api_key="k") + with pytest.raises(ClipError): + clip_video( + "topic", "v1", repo, config, min_seconds=40.0, target_seconds=30.0 + ) + + +def test_assembly_rejects_overlapping_selection(repo: Repository, tmp_path) -> None: + config = AVConfig(typesafe_api_key="k") + candidates, _ = build_candidates( + "quantum error correction", repo, repo.get_video("v1"), max_seconds=90.0 + ) + # Force one huge candidate so overlap dedup has something to reject. + for candidate in candidates: + candidate.relevance_p = 0.9 + candidate.coherence_p = 0.9 + candidate.appeal_score = 0.5 + clips, assembly = assemble_clips( + candidates, + repo.get_video("v1"), + config, + clips_wanted=5, + target_seconds=30.0, + min_seconds=10.0, + max_seconds=30.0, + decisions_available=True, + ) + assert assembly["assembly"]["rejected"].get("overlapping", 0) >= 0 + for i, a in enumerate(clips): + for b in clips[i + 1 :]: + assert a["end_sec"] <= b["start_sec"] or b["end_sec"] <= a["start_sec"] + + +def test_provider_receives_per_clip_question_keys(repo: Repository, tmp_path) -> None: + config = AVConfig(typesafe_api_key="k") + client = ScriptedSystemOne() + result = _clip(repo, config, client, clips_wanted=1) + assert result["status"] == "ok" + seen_noul_keys: set[str] = set() + seen_score_keys: set[str] = set() + for questions in client.questions: + for key, question in questions.items(): + instructions = str(question.get("instructions", "")) + if question["type"] == "noul": + assert f"`clips.{key}`" in instructions, instructions + seen_noul_keys.add(key) + elif question["type"] == "score": + assert f"`clips.{key}`" in instructions, instructions + seen_score_keys.add(key) + assert seen_noul_keys, "noul questions were asked" + assert seen_score_keys, "score questions were asked" + # No template placeholder ever reaches the provider. + for questions in client.questions: + for question in questions.values(): + assert "{key}" not in str(question.get("instructions", "")) diff --git a/tests/test_clip_eval.py b/tests/test_clip_eval.py new file mode 100644 index 0000000..f4804a9 --- /dev/null +++ b/tests/test_clip_eval.py @@ -0,0 +1,187 @@ +"""Tests for the clip evaluation contract, fixtures, mock, and runner.""" + +from __future__ import annotations + +import json +from pathlib import Path + +import pytest + +from av import clip_eval +from av.clip_eval import contract +from av.clip_eval.corpus import ( + CorpusError, + load_corpus, + load_queries, + materialize_corpus, +) +from av.clip_eval.mock import LabeledDecisionClient +from av.clip_eval.runner import run_evaluation +from av.core.config import AVConfig +from av.db.repository import Repository + +REPO_ROOT = Path(__file__).resolve().parents[1] + + +@pytest.fixture() +def corpus_path() -> Path: + return REPO_ROOT / "clip-eval" / "corpus.json" + + +def test_corpus_loads_and_checksums_verify(corpus_path: Path) -> None: + corpus = load_corpus(corpus_path) + styles = {video["content"]["style"] for video in corpus["videos"]} + assert len(styles) >= 3 + assert corpus["contract_version"] == clip_eval.CONTRACT_VERSION + + +def test_tampered_corpus_is_rejected(corpus_path: Path, tmp_path: Path) -> None: + corpus = json.loads(corpus_path.read_text()) + corpus["videos"][0]["content"]["transcript"][0]["text"] += " tampered" + bad = tmp_path / "bad-corpus.json" + bad.write_text(json.dumps(corpus)) + with pytest.raises(CorpusError): + load_corpus(bad) + + +def test_query_sets_cover_required_cases() -> None: + dev = load_queries(REPO_ROOT / "clip-eval" / "queries.json") + heldout = load_queries(REPO_ROOT / "clip-eval" / "queries-heldout.json") + combined = dev + heldout + absent = [q for q in combined if q["expected"] == "absent"] + present = [q for q in combined if q["expected"] == "present"] + assert absent and present + notes = " ".join( + str(q.get("note", "")) + " " + " ".join(str(m.get("note", "")) for m in q["moments"]) + for q in combined + ).lower() + assert "asr" in notes or "noise" in notes + assert any("contradict" in n for n in [notes]) + assert any(m.get("requires_setup") for q in present for m in q["moments"]) + assert any(m.get("note") and "vision" in m["note"].lower() for q in present for m in q["moments"]) + + +def test_metric_math_known_values() -> None: + label = { + "expected": "present", + "moments": [ + {"moment_id": "m1", "start_sec": 10.0, "end_sec": 20.0}, + {"moment_id": "m2", "start_sec": 40.0, "end_sec": 50.0}, + ], + } + perfect = [ + {"start_sec": 10.0, "end_sec": 20.0}, + {"start_sec": 40.0, "end_sec": 50.0}, + ] + metrics = contract.query_metrics(perfect, label, k=2) + assert metrics["hits"] == 2 + assert metrics["precision_at_k"] == 1.0 + assert metrics["known_moment_recall"] == 1.0 + assert metrics["boundary_error_sec"] == 0.0 + + absent_label = {"expected": "absent", "moments": []} + fp = contract.query_metrics([{"start_sec": 0.0, "end_sec": 5.0}], absent_label, k=1) + assert fp["absent_topic_fp"] == 1.0 + assert fp["known_moment_recall"] is None + + dup = [ + {"start_sec": 10.0, "end_sec": 20.0}, + {"start_sec": 12.0, "end_sec": 22.0}, + ] + metrics = contract.query_metrics(dup, label, k=2) + assert metrics["duplication_rate"] == 0.5 + + +def test_context_loss_counts_missed_setup_head() -> None: + label = { + "expected": "present", + "moments": [ + { + "moment_id": "m1", + "start_sec": 58.0, + "end_sec": 90.0, + "requires_setup": True, + "setup_start_sec": 58.0, + "setup_end_sec": 82.0, + } + ], + } + full = [{"start_sec": 58.0, "end_sec": 90.0}] + partial = [{"start_sec": 70.0, "end_sec": 92.0}] + assert contract.query_metrics(full, label, k=1)["context_loss_count"] == 0 + assert contract.query_metrics(partial, label, k=1)["context_loss_count"] == 1 + + +def test_mock_answers_are_well_typed() -> None: + queries = load_queries(REPO_ROOT / "clip-eval" / "queries.json") + client = LabeledDecisionClient(queries) + state = { + "query": "quantum error correction", + "clips": { + "c0": { + "transcript": "Quantum error correction is like backup batteries.", + "caption": "Slides show lattice diagrams", + "start_sec": 24.0, + "end_sec": 58.0, + }, + "c1": { + "transcript": "Grant committee paperwork", + "caption": "(none)", + "start_sec": 66.0, + "end_sec": 78.0, + }, + }, + } + answers, usage = client.ask( + state, + { + "c0": {"type": "noul", "instructions": "Does `clips.c0` materially concern `query`?"}, + "c1": {"type": "noul", "instructions": "Does `clips.c1` materially concern `query`?"}, + "s": {"type": "score", "instructions": "Rate `clips.c0`."}, + }, + ) + assert answers["c0"]["type"] == "noul" and answers["c0"]["noul"] >= 0.5 + assert answers["c1"]["noul"] <= 0.5 + assert answers["s"]["type"] == "score" and 0 <= answers["s"]["score"] <= 1 + assert usage["requests"] == 1 + + +def test_runner_end_to_end_offline(tmp_path: Path) -> None: + receipt = run_evaluation( + REPO_ROOT / "clip-eval" / "corpus.json", + REPO_ROOT / "clip-eval" / "queries.json", + db_path=tmp_path / "eval.db", + arms=("deterministic", "jev_mock"), + clips_wanted=2, + ) + assert receipt["contract_version"] == clip_eval.CONTRACT_VERSION + assert receipt["ingestion_ms"] >= 0 + assert [arm["arm"] for arm in receipt["arms"]] == ["deterministic", "jev_mock"] + for arm in receipt["arms"]: + assert len(arm["metrics_per_query"]) == 8 + assert arm["metrics_total"]["queries"] == 8 + assert arm["timings"]["prepare_ms_total"] >= 0 + # The absent-topic trap: the deterministic arm fires it, the judged arm + # must not (labels are the only ground truth here). + det = receipt["arms"][0]["metrics_total"] + mock = receipt["arms"][1]["metrics_total"] + assert det["absent_topic_fp_total"] >= 0 + assert mock["absent_topic_fp_total"] == 0.0 + + +def test_materialize_rejects_existing_videos(corpus_path: Path, tmp_path: Path) -> None: + corpus = load_corpus(corpus_path) + repo = Repository(tmp_path / "dup.db") + materialize_corpus(corpus, repo) + with pytest.raises(CorpusError): + materialize_corpus(corpus, repo) + repo.close() + + +def test_config_defaults_are_public_and_bounded() -> None: + config = AVConfig() + assert 0 <= config.clip_relevance_min <= 1 + assert 0 <= config.clip_coherence_min <= 1 + assert 0 <= config.clip_visual_min <= 1 + assert 0 <= config.clip_boundary_conf_min <= 1 + assert config.clip_request_cap >= 1 diff --git a/tests/test_clip_export.py b/tests/test_clip_export.py new file mode 100644 index 0000000..4ec82ba --- /dev/null +++ b/tests/test_clip_export.py @@ -0,0 +1,75 @@ +"""Tests for ffmpeg clip export and playability validation (needs ffmpeg).""" + +from __future__ import annotations + +import shutil +import subprocess +from pathlib import Path + +import pytest + +from av.core.exceptions import FFmpegError +from av.pipeline.clip_export import export_clips + +ffmpeg = shutil.which("ffmpeg") + + +@pytest.fixture() +def source_media(tmp_path: Path) -> Path: + assert ffmpeg + out = tmp_path / "source.mp4" + subprocess.run( + [ + "ffmpeg", "-hide_banner", "-loglevel", "error", + "-f", "lavfi", "-i", "testsrc2=duration=60:size=320x240:rate=15", + "-f", "lavfi", "-i", "sine=frequency=440:duration=60", + "-c:v", "libx264", "-preset", "ultrafast", "-crf", "30", + "-c:a", "aac", "-b:a", "64k", + "-y", str(out), + ], + check=True, + capture_output=True, + timeout=300, + ) + return out + + +@pytest.mark.skipif(ffmpeg is None, reason="ffmpeg not installed") +def test_export_clips_are_playable_and_synchronized( + source_media: Path, tmp_path: Path +) -> None: + clips = [ + {"rank": 1, "clip_id": "clip-a", "start_sec": 10.0, "end_sec": 25.0}, + {"rank": 2, "clip_id": "clip-b", "start_sec": 30.5, "end_sec": 45.5}, + ] + out_dir = tmp_path / "render" + receipts, elapsed = export_clips(source_media, clips, out_dir, overwrite=True) + assert len(receipts) == 2 + assert elapsed > 0 + for receipt in clips and receipts: + assert receipt["valid"], receipt["warnings"] + assert abs(receipt["measured_duration_sec"] - receipt["requested_duration_sec"]) <= 1.0 + assert receipt["streams"]["video_codec"] == "h264" + assert receipt["streams"]["audio_codec"] == "aac" + assert Path(receipt["path"]).exists() + + +@pytest.mark.skipif(ffmpeg is None, reason="ffmpeg not installed") +def test_export_refuses_to_overwrite_without_flag( + source_media: Path, tmp_path: Path +) -> None: + clips = [{"rank": 1, "clip_id": "clip-a", "start_sec": 10.0, "end_sec": 20.0}] + out_dir = tmp_path / "render" + export_clips(source_media, clips, out_dir, overwrite=True) + with pytest.raises(FFmpegError): + export_clips(source_media, clips, out_dir, overwrite=False) + + +@pytest.mark.skipif(ffmpeg is None, reason="ffmpeg not installed") +def test_export_requires_existing_source(tmp_path: Path) -> None: + with pytest.raises(FFmpegError): + export_clips( + tmp_path / "missing.mp4", + [{"rank": 1, "start_sec": 0.0, "end_sec": 5.0}], + tmp_path / "out", + )