diff --git a/solstone/think/pipeline_health.py b/solstone/think/pipeline_health.py index ca4c587f4..607d1c6a3 100644 --- a/solstone/think/pipeline_health.py +++ b/solstone/think/pipeline_health.py @@ -9,7 +9,7 @@ import json import logging import os from dataclasses import dataclass -from datetime import datetime +from datetime import datetime, timedelta from pathlib import Path from solstone.think.cluster import cluster_segments @@ -29,7 +29,7 @@ logger = logging.getLogger(__name__) # Test indirection: tests monkeypatch this for time-sensitive branches. _now = datetime.now -_MODES = ("segment", "daily", "activity", "weekly", "flush") +_MODES = ("segment", "daily", "activity", "weekly", "flush", "cadence") _FAILED_LIST_CAP = 20 SEGMENT_FLOOR_TALENTS: tuple[str, ...] = ("entities", "documents") STUCK_FAIL_THRESHOLD = 3 @@ -111,6 +111,14 @@ class TerminalState: model: str | None +@dataclass(frozen=True) +class CompletionsSince: + """Completed segment/activity units newer than a timestamp, for cadence.""" + + segments: tuple[dict, ...] + activities: tuple[dict, ...] + + @dataclass(frozen=True) class DeterministicFailure: """A daily unit whose latest terminal is a deterministic crash.""" @@ -486,6 +494,61 @@ def read_completed_units(day: str) -> set[tuple[str, str, str | None]]: } +def read_completed_since(day: str, since_ms: int) -> CompletionsSince: + """Return unique completed segment/activity units newer than since_ms. + + Scans ``day`` and the prior day because post-midnight completions can + reference the previous day's health dir. Projects ``read_terminal_states`` + from per-talent identities to unique segment/activity units, each tagged + with the newest completion ts. + + This function does not create, modify, or delete journal state. + """ + prev = (datetime.strptime(day, "%Y%m%d") - timedelta(days=1)).strftime("%Y%m%d") + seg_max: dict[tuple[str | None, str], int] = {} + act_max: dict[tuple[str | None, str], int] = {} + + for scan_day in (day, prev): + for unit, state in read_terminal_states(scan_day).items(): + if state.latest_event != TERMINAL_COMPLETE or state.latest_ts <= since_ms: + continue + + if unit.segment: + seg_key = (unit.stream, unit.segment) + seg_max[seg_key] = max(seg_max.get(seg_key, 0), state.latest_ts) + elif unit.activity: + act_key = (unit.facet, unit.activity) + act_max[act_key] = max(act_max.get(act_key, 0), state.latest_ts) + + segments = tuple( + sorted( + ( + {"stream": stream, "segment": segment, "ts": ts} + for (stream, segment), ts in seg_max.items() + ), + key=lambda item: ( + item["ts"], + item["stream"] or "", + item["segment"], + ), + ) + ) + activities = tuple( + sorted( + ( + {"facet": facet, "activity": activity, "ts": ts} + for (facet, activity), ts in act_max.items() + ), + key=lambda item: ( + item["ts"], + item["facet"] or "", + item["activity"], + ), + ) + ) + return CompletionsSince(segments=segments, activities=activities) + + def read_daily_deterministic_failures( day: str, ) -> dict[tuple[str, str | None], DeterministicFailure]: diff --git a/solstone/think/thinking.py b/solstone/think/thinking.py index aaa9f8370..7939bf407 100644 --- a/solstone/think/thinking.py +++ b/solstone/think/thinking.py @@ -46,6 +46,7 @@ from solstone.think.pipeline_health import ( SEGMENT_FLOOR_TALENTS, DeterministicFailure, classify_segment_completion, + read_completed_since, read_completed_units, read_daily_deterministic_failures, read_segment_progress, @@ -123,6 +124,26 @@ def _jsonl_log(event: str, **fields) -> None: _jsonl.log(event, **fields) +def load_cadence_state() -> dict[str, int]: + """Read health/cadence.json per-talent last-run timestamps.""" + path = Path(get_journal()) / "health" / "cadence.json" + if not path.exists(): + return {} + try: + with path.open("r", encoding="utf-8") as handle: + data = json.load(handle) + return data if isinstance(data, dict) else {} + except (json.JSONDecodeError, OSError) as exc: + logging.warning("Failed to load cadence state: %s", exc) + return {} + + +def save_cadence_state(state: dict[str, int]) -> None: + """Persist cadence state to health/cadence.json atomically.""" + path = Path(get_journal()) / "health" / "cadence.json" + atomic_replace(path, json.dumps(state, indent=2)) + + def _provider_model_fields(use_id: str) -> dict[str, str | None]: provider, model, reason_code = read_use_provider_model_reason(use_id) fields = {"provider": provider, "model": model} @@ -2081,6 +2102,119 @@ def run_weekly_prompts( return (total_success, total_failed, all_failed_names) +def run_cadence_prompts( + day: str, + refresh: bool, + verbose: bool, + max_concurrency: int = 2, + stream: str | None = None, + timeout: int | None = 610, +) -> tuple[int, int, list[str]]: + """Run cadence-scheduled prompts whose completion gate is open.""" + all_prompts = get_talent_configs(schedule="cadence") + if not all_prompts: + logging.info("cadence: no cadence talents configured") + return (0, 0, []) + + cadence_state = load_cadence_state() + dirty = False + total_success = 0 + total_failed = 0 + failed_names: list[str] = [] + fired = 0 + skipped = 0 + + for name, config in sorted( + all_prompts.items(), key=lambda item: (item[1]["priority"], item[0]) + ): + now = now_ms() + cadence_minutes = config.get("cadence_minutes", 5) + last = cadence_state.get(name) + if last is not None and now - last < cadence_minutes * 60_000: + _log_skip( + name, + "interval_not_elapsed", + f"{(now - last) // 1000}s since last < {cadence_minutes}m", + mode="cadence", + day=day, + ) + skipped += 1 + continue + + since_ms = last or 0 + window = read_completed_since(day, since_ms) + if not window.segments and not window.activities: + _log_skip( + name, + "no_new_work", + "no segment/activity completed since last cadence run", + mode="cadence", + day=day, + ) + skipped += 1 + continue + + is_generate = config["type"] == "generate" + request_config: dict = { + "day": day, + "schedule": "cadence", + "env": {"SOL_DAY": day}, + "cadence_window": { + "since_ms": since_ms, + "segments": list(window.segments), + "activities": list(window.activities), + }, + } + _apply_output_persistence(request_config, config, force_refresh=refresh) + prompt = "" if is_generate else f"Running cadence task for {iso_date(day)}." + + use_id = _cortex_request_with_retry( + prompt=prompt, + name=name, + config=request_config, + ) + if use_id is None: + _log_skip( + name, + "send_failed", + f"All cortex request attempts failed for {name}", + mode="cadence", + day=day, + ) + total_failed += 1 + failed_names.append(f"{name} (send)") + continue + + emit("talent_started", mode="cadence", day=day, name=name, use_id=use_id) + _jsonl_log( + "talent.dispatch", + mode="cadence", + day=day, + name=name, + use_id=use_id, + ) + s, f, fn = _drain_priority_batch( + [(use_id, name, config, None)], "cadence", day, None, stream, timeout + ) + total_success += s + total_failed += f + failed_names.extend(fn) + if s == 1 and f == 0: + cadence_state[name] = now + dirty = True + fired += 1 + + if dirty: + save_cadence_state(cadence_state) + logging.info( + "cadence: %d fired, %d skipped (no new work or interval), %d failed", + fired, + skipped, + total_failed, + ) + return (total_success, total_failed, failed_names) + + def run_activity_prompts( day: str, activity_id: str, @@ -2699,6 +2833,7 @@ def dry_run( refresh: bool = False, stream: str | None = None, weekly: bool = False, + cadence: bool = False, ) -> None: """Print what think would execute without spawning any agents.""" day_formatted = iso_date(day) @@ -2768,6 +2903,36 @@ def dry_run( _print_prompt_table(all_prompts, day, refresh=refresh, stream=stream) return + if cadence: + all_prompts = get_talent_configs(schedule="cadence") + print(f"Day {day_formatted} — cadence agents\n") + if not all_prompts: + print("No prompts for schedule: cadence") + return + cadence_state = load_cadence_state() + now = now_ms() + for name, config in sorted( + all_prompts.items(), key=lambda item: (item[1]["priority"], item[0]) + ): + cadence_minutes = config.get("cadence_minutes", 5) + last = cadence_state.get(name) + if last is not None and now - last < cadence_minutes * 60_000: + print( + f" skip {name} — interval not elapsed " + f"({(now - last) // 1000}s < {cadence_minutes}m)" + ) + continue + window = read_completed_since(day, last or 0) + count = len(window.segments) + len(window.activities) + if count == 0: + print(f" no-op {name} — no new work since last cadence run") + else: + print( + f" fire {name} — window: {len(window.segments)} segment(s), " + f"{len(window.activities)} activity(ies)" + ) + return + if segments: segs = cluster_segments(day) if not segs: @@ -2994,7 +3159,7 @@ def parse_args() -> argparse.ArgumentParser: ) parser.add_argument( "--day", - help="Day folder in YYYYMMDD format (defaults to yesterday)", + help="Day folder in YYYYMMDD format (defaults to yesterday, or today with --cadence)", ) parser.add_argument( "--segment", @@ -3090,6 +3255,11 @@ def parse_args() -> argparse.ArgumentParser: action="store_true", help="Run weekly-scheduled agents (incompatible with --segment, --segments, --activity, --flush)", ) + parser.add_argument( + "--cadence", + action="store_true", + help="Run cadence-scheduled agents on completed segments/activities (incompatible with --segment, --segments, --activity, --flush, --weekly)", + ) parser.add_argument( "--dry-run", action="store_true", @@ -3123,6 +3293,8 @@ def main() -> None: incompatible.append("--flush") if args.segments: incompatible.append("--segments") + if args.cadence: + incompatible.append("--cadence") if incompatible: parser.error(f"--updated is incompatible with {', '.join(incompatible)}") today = date.today().strftime("%Y%m%d") @@ -3132,7 +3304,11 @@ def main() -> None: day = args.day if day is None: - day = (datetime.now() - timedelta(days=1)).strftime("%Y%m%d") + day = ( + date.today().strftime("%Y%m%d") + if args.cadence + else (datetime.now() - timedelta(days=1)).strftime("%Y%m%d") + ) day_dir = day_path(day) if not day_dir.is_dir(): @@ -3173,6 +3349,13 @@ def main() -> None: "--weekly is incompatible with --segment, --segments, --activity, and --flush" ) + if args.cadence and ( + args.segment or args.segments or args.activity or args.flush or args.weekly + ): + parser.error( + "--cadence is incompatible with --segment, --segments, --activity, --flush, and --weekly" + ) + if args.dry_run: dry_run( day, @@ -3184,6 +3367,7 @@ def main() -> None: refresh=args.refresh, stream=args.stream, weekly=args.weekly, + cadence=args.cadence, ) sys.exit(0) @@ -3195,11 +3379,17 @@ def main() -> None: _run_mode = "segment" elif args.weekly: _run_mode = "weekly" + elif args.cadence: + _run_mode = "cadence" elif args.segment: _run_mode = "segment" else: _run_mode = "daily" + if args.cadence and not get_talent_configs(schedule="cadence"): + logging.info("cadence: no cadence talents configured") + sys.exit(0) + _run_ref = str(now_ms()) _run_start_time = time.time() _run_result = {"success": 0, "failed": 0} @@ -3354,6 +3544,31 @@ def main() -> None: sys.exit(1) sys.exit(0) + # Handle cadence mode — dispatch only agents whose completion gate is open + if args.cadence: + success_count, fail_count, failed_names = run_cadence_prompts( + day=day, + refresh=args.refresh, + verbose=args.verbose, + max_concurrency=args.jobs, + stream=args.stream, + ) + + duration_ms = int((time.time() - start_time) * 1000) + logging.info( + f"Cadence think completed in {duration_ms}ms: " + f"{success_count} succeeded, {fail_count} failed" + ) + day_log(day, f"think --cadence failed={fail_count}") + _run_result["success"] = success_count + _run_result["failed"] = fail_count + + if fail_count > 0: + names = ", ".join(failed_names) + logging.error(f"{fail_count} cadence prompt(s) failed: {names}") + sys.exit(1) + sys.exit(0) + # PRE-PHASE: Run sense repair (daily only) if not args.segment: logging.info("Running pre-phase: sense repair") diff --git a/tests/test_pipeline_health.py b/tests/test_pipeline_health.py index f6acd2abf..04fefcf65 100644 --- a/tests/test_pipeline_health.py +++ b/tests/test_pipeline_health.py @@ -16,9 +16,11 @@ import pytest from solstone.think.pipeline_health import ( STUCK_FAIL_THRESHOLD, + CompletionsSince, TerminalUnit, pipeline_status_message, read_backlog_view, + read_completed_since, read_completed_units, read_daily_deterministic_failures, read_day_stuck, @@ -576,6 +578,151 @@ def test_read_completed_units_returns_old_daily_tuple_shape_and_filters_scoped_u assert all(isinstance(unit, tuple) and len(unit) == 3 for unit in completed) +def test_read_completed_since_missing_health_dirs(pipeline_journal): + assert read_completed_since("20990210", 0) == CompletionsSince((), ()) + + +def test_read_completed_since_dedups_segment_completions_at_max_ts(pipeline_journal): + day = "20990211" + base = pipeline_journal / "chronicle" / day / "health" + _write_jsonl( + base / "001_segment.jsonl", + [ + _complete("090000_300", "entities", 100, stream="default"), + _complete("090000_300", "sense", 150, stream="default"), + _complete("090000_300", "documents", 125, stream="default"), + ], + ) + + assert read_completed_since(day, 99).segments == ( + {"stream": "default", "segment": "090000_300", "ts": 150}, + ) + + +def test_read_completed_since_excludes_ts_at_or_before_since(pipeline_journal): + day = "20990212" + base = pipeline_journal / "chronicle" / day / "health" + _write_jsonl( + base / "001_segment.jsonl", + [_complete("090000_300", "entities", 100, stream="default")], + ) + + assert read_completed_since(day, 100) == CompletionsSince((), ()) + + +def test_read_completed_since_projects_activity_units(pipeline_journal): + day = "20990213" + base = pipeline_journal / "chronicle" / day / "health" + _write_jsonl( + base / "001_activity.jsonl", + [ + { + "event": "talent.complete", + "ts": 101, + "mode": "activity", + "name": "summary", + "facet": "work", + "activity": "meeting_090000_300", + } + ], + ) + + assert read_completed_since(day, 100).activities == ( + {"facet": "work", "activity": "meeting_090000_300", "ts": 101}, + ) + + +def test_read_completed_since_includes_prior_day_completions(pipeline_journal): + day = "20990214" + prev_day = "20990213" + base = pipeline_journal / "chronicle" / prev_day / "health" + _write_jsonl( + base / "001_segment.jsonl", + [_complete("235900_300", "entities", 200, stream="default")], + ) + + assert read_completed_since(day, 199).segments == ( + {"stream": "default", "segment": "235900_300", "ts": 200}, + ) + + +def test_read_completed_since_excludes_units_whose_latest_terminal_failed( + pipeline_journal, +): + day = "20990215" + base = pipeline_journal / "chronicle" / day / "health" + _write_jsonl( + base / "001_segment.jsonl", + [ + _complete("090000_300", "entities", 100, stream="default"), + _fail("090000_300", "entities", 110, stream="default"), + ], + ) + + assert read_completed_since(day, 99) == CompletionsSince((), ()) + + +def test_read_completed_since_ignores_cadence_own_terminal_record(pipeline_journal): + day = "20990216" + base = pipeline_journal / "chronicle" / day / "health" + _write_jsonl( + base / "001_cadence.jsonl", + [ + { + "event": "talent.complete", + "ts": 100, + "mode": "cadence", + "name": "pulse", + "segment": None, + "activity": None, + } + ], + ) + + assert read_completed_since(day, 99) == CompletionsSince((), ()) + + +def test_read_completed_since_sorts_outputs_deterministically(pipeline_journal): + day = "20990217" + base = pipeline_journal / "chronicle" / day / "health" + _write_jsonl( + base / "001_units.jsonl", + [ + _complete("z_seg", "entities", 200, stream="zeta"), + _complete("a_seg", "entities", 100, stream="zeta"), + _complete("b_seg", "entities", 100, stream="alpha"), + { + "event": "talent.complete", + "ts": 150, + "mode": "activity", + "name": "summary", + "facet": "work", + "activity": "z_activity", + }, + { + "event": "talent.complete", + "ts": 150, + "mode": "activity", + "name": "summary", + "facet": "home", + "activity": "a_activity", + }, + ], + ) + + completions = read_completed_since(day, 99) + + assert completions.segments == ( + {"stream": "alpha", "segment": "b_seg", "ts": 100}, + {"stream": "zeta", "segment": "a_seg", "ts": 100}, + {"stream": "zeta", "segment": "z_seg", "ts": 200}, + ) + assert completions.activities == ( + {"facet": "home", "activity": "a_activity", "ts": 150}, + {"facet": "work", "activity": "z_activity", "ts": 150}, + ) + + def test_empty_day_is_healthy(pipeline_journal): summary = summarize_pipeline_day("20260101") diff --git a/tests/test_think_cadence.py b/tests/test_think_cadence.py new file mode 100644 index 000000000..3e4688f92 --- /dev/null +++ b/tests/test_think_cadence.py @@ -0,0 +1,258 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Tests for cadence-scheduled think dispatch.""" + +from __future__ import annotations + +import importlib +import json +import time +from pathlib import Path + +import pytest + +from solstone.think.pipeline_health import CompletionsSince + +DAY = "20990302" +NOW = 1_800_000_000_000 + + +@pytest.fixture +def cadence_runtime(tmp_path, monkeypatch): + journal = tmp_path / "journal" + (journal / "health").mkdir(parents=True) + monkeypatch.setenv("SOLSTONE_JOURNAL", str(journal)) + mod = importlib.import_module("solstone.think.thinking") + monkeypatch.setattr(mod, "emit", lambda *args, **kwargs: None) + monkeypatch.setattr(mod, "_jsonl_log", lambda *args, **kwargs: None) + return mod, journal + + +def _cadence_config(**extra) -> dict: + config = { + "type": "generate", + "priority": 1, + "output": "md", + "schedule": "cadence", + } + config.update(extra) + return config + + +def _window( + *, + segments: tuple[dict, ...] = ( + {"stream": "default", "segment": "090000_300", "ts": NOW - 120_000}, + ), + activities: tuple[dict, ...] = (), +) -> CompletionsSince: + return CompletionsSince(segments=segments, activities=activities) + + +def _write_cadence_state(journal: Path, state: dict[str, int]) -> None: + path = journal / "health" / "cadence.json" + path.write_text(json.dumps(state), encoding="utf-8") + + +def test_run_cadence_prompts_zero_talents_exits_without_state_write( + cadence_runtime, monkeypatch +): + mod, journal = cadence_runtime + monkeypatch.setattr(mod, "get_talent_configs", lambda schedule: {}) + monkeypatch.setattr( + mod, + "_cortex_request_with_retry", + lambda **kwargs: pytest.fail("cadence should not dispatch"), + ) + + assert mod.run_cadence_prompts(DAY, refresh=False, verbose=False) == (0, 0, []) + assert not (journal / "health" / "cadence.json").exists() + + +def test_run_cadence_prompts_fires_when_window_has_new_work( + cadence_runtime, monkeypatch +): + mod, journal = cadence_runtime + last = NOW - 600_000 + requests: list[dict] = [] + _write_cadence_state(journal, {"talentA": last}) + monkeypatch.setattr( + mod, + "get_talent_configs", + lambda schedule: {"talentA": _cadence_config()}, + ) + monkeypatch.setattr(mod, "now_ms", lambda: NOW) + + def fake_completed_since(day: str, since_ms: int) -> CompletionsSince: + assert day == DAY + assert since_ms == last + return _window() + + def fake_request(**kwargs): + requests.append(kwargs) + return "use-1" + + monkeypatch.setattr(mod, "read_completed_since", fake_completed_since) + monkeypatch.setattr(mod, "_cortex_request_with_retry", fake_request) + monkeypatch.setattr( + mod, + "_drain_priority_batch", + lambda spawned, *_args: (1, 0, []), + ) + + assert mod.run_cadence_prompts(DAY, refresh=False, verbose=False) == (1, 0, []) + + assert len(requests) == 1 + request_config = requests[0]["config"] + assert request_config["schedule"] == "cadence" + assert request_config["cadence_window"]["since_ms"] == last + assert request_config["cadence_window"]["segments"] + assert mod.load_cadence_state()["talentA"] == NOW + + +def test_run_cadence_prompts_noops_without_new_work(cadence_runtime, monkeypatch): + mod, journal = cadence_runtime + last = NOW - 600_000 + save_calls: list[dict] = [] + _write_cadence_state(journal, {"talentA": last}) + monkeypatch.setattr( + mod, + "get_talent_configs", + lambda schedule: {"talentA": _cadence_config()}, + ) + monkeypatch.setattr(mod, "now_ms", lambda: NOW) + monkeypatch.setattr( + mod, + "read_completed_since", + lambda day, since_ms: CompletionsSince((), ()), + ) + monkeypatch.setattr( + mod, + "_cortex_request_with_retry", + lambda **kwargs: pytest.fail("cadence should not dispatch"), + ) + monkeypatch.setattr( + mod, "save_cadence_state", lambda state: save_calls.append(state) + ) + + assert mod.run_cadence_prompts(DAY, refresh=False, verbose=False) == (0, 0, []) + assert save_calls == [] + assert json.loads((journal / "health" / "cadence.json").read_text()) == { + "talentA": last + } + + +def test_run_cadence_prompts_respects_per_talent_interval(cadence_runtime, monkeypatch): + mod, journal = cadence_runtime + _write_cadence_state(journal, {"talentA": NOW - 600_000}) + monkeypatch.setattr( + mod, + "get_talent_configs", + lambda schedule: {"talentA": _cadence_config(cadence_minutes=30)}, + ) + monkeypatch.setattr(mod, "now_ms", lambda: NOW) + monkeypatch.setattr( + mod, + "read_completed_since", + lambda day, since_ms: pytest.fail("interval gate should run first"), + ) + monkeypatch.setattr( + mod, + "_cortex_request_with_retry", + lambda **kwargs: pytest.fail("cadence should not dispatch"), + ) + + assert mod.run_cadence_prompts(DAY, refresh=False, verbose=False) == (0, 0, []) + + +def test_run_cadence_prompts_writes_back_only_on_success(cadence_runtime, monkeypatch): + mod, journal = cadence_runtime + last = NOW - 600_000 + save_calls: list[dict] = [] + _write_cadence_state(journal, {"talentA": last}) + monkeypatch.setattr( + mod, + "get_talent_configs", + lambda schedule: {"talentA": _cadence_config()}, + ) + monkeypatch.setattr(mod, "now_ms", lambda: NOW) + monkeypatch.setattr(mod, "read_completed_since", lambda day, since_ms: _window()) + monkeypatch.setattr(mod, "_cortex_request_with_retry", lambda **kwargs: "use-fail") + monkeypatch.setattr( + mod, + "_drain_priority_batch", + lambda spawned, *_args: (0, 1, ["talentA (error)"]), + ) + monkeypatch.setattr( + mod, "save_cadence_state", lambda state: save_calls.append(state) + ) + + assert mod.run_cadence_prompts(DAY, refresh=False, verbose=False) == ( + 0, + 1, + ["talentA (error)"], + ) + assert save_calls == [] + assert json.loads((journal / "health" / "cadence.json").read_text()) == { + "talentA": last + } + + +def test_run_cadence_prompts_missing_state_treats_talent_as_never_run( + cadence_runtime, monkeypatch +): + mod, journal = cadence_runtime + requests: list[dict] = [] + monkeypatch.setattr( + mod, + "get_talent_configs", + lambda schedule: {"talentA": _cadence_config()}, + ) + monkeypatch.setattr(mod, "now_ms", lambda: NOW) + monkeypatch.setattr(mod, "read_completed_since", lambda day, since_ms: _window()) + monkeypatch.setattr( + mod, + "_cortex_request_with_retry", + lambda **kwargs: requests.append(kwargs) or "use-1", + ) + monkeypatch.setattr( + mod, + "_drain_priority_batch", + lambda spawned, *_args: (1, 0, []), + ) + + assert mod.run_cadence_prompts(DAY, refresh=False, verbose=False) == (1, 0, []) + assert requests[0]["config"]["cadence_window"]["since_ms"] == 0 + assert mod.load_cadence_state()["talentA"] == NOW + assert (journal / "health" / "cadence.json").exists() + + +def test_run_cadence_prompts_writeback_uses_window_read_time( + cadence_runtime, monkeypatch +): + mod, _journal = cadence_runtime + calls = 0 + monkeypatch.setattr( + mod, + "get_talent_configs", + lambda schedule: {"talentA": _cadence_config()}, + ) + + def fake_now_ms() -> int: + nonlocal calls + calls += 1 + return NOW if calls == 1 else NOW + 999_999 + + def fake_drain(spawned, *_args): + assert mod.now_ms() > NOW + assert time.time() > 0 + return (1, 0, []) + + monkeypatch.setattr(mod, "now_ms", fake_now_ms) + monkeypatch.setattr(mod, "read_completed_since", lambda day, since_ms: _window()) + monkeypatch.setattr(mod, "_cortex_request_with_retry", lambda **kwargs: "use-1") + monkeypatch.setattr(mod, "_drain_priority_batch", fake_drain) + + assert mod.run_cadence_prompts(DAY, refresh=False, verbose=False) == (1, 0, []) + assert mod.load_cadence_state()["talentA"] == NOW diff --git a/tests/test_think_dry_run.py b/tests/test_think_dry_run.py index 5fc0ec760..4370d7a76 100644 --- a/tests/test_think_dry_run.py +++ b/tests/test_think_dry_run.py @@ -5,6 +5,8 @@ import importlib +from solstone.think.pipeline_health import CompletionsSince + def test_dry_run_daily(journal_copy, capsys): """Dry-run daily mode prints prompts without spawning agents.""" @@ -73,3 +75,73 @@ def test_dry_run_no_callosum(journal_copy, monkeypatch, capsys): mod.dry_run("20240101") assert mod._callosum is None monkeypatch.setattr(mod, "_callosum", prev) + + +def test_dry_run_cadence_reports_gate_decisions(tmp_path, monkeypatch, capsys): + """Dry-run cadence mode reports fire/no-op/skip decisions without writes.""" + mod = importlib.import_module("solstone.think.thinking") + journal = tmp_path / "journal" + (journal / "health").mkdir(parents=True) + monkeypatch.setenv("SOLSTONE_JOURNAL", str(journal)) + now = 1_800_000_000_000 + state = { + "fire": now - 600_000, + "no_work": now - 700_000, + "wait": now - 60_000, + } + + monkeypatch.setattr( + mod, + "get_talent_configs", + lambda schedule: { + "fire": {"type": "generate", "priority": 1, "schedule": "cadence"}, + "no_work": {"type": "generate", "priority": 2, "schedule": "cadence"}, + "wait": { + "type": "generate", + "priority": 3, + "schedule": "cadence", + "cadence_minutes": 5, + }, + }, + ) + monkeypatch.setattr(mod, "load_cadence_state", lambda: state) + monkeypatch.setattr(mod, "now_ms", lambda: now) + + def fake_completed_since(day: str, since_ms: int) -> CompletionsSince: + if since_ms == state["fire"]: + return CompletionsSince( + segments=( + {"stream": "default", "segment": "090000_300", "ts": now - 1}, + ), + activities=(), + ) + return CompletionsSince((), ()) + + monkeypatch.setattr(mod, "read_completed_since", fake_completed_since) + + mod.dry_run("20990302", cadence=True) + + out = capsys.readouterr().out + assert "2099-03-02" in out + assert "cadence agents" in out + assert "fire fire" in out + assert "window: 1 segment(s), 0 activity(ies)" in out + assert "no-op no_work" in out + assert "skip wait" in out + assert not (journal / "health" / "cadence.json").exists() + + +def test_dry_run_cadence_zero_talents(tmp_path, monkeypatch, capsys): + """Dry-run cadence mode reports empty cadence schedules without writes.""" + mod = importlib.import_module("solstone.think.thinking") + journal = tmp_path / "journal" + (journal / "health").mkdir(parents=True) + monkeypatch.setenv("SOLSTONE_JOURNAL", str(journal)) + monkeypatch.setattr(mod, "get_talent_configs", lambda schedule: {}) + + mod.dry_run("20990302", cadence=True) + + out = capsys.readouterr().out + assert "cadence agents" in out + assert "No prompts for schedule: cadence" in out + assert not (journal / "health" / "cadence.json").exists()