From c038ce35336f590f1fa352af774e65cd0ecc2fcf Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Fri, 3 Jul 2026 16:04:47 -0600 Subject: [PATCH] fix(retention): gate raw-media purge on derived per-file state, fail closed A failed or empty extraction also writes the extraction JSONL, so the old presence-only completion gate treated a failed audio transcription as complete and deleted the only copy of the raw media. Because the gate is whole-segment, a failed audio extraction also exposed the segment's video. Re-key eligibility on derive_modality_state per extraction file. Fail closed on failed, malformed, and unreadable files. The gate computes has_chunks only from rows after the header so a stray start or timestamp merged via SEGMENT_META cannot mask a failure. Legacy chunk-bearing files with no processing record stay eligible. Add the adjacent honesty fixes: a distinct segments_blocked_failed counter and per-segment detail surfaced in the result, CLI, settings route, and logs; per-segment pruning-runs audit writes that survive a mid-loop crash, with AuditOutcome captured and surfaced; retention.log now also writes for real runs that block but delete nothing, while dry-run still writes nothing. Correct the docstring's false 7-day default claim: the real default is mode="keep". Regenerate docs/deletion-sites-inventory.md with current line refs and the omitted log_retention.py rows, and update logs.md for the new audit record shape. Co-Authored-By: Claude Opus 4.8 (1M context) --- docs/deletion-sites-inventory.md | 18 +- solstone/apps/settings/routes.py | 2 + solstone/talent/journal/references/logs.md | 3 +- solstone/think/retention.py | 347 ++++++++++----- solstone/think/tools/call.py | 12 + tests/test_retention.py | 472 ++++++++++++++++++--- 6 files changed, 680 insertions(+), 174 deletions(-) diff --git a/docs/deletion-sites-inventory.md b/docs/deletion-sites-inventory.md index b52a0d32f..500213d5e 100644 --- a/docs/deletion-sites-inventory.md +++ b/docs/deletion-sites-inventory.md @@ -16,11 +16,11 @@ Inventory of every non-test, non-scratch, non-atomic-tmp destructive removal (`s - Reference model: `solstone/think/retention.py` - scope-narrow docstring at `:4-19` - - completion check at `:73-115` - - per-file stream-hashed SHA-256 at `:416-422` - - dry-run support at `:349-369`, `:427-429`, `:450-451` - - narrow exception handling at `:378-381`, `:416-429` - - retention log at `:456-472` + - deletion gate at `:75-253` + - per-file stream-hashed SHA-256 at `:563-576` + - dry-run support at `:478-498`, `:577-579`, `:599` + - failed-extraction block at `:523-540` + - retention log at `:624-625`, `:664-683` - Write-owner table pointer: `CLAUDE.md` / `AGENTS.md` §7 L2 - Importer convention: importers audit destructive operations via `log_app_action(app='import', ...)` per repo convention (`solstone/think/importers/journal_source_cli.py:40, 75, 230, 250`). - Raw grep noise removed manually: nested app test hit at `solstone/apps/observer/tests/test_routes.py:1008` and regex literals in `scripts/check_layer_hygiene.py:54-55`. @@ -35,7 +35,13 @@ Inventory of every non-test, non-scratch, non-atomic-tmp destructive removal (`s | file:line | target | trigger | path validation | audit log | dry-run | class | why | | --- | --- | --- | --- | --- | --- | --- | --- | -| `solstone/think/retention.py:428` | raw media files in completed segments | retention purge on eligible segments | `is_segment_complete()` plus retention-policy eligibility | `_write_retention_log()` to `health/retention.log` | yes | `✅` | reference template for this sweep | +| `solstone/think/retention.py:578` | raw media files in gate-eligible segments | retention purge on eligible segments | `resolve_segment_gate()` plus retention-policy eligibility | per-segment `write_prune_audit()` record plus `_write_retention_log()` summary | yes | `✅` | reference template for this sweep | + +## think/log_retention + +| file:line | target | trigger | path validation | audit log | dry-run | class | why | +| --- | --- | --- | --- | --- | --- | --- | --- | +| `solstone/think/log_retention.py:682,684,686` | dated operational log/cache files and dirs from the retention allowlist | operational log/cache pruning | class-specific scanners feed `_delete_target()` only with dated allowlist paths under the journal | `write_prune_audit()` plus structured result errors | yes | `✅` | declared operational-log/cache pruning owner with dry-run and audit outcome surfacing | ## think/entities diff --git a/solstone/apps/settings/routes.py b/solstone/apps/settings/routes.py index 6b50c61c7..3ce2faa81 100644 --- a/solstone/apps/settings/routes.py +++ b/solstone/apps/settings/routes.py @@ -2077,6 +2077,8 @@ def run_purge() -> Any: "segments_processed": result.segments_processed, "segments_skipped_incomplete": result.segments_skipped_incomplete, "segments_skipped_policy": result.segments_skipped_policy, + "segments_blocked_failed": result.segments_blocked_failed, + "partial_error": result.partial_error, "dry_run": dry_run, } diff --git a/solstone/talent/journal/references/logs.md b/solstone/talent/journal/references/logs.md index f014e059f..3e4049605 100644 --- a/solstone/talent/journal/references/logs.md +++ b/solstone/talent/journal/references/logs.md @@ -141,6 +141,7 @@ The `health/` directory contains log files for long-running services. **Files:** - `health/.log` – log output for each service (e.g., `observe.log`, `cortex.log`, `convey.log`) -- `health/retention.log` – JSONL log of retention purge operations with timestamps, files deleted, bytes freed, and per-segment details +- `health/retention.log` – JSONL run summaries for real retention purge operations, including files deleted, bytes freed, per-segment deletion details, `segments_blocked_failed`, `blocked_failed_details`, and `partial_error` +- `health/pruning-runs/YYYYMMDD.jsonl` – per-deleted-segment raw-media prune audit records, one JSONL line per segment with day, stream, segment, per-file name/bytes/SHA-256 hash, bytes freed, and `processed_at`; dry-runs and zero-deletion blocked-only runs do not write pruning-run records These logs are useful for debugging service issues. See [DOCTOR.md](../../../docs/DOCTOR.md) for diagnostics and troubleshooting guidance. diff --git a/solstone/think/retention.py b/solstone/think/retention.py index a67dd0bfc..760151033 100644 --- a/solstone/think/retention.py +++ b/solstone/think/retention.py @@ -5,14 +5,14 @@ Manages the lifecycle of raw media files (layer 1 captures) in journal segments. Three retention modes: -- keep: retain raw media indefinitely -- days: delete raw media after N days, once processing is complete (default: 7) +- keep: retain raw media indefinitely and purge nothing (default) +- days: delete raw media after explicit N days, once processing is complete - processed: delete raw media as soon as processing completes Scope: raw media ONLY. Chronicle JSONL, derived outputs, talents/ directories, and all other journal content persist indefinitely and are never touched by -retention. Do not extrapolate the "7 days" default to any other data — it is -specific to raw media purging. +retention. Days mode requires an explicit days value from config or CLI override; +without one, it purges nothing. Safety invariant: never delete raw media from segments that haven't finished processing. All completion checks must pass before any deletion. @@ -29,10 +29,11 @@ from datetime import datetime from pathlib import Path from typing import Any +from solstone.think.data_state import DataState, derive_modality_state from solstone.think.media import AUDIO_EXTENSIONS as RAW_AUDIO_EXTENSIONS from solstone.think.media import MEDIA_EXTENSIONS as RAW_MEDIA_EXTENSIONS from solstone.think.media import VIDEO_EXTENSIONS as RAW_VIDEO_EXTENSIONS -from solstone.think.pruning_audit import write_prune_audit +from solstone.think.pruning_audit import AuditOutcome, write_prune_audit from solstone.think.utils import day_dirs, get_journal, iter_segments logger = logging.getLogger(__name__) @@ -71,72 +72,185 @@ def get_raw_media_files(segment_path: Path) -> list[Path]: # --------------------------------------------------------------------------- -def is_segment_complete(segment_path: Path) -> bool: - """Check if a segment has finished all processing. +@dataclass +class SegmentGate: + """Deletion gate verdict for one segment.""" + + verdict: str + failed_files: dict[str, str] = field(default_factory=dict) + completion_files: list[Path] = field(default_factory=list) + + +def _matching_extraction_files( + segment_path: Path, patterns: tuple[str, ...] +) -> list[Path]: + return sorted( + { + path + for pattern in patterns + for path in segment_path.glob(pattern) + if path.is_file() + } + ) - Completion checks (ALL must pass): - 1. No _active.jsonl files in talents/ - 2. audio.jsonl exists if any audio raw media was captured - 3. screen.jsonl exists if any video raw media was captured - 4. talents/speaker_labels.json exists if embeddings (.npz) are present - """ + +def _derive_extraction_file_state( + seg_path: Path, + jsonl_path: Path, + *, + modality: str, + marker_key: str, + has_raw: bool, +) -> str: + """Read one extraction file strictly enough for irreversible deletion.""" + try: + lines: list[str] = [] + with jsonl_path.open("r", encoding="utf-8") as handle: + for line in handle: + if not line.strip(): + continue + lines.append(line) + if len(lines) == 2: + break + except OSError: + return "malformed" + + if not lines: + return "malformed" + + try: + header = json.loads(lines[0]) + except (json.JSONDecodeError, ValueError): + return "malformed" + if not isinstance(header, dict): + return "malformed" + + record = header.get("_solstone_processing") + if not isinstance(record, dict): + record = None + + # Row 1 is always the metadata header. A stray marker key merged through + # SEGMENT_META must never make a header-only file look chunk-bearing when + # retention is deciding whether raw media can be irreversibly deleted. + has_chunks = False + if len(lines) == 2: + try: + first_chunk = json.loads(lines[1]) + except (json.JSONDecodeError, ValueError): + first_chunk = None + has_chunks = isinstance(first_chunk, dict) and marker_key in first_chunk + + return derive_modality_state( + seg_path, + modality, + has_chunks=has_chunks, + has_jsonl=True, + has_raw=has_raw, + record=record, + ) + + +def _collect_extraction_states( + segment_path: Path, + paths: list[Path], + *, + modality: str, + marker_key: str, + has_raw: bool, + failed_files: dict[str, str], + completion_files: list[Path], +) -> bool: + incomplete = False + for path in paths: + state = _derive_extraction_file_state( + segment_path, + path, + modality=modality, + marker_key=marker_key, + has_raw=has_raw, + ) + if state == "malformed": + failed_files[path.name] = "malformed" + elif state == DataState.FAILED.value: + failed_files[path.name] = DataState.FAILED.value + elif state in {DataState.PENDING.value, DataState.ANALYZING.value}: + incomplete = True + elif state in {DataState.ANALYZED.value, DataState.EMPTY.value}: + completion_files.append(path) + else: # pragma: no cover - derive_modality_state owns this closed vocabulary. + raise RuntimeError(f"unexpected {modality} data state for {path}: {state}") + return incomplete + + +def resolve_segment_gate(segment_path: Path) -> SegmentGate: + """Resolve whether a segment's raw media is safe to purge.""" agents_dir = segment_path / "talents" + incomplete = False + failed_files: dict[str, str] = {} + completion_files: list[Path] = [] - # Check 1: no active agent files if agents_dir.is_dir(): for f in agents_dir.iterdir(): if f.is_file() and f.name.endswith("_active.jsonl"): - return False + incomplete = True + break files = [f for f in segment_path.iterdir() if f.is_file()] - file_names = {f.name for f in files} file_suffixes = {f.suffix.lower() for f in files} + has_audio_raw = bool(file_suffixes & RAW_AUDIO_EXTENSIONS) + has_video_raw = bool(file_suffixes & RAW_VIDEO_EXTENSIONS) + audio_extracts = _matching_extraction_files( + segment_path, ("audio.jsonl", "*_audio.jsonl") + ) + screen_extracts = _matching_extraction_files( + segment_path, ("screen.jsonl", "*_screen.jsonl") + ) - # Check 2: audio transcript exists if audio was captured - if file_suffixes & RAW_AUDIO_EXTENSIONS: - has_audio_extract = "audio.jsonl" in file_names or any( - n.endswith("_audio.jsonl") for n in file_names - ) - if not has_audio_extract: - return False + if has_audio_raw and not audio_extracts: + incomplete = True - # Check 3: screen extract exists if video was captured - if file_suffixes & RAW_VIDEO_EXTENSIONS: - has_screen_extract = "screen.jsonl" in file_names or any( - n.endswith("_screen.jsonl") for n in file_names - ) - if not has_screen_extract: - return False + # monitor_*_diff.png files are raw media but have no extraction record. They + # ride the whole-segment gate and are only deleted when audio/video checks pass. + if has_video_raw and not screen_extracts: + incomplete = True - # Check 4: speaker labels exist if embeddings are present + speaker_labels = agents_dir / "speaker_labels.json" if ".npz" in file_suffixes: - if not agents_dir.is_dir() or not (agents_dir / "speaker_labels.json").exists(): - return False - - return True - - -def _get_completion_files(segment_path: Path) -> list[Path]: - """Return existing completion-indicating files for a segment.""" - completion_files: list[Path] = [] - - for name in ("audio.jsonl", "screen.jsonl"): - path = segment_path / name - if path.exists(): - completion_files.append(path) - - completion_files.extend( - path - for pattern in ("*_audio.jsonl", "*_screen.jsonl") - for path in segment_path.glob(pattern) - if path.is_file() + if not agents_dir.is_dir() or not speaker_labels.exists(): + incomplete = True + + incomplete = ( + _collect_extraction_states( + segment_path, + audio_extracts, + modality="audio", + marker_key="start", + has_raw=has_audio_raw, + failed_files=failed_files, + completion_files=completion_files, + ) + or incomplete + ) + incomplete = ( + _collect_extraction_states( + segment_path, + screen_extracts, + modality="screen", + marker_key="timestamp", + has_raw=has_video_raw, + failed_files=failed_files, + completion_files=completion_files, + ) + or incomplete ) - speaker_labels = segment_path / "talents" / "speaker_labels.json" + if failed_files: + return SegmentGate("failed", failed_files=failed_files) if speaker_labels.exists(): completion_files.append(speaker_labels) - - return completion_files + if incomplete: + return SegmentGate("incomplete", completion_files=completion_files) + return SegmentGate("eligible", completion_files=completion_files) # --------------------------------------------------------------------------- @@ -355,6 +469,9 @@ class PurgeResult: segments_processed: int = 0 segments_skipped_incomplete: int = 0 segments_skipped_policy: int = 0 + segments_blocked_failed: int = 0 + blocked_failed_details: list[dict[str, Any]] = field(default_factory=list) + partial_error: bool = False details: list[dict[str, Any]] = field(default_factory=list) @@ -403,13 +520,34 @@ def purge( result.segments_processed += 1 + gate = resolve_segment_gate(seg_path) + if gate.verdict == "failed": + result.segments_blocked_failed += 1 + blocked_detail = { + "day": day_name, + "stream": stream_name, + "segment": seg_key, + "files": gate.failed_files, + } + result.blocked_failed_details.append(blocked_detail) + logger.warning( + "retention: blocked purge - extraction failed: %s/%s/%s files=%s", + day_name, + stream_name, + seg_key, + gate.failed_files, + ) + continue + # Safety invariant: never delete from incomplete segments - if not is_segment_complete(seg_path): + if gate.verdict == "incomplete": result.segments_skipped_incomplete += 1 logger.debug( "Skipping incomplete: %s/%s/%s", day_name, stream_name, seg_key ) continue + if gate.verdict != "eligible": + raise RuntimeError(f"unexpected retention gate verdict: {gate.verdict}") # Check eligibility if older_than_days is not None: @@ -440,10 +578,9 @@ def purge( f.unlink() logger.info("Deleted: %s (%s)", f, _human_bytes(size)) - completion_files = _get_completion_files(seg_path) processed_at = None - if completion_files: - latest_mtime = max(f.stat().st_mtime for f in completion_files) + if gate.completion_files: + latest_mtime = max(f.stat().st_mtime for f in gate.completion_files) processed_at = datetime.fromtimestamp(latest_mtime).isoformat() result.files_deleted += len(raw_files) @@ -459,67 +596,68 @@ def purge( } ) - if not dry_run and result.files_deleted > 0: + if not dry_run and segment_files: + outcome = _write_segment_prune_audit( + journal_path, + day_name, + stream_name, + seg_key, + segment_files, + segment_bytes, + processed_at, + older_than_days, + ) + for audit_day, error in outcome.per_day_failures.items(): + result.partial_error = True + logger.warning( + "retention: failed to append pruning task log for %s: %s", + audit_day, + error, + ) + if outcome.global_record_error is not None: + result.partial_error = True + logger.warning( + "retention: failed to append pruning run record: %s", + outcome.global_record_error, + ) + + if not dry_run and (result.files_deleted > 0 or result.segments_blocked_failed > 0): _write_retention_log(journal_path, result) - _write_retention_prune_audit(journal_path, result, older_than_days) return result -def _raw_media_audit_by_day(result: PurgeResult) -> dict[str, dict[str, int]]: - by_day: dict[str, dict[str, int]] = {} - for detail in result.details: - day = str(detail["day"]) - entry = by_day.setdefault( - day, - {"files_deleted": 0, "bytes_freed": 0, "segments": 0}, - ) - entry["files_deleted"] += len(detail.get("files", [])) - entry["bytes_freed"] += int(detail.get("bytes_freed", 0)) - entry["segments"] += 1 - return by_day - - -def _raw_media_audit_day_messages( - by_day: dict[str, dict[str, int]], -) -> dict[str, str]: - return { - day: ( - "raw-media retention: pruned " - f"{stats['files_deleted']} raw media file(s) " - f"({_human_bytes(stats['bytes_freed'])}) from " - f"{stats['segments']} segment(s) for this day" - ) - for day, stats in sorted(by_day.items()) - } - - -def _write_retention_prune_audit( +def _write_segment_prune_audit( journal_path: Path, - result: PurgeResult, + day: str, + stream: str, + segment: str, + files: list[dict[str, Any]], + bytes_freed: int, + processed_at: str | None, older_than_days: int | None, -) -> None: - by_day = _raw_media_audit_by_day(result) +) -> AuditOutcome: run_record = { "timestamp": datetime.now().isoformat(), "kind": "raw_media", "dry_run": False, "days": older_than_days, - "by_day": by_day, - "totals": { - "files_deleted": result.files_deleted, - "bytes_freed": result.bytes_freed, - "segments_processed": result.segments_processed, - "segments_skipped_incomplete": result.segments_skipped_incomplete, - }, - "details": result.details, - "errors": [], + "day": day, + "stream": stream, + "segment": segment, + "files": files, + "bytes_freed": bytes_freed, + "processed_at": processed_at, } - write_prune_audit( + message = ( + f"raw-media retention: pruned {len(files)} raw media file(s) " + f"({_human_bytes(bytes_freed)}) from segment {stream}/{segment}" + ) + return write_prune_audit( journal_path, kind="raw_media", run_record=run_record, - per_day_messages=_raw_media_audit_day_messages(by_day), + per_day_messages={day: message}, ) @@ -535,6 +673,9 @@ def _write_retention_log(journal_path: Path, result: PurgeResult) -> None: "bytes_freed": result.bytes_freed, "segments_processed": result.segments_processed, "segments_skipped_incomplete": result.segments_skipped_incomplete, + "segments_blocked_failed": result.segments_blocked_failed, + "blocked_failed_details": result.blocked_failed_details, + "partial_error": result.partial_error, "details": result.details, } diff --git a/solstone/think/tools/call.py b/solstone/think/tools/call.py index a671f6a22..db2141bf2 100644 --- a/solstone/think/tools/call.py +++ b/solstone/think/tools/call.py @@ -925,6 +925,18 @@ def purge( "(not yet eligible under retention policy)." ) + if result.segments_blocked_failed: + typer.echo( + f"Blocked {result.segments_blocked_failed} segments " + "(extraction failed - raw media preserved)." + ) + for detail in result.blocked_failed_details: + failed_names = ", ".join(sorted(detail["files"])) + typer.echo( + f" {detail['day']}/{detail['stream']}/{detail['segment']}: " + f"{failed_names}" + ) + @retention_app.command() def config( diff --git a/tests/test_retention.py b/tests/test_retention.py index b311a07b3..587e43f97 100644 --- a/tests/test_retention.py +++ b/tests/test_retention.py @@ -8,7 +8,12 @@ import json import os import shutil from datetime import datetime +from pathlib import Path +import pytest +from typer.testing import CliRunner + +from solstone.observe.processing_record import STATE_EMPTY, STATE_FAILED from solstone.think.retention import ( RetentionConfig, RetentionPolicy, @@ -17,10 +22,63 @@ from solstone.think.retention import ( check_storage_health, get_raw_media_files, is_raw_media, - is_segment_complete, load_retention_config, purge, + resolve_segment_gate, ) +from solstone.think.tools.call import retention_app + + +def _write_jsonl(path: Path, *records: dict) -> None: + path.write_text( + "".join(json.dumps(record) + "\n" for record in records), + encoding="utf-8", + ) + + +def _processing_record(state: str) -> dict: + return { + "schema": "solstone.processing.v1", + "state": state, + "reason_code": "test", + "handler": "test", + "attempted_at": "2026-01-01T00:00:00Z", + "input_size": 1, + } + + +def _write_audio_success(path: Path, raw: str = "audio.flac") -> None: + _write_jsonl(path, {"raw": raw}, {"start": "00:00:00", "text": "ok"}) + + +def _write_screen_success(path: Path, raw: str = "screen.webm") -> None: + _write_jsonl(path, {"raw": raw}, {"timestamp": 0.0, "content": {}}) + + +def _write_processing_header(path: Path, raw: str, state: str) -> None: + _write_jsonl(path, {"raw": raw, "_solstone_processing": _processing_record(state)}) + + +def _install_test_journal( + journal: Path, monkeypatch, fixed_now: datetime | None = None +): + monkeypatch.setenv("SOLSTONE_JOURNAL", str(journal)) + fixed_now = fixed_now or datetime(2026, 4, 15) + + class FixedDateTime(datetime): + @classmethod + def now(cls, tz=None): + if tz is not None: + return fixed_now.replace(tzinfo=tz) + return fixed_now + + monkeypatch.setattr("solstone.think.retention.datetime", FixedDateTime) + monkeypatch.setattr("solstone.think.pruning_audit.datetime", FixedDateTime) + + import solstone.think.utils as think_utils + + think_utils._journal_path_cache = None + # --------------------------------------------------------------------------- # is_raw_media @@ -88,7 +146,7 @@ class TestGetRawMediaFiles: # --------------------------------------------------------------------------- -# is_segment_complete +# resolve_segment_gate # --------------------------------------------------------------------------- @@ -106,7 +164,7 @@ def _make_segment( ): """Create a segment directory with specified contents.""" seg = tmp_path / "segment" - seg.mkdir(exist_ok=True) + seg.mkdir(parents=True, exist_ok=True) agents_dir = seg / "talents" agents_dir.mkdir(exist_ok=True) @@ -117,9 +175,9 @@ def _make_segment( if embeddings: (seg / "audio.npz").write_bytes(b"npz") if audio and audio_extract: - (seg / "audio.jsonl").write_text('{"raw":"audio.flac"}\n') + _write_audio_success(seg / "audio.jsonl") if video and screen_extract: - (seg / "screen.jsonl").write_text('{"raw":"screen.webm"}\n') + _write_screen_success(seg / "screen.jsonl", video_name) if embeddings and speaker_labels: (agents_dir / "speaker_labels.json").write_text("{}") if active_agents: @@ -129,40 +187,40 @@ def _make_segment( return seg -class TestIsSegmentComplete: - def test_complete_audio_video(self, tmp_path): +class TestResolveSegmentGate: + def test_complete_audio_video_is_eligible(self, tmp_path): seg = _make_segment(tmp_path, audio=True, video=True, embeddings=True) - assert is_segment_complete(seg) + assert resolve_segment_gate(seg).verdict == "eligible" - def test_complete_audio_only(self, tmp_path): + def test_complete_audio_only_is_eligible(self, tmp_path): seg = _make_segment(tmp_path, audio=True) - assert is_segment_complete(seg) + assert resolve_segment_gate(seg).verdict == "eligible" - def test_complete_video_only(self, tmp_path): + def test_complete_video_only_is_eligible(self, tmp_path): seg = _make_segment(tmp_path, video=True) - assert is_segment_complete(seg) + assert resolve_segment_gate(seg).verdict == "eligible" def test_incomplete_missing_audio_extract(self, tmp_path): seg = _make_segment(tmp_path, audio=True, audio_extract=False) - assert not is_segment_complete(seg) + assert resolve_segment_gate(seg).verdict == "incomplete" def test_incomplete_missing_screen_extract(self, tmp_path): seg = _make_segment(tmp_path, video=True, screen_extract=False) - assert not is_segment_complete(seg) + assert resolve_segment_gate(seg).verdict == "incomplete" def test_incomplete_missing_screen_extract_for_mp4(self, tmp_path): seg = _make_segment( tmp_path, video=True, video_name="screen.mp4", screen_extract=False ) - assert not is_segment_complete(seg) + assert resolve_segment_gate(seg).verdict == "incomplete" def test_complete_mp4_with_screen_extract(self, tmp_path): seg = _make_segment(tmp_path, video=True, video_name="screen.mp4") - assert is_segment_complete(seg) + assert resolve_segment_gate(seg).verdict == "eligible" def test_incomplete_missing_speaker_labels(self, tmp_path): seg = _make_segment(tmp_path, audio=True, embeddings=True, speaker_labels=False) - assert not is_segment_complete(seg) + assert resolve_segment_gate(seg).verdict == "incomplete" def test_complete_with_stub_speaker_labels(self, tmp_path): """Stub speaker_labels.json (skipped=True, labels=[]) unblocks retention.""" @@ -171,26 +229,206 @@ class TestIsSegmentComplete: stub.write_text( json.dumps({"labels": [], "skipped": True, "reason": "no_owner_centroid"}) ) - assert is_segment_complete(seg) + assert resolve_segment_gate(seg).verdict == "eligible" def test_incomplete_active_agents(self, tmp_path): seg = _make_segment(tmp_path, audio=True, active_agents=True) - assert not is_segment_complete(seg) + assert resolve_segment_gate(seg).verdict == "incomplete" - def test_no_raw_media_is_complete(self, tmp_path): - """Segment with only derived content is considered complete.""" + def test_no_raw_media_is_eligible_with_terminal_extract(self, tmp_path): + """Segment with only terminal derived content is considered eligible.""" seg = tmp_path / "segment" seg.mkdir() - (seg / "audio.jsonl").write_text("transcript") + _write_audio_success(seg / "audio.jsonl") (seg / "stream.json").write_text("{}") - assert is_segment_complete(seg) + assert resolve_segment_gate(seg).verdict == "eligible" def test_no_agents_dir_is_ok(self, tmp_path): """No agents/ directory = no active agents = passes check 1.""" seg = tmp_path / "segment" seg.mkdir() (seg / "stream.json").write_text("{}") - assert is_segment_complete(seg) + assert resolve_segment_gate(seg).verdict == "eligible" + + def test_failed_audio_record_blocks(self, tmp_path): + seg = _make_segment(tmp_path, audio=True, audio_extract=False) + _write_processing_header(seg / "audio.jsonl", "audio.flac", STATE_FAILED) + + gate = resolve_segment_gate(seg) + + assert gate.verdict == "failed" + assert gate.failed_files == {"audio.jsonl": "failed"} + + def test_corrupt_analyzing_marker_blocks(self, tmp_path): + seg = _make_segment(tmp_path, audio=True, audio_extract=False) + _write_jsonl(seg / "audio.jsonl", {"raw": "audio.flac"}) + (seg / ".analyzing_audio").write_text("not json", encoding="utf-8") + + gate = resolve_segment_gate(seg) + + assert gate.verdict == "failed" + assert gate.failed_files == {"audio.jsonl": "failed"} + + def test_stale_analyzing_marker_blocks(self, tmp_path): + seg = _make_segment(tmp_path, audio=True, audio_extract=False) + _write_jsonl(seg / "audio.jsonl", {"raw": "audio.flac"}) + marker = seg / ".analyzing_audio" + marker.write_text( + '{"started_at": "2026-04-15T10:00:00Z", "modality": "audio"}\n', + encoding="utf-8", + ) + os.utime(marker, (0, 0)) + + gate = resolve_segment_gate(seg) + + assert gate.verdict == "failed" + assert gate.failed_files == {"audio.jsonl": "failed"} + + def test_failed_marker_blocks(self, tmp_path): + seg = _make_segment(tmp_path, audio=True, audio_extract=False) + _write_jsonl(seg / "audio.jsonl", {"raw": "audio.flac"}) + (seg / ".analyze_failed_audio").write_text( + '{"reason": "test"}\n', encoding="utf-8" + ) + + gate = resolve_segment_gate(seg) + + assert gate.verdict == "failed" + assert gate.failed_files == {"audio.jsonl": "failed"} + + def test_empty_audio_record_is_eligible(self, tmp_path): + seg = _make_segment(tmp_path, audio=True, audio_extract=False) + _write_processing_header(seg / "audio.jsonl", "audio.flac", STATE_EMPTY) + + gate = resolve_segment_gate(seg) + + assert gate.verdict == "eligible" + assert seg / "audio.jsonl" in gate.completion_files + + def test_legacy_chunk_bearing_no_record_is_eligible(self, tmp_path): + seg = _make_segment(tmp_path, audio=True) + + assert resolve_segment_gate(seg).verdict == "eligible" + + def test_audio_header_start_poisoning_does_not_mask_failure(self, tmp_path): + seg = _make_segment(tmp_path, audio=True, audio_extract=False) + header = { + "raw": "audio.flac", + "start": "not-a-chunk", + "_solstone_processing": _processing_record(STATE_FAILED), + } + _write_jsonl(seg / "audio.jsonl", header) + + gate = resolve_segment_gate(seg) + + assert gate.verdict == "failed" + assert gate.failed_files == {"audio.jsonl": "failed"} + + def test_screen_header_timestamp_poisoning_does_not_mask_failure(self, tmp_path): + seg = _make_segment(tmp_path, video=True, screen_extract=False) + header = { + "raw": "screen.webm", + "timestamp": "not-a-chunk", + "_solstone_processing": _processing_record(STATE_FAILED), + } + _write_jsonl(seg / "screen.jsonl", header) + + gate = resolve_segment_gate(seg) + + assert gate.verdict == "failed" + assert gate.failed_files == {"screen.jsonl": "failed"} + + def test_audio_analyzed_sibling_does_not_mask_failed_sibling(self, tmp_path): + seg = _make_segment(tmp_path, audio=True) + (seg / "meeting_audio.flac").write_bytes(b"audio") + _write_processing_header( + seg / "meeting_audio.jsonl", "meeting_audio.flac", STATE_FAILED + ) + + gate = resolve_segment_gate(seg) + + assert gate.verdict == "failed" + assert gate.failed_files == {"meeting_audio.jsonl": "failed"} + + def test_screen_analyzed_sibling_does_not_mask_failed_sibling(self, tmp_path): + seg = _make_segment(tmp_path, video=True, screen_extract=False) + (seg / "left_screen.webm").write_bytes(b"video") + (seg / "right_screen.webm").write_bytes(b"video") + _write_screen_success(seg / "left_screen.jsonl", "left_screen.webm") + _write_processing_header( + seg / "right_screen.jsonl", "right_screen.webm", STATE_FAILED + ) + + gate = resolve_segment_gate(seg) + + assert gate.verdict == "failed" + assert gate.failed_files == {"right_screen.jsonl": "failed"} + + def test_pending_and_analyzing_are_incomplete_not_blocked(self, tmp_path): + pending = _make_segment(tmp_path / "pending", audio=True, audio_extract=False) + _write_jsonl(pending / "audio.jsonl", {"raw": "audio.flac"}) + analyzing = _make_segment( + tmp_path / "analyzing", audio=True, audio_extract=False + ) + _write_jsonl(analyzing / "audio.jsonl", {"raw": "audio.flac"}) + (analyzing / ".analyzing_audio").write_text( + '{"started_at": "2026-04-15T10:00:00Z", "modality": "audio"}\n', + encoding="utf-8", + ) + + assert resolve_segment_gate(pending).verdict == "incomplete" + assert resolve_segment_gate(analyzing).verdict == "incomplete" + + def test_zero_nonblank_extraction_blocks_as_malformed(self, tmp_path): + seg = _make_segment(tmp_path, audio=True, audio_extract=False) + (seg / "audio.jsonl").write_text("\n\n", encoding="utf-8") + + gate = resolve_segment_gate(seg) + + assert gate.verdict == "failed" + assert gate.failed_files == {"audio.jsonl": "malformed"} + + def test_invalid_json_header_blocks_as_malformed(self, tmp_path): + seg = _make_segment(tmp_path, audio=True, audio_extract=False) + (seg / "audio.jsonl").write_text("not json\n", encoding="utf-8") + + gate = resolve_segment_gate(seg) + + assert gate.verdict == "failed" + assert gate.failed_files == {"audio.jsonl": "malformed"} + + def test_non_object_header_blocks_as_malformed(self, tmp_path): + seg = _make_segment(tmp_path, audio=True, audio_extract=False) + (seg / "audio.jsonl").write_text('"header"\n', encoding="utf-8") + + gate = resolve_segment_gate(seg) + + assert gate.verdict == "failed" + assert gate.failed_files == {"audio.jsonl": "malformed"} + + def test_unreadable_extraction_blocks_as_malformed(self, tmp_path, monkeypatch): + seg = _make_segment(tmp_path, audio=True, audio_extract=False) + target = seg / "audio.jsonl" + _write_audio_success(target) + original_open = Path.open + + def blocked_open(path, *args, **kwargs): + if path == target: + raise OSError("blocked") + return original_open(path, *args, **kwargs) + + monkeypatch.setattr(Path, "open", blocked_open) + + gate = resolve_segment_gate(seg) + + assert gate.verdict == "failed" + assert gate.failed_files == {"audio.jsonl": "malformed"} + + def test_monitor_diff_rides_segment_gate(self, tmp_path): + seg = _make_segment(tmp_path, audio=True) + (seg / "monitor_1_diff.png").write_bytes(b"diff") + + assert resolve_segment_gate(seg).verdict == "eligible" # --------------------------------------------------------------------------- @@ -298,14 +536,14 @@ class TestPurge: day1 = journal / "chronicle" / "20260115" / "default" / "100000_300" day1.mkdir(parents=True) (day1 / "audio.flac").write_bytes(b"x" * 1000) - (day1 / "audio.jsonl").write_text('{"raw":"audio.flac"}\n') + _write_audio_success(day1 / "audio.jsonl") (day1 / "stream.json").write_text('{"stream":"default"}') (day1 / "talents").mkdir() day1b = journal / "chronicle" / "20260115" / "plaud" / "103000_300" day1b.mkdir(parents=True) (day1b / "audio.m4a").write_bytes(b"x" * 500) - (day1b / "audio.jsonl").write_text('{"raw":"audio.m4a"}\n') + _write_audio_success(day1b / "audio.jsonl", "audio.m4a") (day1b / "stream.json").write_text('{"stream":"plaud"}') (day1b / "talents").mkdir() @@ -313,7 +551,7 @@ class TestPurge: day2 = journal / "chronicle" / "20260401" / "default" / "120000_300" day2.mkdir(parents=True) (day2 / "audio.flac").write_bytes(b"x" * 800) - (day2 / "audio.jsonl").write_text('{"raw":"audio.flac"}\n') + _write_audio_success(day2 / "audio.jsonl") (day2 / "stream.json").write_text('{"stream":"default"}') (day2 / "talents").mkdir() @@ -323,23 +561,7 @@ class TestPurge: (day3 / "audio.flac").write_bytes(b"x" * 600) (day3 / "stream.json").write_text('{"stream":"default"}') - monkeypatch.setenv("SOLSTONE_JOURNAL", str(journal)) - fixed_now = datetime(2026, 4, 15) - - class FixedDateTime(datetime): - @classmethod - def now(cls, tz=None): - if tz is not None: - return fixed_now.replace(tzinfo=tz) - return fixed_now - - monkeypatch.setattr("solstone.think.retention.datetime", FixedDateTime) - monkeypatch.setattr("solstone.think.pruning_audit.datetime", FixedDateTime) - # Clear cached journal path - import solstone.think.utils as think_utils - - think_utils._journal_path_cache = None - + _install_test_journal(journal, monkeypatch) return journal def test_dry_run(self, tmp_path, monkeypatch): @@ -356,7 +578,6 @@ class TestPurge: assert ( journal / "chronicle" / "20260115" / "plaud" / "103000_300" / "audio.m4a" ).exists() - # No retention log for dry run assert not (journal / "health" / "retention.log").exists() assert not (journal / "health" / "pruning-runs").exists() @@ -387,29 +608,76 @@ class TestPurge: assert (journal / "health" / "retention.log").exists() run_log = journal / "health" / "pruning-runs" / "20260415.jsonl" assert run_log.exists() - run_record = json.loads(run_log.read_text(encoding="utf-8").strip()) - assert run_record["kind"] == "raw_media" - assert run_record["dry_run"] is False - assert run_record["days"] == 30 - assert run_record["by_day"] == { - "20260115": { - "files_deleted": 2, - "bytes_freed": 1500, - "segments": 2, - } - } - assert run_record["totals"]["files_deleted"] == 2 - assert run_record["totals"]["bytes_freed"] == 1500 - assert run_record["totals"]["segments_skipped_incomplete"] == 1 - assert run_record["details"] == result.details - assert run_record["errors"] == [] + records = [ + json.loads(line) + for line in run_log.read_text(encoding="utf-8").splitlines() + if line.strip() + ] + assert len(records) == 2 + by_segment = {record["segment"]: record for record in records} + default_record = by_segment["100000_300"] + plaud_record = by_segment["103000_300"] + assert default_record["kind"] == "raw_media" + assert default_record["dry_run"] is False + assert default_record["days"] == 30 + assert default_record["day"] == "20260115" + assert default_record["stream"] == "default" + assert default_record["files"][0]["name"] == "audio.flac" + assert default_record["files"][0]["bytes"] == 1000 + assert default_record["bytes_freed"] == 1000 + assert isinstance(default_record["processed_at"], str) + assert plaud_record["stream"] == "plaud" + assert plaud_record["files"][0]["name"] == "audio.m4a" + assert plaud_record["files"][0]["bytes"] == 500 + assert plaud_record["bytes_freed"] == 500 task_log = journal / "chronicle" / "20260115" / "task_log.txt" task_line = task_log.read_text(encoding="utf-8") assert ( - "raw-media retention: pruned 2 raw media file(s) (1.5 KB) " - "from 2 segment(s) for this day" + "raw-media retention: pruned 1 raw media file(s) (1000 B) " + "from segment default/100000_300" + ) in task_line + assert ( + "raw-media retention: pruned 1 raw media file(s) (500 B) " + "from segment plaud/103000_300" ) in task_line + def test_mid_purge_interruption_preserves_prior_segment_audit( + self, tmp_path, monkeypatch + ): + journal = self._setup_journal(tmp_path, monkeypatch) + first_raw = ( + journal / "chronicle" / "20260115" / "default" / "100000_300" / "audio.flac" + ) + second_raw = ( + journal / "chronicle" / "20260115" / "plaud" / "103000_300" / "audio.m4a" + ) + original_unlink = Path.unlink + + def fail_second_unlink(path, *args, **kwargs): + if path == second_raw: + raise OSError("interrupted") + return original_unlink(path, *args, **kwargs) + + monkeypatch.setattr(Path, "unlink", fail_second_unlink) + + with pytest.raises(OSError, match="interrupted"): + purge(older_than_days=30, dry_run=False) + + assert not first_raw.exists() + assert second_raw.exists() + + run_log = journal / "health" / "pruning-runs" / "20260415.jsonl" + assert run_log.exists() + records = [ + json.loads(line) + for line in run_log.read_text(encoding="utf-8").splitlines() + if line.strip() + ] + assert len(records) == 1 + assert records[0]["stream"] == "default" + assert records[0]["segment"] == "100000_300" + assert records[0]["files"][0]["name"] == "audio.flac" + def test_skips_incomplete(self, tmp_path, monkeypatch): self._setup_journal(tmp_path, monkeypatch) @@ -442,6 +710,82 @@ class TestPurge: assert result.files_deleted == 1 assert result.details[0]["stream"] == "plaud" + def test_failed_extraction_blocks_purge_and_preserves_raw( + self, tmp_path, monkeypatch, caplog + ): + journal = tmp_path / "journal" + segment = journal / "chronicle" / "20260115" / "default" / "100000_300" + segment.mkdir(parents=True) + (segment / "audio.flac").write_bytes(b"x" * 1000) + _write_processing_header(segment / "audio.jsonl", "audio.flac", STATE_FAILED) + (segment / "stream.json").write_text('{"stream":"default"}') + (segment / "talents").mkdir() + _install_test_journal(journal, monkeypatch) + + with caplog.at_level("WARNING", logger="solstone.think.retention"): + result = purge(older_than_days=30, dry_run=False) + + assert result.files_deleted == 0 + assert result.segments_skipped_incomplete == 0 + assert result.segments_blocked_failed == 1 + assert result.blocked_failed_details == [ + { + "day": "20260115", + "stream": "default", + "segment": "100000_300", + "files": {"audio.jsonl": "failed"}, + } + ] + assert (segment / "audio.flac").exists() + assert "blocked purge" in caplog.text + + log_entry = json.loads( + (journal / "health" / "retention.log").read_text(encoding="utf-8") + ) + assert log_entry["segments_blocked_failed"] == 1 + assert log_entry["blocked_failed_details"] == result.blocked_failed_details + assert log_entry["partial_error"] is False + assert not (journal / "health" / "pruning-runs").exists() + + def test_audit_partial_error_is_surfaced(self, tmp_path, monkeypatch): + journal = self._setup_journal(tmp_path, monkeypatch) + + def fail_day_log(day, message): + raise OSError(f"blocked {day}") + + monkeypatch.setattr( + "solstone.think.pruning_audit.day_log_checked", fail_day_log + ) + + result = purge(older_than_days=30, dry_run=False) + + assert result.files_deleted == 2 + assert result.partial_error is True + log_entry = json.loads( + (journal / "health" / "retention.log").read_text(encoding="utf-8") + ) + assert log_entry["partial_error"] is True + + def test_cli_dry_run_surfaces_blocked_failed(self, tmp_path, monkeypatch): + journal = tmp_path / "journal" + segment = journal / "chronicle" / "20260115" / "default" / "100000_300" + segment.mkdir(parents=True) + (segment / "audio.flac").write_bytes(b"x" * 1000) + _write_processing_header(segment / "audio.jsonl", "audio.flac", STATE_FAILED) + (segment / "stream.json").write_text('{"stream":"default"}') + (segment / "talents").mkdir() + _install_test_journal(journal, monkeypatch) + + result = CliRunner().invoke( + retention_app, ["purge", "--older-than", "30d", "--dry-run"] + ) + + assert result.exit_code == 0 + assert "Blocked 1 segments" in result.output + assert "20260115/default/100000_300: audio.jsonl" in result.output + assert not (journal / "health" / "retention.log").exists() + assert not (journal / "health" / "pruning-runs").exists() + class TestPurgeProvenance: def _setup_journal(self, tmp_path, monkeypatch): @@ -502,7 +846,7 @@ class TestPurgeProvenance: alternate_audio_jsonl = segment / "meeting_audio.jsonl" speaker_labels = segment / "talents" / "speaker_labels.json" - alternate_audio_jsonl.write_text('{"raw":"audio.flac"}\n') + _write_audio_success(alternate_audio_jsonl, "meeting_audio.flac") speaker_labels.write_text("{}") older_ts = datetime(2026, 1, 15, 10, 0, 0).timestamp() -- 2.51.2