diff --git a/solstone/think/cogitate_policy.py b/solstone/think/cogitate_policy.py index f90c11ea8..ec6bb473c 100644 --- a/solstone/think/cogitate_policy.py +++ b/solstone/think/cogitate_policy.py @@ -33,10 +33,9 @@ MAX_TURNS_HEADROOM = 2 # output rate ($2.50 / 1M tokens) for ALL fresh non-cache tokens so the estimate # errs high and the ceiling trips early rather than late. _FALLBACK_USD_PER_TOKEN = 0.0000025 -DETERMINISTIC_FAILURE_THRESHOLD = 2 DEFAULT_READ_CALL_BUDGET = 200 # Reason codes for content-deterministic crashes and high-recurrence stochastic -# failures we decline to auto-retry past the threshold. +# failures we decline to auto-retry past their caps. DETERMINISTIC_FAILURE_REASON_CODES = frozenset( { "agent_stuck", @@ -48,6 +47,31 @@ DETERMINISTIC_FAILURE_REASON_CODES = frozenset( "wall_clock_exceeded", } ) +# Single source of truth for deterministic failure caps consumed by +# thinking._check_daily_skip and thinking.evaluate_daily_completion. The +# judgement shape is the latest deterministic reason plus the total in-set +# deterministic count from pipeline_health.read_daily_deterministic_failures. +# Scope-provided calibration: schema_invalid measured 24.3% per-call failure on +# the affected talent (87 complete / 28 schema_invalid since the local cutover), +# with same-day fail-then-pass observed on 20260723 for entity_observer:vconic +# (failed 00:24, completed 00:36). The other reasons are unmeasured and kept +# deliberately tight. +DETERMINISTIC_FAILURE_CAPS: dict[str, int] = { + "agent_stuck": 2, + "context_window_exceeded": 2, + "max_turns_exhausted": 2, + "no_output": 2, + "schema_invalid": 3, + "token_budget_exceeded": 2, + "wall_clock_exceeded": 2, +} + + +def failure_capped(reason_code: str | None, count: int) -> bool: + """Return True when a deterministic failure count reaches its cap.""" + cap = DETERMINISTIC_FAILURE_CAPS.get(reason_code or "") + return cap is not None and count >= cap + _JOURNAL_COMMANDS = {"identity", "health", "talent"} _SHELL_OPERATOR_CHARS = frozenset("();<>|&") diff --git a/solstone/think/doctor.py b/solstone/think/doctor.py index 7af07983b..b105a8029 100644 --- a/solstone/think/doctor.py +++ b/solstone/think/doctor.py @@ -904,6 +904,7 @@ def journal_caught_up_check(args: Args) -> CheckResult: check = JOURNAL_CAUGHT_UP_CHECK try: from solstone.think.pipeline_health import ( + BACKLOG_STATE_COMPLETE, BACKLOG_STATE_UNKNOWN, read_backlog_view, ) @@ -926,6 +927,18 @@ def journal_caught_up_check(args: Args) -> CheckResult: _CAUGHT_UP_CANT_TELL_FIX, ) + capped_complete_days = sum( + 1 + for day in view.days + if day.state == BACKLOG_STATE_COMPLETE and day.capped_daily_unit_count > 0 + ) + if view.pending_days == 0 and view.stuck_days == 0 and capped_complete_days > 0: + return make_result( + check, + "ok", + f"caught up; {capped_complete_days} day(s) completed with capped daily unit(s)", + ) + if view.pending_days == 0 and view.stuck_days == 0: return make_result(check, "ok", "caught up") diff --git a/solstone/think/journal_stats.py b/solstone/think/journal_stats.py index a9e13c8f8..61c88bd80 100644 --- a/solstone/think/journal_stats.py +++ b/solstone/think/journal_stats.py @@ -100,6 +100,9 @@ def _serialize_backlog_day(day: BacklogDay) -> dict: data["segment_repair_cleared"] = day.segment_repair_cleared if day.segment_repair_remaining is not None: data["segment_repair_remaining"] = day.segment_repair_remaining + if day.capped_daily_unit_count > 0: + data["capped_daily_unit_count"] = day.capped_daily_unit_count + data["capped_daily_unit"] = day.capped_daily_unit return data diff --git a/solstone/think/pipeline_health.py b/solstone/think/pipeline_health.py index 80b000a81..806450490 100644 --- a/solstone/think/pipeline_health.py +++ b/solstone/think/pipeline_health.py @@ -18,7 +18,10 @@ from solstone.think.catchup_state import ( read_segment_repair_summary, ) from solstone.think.cluster import cluster_segments -from solstone.think.cogitate_policy import DETERMINISTIC_FAILURE_REASON_CODES +from solstone.think.cogitate_policy import ( + DETERMINISTIC_FAILURE_REASON_CODES, + failure_capped, +) from solstone.think.data_state import DataState from solstone.think.utils import ( DEFAULT_STREAM, @@ -237,6 +240,8 @@ class BacklogDay: segment_repair_bounded: bool | None = None segment_repair_cleared: int | None = None segment_repair_remaining: int | None = None + capped_daily_unit_count: int = 0 + capped_daily_unit: dict[str, object] | None = None @dataclass(frozen=True) @@ -1401,6 +1406,28 @@ def _non_segment_failed_units( return tuple(why) +def _capped_daily_complete_fields(day: str) -> dict[str, object]: + capped = [ + { + "name": name, + "facet": facet, + "reason_code": failure.reason_code, + "count": failure.count, + } + for (name, facet), failure in sorted( + read_daily_deterministic_failures(day).items(), + key=lambda item: (item[0][0], item[0][1] or ""), + ) + if failure_capped(failure.reason_code, failure.count) + ] + if not capped: + return {} + return { + "capped_daily_unit_count": len(capped), + "capped_daily_unit": capped[0], + } + + def _complete_backlog_day(day: str) -> BacklogDay: return BacklogDay( day=day, @@ -1414,6 +1441,7 @@ def _complete_backlog_day(day: str) -> BacklogDay: provider=None, model=None, error=None, + **_capped_daily_complete_fields(day), ) @@ -1486,6 +1514,7 @@ def _backlog_day_for_complete(day: str, repair: dict | None) -> BacklogDay: provider=None, model=None, error=error, + **_capped_daily_complete_fields(day), **_segment_repair_fields(repair), ) diff --git a/solstone/think/thinking.py b/solstone/think/thinking.py index 5359af713..5c4c1926a 100644 --- a/solstone/think/thinking.py +++ b/solstone/think/thinking.py @@ -19,6 +19,7 @@ import sys import threading import time from concurrent.futures import ThreadPoolExecutor, as_completed +from dataclasses import dataclass from datetime import date, datetime, timedelta, timezone from pathlib import Path from typing import Any @@ -39,7 +40,7 @@ from solstone.think.catchup_state import ( ) from solstone.think.change_detection import detect_segment_change, resolve_predecessor from solstone.think.cluster import cluster_segments, read_segment_data_state -from solstone.think.cogitate_policy import DETERMINISTIC_FAILURE_THRESHOLD +from solstone.think.cogitate_policy import failure_capped from solstone.think.cortex_client import ( PATIENT_CLAIM_WINDOWS, CortexNotClaimed, @@ -906,6 +907,22 @@ def check_callosum_available() -> bool: _SKIPPED: object = object() +@dataclass(frozen=True) +class CappedDailyUnit: + name: str + facet: str | None + reason_code: str + count: int + + +@dataclass(frozen=True) +class DailyCompletionVerdict: + complete: bool + daily_units_terminal: bool + segment_blockers: tuple[dict[str, str], ...] + capped_daily_units: tuple[CappedDailyUnit, ...] + + class _NotClaimed: __slots__ = ("use_id",) @@ -1267,11 +1284,101 @@ def _check_daily_skip( return (True, "already_complete") if not retry_on_deterministic_failure: failure = deterministic_failures.get((name, facet)) - if failure is not None and failure.count >= DETERMINISTIC_FAILURE_THRESHOLD: + # Dispatch uses cogitate_policy.failure_capped, as does + # evaluate_daily_completion. retry_on_deterministic_failure affects + # dispatch only; completed day finalization ends retries, so forcing + # convergence-by-retry with this flag is unsupported. + if failure is not None and failure_capped(failure.reason_code, failure.count): return (True, "deterministic_failure_no_retry") return (False, None) +def evaluate_daily_completion( + applicable_units: set[tuple[str, str | None]], + completed_units: set[tuple[str, str, str | None]], + deterministic_failures: dict[tuple[str, str | None], DeterministicFailure], + segment_blockers: list[dict[str, str]] | tuple[dict[str, str], ...], +) -> DailyCompletionVerdict: + """Return the terminal daily completion verdict without journal writes.""" + capped_daily_units: list[CappedDailyUnit] = [] + daily_units_terminal = True + + # Invariant shared with _check_daily_skip via cogitate_policy.failure_capped: + # a unit's failure cap is terminal; the same predicate that stops dispatch + # marks the unit terminal-degraded for day completion. Day completion means + # all applicable units terminal, not all units succeeded. Degradation is + # terminal but visible. This deliberately ignores the dispatch retry override. + for name, facet in sorted( + applicable_units, key=lambda unit: (unit[0], unit[1] or "") + ): + if ("daily", name, facet) in completed_units: + continue + failure = deterministic_failures.get((name, facet)) + if failure is not None and failure_capped(failure.reason_code, failure.count): + capped_daily_units.append( + CappedDailyUnit( + name=name, + facet=facet, + reason_code=failure.reason_code, + count=failure.count, + ) + ) + continue + daily_units_terminal = False + + blockers = tuple(segment_blockers) + return DailyCompletionVerdict( + complete=daily_units_terminal and not blockers, + daily_units_terminal=daily_units_terminal, + segment_blockers=blockers, + capped_daily_units=tuple(capped_daily_units), + ) + + +def _capped_daily_unit_payload(unit: CappedDailyUnit) -> dict[str, object]: + return { + "name": unit.name, + "facet": unit.facet, + "reason_code": unit.reason_code, + "count": unit.count, + } + + +def finalize_day_completion( + day: str, verdict: DailyCompletionVerdict +) -> dict[str, object]: + """Write or withhold the daily marker and return daily_complete extras.""" + capped_payload = [ + _capped_daily_unit_payload(unit) for unit in verdict.capped_daily_units + ] + payload_fragment: dict[str, object] = {} + if capped_payload: + payload_fragment["capped_daily_units"] = capped_payload + + if verdict.complete: + health_dir = day_path(day) / "health" + health_dir.mkdir(parents=True, exist_ok=True) + (health_dir / "daily.updated").touch() + if capped_payload: + logging.info( + "Day %s complete with capped daily unit(s); wrote daily.updated: capped_daily_units=%s", + day, + capped_payload, + ) + else: + logging.info("Day %s fully complete; wrote daily.updated", day) + else: + logging.info( + "Day %s withholding daily.updated: daily_units_terminal=%s segment_blockers=%s capped_daily_units=%s", + day, + verdict.daily_units_terminal, + list(verdict.segment_blockers), + capped_payload, + ) + + return payload_fragment + + def run_segment_sense( day: str, segment: str, @@ -4704,32 +4811,22 @@ def main() -> None: ) # Touch daily.updated marker only after daily and segment work completes. + completion_payload_fragment: dict[str, object] = {} try: completed = read_completed_units(day) - daily_done = all( - ("daily", name, facet) in completed - for name, facet in applicable_units - ) + deterministic_failures = read_daily_deterministic_failures(day) segments = cluster_segments(day) progress = read_segment_progress(day) completion = classify_segment_completion(segments, progress) blocked_after_cycle = blocked_segment_keys(segments, progress) - blockers = completion.blockers - - if daily_done and not blockers: - health_dir = day_path(day) / "health" - health_dir.mkdir(parents=True, exist_ok=True) - (health_dir / "daily.updated").touch() - logging.info("Day %s fully complete; wrote daily.updated", day) - else: - logging.info( - "Day %s withholding daily.updated: " - "daily_units_complete=%s segment_blockers=%s", - day, - daily_done, - blockers, - ) + verdict = evaluate_daily_completion( + applicable_units, + completed, + deterministic_failures, + completion.blockers, + ) + completion_payload_fragment = finalize_day_completion(day, verdict) except Exception: logging.warning("Failed to update daily marker", exc_info=True) @@ -4749,6 +4846,8 @@ def main() -> None: ) if remaining_cycle is not None: daily_complete_payload["remaining"] = remaining_cycle + if completion_payload_fragment: + daily_complete_payload.update(completion_payload_fragment) # Set first_daily_ready awareness flag after first daily analysis try: diff --git a/tests/test_cogitate_policy.py b/tests/test_cogitate_policy.py index f6050ff60..c55b177a5 100644 --- a/tests/test_cogitate_policy.py +++ b/tests/test_cogitate_policy.py @@ -40,6 +40,23 @@ def test_resolve_read_scope_span_is_inclusive(): ) == ["chronicle/20260425", "chronicle/20260426", "chronicle/20260427"] +def test_failure_capped_schema_invalid_cap_is_three(): + assert cogitate_policy.failure_capped("schema_invalid", 2) is False + assert cogitate_policy.failure_capped("schema_invalid", 3) is True + + +def test_failure_capped_default_deterministic_cap_is_two(): + assert cogitate_policy.failure_capped("context_window_exceeded", 1) is False + assert cogitate_policy.failure_capped("context_window_exceeded", 2) is True + + +def test_deterministic_failure_caps_cover_reason_codes_exactly(): + assert ( + set(cogitate_policy.DETERMINISTIC_FAILURE_CAPS) + == cogitate_policy.DETERMINISTIC_FAILURE_REASON_CODES + ) + + def test_policy_denies_write_tools(tmp_path): policy = _policy(tmp_path) diff --git a/tests/test_journal_caught_up.py b/tests/test_journal_caught_up.py index 337376e7d..6dbb241be 100644 --- a/tests/test_journal_caught_up.py +++ b/tests/test_journal_caught_up.py @@ -229,6 +229,62 @@ def test_journal_caught_up_ok_for_all_complete_view(doctor, monkeypatch): assert result.detail == "caught up" +def test_journal_caught_up_ok_with_capped_complete_days_detail(doctor, monkeypatch): + day = BacklogDay( + day="20200229", + state=BACKLOG_STATE_COMPLETE, + segments=0, + units=0, + not_sensed=0, + why=(), + reason=None, + reason_code=None, + provider=None, + model=None, + error=None, + capped_daily_unit_count=1, + capped_daily_unit={ + "name": "entities:entity_observer", + "facet": "vconic", + "reason_code": "context_window_exceeded", + "count": 2, + }, + ) + view = BacklogView( + window=30, + days=(day,), + pending_days=0, + stuck_days=0, + oldest_pending_day=None, + errors=(), + ) + monkeypatch.setattr(pipeline_health, "read_backlog_view", lambda: view) + + result = doctor.journal_caught_up_check(args(doctor)) + + assert result.status == "ok" + assert result.detail == "caught up; 1 day(s) completed with capped daily unit(s)" + + +def test_journal_caught_up_plain_caught_up_without_capped_complete_days( + doctor, monkeypatch +): + view = BacklogView( + window=30, + days=(backlog_day("20200229", BACKLOG_STATE_COMPLETE),), + pending_days=0, + stuck_days=0, + oldest_pending_day=None, + errors=(), + ) + monkeypatch.setattr(pipeline_health, "read_backlog_view", lambda: view) + + result = doctor.journal_caught_up_check(args(doctor)) + + assert result.status == "ok" + assert result.detail == "caught up" + + def test_journal_caught_up_reports_backoff_stuck_day(doctor, tmp_path, monkeypatch): journal = tmp_path / "journal" day = "20990401" diff --git a/tests/test_journal_stats.py b/tests/test_journal_stats.py index 81712b7b2..77a82f10a 100644 --- a/tests/test_journal_stats.py +++ b/tests/test_journal_stats.py @@ -214,6 +214,58 @@ def test_serialize_backlog_day_includes_segment_repair_fields_only_when_present( assert progressing_data["segment_repair_remaining"] == 5 +def test_serialize_backlog_day_includes_capped_daily_fields_only_when_present(): + stats_mod = importlib.import_module("solstone.think.journal_stats") + health_mod = importlib.import_module("solstone.think.pipeline_health") + + clean = health_mod.BacklogDay( + day="20240101", + state=health_mod.BACKLOG_STATE_COMPLETE, + segments=0, + units=0, + not_sensed=0, + why=(), + reason=None, + reason_code=None, + provider=None, + model=None, + error=None, + ) + capped = health_mod.BacklogDay( + day="20240102", + state=health_mod.BACKLOG_STATE_COMPLETE, + segments=0, + units=0, + not_sensed=0, + why=(), + reason=None, + reason_code=None, + provider=None, + model=None, + error=None, + capped_daily_unit_count=1, + capped_daily_unit={ + "name": "entities:entity_observer", + "facet": "vconic", + "reason_code": "context_window_exceeded", + "count": 2, + }, + ) + + clean_data = stats_mod._serialize_backlog_day(clean) + capped_data = stats_mod._serialize_backlog_day(capped) + + assert "capped_daily_unit_count" not in clean_data + assert "capped_daily_unit" not in clean_data + assert capped_data["capped_daily_unit_count"] == 1 + assert capped_data["capped_daily_unit"] == { + "name": "entities:entity_observer", + "facet": "vconic", + "reason_code": "context_window_exceeded", + "count": 2, + } + + def test_scan_day(tmp_path, monkeypatch): stats_mod = importlib.import_module("solstone.think.journal_stats") journal = tmp_path diff --git a/tests/test_pipeline_health.py b/tests/test_pipeline_health.py index 49e375e9a..af0ddcf7b 100644 --- a/tests/test_pipeline_health.py +++ b/tests/test_pipeline_health.py @@ -395,6 +395,60 @@ def test_read_backlog_view_suppresses_stale_segment_repair_fingerprint( assert view.stuck_days == 0 +def test_backlog_complete_day_reports_capped_daily_unit_without_pending( + pipeline_journal, +): + capped_day = "20990310" + clean_day = "20990311" + capped_base = pipeline_journal / "chronicle" / capped_day / "health" + clean_base = pipeline_journal / "chronicle" / clean_day / "health" + _write_jsonl( + capped_base / "001_daily.jsonl", + [ + { + "event": "talent.fail", + "ts": 1, + "mode": "daily", + "name": "entities:entity_observer", + "facet": "vconic", + "reason_code": "context_window_exceeded", + }, + { + "event": "talent.fail", + "ts": 2, + "mode": "daily", + "name": "entities:entity_observer", + "facet": "vconic", + "reason_code": "context_window_exceeded", + }, + ], + ) + _write_jsonl( + clean_base / "001_daily.jsonl", + [{"event": "talent.complete", "ts": 1, "mode": "daily", "name": "alpha"}], + ) + for day in (capped_day, clean_day): + _touch_marker(pipeline_journal, day, "stream.updated", mtime_ms=1000) + _touch_marker(pipeline_journal, day, "daily.updated", mtime_ms=2000) + + view = read_backlog_view(window=2) + by_day = {item.day: item for item in view.days} + + assert by_day[capped_day].state == BACKLOG_STATE_COMPLETE + assert by_day[capped_day].capped_daily_unit_count == 1 + assert by_day[capped_day].capped_daily_unit == { + "name": "entities:entity_observer", + "facet": "vconic", + "reason_code": "context_window_exceeded", + "count": 2, + } + assert by_day[clean_day].state == BACKLOG_STATE_COMPLETE + assert by_day[clean_day].capped_daily_unit_count == 0 + assert by_day[clean_day].capped_daily_unit is None + assert view.pending_days == 0 + assert view.stuck_days == 0 + + def test_read_completed_units_terminal_presence(pipeline_journal): day = "20990202" base = pipeline_journal / "chronicle" / day / "health" diff --git a/tests/test_think_daily_idempotency.py b/tests/test_think_daily_idempotency.py index c46418f5a..f5b02861c 100644 --- a/tests/test_think_daily_idempotency.py +++ b/tests/test_think_daily_idempotency.py @@ -7,10 +7,19 @@ from __future__ import annotations import importlib import json +import os from pathlib import Path import pytest +from solstone.think.cogitate_policy import DETERMINISTIC_FAILURE_CAPS +from solstone.think.pipeline_health import ( + DeterministicFailure, + read_completed_units, + read_daily_deterministic_failures, +) +from solstone.think.utils import updated_days + DAY = "20990301" @@ -172,6 +181,170 @@ def test_check_daily_skip_has_no_freshness_inputs(): assert "from_scratch" in names +@pytest.mark.parametrize("retry_on_deterministic_failure", [False, True]) +@pytest.mark.parametrize("at_cap", [False, True]) +@pytest.mark.parametrize( + ("reason_code", "cap"), sorted(DETERMINISTIC_FAILURE_CAPS.items()) +) +def test_daily_skip_and_completion_cap_predicates_match( + retry_on_deterministic_failure, + at_cap, + reason_code, + cap, +): + mod = importlib.import_module("solstone.think.thinking") + count = cap if at_cap else cap - 1 + deterministic_failures = { + ("beta", None): DeterministicFailure(count=count, reason_code=reason_code) + } + + skip, _reason = mod._check_daily_skip( + "beta", + None, + mode="daily", + completed=set(), + deterministic_failures=deterministic_failures, + retry_on_deterministic_failure=retry_on_deterministic_failure, + ) + verdict = mod.evaluate_daily_completion( + {("beta", None)}, + set(), + deterministic_failures, + [], + ) + terminal_degraded = bool(verdict.capped_daily_units) + + if retry_on_deterministic_failure: + assert skip is False + assert terminal_degraded is at_cap + else: + assert skip is terminal_degraded + assert not (skip and not terminal_degraded) + + +def test_evaluate_daily_completion_terminal_cases(): + mod = importlib.import_module("solstone.think.thinking") + + complete_and_capped = mod.evaluate_daily_completion( + {("alpha", None), ("beta", None)}, + {("daily", "alpha", None)}, + { + ("beta", None): DeterministicFailure( + count=2, reason_code="context_window_exceeded" + ) + }, + [], + ) + assert complete_and_capped.complete is True + assert complete_and_capped.daily_units_terminal is True + assert complete_and_capped.capped_daily_units == ( + mod.CappedDailyUnit( + name="beta", + facet=None, + reason_code="context_window_exceeded", + count=2, + ), + ) + + below_cap = mod.evaluate_daily_completion( + {("beta", None)}, + set(), + { + ("beta", None): DeterministicFailure( + count=1, reason_code="context_window_exceeded" + ) + }, + [], + ) + assert below_cap.complete is False + assert below_cap.capped_daily_units == () + + +def test_evaluate_daily_completion_transient_latest_stays_incomplete(journal_copy): + mod = importlib.import_module("solstone.think.thinking") + day = "20990319" + _prepare_main_day(journal_copy, day) + _write_health( + journal_copy, + day, + "001_daily.jsonl", + [ + _complete("alpha"), + _fail("beta", ts=1, reason_code="context_window_exceeded"), + _fail("beta", ts=2, reason_code="context_window_exceeded"), + _fail("beta", ts=3, reason_code="schema_invalid"), + _fail("beta", ts=4, reason_code="token_budget_exceeded"), + _fail("beta", ts=5, reason_code="provider_transient"), + ], + ) + + completed = read_completed_units(day) + deterministic_failures = read_daily_deterministic_failures(day) + + assert ("beta", None) not in deterministic_failures + transient_latest = mod.evaluate_daily_completion( + {("alpha", None), ("beta", None)}, + completed, + deterministic_failures, + [], + ) + assert transient_latest.complete is False + assert transient_latest.capped_daily_units == () + + +def test_evaluate_daily_completion_dispatch_without_terminal_stays_incomplete( + journal_copy, +): + mod = importlib.import_module("solstone.think.thinking") + day = "20990320" + _prepare_main_day(journal_copy, day) + _write_health( + journal_copy, + day, + "001_daily.jsonl", + [ + _complete("alpha"), + {"event": "talent.dispatch", "ts": 2, "mode": "daily", "name": "beta"}, + ], + ) + + completed = read_completed_units(day) + deterministic_failures = read_daily_deterministic_failures(day) + + assert ("daily", "beta", None) not in completed + assert ("beta", None) not in deterministic_failures + dispatched_without_terminal = mod.evaluate_daily_completion( + {("alpha", None), ("beta", None)}, + completed, + deterministic_failures, + [], + ) + assert dispatched_without_terminal.complete is False + assert dispatched_without_terminal.capped_daily_units == () + + +def test_evaluate_daily_completion_withholds_with_segment_blockers(): + mod = importlib.import_module("solstone.think.thinking") + + verdict = mod.evaluate_daily_completion( + {("beta", None)}, + set(), + { + ("beta", None): DeterministicFailure( + count=2, reason_code="context_window_exceeded" + ) + }, + [{"segment": "090000_300", "dimension": "not_thought", "detail": "floor"}], + ) + + assert verdict.complete is False + assert verdict.daily_units_terminal is True + assert verdict.capped_daily_units + assert verdict.segment_blockers == ( + {"segment": "090000_300", "dimension": "not_thought", "detail": "floor"}, + ) + + def test_run_daily_prompts_skips_all_completed_units(daily_journal, monkeypatch): mod = importlib.import_module("solstone.think.thinking") _write_health( @@ -324,6 +497,68 @@ def test_run_daily_prompts_skips_two_deterministic_failures(daily_journal, monke ) +def test_run_daily_prompts_schema_invalid_retries_until_third_failure( + daily_journal, monkeypatch +): + mod = importlib.import_module("solstone.think.thinking") + _write_health( + daily_journal, + DAY, + "001_daily.jsonl", + [ + _fail("alpha", ts=1, reason_code="schema_invalid"), + _fail("alpha", ts=2, reason_code="schema_invalid"), + ], + ) + dispatched: list[tuple[str, dict]] = [] + _install_daily_mocks(monkeypatch, mod, _single_configs("alpha"), dispatched) + + _run_daily_with_writer(mod, daily_journal, DAY, "002_daily.jsonl") + + assert [name for name, _config in dispatched] == ["alpha"] + dispatched.clear() + _write_health( + daily_journal, + DAY, + "003_daily.jsonl", + [_fail("alpha", ts=3, reason_code="schema_invalid")], + ) + + _run_daily_with_writer(mod, daily_journal, DAY, "004_daily.jsonl") + + assert dispatched == [] + skips = _skip_events(daily_journal, DAY, "004_daily.jsonl") + assert skips[0]["reason"] == "deterministic_failure_no_retry" + assert skips[0]["detail"] == ( + "3 same-day deterministic failures (schema_invalid); not re-dispatching" + ) + + +def test_mixed_deterministic_reason_uses_latest_reason_cap(daily_journal, monkeypatch): + mod = importlib.import_module("solstone.think.thinking") + _write_health( + daily_journal, + DAY, + "001_daily.jsonl", + [ + _fail("alpha", ts=1, reason_code="context_window_exceeded"), + _fail("alpha", ts=2, reason_code="schema_invalid"), + _fail("beta", ts=1, reason_code="schema_invalid"), + _fail("beta", ts=2, reason_code="context_window_exceeded"), + ], + ) + dispatched: list[tuple[str, dict]] = [] + _install_daily_mocks(monkeypatch, mod, _single_configs("alpha", "beta"), dispatched) + + _run_daily_with_writer(mod, daily_journal, DAY, "002_daily.jsonl") + + assert [name for name, _config in dispatched] == ["alpha"] + skips = _skip_events(daily_journal, DAY, "002_daily.jsonl") + assert [(event["name"], event["reason"]) for event in skips] == [ + ("beta", "deterministic_failure_no_retry") + ] + + def test_run_daily_prompts_reruns_one_deterministic_failure(daily_journal, monkeypatch): mod = importlib.import_module("solstone.think.thinking") _write_health( @@ -669,6 +904,84 @@ def test_main_withholds_daily_marker_when_no_output_failure_is_incomplete( assert not (health / "daily.updated").exists() +def test_main_writes_daily_marker_when_capped_unit_terminal_and_payload_includes_capped_units( + journal_copy, monkeypatch +): + mod = importlib.import_module("solstone.think.thinking") + day = "20990316" + health = _prepare_main_day(journal_copy, day) + _write_health( + journal_copy, + day, + "001_daily.jsonl", + [ + _complete("alpha"), + _fail("beta", ts=1, reason_code="context_window_exceeded"), + _fail("beta", ts=2, reason_code="context_window_exceeded"), + ], + ) + _patch_main(monkeypatch, mod, {("alpha", None), ("beta", None)}) + emitted: list[tuple[str, dict]] = [] + monkeypatch.setattr( + mod, + "emit", + lambda event, **fields: emitted.append((event, fields)), + ) + monkeypatch.setattr("sys.argv", ["sol think", "--day", day]) + + mod.main() + + assert (health / "daily.updated").exists() + daily_complete = next( + fields for event, fields in emitted if event == "daily_complete" + ) + assert daily_complete["capped_daily_units"] == [ + { + "name": "beta", + "facet": None, + "reason_code": "context_window_exceeded", + "count": 2, + } + ] + + +def test_degraded_completion_clears_after_later_complete(journal_copy): + mod = importlib.import_module("solstone.think.thinking") + day = "20990317" + _prepare_main_day(journal_copy, day) + _write_health( + journal_copy, + day, + "001_daily.jsonl", + [ + _complete("alpha"), + _fail("beta", ts=1, reason_code="context_window_exceeded"), + _fail("beta", ts=2, reason_code="context_window_exceeded"), + ], + ) + + capped = mod.evaluate_daily_completion( + {("alpha", None), ("beta", None)}, + read_completed_units(day), + read_daily_deterministic_failures(day), + [], + ) + + assert capped.complete is True + assert capped.capped_daily_units + + _write_health(journal_copy, day, "002_daily.jsonl", [_complete("beta", ts=3)]) + cleared = mod.evaluate_daily_completion( + {("alpha", None), ("beta", None)}, + read_completed_units(day), + read_daily_deterministic_failures(day), + [], + ) + + assert cleared.complete is True + assert cleared.capped_daily_units == () + + def test_main_ignores_not_applicable_incomplete_units(journal_copy, monkeypatch): mod = importlib.import_module("solstone.think.thinking") day = "20990312" @@ -710,6 +1023,49 @@ def test_main_does_not_force_refresh_from_stream_marker(journal_copy, monkeypatc ] +def test_capped_terminal_completion_clears_updated_days_until_stream_newer( + journal_copy, monkeypatch +): + mod = importlib.import_module("solstone.think.thinking") + day = "20990318" + health = _prepare_main_day(journal_copy, day) + _write_health( + journal_copy, + day, + "001_daily.jsonl", + [ + _complete("alpha"), + _fail("beta", ts=1, reason_code="context_window_exceeded"), + _fail("beta", ts=2, reason_code="context_window_exceeded"), + ], + ) + (health / "stream.updated").touch() + dispatched: list[tuple[str, dict]] = [] + _install_daily_mocks(monkeypatch, mod, _single_configs("alpha", "beta"), dispatched) + monkeypatch.setattr( + mod, "run_bounded_phase", lambda *_args, **_kwargs: (True, False) + ) + monkeypatch.setattr(mod, "run_queued_command", lambda *_args, **_kwargs: True) + monkeypatch.setattr("sys.argv", ["sol think", "--day", day]) + + assert day in updated_days() + mod.main() + + assert dispatched == [] + assert (health / "daily.updated").exists() + assert day not in updated_days() + + os.utime(health / "daily.updated", (1000, 1000)) + os.utime(health / "stream.updated", (1010, 1010)) + assert day in updated_days() + monkeypatch.setattr("sys.argv", ["sol think", "--day", day]) + + mod.main() + + assert dispatched == [] + assert day not in updated_days() + + def test_main_passes_from_scratch_to_daily_prompts(journal_copy, monkeypatch): mod = importlib.import_module("solstone.think.thinking") day = "20990315"