diff --git a/solstone/think/catchup_state.py b/solstone/think/catchup_state.py index 8867a9a6f..93ac51a3e 100644 --- a/solstone/think/catchup_state.py +++ b/solstone/think/catchup_state.py @@ -313,6 +313,48 @@ def day_eligible_to_drain(day: str, kind: str) -> bool: return read_raw_input_fingerprint(day) != record.get("fingerprint") +def read_drain_hold_retry_at(day: str) -> float | None: + """Return the retry time only when ``day`` is positively held from draining. + + ``None`` means "not provably held"; callers must fall through to status-quo + behavior. If any gated record is active, active dominates and this returns + ``None`` regardless of any other held record. + """ + state = read_catchup_state() + entries = state["entries"] + records = [] + for kind in (KIND_DAILY_CATCHUP, KIND_SEGMENT_REPAIR): + record = entries.get(_key(day, kind)) + if isinstance(record, dict): + records.append(record) + + if any(record.get("active") for record in records): + return None + + now = time.time() + candidates = [] + for record in records: + try: + retry_at = float(record.get("next_retry_at") or 0) + except (TypeError, ValueError): + continue + if now < retry_at: + candidates.append((record, retry_at)) + + if not candidates: + return None + + fingerprint = read_raw_input_fingerprint(day) + held_retries = [ + retry_at + for record, retry_at in candidates + if record.get("fingerprint") == fingerprint + ] + if not held_retries: + return None + return max(held_retries) + + def read_backoff_summary(day: str) -> dict | None: record = read_day_record(day, KIND_DAILY_CATCHUP) if record is None or record.get("entered_backoff_at") is None: diff --git a/solstone/think/reprocess.py b/solstone/think/reprocess.py index 6ed633a6c..95822d8bd 100644 --- a/solstone/think/reprocess.py +++ b/solstone/think/reprocess.py @@ -12,6 +12,7 @@ from datetime import date, datetime, timedelta from enum import Enum from solstone.think.callosum import callosum_send +from solstone.think.catchup_state import read_drain_hold_retry_at from solstone.think.cluster import cluster_segments from solstone.think.streams import touch_stream_health_marker from solstone.think.utils import ( @@ -38,6 +39,7 @@ class ReprocessCode(Enum): FROM_SCRATCH_SUBMITTED = "from_scratch_submitted" MARK_UPDATED_SUBMITTED = "mark_updated_submitted" ALREADY_COMPLETE = "already_complete" + HELD_BY_BACKOFF = "held_by_backoff" PROCESS_NOW_SUBMITTED = "process_now_submitted" UNREACHABLE = "unreachable" @@ -45,6 +47,7 @@ class ReprocessCode(Enum): @dataclass(frozen=True) class ReprocessOutcome: code: ReprocessCode + when: str | None = None @dataclass(frozen=True) @@ -65,6 +68,9 @@ _CLI_STDOUT = { ReprocessCode.ALREADY_COMPLETE: ( "day {day} already complete; use --from-scratch to force a full re-run" ), + ReprocessCode.HELD_BY_BACKOFF: ( + "day {day} is held until {when}; use --from-scratch to start it over now" + ), } _CLI_STDERR = { @@ -77,6 +83,23 @@ _CLI_STDERR = { } +def _format_retry_clock(value: datetime) -> str: + return value.strftime("%I:%M%p").lstrip("0").lower() + + +def _format_retry_when(epoch_seconds: float) -> str: + retry = datetime.fromtimestamp(epoch_seconds) + retry_day = retry.date() + today = date.today() + if retry_day == today: + label = "today" + elif retry_day == today + timedelta(days=1): + label = "tomorrow" + else: + label = f"{retry:%b}".lower() + f" {retry.day}" + return f"{label} at {_format_retry_clock(retry)}" + + def reprocess_day(day: str, flavor: str) -> ReprocessOutcome: if not DATE_RE.fullmatch(day): return ReprocessOutcome(ReprocessCode.MALFORMED_DAY) @@ -116,6 +139,12 @@ def reprocess_day(day: str, flavor: str) -> ReprocessOutcome: if day_is_complete(day): return ReprocessOutcome(ReprocessCode.ALREADY_COMPLETE) + retry_at = read_drain_hold_retry_at(day) + if retry_at is not None: + return ReprocessOutcome( + ReprocessCode.HELD_BY_BACKOFF, when=_format_retry_when(retry_at) + ) + ok = callosum_send("supervisor", "drain", day=day) return ReprocessOutcome( ReprocessCode.PROCESS_NOW_SUBMITTED if ok else ReprocessCode.UNREACHABLE @@ -274,7 +303,7 @@ def main() -> None: outcome = reprocess_day(args.day, flavor) code = outcome.code if code in _CLI_STDOUT: - print(_CLI_STDOUT[code].format(day=args.day)) + print(_CLI_STDOUT[code].format(day=args.day, when=outcome.when)) return print(_CLI_STDERR[code].format(day=args.day), file=sys.stderr) diff --git a/tests/test_reprocess.py b/tests/test_reprocess.py index a8011f9fc..3171c4cf6 100644 --- a/tests/test_reprocess.py +++ b/tests/test_reprocess.py @@ -4,14 +4,16 @@ from __future__ import annotations import importlib +import json import os -from datetime import date, timedelta +import time +from datetime import date, datetime, timedelta from pathlib import Path from unittest.mock import Mock import pytest -from solstone.think import reprocess +from solstone.think import catchup_state, reprocess from tests.helpers.module_mocks import capturing_thread_constructor, module_mock DAY = "20250115" @@ -52,6 +54,68 @@ def _touch_marker(journal: Path, day: str, name: str, ns: int) -> Path: return marker +def _write_catchup_record( + journal: Path, + day: str, + kind: str, + *, + fingerprint: str | None, + next_retry_at: float, + active: dict | None = None, +) -> None: + state_path = journal / "health" / "catchup-state.json" + state_path.parent.mkdir(parents=True, exist_ok=True) + state_path.write_text( + json.dumps( + { + "version": catchup_state.STATE_VERSION, + "entries": { + f"{day}:{kind}": { + "day": day, + "command_kind": kind, + "attempts": 3, + "consecutive_non_completion": 3, + "last_attempt_at": 0, + "last_outcome": "timeout", + "next_retry_at": next_retry_at, + "entered_backoff_at": next_retry_at - 600, + "notified_at": next_retry_at - 600, + "fingerprint": fingerprint, + "active": active, + "reason_code": None, + "timeout_seconds": None, + "bounded": None, + "cleared": None, + "remaining": None, + "exit_reason": None, + "daily_progress": None, + } + }, + } + ), + encoding="utf-8", + ) + + +def test_format_retry_when_local_day_labels(): + today = date.today() + tomorrow = today + timedelta(days=1) + later = today + timedelta(days=2) + + same_day = datetime(today.year, today.month, today.day, 20, 42) + next_day = datetime(tomorrow.year, tomorrow.month, tomorrow.day, 20, 42) + midnight = datetime(tomorrow.year, tomorrow.month, tomorrow.day, 0, 3) + beyond_tomorrow = datetime(later.year, later.month, later.day, 20, 42) + + assert reprocess._format_retry_when(same_day.timestamp()) == "today at 8:42pm" + assert reprocess._format_retry_when(next_day.timestamp()) == "tomorrow at 8:42pm" + assert reprocess._format_retry_when(midnight.timestamp()) == "tomorrow at 12:03am" + assert ( + reprocess._format_retry_when(beyond_tomorrow.timestamp()) + == f"{beyond_tomorrow:%b}".lower() + f" {later.day} at 8:42pm" + ) + + class _ActiveTaskStub: def __init__(self, cmd): self.cmd = cmd @@ -151,6 +215,38 @@ def test_process_now_complete_day_is_noop_and_preserves_markers( assert (stream.stat().st_mtime_ns, daily.stat().st_mtime_ns) == before +def test_process_now_held_by_backoff_prints_stdout_and_exits_zero( + tmp_path, monkeypatch, capsys +): + journal = tmp_path / "journal" + segment = _seed_segment(journal) + (segment / "audio.jsonl").write_text("one\n", encoding="utf-8") + _touch_marker(journal, DAY, "stream.updated", 2_000_000_000) + monkeypatch.setenv("SOLSTONE_JOURNAL", str(journal)) + fingerprint = catchup_state.read_raw_input_fingerprint(DAY) + retry_at = time.time() + 3600 + _write_catchup_record( + journal, + DAY, + catchup_state.KIND_DAILY_CATCHUP, + fingerprint=fingerprint, + next_retry_at=retry_at, + ) + send = Mock(return_value=True) + monkeypatch.setattr(reprocess, "callosum_send", send) + + code, out, err = _invoke_reprocess(monkeypatch, capsys, journal, DAY) + + assert code == 0 + assert ( + out == "day " + f"{DAY} is held until {reprocess._format_retry_when(retry_at)}; " + "use --from-scratch to start it over now\n" + ) + assert err == "" + send.assert_not_called() + + def test_from_scratch_sends_request_and_preserves_marker(tmp_path, monkeypatch, capsys): journal = tmp_path / "journal" _seed_segment(journal) @@ -460,6 +556,111 @@ def test_reprocess_day_process_now_submitted(tmp_path, monkeypatch): send.assert_called_once_with("supervisor", "drain", day=DAY) +def test_reprocess_day_process_now_held_by_daily_backoff_returns_held_without_send( + tmp_path, monkeypatch +): + journal = tmp_path / "journal" + segment = _seed_segment(journal) + (segment / "audio.jsonl").write_text("one\n", encoding="utf-8") + _touch_marker(journal, DAY, "stream.updated", 2_000_000_000) + monkeypatch.setenv("SOLSTONE_JOURNAL", str(journal)) + fingerprint = catchup_state.read_raw_input_fingerprint(DAY) + retry_at = time.time() + 3600 + _write_catchup_record( + journal, + DAY, + catchup_state.KIND_DAILY_CATCHUP, + fingerprint=fingerprint, + next_retry_at=retry_at, + ) + send = Mock(return_value=True) + monkeypatch.setattr(reprocess, "callosum_send", send) + + outcome = reprocess.reprocess_day(DAY, reprocess.FLAVOR_PROCESS_NOW) + + assert (outcome.code.value, send.call_args_list) == ("held_by_backoff", []) + assert outcome.when + assert outcome.when == reprocess._format_retry_when(retry_at) + + +def test_reprocess_day_active_backoff_record_falls_through_and_sends( + tmp_path, monkeypatch +): + journal = tmp_path / "journal" + segment = _seed_segment(journal) + (segment / "audio.jsonl").write_text("one\n", encoding="utf-8") + _touch_marker(journal, DAY, "stream.updated", 2_000_000_000) + monkeypatch.setenv("SOLSTONE_JOURNAL", str(journal)) + fingerprint = catchup_state.read_raw_input_fingerprint(DAY) + _write_catchup_record( + journal, + DAY, + catchup_state.KIND_DAILY_CATCHUP, + fingerprint=fingerprint, + next_retry_at=time.time() + 3600, + active={"ref": "daily", "started_at": 1}, + ) + send = Mock(return_value=True) + monkeypatch.setattr(reprocess, "callosum_send", send) + + outcome = reprocess.reprocess_day(DAY, reprocess.FLAVOR_PROCESS_NOW) + + assert outcome.code is reprocess.ReprocessCode.PROCESS_NOW_SUBMITTED + send.assert_called_once_with("supervisor", "drain", day=DAY) + + +def test_reprocess_day_process_now_held_by_segment_repair_backoff( + tmp_path, monkeypatch +): + journal = tmp_path / "journal" + segment = _seed_segment(journal) + (segment / "audio.jsonl").write_text("one\n", encoding="utf-8") + _touch_marker(journal, DAY, "stream.updated", 2_000_000_000) + monkeypatch.setenv("SOLSTONE_JOURNAL", str(journal)) + fingerprint = catchup_state.read_raw_input_fingerprint(DAY) + retry_at = time.time() + 7200 + _write_catchup_record( + journal, + DAY, + catchup_state.KIND_SEGMENT_REPAIR, + fingerprint=fingerprint, + next_retry_at=retry_at, + ) + send = Mock(return_value=True) + monkeypatch.setattr(reprocess, "callosum_send", send) + + outcome = reprocess.reprocess_day(DAY, reprocess.FLAVOR_PROCESS_NOW) + + assert (outcome.code.value, send.call_args_list) == ("held_by_backoff", []) + assert outcome.when + assert outcome.when == reprocess._format_retry_when(retry_at) + + +def test_reprocess_day_backoff_fingerprint_change_sends(tmp_path, monkeypatch): + journal = tmp_path / "journal" + segment = _seed_segment(journal) + raw = segment / "audio.jsonl" + raw.write_text("one\n", encoding="utf-8") + _touch_marker(journal, DAY, "stream.updated", 2_000_000_000) + monkeypatch.setenv("SOLSTONE_JOURNAL", str(journal)) + fingerprint = catchup_state.read_raw_input_fingerprint(DAY) + _write_catchup_record( + journal, + DAY, + catchup_state.KIND_DAILY_CATCHUP, + fingerprint=fingerprint, + next_retry_at=time.time() + 3600, + ) + raw.write_text("two\n", encoding="utf-8") + send = Mock(return_value=True) + monkeypatch.setattr(reprocess, "callosum_send", send) + + outcome = reprocess.reprocess_day(DAY, reprocess.FLAVOR_PROCESS_NOW) + + assert outcome.code is reprocess.ReprocessCode.PROCESS_NOW_SUBMITTED + send.assert_called_once_with("supervisor", "drain", day=DAY) + + def test_reprocess_day_from_scratch_submitted_before_complete_check( tmp_path, monkeypatch ): diff --git a/tests/test_supervisor_schedule.py b/tests/test_supervisor_schedule.py index 2c6dbf856..7c663b593 100644 --- a/tests/test_supervisor_schedule.py +++ b/tests/test_supervisor_schedule.py @@ -3,14 +3,18 @@ """Test supervisor daily scheduling functionality.""" +import json import logging import os +import time from datetime import date from unittest.mock import MagicMock, Mock, call import pytest import solstone.think.supervisor as mod +from solstone.think import catchup_state +from solstone.think import utils as think_utils def _daily_think_calls(days): @@ -205,6 +209,72 @@ def test_force_day_blocked_when_not_eligible(mock_callosum, monkeypatch, submit_ assert submit_mock.call_args_list == _daily_think_calls(submitted) +def test_force_day_blocked_by_real_catchup_state_when_not_eligible( + tmp_path, mock_callosum, monkeypatch, submit_mock +): + day = "20250102" + journal = tmp_path / "journal" + monkeypatch.setenv("SOLSTONE_JOURNAL", str(journal)) + think_utils._journal_path_cache = None + config = journal / "config" / "journal.json" + config.parent.mkdir(parents=True) + config.write_text( + json.dumps( + { + "setup": {"completed_at": 1700000000000}, + "providers": {"active": {"provider": "openai"}}, + } + ), + encoding="utf-8", + ) + segment = journal / "chronicle" / day / "default" / "120000_300" + segment.mkdir(parents=True) + (segment / "audio.jsonl").write_text("one\n", encoding="utf-8") + marker = journal / "chronicle" / day / "health" / "stream.updated" + marker.parent.mkdir(parents=True) + marker.touch() + os.utime(marker, ns=(2_000_000_000, 2_000_000_000)) + fingerprint = catchup_state.read_raw_input_fingerprint(day) + retry_at = time.time() + 3600 + state_path = journal / "health" / "catchup-state.json" + state_path.parent.mkdir(parents=True) + state_path.write_text( + json.dumps( + { + "version": catchup_state.STATE_VERSION, + "entries": { + f"{day}:{catchup_state.KIND_DAILY_CATCHUP}": { + "day": day, + "command_kind": catchup_state.KIND_DAILY_CATCHUP, + "attempts": 3, + "consecutive_non_completion": 3, + "last_attempt_at": 0, + "last_outcome": "timeout", + "next_retry_at": retry_at, + "entered_backoff_at": retry_at - 600, + "notified_at": retry_at - 600, + "fingerprint": fingerprint, + "active": None, + "reason_code": None, + "timeout_seconds": None, + "bounded": None, + "cleared": None, + "remaining": None, + "exit_reason": None, + "daily_progress": None, + } + }, + } + ), + encoding="utf-8", + ) + + submitted = mod.run_catchup_drain(force_days={day}) + + assert submitted == [] + submit_mock.assert_not_called() + + def test_force_day_drains_when_eligible(mock_callosum, monkeypatch, submit_mock): monkeypatch.setattr( mod,