From ffbae781a949be3cda7741767ce100c8eae2dc83 Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Sun, 26 Jul 2026 11:30:50 -0600 Subject: [PATCH] fix(retention): keep talent day indexes holding recent runs Delete a filename-old talent day index only when it has at least one valid row and every valid row's local execution date is older than the retention cutoff. Recent-filename indexes are never opened. Unreadable or malformed old indexes are preserved and recorded via unreadable_talent_index or malformed_talent_index_row. Empty old indexes are preserved as a plain skip. Map row ts values to local dates to match _parse_epoch_ms and the existing retention cutoff convention. --- solstone/think/log_retention.py | 119 +++++++++++++++++++++++++++++++- tests/test_log_retention.py | 28 +++++++- 2 files changed, 142 insertions(+), 5 deletions(-) diff --git a/solstone/think/log_retention.py b/solstone/think/log_retention.py index 2fedb0512..bae0d617a 100644 --- a/solstone/think/log_retention.py +++ b/solstone/think/log_retention.py @@ -5,6 +5,7 @@ from __future__ import annotations +import json import shutil from dataclasses import dataclass, field from datetime import date, datetime, timedelta @@ -111,6 +112,14 @@ class PruneResult: return messages +@dataclass(frozen=True) +class _TalentDayIndexInspection: + delete: bool + reason: str | None = None + message: str | None = None + hint: str | None = None + + def load_log_retention_config() -> LogRetentionConfig: """Load journal log/cache retention config with per-field defaults.""" config = get_config() @@ -389,16 +398,122 @@ def _scan_talent_day_index( path=target, result=result, ) - if day is not None and day < cutoff: + if day is None: + continue + if day >= cutoff: + continue + + formatted_day = _format_day(day) + inspection = _inspect_old_talent_day_index(target, cutoff) + if inspection.reason is not None: + _record_skip_error( + journal_path, + result, + "talent_day_index", + target, + formatted_day, + inspection.reason, + inspection.message or "talent day-index skipped during pruning", + inspection.hint, + ) + continue + if inspection.delete: _delete_target( journal_path, result, "talent_day_index", target, - _format_day(day), + formatted_day, target_kind="file", dry_run=dry_run, ) + continue + _mark_skipped(result, "talent_day_index") + + +def _inspect_old_talent_day_index(path: Path, cutoff: date) -> _TalentDayIndexInspection: + valid_row_days: list[date] = [] + try: + with path.open(encoding="utf-8") as handle: + for line_number, raw_line in enumerate(handle, start=1): + line = raw_line.strip() + if not line: + continue + try: + payload = json.loads(line) + except json.JSONDecodeError: + return _malformed_talent_index_row( + path, + f"line {line_number}: malformed JSON", + ) + if not isinstance(payload, dict): + return _malformed_talent_index_row( + path, + f"line {line_number}: non-object row", + ) + if "ts" not in payload: + return _malformed_talent_index_row( + path, + f"line {line_number}: missing ts", + ) + timestamp = payload.get("ts") + if isinstance(timestamp, bool) or not isinstance(timestamp, int): + return _malformed_talent_index_row( + path, + f"line {line_number}: non-integer ts", + ) + row_day = _talent_day_index_row_local_date(timestamp) + if row_day is None: + return _malformed_talent_index_row( + path, + f"line {line_number}: invalid ts", + ) + valid_row_days.append(row_day) + except FileNotFoundError: + return _TalentDayIndexInspection( + delete=False, + reason="missing_file_race", + message="talent day-index disappeared before pruning could inspect rows", + hint=None, + ) + except UnicodeDecodeError: + return _TalentDayIndexInspection( + delete=False, + reason="unreadable_talent_index", + message=f"failed to read talent day-index as UTF-8: {path.name}", + hint=None, + ) + except OSError as exc: + return _TalentDayIndexInspection( + delete=False, + reason="unreadable_talent_index", + message=f"failed to read talent day-index before pruning: {exc}", + hint=DELETE_FAILED_HINT, + ) + + if not valid_row_days: + return _TalentDayIndexInspection(delete=False) + if all(row_day < cutoff for row_day in valid_row_days): + return _TalentDayIndexInspection(delete=True) + return _TalentDayIndexInspection(delete=False) + + +def _malformed_talent_index_row(path: Path, detail: str) -> _TalentDayIndexInspection: + return _TalentDayIndexInspection( + delete=False, + reason="malformed_talent_index_row", + message=f"malformed talent day-index row in {path.name}: {detail}", + hint=None, + ) + + +def _talent_day_index_row_local_date(timestamp_ms: int) -> date | None: + # Log-retention cutoffs are local dates; match _parse_epoch_ms so a row and + # a run-log filename on the same timestamp prune on the same local day. + try: + return datetime.fromtimestamp(timestamp_ms / 1000).date() + except (OSError, OverflowError, ValueError): + return None def _scan_cogitate_history_cache( diff --git a/tests/test_log_retention.py b/tests/test_log_retention.py index d874935c9..c5e609696 100644 --- a/tests/test_log_retention.py +++ b/tests/test_log_retention.py @@ -341,6 +341,10 @@ def test_ac7_chronicle_allowlist_and_task_log_day_routing(journal): def test_ac8_ac9_talent_logs_indexes_and_malformed_names(journal): old_day = _day(31) + older_day = _day(32) + malformed_day = _day(33) + empty_day = _day(34) + invalid_utf8_day = _day(35) recent_day = _day(1) talent_dir = journal / "talents" / "default" old_run = _write(talent_dir / f"{_epoch_ms(old_day)}.jsonl") @@ -352,14 +356,26 @@ def test_ac8_ac9_talent_logs_indexes_and_malformed_names(journal): log_link = journal / "talents" / "linked.log" log_link.parent.mkdir(parents=True, exist_ok=True) log_link.symlink_to(log_target) - old_index = _write(journal / "talents" / f"{old_day}.jsonl") - recent_index = _write(journal / "talents" / f"{recent_day}.jsonl") + old_all_valid_old_index = _write( + journal / "talents" / f"{old_day}.jsonl", + json.dumps({"ts": int(_epoch_ms(old_day)), "status": "completed"}) + "\n", + ) + old_recent_ts_index = _write( + journal / "talents" / f"{older_day}.jsonl", + json.dumps({"ts": int(_epoch_ms(recent_day)), "status": "completed"}) + "\n", + ) + old_malformed_index = _write(journal / "talents" / f"{malformed_day}.jsonl") + old_empty_index = _write(journal / "talents" / f"{empty_day}.jsonl", "\n\n") + old_invalid_utf8_index = journal / "talents" / f"{invalid_utf8_day}.jsonl" + old_invalid_utf8_index.parent.mkdir(parents=True, exist_ok=True) + old_invalid_utf8_index.write_bytes(b"\xff") + recent_index = _write(journal / "talents" / f"{recent_day}.jsonl", "x") bad_index = _write(journal / "talents" / "20260231.jsonl") result = prune(config=LogRetentionConfig(days=30)) assert not old_run.exists() - assert not old_index.exists() + assert not old_all_valid_old_index.exists() for protected in ( recent_run, active, @@ -367,6 +383,10 @@ def test_ac8_ac9_talent_logs_indexes_and_malformed_names(journal): regular_log, log_link, log_target, + old_recent_ts_index, + old_malformed_index, + old_empty_index, + old_invalid_utf8_index, recent_index, bad_index, ): @@ -374,6 +394,8 @@ def test_ac8_ac9_talent_logs_indexes_and_malformed_names(journal): reasons = {error["reason"] for error in result.errors} assert "malformed_talent_timestamp" in reasons assert "malformed_date" in reasons + assert "malformed_talent_index_row" in reasons + assert "unreadable_talent_index" in reasons def test_ac10_ac11_symlink_unlinked_only_and_cache_mtime_day_logged(journal): -- 2.51.2