From abb6c28d9cd989c46f5ccd2ddfb2feda16389cfd Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Mon, 16 Mar 2026 14:40:57 -0600 Subject: [PATCH] =?UTF-8?q?Add=20media=20retention=20service=20=E2=80=94?= =?UTF-8?q?=20configurable=20raw=20media=20lifecycle?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Implements wave 1 of the media retention spec: retention config, service, and CLI commands. Three composable modes (keep, days, processed) with per-stream overrides and retroactive purge. Safety invariant: never deletes raw media from incomplete segments. Completion checks verify all processing is done before any deletion. New files: - think/retention.py — core retention service - tests/test_retention.py — 34 tests covering all retention logic CLI commands: - sol call journal retention purge --older-than 30d [--stream] [--dry-run] - sol call journal storage-summary [--json] Co-Authored-By: Claude Opus 4.6 (1M context) --- tests/test_retention.py | 370 +++++++++++++++++++++++++++++++++++++ think/journal_default.json | 5 + think/retention.py | 355 +++++++++++++++++++++++++++++++++++ think/tools/call.py | 101 ++++++++++ 4 files changed, 831 insertions(+) create mode 100644 tests/test_retention.py create mode 100644 think/retention.py diff --git a/tests/test_retention.py b/tests/test_retention.py new file mode 100644 index 000000000..e87510118 --- /dev/null +++ b/tests/test_retention.py @@ -0,0 +1,370 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Tests for think.retention — media retention service.""" + +import json +from pathlib import Path + +import pytest + +from think.retention import ( + RetentionConfig, + RetentionPolicy, + StorageSummary, + _human_bytes, + get_raw_media_files, + is_raw_media, + is_segment_complete, + load_retention_config, + purge, +) + + +# --------------------------------------------------------------------------- +# is_raw_media +# --------------------------------------------------------------------------- + + +class TestIsRawMedia: + def test_audio_extensions(self, tmp_path): + for ext in (".flac", ".opus", ".ogg", ".m4a"): + p = tmp_path / f"audio{ext}" + p.touch() + assert is_raw_media(p), f"{ext} should be raw media" + + def test_video_extensions(self, tmp_path): + for ext in (".webm", ".mov"): + p = tmp_path / f"screen{ext}" + p.touch() + assert is_raw_media(p), f"{ext} should be raw media" + + def test_monitor_diff_png(self, tmp_path): + p = tmp_path / "monitor_1_diff.png" + p.touch() + assert is_raw_media(p) + + p2 = tmp_path / "monitor_2_diff.png" + p2.touch() + assert is_raw_media(p2) + + def test_not_raw_media(self, tmp_path): + for name in ( + "audio.jsonl", + "screen.jsonl", + "stream.json", + "speaker_labels.json", + "audio.npz", + "summary.md", + "regular.png", + ): + p = tmp_path / name + p.touch() + assert not is_raw_media(p), f"{name} should NOT be raw media" + + +# --------------------------------------------------------------------------- +# get_raw_media_files +# --------------------------------------------------------------------------- + + +class TestGetRawMediaFiles: + def test_returns_only_raw(self, tmp_path): + (tmp_path / "audio.flac").write_bytes(b"x" * 100) + (tmp_path / "screen.webm").write_bytes(b"x" * 200) + (tmp_path / "audio.jsonl").write_text("transcript") + (tmp_path / "stream.json").write_text("{}") + + raw = get_raw_media_files(tmp_path) + names = {f.name for f in raw} + assert names == {"audio.flac", "screen.webm"} + + def test_empty_dir(self, tmp_path): + assert get_raw_media_files(tmp_path) == [] + + def test_nonexistent_dir(self, tmp_path): + assert get_raw_media_files(tmp_path / "nope") == [] + + +# --------------------------------------------------------------------------- +# is_segment_complete +# --------------------------------------------------------------------------- + + +def _make_segment(tmp_path, *, audio=False, video=False, embeddings=False, + audio_extract=True, screen_extract=True, speaker_labels=True, + active_agents=False): + """Create a segment directory with specified contents.""" + seg = tmp_path / "segment" + seg.mkdir(exist_ok=True) + agents_dir = seg / "agents" + agents_dir.mkdir(exist_ok=True) + + if audio: + (seg / "audio.flac").write_bytes(b"audio") + if video: + (seg / "screen.webm").write_bytes(b"video") + if embeddings: + (seg / "audio.npz").write_bytes(b"npz") + if audio and audio_extract: + (seg / "audio.jsonl").write_text('{"raw":"audio.flac"}\n') + if video and screen_extract: + (seg / "screen.jsonl").write_text('{"raw":"screen.webm"}\n') + if embeddings and speaker_labels: + (agents_dir / "speaker_labels.json").write_text("{}") + if active_agents: + (agents_dir / "1234_active.jsonl").write_text("{}") + + (seg / "stream.json").write_text('{"stream":"default"}') + return seg + + +class TestIsSegmentComplete: + def test_complete_audio_video(self, tmp_path): + seg = _make_segment(tmp_path, audio=True, video=True, embeddings=True) + assert is_segment_complete(seg) + + def test_complete_audio_only(self, tmp_path): + seg = _make_segment(tmp_path, audio=True) + assert is_segment_complete(seg) + + def test_complete_video_only(self, tmp_path): + seg = _make_segment(tmp_path, video=True) + assert is_segment_complete(seg) + + 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) + + 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) + + 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) + + def test_incomplete_active_agents(self, tmp_path): + seg = _make_segment(tmp_path, audio=True, active_agents=True) + assert not is_segment_complete(seg) + + def test_no_raw_media_is_complete(self, tmp_path): + """Segment with only derived content is considered complete.""" + seg = tmp_path / "segment" + seg.mkdir() + (seg / "audio.jsonl").write_text("transcript") + (seg / "stream.json").write_text("{}") + assert is_segment_complete(seg) + + 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) + + +# --------------------------------------------------------------------------- +# RetentionPolicy +# --------------------------------------------------------------------------- + + +class TestRetentionPolicy: + def test_keep_never_eligible(self): + p = RetentionPolicy(mode="keep") + assert not p.is_eligible(0) + assert not p.is_eligible(365) + + def test_processed_always_eligible(self): + p = RetentionPolicy(mode="processed") + assert p.is_eligible(0) + assert p.is_eligible(1) + + def test_days_threshold(self): + p = RetentionPolicy(mode="days", days=30) + assert not p.is_eligible(29) + assert p.is_eligible(30) + assert p.is_eligible(31) + + def test_days_no_value(self): + p = RetentionPolicy(mode="days", days=None) + assert not p.is_eligible(100) + + +class TestRetentionConfig: + def test_default_policy(self): + cfg = RetentionConfig() + assert cfg.policy_for_stream("default").mode == "keep" + + def test_per_stream_override(self): + cfg = RetentionConfig( + default=RetentionPolicy(mode="keep"), + per_stream={ + "archon.plaud": RetentionPolicy(mode="days", days=7), + }, + ) + assert cfg.policy_for_stream("archon.plaud").mode == "days" + assert cfg.policy_for_stream("archon.plaud").days == 7 + assert cfg.policy_for_stream("default").mode == "keep" + + +# --------------------------------------------------------------------------- +# load_retention_config +# --------------------------------------------------------------------------- + + +class TestLoadRetentionConfig: + def test_default_config(self, monkeypatch): + monkeypatch.setattr("think.utils.get_config", lambda: {}) + cfg = load_retention_config() + assert cfg.default.mode == "keep" + assert cfg.per_stream == {} + + def test_custom_config(self, monkeypatch): + monkeypatch.setattr( + "think.utils.get_config", + lambda: { + "retention": { + "raw_media": "days", + "raw_media_days": 30, + "per_stream": { + "default": {"raw_media": "processed"}, + }, + } + }, + ) + cfg = load_retention_config() + assert cfg.default.mode == "days" + assert cfg.default.days == 30 + assert cfg.per_stream["default"].mode == "processed" + + +# --------------------------------------------------------------------------- +# purge +# --------------------------------------------------------------------------- + + +class TestPurge: + def _setup_journal(self, tmp_path, monkeypatch): + """Create a journal structure with test segments.""" + journal = tmp_path / "journal" + + # Day 1: 60 days old — two complete segments + day1 = journal / "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') + (day1 / "stream.json").write_text('{"stream":"default"}') + (day1 / "agents").mkdir() + + day1b = journal / "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') + (day1b / "stream.json").write_text('{"stream":"plaud"}') + (day1b / "agents").mkdir() + + # Day 2: 10 days old — one complete segment + day2 = journal / "20260306" / "default" / "120000_300" + day2.mkdir(parents=True) + (day2 / "audio.flac").write_bytes(b"x" * 800) + (day2 / "audio.jsonl").write_text('{"raw":"audio.flac"}\n') + (day2 / "stream.json").write_text('{"stream":"default"}') + (day2 / "agents").mkdir() + + # Day 3: incomplete segment (no audio.jsonl) + day3 = journal / "20260101" / "default" / "140000_300" + day3.mkdir(parents=True) + (day3 / "audio.flac").write_bytes(b"x" * 600) + (day3 / "stream.json").write_text('{"stream":"default"}') + + monkeypatch.setenv("JOURNAL_PATH", str(journal)) + # Clear cached journal path + import think.utils + think.utils._journal_path_cache = None + + return journal + + def test_dry_run(self, tmp_path, monkeypatch): + journal = self._setup_journal(tmp_path, monkeypatch) + + result = purge(older_than_days=30, dry_run=True) + + # Should report but not delete + assert result.files_deleted == 2 # day1 default + plaud + assert result.bytes_freed == 1500 + assert (journal / "20260115" / "default" / "100000_300" / "audio.flac").exists() + assert (journal / "20260115" / "plaud" / "103000_300" / "audio.m4a").exists() + # No retention log for dry run + assert not (journal / "health" / "retention.log").exists() + + def test_actual_purge(self, tmp_path, monkeypatch): + journal = self._setup_journal(tmp_path, monkeypatch) + + result = purge(older_than_days=30, dry_run=False) + + assert result.files_deleted == 2 + # Files should be gone + assert not (journal / "20260115" / "default" / "100000_300" / "audio.flac").exists() + assert not (journal / "20260115" / "plaud" / "103000_300" / "audio.m4a").exists() + # Derived content preserved + assert (journal / "20260115" / "default" / "100000_300" / "audio.jsonl").exists() + # Retention log written + assert (journal / "health" / "retention.log").exists() + + def test_skips_incomplete(self, tmp_path, monkeypatch): + self._setup_journal(tmp_path, monkeypatch) + + result = purge(older_than_days=0, dry_run=True) + + # Day3 segment should be skipped (incomplete) + assert result.segments_skipped_incomplete == 1 + + def test_stream_filter(self, tmp_path, monkeypatch): + self._setup_journal(tmp_path, monkeypatch) + + result = purge(older_than_days=30, stream_filter="plaud", dry_run=True) + + assert result.files_deleted == 1 + assert result.details[0]["stream"] == "plaud" + + def test_policy_based_purge(self, tmp_path, monkeypatch): + self._setup_journal(tmp_path, monkeypatch) + + config = RetentionConfig( + default=RetentionPolicy(mode="keep"), + per_stream={ + "plaud": RetentionPolicy(mode="days", days=7), + }, + ) + + result = purge(dry_run=True, config=config) + + # Only plaud segment (60 days old) should be eligible + assert result.files_deleted == 1 + assert result.details[0]["stream"] == "plaud" + + +# --------------------------------------------------------------------------- +# _human_bytes +# --------------------------------------------------------------------------- + + +class TestHumanBytes: + def test_bytes(self): + assert _human_bytes(0) == "0 B" + assert _human_bytes(512) == "512 B" + + def test_kilobytes(self): + assert _human_bytes(1024) == "1.0 KB" + + def test_megabytes(self): + assert _human_bytes(1024 * 1024) == "1.0 MB" + + def test_gigabytes(self): + assert _human_bytes(1024 ** 3) == "1.0 GB" + + def test_large(self): + result = _human_bytes(12_400_000_000) + assert "GB" in result diff --git a/think/journal_default.json b/think/journal_default.json index bd5da666a..83723f697 100644 --- a/think/journal_default.json +++ b/think/journal_default.json @@ -31,5 +31,10 @@ "name_status": "default", "named_date": null, "proposal_count": 0 + }, + "retention": { + "raw_media": "keep", + "raw_media_days": null, + "per_stream": {} } } diff --git a/think/retention.py b/think/retention.py new file mode 100644 index 000000000..52d605265 --- /dev/null +++ b/think/retention.py @@ -0,0 +1,355 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Media retention service for solstone journals. + +Manages the lifecycle of raw media files (layer 1 captures) in journal segments. +Three retention modes: +- keep: retain raw media indefinitely (default) +- days: delete raw media after N days, once processing is complete +- processed: delete raw media as soon as processing completes + +Safety invariant: never delete raw media from segments that haven't finished +processing. All completion checks must pass before any deletion. +""" + +from __future__ import annotations + +import json +import logging +from dataclasses import dataclass, field +from datetime import datetime +from pathlib import Path +from typing import Any + +from think.utils import day_dirs, get_journal, iter_segments + +logger = logging.getLogger(__name__) + +# --------------------------------------------------------------------------- +# Raw media file identification +# --------------------------------------------------------------------------- + +RAW_AUDIO_EXTENSIONS = frozenset({".flac", ".opus", ".ogg", ".m4a"}) +RAW_VIDEO_EXTENSIONS = frozenset({".webm", ".mov"}) +RAW_MEDIA_EXTENSIONS = RAW_AUDIO_EXTENSIONS | RAW_VIDEO_EXTENSIONS + + +def is_raw_media(path: Path) -> bool: + """Check if a file is raw media (layer 1 capture). + + Raw media: *.flac, *.opus, *.ogg, *.m4a (audio), + *.webm, *.mov (video), monitor_*_diff.png (screen diffs). + """ + if path.suffix.lower() in RAW_MEDIA_EXTENSIONS: + return True + if ( + path.suffix.lower() == ".png" + and path.name.startswith("monitor_") + and "_diff" in path.name + ): + return True + return False + + +def get_raw_media_files(segment_path: Path) -> list[Path]: + """Return all raw media files in a segment directory.""" + if not segment_path.is_dir(): + return [] + return [f for f in segment_path.iterdir() if f.is_file() and is_raw_media(f)] + + +# --------------------------------------------------------------------------- +# Completion detection (safety invariant) +# --------------------------------------------------------------------------- + + +def is_segment_complete(segment_path: Path) -> bool: + """Check if a segment has finished all processing. + + Completion checks (ALL must pass): + 1. No _active.jsonl files in agents/ + 2. audio.jsonl exists if any audio raw media was captured + 3. screen.jsonl exists if any video raw media was captured + 4. agents/speaker_labels.json exists if embeddings (.npz) are present + """ + agents_dir = segment_path / "agents" + + # 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 + + 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} + + # 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 + + # 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 + + # Check 4: speaker labels exist if embeddings are present + if ".npz" in file_suffixes: + if not agents_dir.is_dir() or not (agents_dir / "speaker_labels.json").exists(): + return False + + return True + + +# --------------------------------------------------------------------------- +# Retention configuration +# --------------------------------------------------------------------------- + + +@dataclass +class RetentionPolicy: + """Retention policy for a single scope (global or per-stream).""" + + mode: str = "keep" # "keep", "days", or "processed" + days: int | None = None + + def is_eligible(self, segment_age_days: int) -> bool: + """Check if a segment's raw media should be purged under this policy.""" + if self.mode == "keep": + return False + if self.mode == "processed": + return True + if self.mode == "days" and self.days is not None: + return segment_age_days >= self.days + return False + + +@dataclass +class RetentionConfig: + """Retention configuration from journal.json.""" + + default: RetentionPolicy = field(default_factory=RetentionPolicy) + per_stream: dict[str, RetentionPolicy] = field(default_factory=dict) + + def policy_for_stream(self, stream: str) -> RetentionPolicy: + """Return the effective policy for a stream.""" + return self.per_stream.get(stream, self.default) + + +def load_retention_config() -> RetentionConfig: + """Load retention configuration from journal.json.""" + from think.utils import get_config + + config = get_config() + retention = config.get("retention", {}) + + mode = retention.get("raw_media", "keep") + days = retention.get("raw_media_days") + default = RetentionPolicy(mode=mode, days=days) + + per_stream: dict[str, RetentionPolicy] = {} + for stream_name, stream_config in retention.get("per_stream", {}).items(): + per_stream[stream_name] = RetentionPolicy( + mode=stream_config.get("raw_media", mode), + days=stream_config.get("raw_media_days", days), + ) + + return RetentionConfig(default=default, per_stream=per_stream) + + +# --------------------------------------------------------------------------- +# Storage summary +# --------------------------------------------------------------------------- + + +def _human_bytes(size: int) -> str: + """Format byte count as human-readable string.""" + n = float(size) + for unit in ("B", "KB", "MB", "GB", "TB"): + if abs(n) < 1024: + if unit == "B": + return f"{int(n)} B" + return f"{n:.1f} {unit}" + n /= 1024 + return f"{n:.1f} PB" + + +@dataclass +class StorageSummary: + """Storage usage summary for a journal.""" + + raw_media_bytes: int = 0 + derived_bytes: int = 0 + total_segments: int = 0 + segments_with_raw: int = 0 + segments_purged: int = 0 + + @property + def raw_media_human(self) -> str: + return _human_bytes(self.raw_media_bytes) + + @property + def derived_human(self) -> str: + return _human_bytes(self.derived_bytes) + + +def compute_storage_summary() -> StorageSummary: + """Compute storage summary across all journal segments.""" + summary = StorageSummary() + + for day_name in sorted(day_dirs().keys()): + for _stream, _seg_key, seg_path in iter_segments(day_name): + summary.total_segments += 1 + + raw_files = get_raw_media_files(seg_path) + summary.raw_media_bytes += sum(f.stat().st_size for f in raw_files) + + if raw_files: + summary.segments_with_raw += 1 + elif (seg_path / "audio.jsonl").exists() or ( + seg_path / "screen.jsonl" + ).exists(): + summary.segments_purged += 1 + + for f in seg_path.rglob("*"): + if f.is_file() and not is_raw_media(f): + summary.derived_bytes += f.stat().st_size + + return summary + + +# --------------------------------------------------------------------------- +# Retention purge +# --------------------------------------------------------------------------- + + +@dataclass +class PurgeResult: + """Result of a purge operation.""" + + files_deleted: int = 0 + bytes_freed: int = 0 + segments_processed: int = 0 + segments_skipped_incomplete: int = 0 + segments_skipped_policy: int = 0 + details: list[dict[str, Any]] = field(default_factory=list) + + +def purge( + *, + older_than_days: int | None = None, + stream_filter: str | None = None, + dry_run: bool = False, + config: RetentionConfig | None = None, +) -> PurgeResult: + """Run retention purge across the journal. + + Parameters + ---------- + older_than_days + Override: purge raw media older than this many days. + If None, uses the configured retention policy. + stream_filter + Only process segments from this stream. + dry_run + If True, report what would be deleted without deleting. + config + Retention config. If None, loads from journal.json. + """ + if config is None: + config = load_retention_config() + + result = PurgeResult() + today = datetime.now().date() + journal_path = Path(get_journal()) + + for day_name in sorted(day_dirs().keys()): + try: + day_date = datetime.strptime(day_name, "%Y%m%d").date() + except ValueError: + continue + age_days = (today - day_date).days + + for stream_name, seg_key, seg_path in iter_segments(day_name): + if stream_filter and stream_name != stream_filter: + continue + + raw_files = get_raw_media_files(seg_path) + if not raw_files: + continue + + result.segments_processed += 1 + + # Safety invariant: never delete from incomplete segments + if not is_segment_complete(seg_path): + result.segments_skipped_incomplete += 1 + logger.debug( + "Skipping incomplete: %s/%s/%s", day_name, stream_name, seg_key + ) + continue + + # Check eligibility + if older_than_days is not None: + eligible = age_days >= older_than_days + else: + policy = config.policy_for_stream(stream_name) + eligible = policy.is_eligible(age_days) + + if not eligible: + result.segments_skipped_policy += 1 + continue + + # Delete raw media + segment_bytes = 0 + segment_files = [] + for f in raw_files: + size = f.stat().st_size + segment_bytes += size + segment_files.append({"name": f.name, "bytes": size}) + if not dry_run: + f.unlink() + logger.info("Deleted: %s (%s)", f, _human_bytes(size)) + + result.files_deleted += len(raw_files) + result.bytes_freed += segment_bytes + result.details.append( + { + "day": day_name, + "stream": stream_name, + "segment": seg_key, + "files": segment_files, + "bytes_freed": segment_bytes, + } + ) + + if not dry_run and result.files_deleted > 0: + _write_retention_log(journal_path, result) + + return result + + +def _write_retention_log(journal_path: Path, result: PurgeResult) -> None: + """Append retention activity to health/retention.log.""" + health_dir = journal_path / "health" + health_dir.mkdir(parents=True, exist_ok=True) + log_path = health_dir / "retention.log" + + entry = { + "timestamp": datetime.now().isoformat(), + "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, + } + + with open(log_path, "a", encoding="utf-8") as f: + f.write(json.dumps(entry) + "\n") diff --git a/think/tools/call.py b/think/tools/call.py index c31f63280..11945893c 100644 --- a/think/tools/call.py +++ b/think/tools/call.py @@ -48,6 +48,8 @@ from think.utils import ( app = typer.Typer(help="Journal search and browsing.") facet_app = typer.Typer(help="Facet management.") app.add_typer(facet_app, name="facet") +retention_app = typer.Typer(help="Media retention management.") +app.add_typer(retention_app, name="retention") @app.command() @@ -667,3 +669,102 @@ def import_detail( pass typer.echo(json.dumps(info, indent=2, default=str)) + + +# ============================================================================ +# Retention Commands +# ============================================================================ + + +def _parse_age(value: str) -> int: + """Parse age string like '30d' or '30' to number of days.""" + value = value.strip().lower() + if value.endswith("d"): + return int(value[:-1]) + return int(value) + + +@retention_app.command() +def purge( + older_than: str | None = typer.Option( + None, "--older-than", help="Age threshold (e.g. 30d, 7d)." + ), + stream: str | None = typer.Option( + None, "--stream", help="Only purge from this stream." + ), + dry_run: bool = typer.Option(False, "--dry-run", help="Show what would be deleted."), +) -> None: + """Purge raw media from completed segments.""" + from think.retention import _human_bytes, load_retention_config, purge as run_purge + + older_than_days = _parse_age(older_than) if older_than else None + config = load_retention_config() + + if dry_run: + typer.echo("DRY RUN — no files will be deleted.\n") + + result = run_purge( + older_than_days=older_than_days, + stream_filter=stream, + dry_run=dry_run, + config=config, + ) + + if result.details: + for detail in result.details: + typer.echo( + f" {detail['day']}/{detail['stream']}/{detail['segment']}: " + f"{len(detail['files'])} files, {_human_bytes(detail['bytes_freed'])}" + ) + typer.echo("") + + action = "Would delete" if dry_run else "Deleted" + typer.echo( + f"{action} {result.files_deleted} files, " + f"freeing {_human_bytes(result.bytes_freed)}" + ) + + if result.segments_skipped_incomplete: + typer.echo( + f"Skipped {result.segments_skipped_incomplete} incomplete segments " + "(processing not finished)." + ) + + if result.segments_skipped_policy: + typer.echo( + f"Skipped {result.segments_skipped_policy} segments " + "(not yet eligible under retention policy)." + ) + + +@app.command(name="storage-summary") +def storage_summary( + json_output: bool = typer.Option(False, "--json", help="Output as JSON."), +) -> None: + """Show journal storage summary.""" + from think.retention import compute_storage_summary + + summary = compute_storage_summary() + + if json_output: + typer.echo( + json.dumps( + { + "raw_media_bytes": summary.raw_media_bytes, + "derived_bytes": summary.derived_bytes, + "total_segments": summary.total_segments, + "segments_with_raw": summary.segments_with_raw, + "segments_purged": summary.segments_purged, + }, + indent=2, + ) + ) + return + + typer.echo(f"Raw media: {summary.raw_media_human}") + typer.echo(f"AI-processed content: {summary.derived_human}") + typer.echo( + f"Segments: {summary.total_segments} total, " + f"{summary.segments_with_raw} with raw media, " + f"{summary.segments_purged} purged" + ) -- 2.51.2