From 8dc6ead23eb63528c2bf858afbbff7cf545f9953 Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Wed, 27 May 2026 18:30:43 -0600 Subject: [PATCH] fix(supervisor): split think queue by mode so daily catchup doesn't block live thinks MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Split journal think task queueing by mode so daily, segment, flush, activity, and weekly work no longer serialize behind one literal think partition. This keeps startup/daily catchup from blocking live segment-thinks. Use `_command_partition` in runner.py as shared source of truth; `TaskQueue.get_command_name` delegates to it. Think partitions use bare names: daily/segment/flush/activity/weekly. The weekly-agents 30m max_runtime now applies only to weekly; before, the broad think key applied it to all journal think runs. Operators with hand-edited schedules keyed against literal `think` will lose that cap and should re-key against the specific think mode. Clean break per CLAUDE.md ยง8: no fallback alias for old `"think"` cmd_name. Allowing two parallel think modes increases LLM load; that is acknowledged and accepted. Co-Authored-By: Claude Opus 4.7 (1M context) --- solstone/think/runner.py | 45 +++-- solstone/think/supervisor.py | 10 +- tests/test_supervisor.py | 311 ++++++++++++++++++++++++++++++++++- 3 files changed, 341 insertions(+), 25 deletions(-) diff --git a/solstone/think/runner.py b/solstone/think/runner.py index fab07b687..97d9c405a 100644 --- a/solstone/think/runner.py +++ b/solstone/think/runner.py @@ -25,6 +25,7 @@ import subprocess import sys import threading import time +from collections.abc import Sequence from dataclasses import dataclass from datetime import datetime from pathlib import Path @@ -198,6 +199,32 @@ class DailyLogWriter: ) +def _command_partition(cmd: Sequence[str]) -> str: + """Return the queue/log partition name for a managed-process cmd. + + Think tasks partition by bare mode name (daily/segment/flush/activity/weekly); + everything else uses sol/journal subcommand or process basename. + """ + if cmd and cmd[0] in ("sol", "journal") and len(cmd) > 1: + name = cmd[1] + if name == "think": + for flag, mode in [ + ("--activity", "activity"), + ("--flush", "flush"), + ("--segments", "segment"), + ("--weekly", "weekly"), + ("--segment", "segment"), + ]: + if flag in cmd: + name = mode + break + else: + name = "daily" + else: + name = Path(cmd[0]).name if cmd else "unknown" + return name + + @dataclass class ManagedProcess: """Subprocess wrapper with automatic output logging and lifecycle management. @@ -281,23 +308,7 @@ class ManagedProcess: this call. Use a long-lived worker thread that blocks in process.wait() for the lifetime of the child. """ - if cmd[0] in ("sol", "journal") and len(cmd) > 1: - name = cmd[1] - if name == "think": - for flag, mode in [ - ("--activity", "activity"), - ("--flush", "flush"), - ("--segments", "segment"), - ("--weekly", "weekly"), - ("--segment", "segment"), - ]: - if flag in cmd: - name = mode - break - else: - name = "daily" - else: - name = Path(cmd[0]).name + name = _command_partition(cmd) # Generate correlation ID (use provided ref, else timestamp) ref = ref if ref else str(now_ms()) diff --git a/solstone/think/supervisor.py b/solstone/think/supervisor.py index a38844faf..527addbc3 100644 --- a/solstone/think/supervisor.py +++ b/solstone/think/supervisor.py @@ -29,6 +29,7 @@ from solstone.think.callosum import CallosumConnection, CallosumServer from solstone.think.maint import run_pending_tasks from solstone.think.readiness import clear_ready, signal_ready from solstone.think.runner import ManagedProcess as RunnerManagedProcess +from solstone.think.runner import _command_partition from solstone.think.sync_check import ( DEFAULT_INTERVAL_SECONDS, SyncCheckSnapshot, @@ -246,13 +247,8 @@ class TaskQueue: @staticmethod def get_command_name(cmd: list[str]) -> str: - """Extract command name from cmd array for queue serialization. - - For 'sol X' or 'journal X' commands, returns X. Otherwise returns cmd[0] basename. - """ - if cmd and cmd[0] in ("sol", "journal") and len(cmd) > 1: - return cmd[1] - return Path(cmd[0]).name if cmd else "unknown" + """Return the canonical queue/log partition for a command.""" + return _command_partition(cmd) def _notify_queue_change(self, cmd_name: str) -> None: """Notify listener of queue state change (called outside lock).""" diff --git a/tests/test_supervisor.py b/tests/test_supervisor.py index d070e2c47..074eac040 100644 --- a/tests/test_supervisor.py +++ b/tests/test_supervisor.py @@ -285,7 +285,7 @@ def test_get_command_name(): # sol X -> X assert get(["sol", "indexer", "--rescan"]) == "indexer" assert get(["sol", "insight", "20240101"]) == "insight" - assert get(["journal", "think", "--day", "20240101"]) == "think" + assert get(["journal", "think", "--day", "20240101"]) == "daily" # Other commands -> basename assert get(["/usr/bin/python", "script.py"]) == "python" @@ -295,6 +295,315 @@ def test_get_command_name(): assert get([]) == "unknown" +@pytest.mark.parametrize( + "cmd", + [ + ["journal", "think", "--day", "20260527"], + ["journal", "think", "--day", "20260527", "--segment", "120000_300"], + [ + "journal", + "think", + "--day", + "20260527", + "--segment", + "120000_300", + "--stream", + "screen", + ], + [ + "journal", + "think", + "--day", + "20260527", + "--segment", + "120000_300", + "--flush", + ], + ["journal", "think", "--day", "20260527", "--segments"], + [ + "journal", + "think", + "--activity", + "activity-id", + "--facet", + "work", + "--day", + "20260527", + ], + ["journal", "think", "--weekly", "-v"], + ["journal", "think"], + ["sol", "indexer", "--rescan"], + ["journal", "sense", "--day", "20260101"], + ], +) +def test_command_partition_matches_task_queue_get_command_name(cmd): + mod = importlib.import_module("solstone.think.supervisor") + runner = importlib.import_module("solstone.think.runner") + + assert runner._command_partition(cmd) == mod.TaskQueue.get_command_name(cmd) + + +def _fresh_task_queue(mod, *, on_queue_change=None): + mod._task_queue = mod.TaskQueue(on_queue_change=on_queue_change) + mod._supervisor_callosum = None + return mod._task_queue + + +def _capture_thread_starts(monkeypatch, mod): + spawned = [] + + def fake_thread_start(self): + spawned.append(self._args) + + monkeypatch.setattr(mod.threading.Thread, "start", fake_thread_start) + return spawned + + +def test_task_queue_daily_and_segment_run_independently(monkeypatch): + mod = importlib.import_module("solstone.think.supervisor") + queue = _fresh_task_queue(mod) + _capture_thread_starts(monkeypatch, mod) + + queue.submit(["journal", "think", "--day", "20260527"], ref="daily-ref") + queue.submit( + ["journal", "think", "--day", "20260527", "--segment", "120000_300"], + ref="segment-ref", + ) + + assert set(queue._running) == {"daily", "segment"} + assert queue._queues == {} + + +def test_task_queue_segment_and_flush_run_independently(monkeypatch): + mod = importlib.import_module("solstone.think.supervisor") + queue = _fresh_task_queue(mod) + _capture_thread_starts(monkeypatch, mod) + + queue.submit( + ["journal", "think", "--day", "20260527", "--segment", "120000_300"], + ref="segment-ref", + ) + queue.submit( + [ + "journal", + "think", + "--day", + "20260527", + "--segment", + "120000_300", + "--flush", + ], + ref="flush-ref", + ) + + assert set(queue._running) == {"segment", "flush"} + assert queue._queues == {} + + +def test_task_queue_daily_and_activity_run_independently(monkeypatch): + mod = importlib.import_module("solstone.think.supervisor") + queue = _fresh_task_queue(mod) + _capture_thread_starts(monkeypatch, mod) + + queue.submit(["journal", "think", "--day", "20260527"], ref="daily-ref") + queue.submit( + [ + "journal", + "think", + "--activity", + "activity-id", + "--facet", + "work", + "--day", + "20260527", + ], + ref="activity-ref", + ) + + assert set(queue._running) == {"daily", "activity"} + assert queue._queues == {} + + +def test_task_queue_daily_and_weekly_run_independently(monkeypatch): + mod = importlib.import_module("solstone.think.supervisor") + queue = _fresh_task_queue(mod) + _capture_thread_starts(monkeypatch, mod) + + queue.submit(["journal", "think", "--day", "20260527"], ref="daily-ref") + queue.submit(["journal", "think", "--weekly", "-v"], ref="weekly-ref") + + assert set(queue._running) == {"daily", "weekly"} + assert queue._queues == {} + + +def test_task_queue_segments_plural_shares_segment_partition(monkeypatch): + mod = importlib.import_module("solstone.think.supervisor") + queue = _fresh_task_queue(mod) + _capture_thread_starts(monkeypatch, mod) + + queue.submit( + ["journal", "think", "--day", "20260527", "--segment", "120000_300"], + ref="segment-ref", + ) + queue.submit( + ["journal", "think", "--day", "20260527", "--segments"], + ref="segments-ref", + ) + + assert set(queue._running) == {"segment"} + assert queue._queues["segment"][0]["refs"] == ["segments-ref"] + assert queue._queues["segment"][0]["cmd"] == [ + "journal", + "think", + "--day", + "20260527", + "--segments", + ] + + +def test_task_queue_flush_flag_precedes_segment_flag(): + runner = importlib.import_module("solstone.think.runner") + + assert ( + runner._command_partition( + ["journal", "think", "--day", "20260527", "--segment", "00", "--flush"] + ) + == "flush" + ) + + +def test_task_queue_within_mode_serialization_segment(monkeypatch): + mod = importlib.import_module("solstone.think.supervisor") + queue = _fresh_task_queue(mod) + spawned = _capture_thread_starts(monkeypatch, mod) + + queue.submit( + ["journal", "think", "--day", "20260527", "--segment", "120000_300"], + ref="first-ref", + ) + queue.submit( + ["journal", "think", "--day", "20260527", "--segment", "120500_300"], + ref="second-ref", + ) + + assert set(queue._running) == {"segment"} + assert queue._running["segment"]["ref"] == "first-ref" + assert set(queue._queues) == {"segment"} + assert queue._queues["segment"][0]["refs"] == ["second-ref"] + + queue._process_next("segment") + + assert queue._running["segment"]["ref"] == "second-ref" + assert queue._queues["segment"] == [] + assert spawned[-1][0] == ["second-ref"] + + +def test_task_queue_dedup_within_segment_partition(monkeypatch): + mod = importlib.import_module("solstone.think.supervisor") + queue = _fresh_task_queue(mod) + _capture_thread_starts(monkeypatch, mod) + cmd = ["journal", "think", "--day", "20260527", "--segment", "120000_300"] + + queue.submit(cmd, ref="first-ref") + queue.submit(cmd, ref="second-ref") + queue.submit(cmd, ref="third-ref") + + assert queue._running["segment"]["ref"] == "first-ref" + assert len(queue._queues["segment"]) == 1 + assert queue._queues["segment"][0]["cmd"] == cmd + assert queue._queues["segment"][0]["refs"] == ["second-ref", "third-ref"] + + +def test_task_queue_stale_thread_reclamation_per_partition(monkeypatch): + mod = importlib.import_module("solstone.think.supervisor") + dead_thread = threading.Thread(target=lambda: None) + dead_thread.start() + dead_thread.join() + + queue = _fresh_task_queue(mod) + _capture_thread_starts(monkeypatch, mod) + + queue.submit(["journal", "think", "--day", "20260527"], ref="old-daily-ref") + queue.submit( + ["journal", "think", "--day", "20260527", "--segment", "120000_300"], + ref="segment-ref", + ) + + queue._running["daily"]["thread"] = dead_thread + + queue.submit(["journal", "think", "--day", "20260528"], ref="new-daily-ref") + + assert queue._running["daily"]["ref"] == "new-daily-ref" + assert queue._running["segment"]["ref"] == "segment-ref" + assert set(queue._running) == {"daily", "segment"} + + +def test_handle_task_request_routes_to_mode_partition(monkeypatch): + mod = importlib.import_module("solstone.think.supervisor") + queue = _fresh_task_queue(mod) + _capture_thread_starts(monkeypatch, mod) + + queue.submit(["journal", "think", "--day", "20260527"], ref="daily-ref") + mod._handle_task_request( + { + "tract": "supervisor", + "event": "request", + "cmd": ["journal", "think", "--day", "20260527", "--segment", "120000_300"], + "ref": "segment-ref", + } + ) + + assert queue._running["daily"]["ref"] == "daily-ref" + assert queue._running["segment"]["ref"] == "segment-ref" + assert queue._queues == {} + + +def test_scheduler_weekly_cap_registers_under_weekly(monkeypatch): + mod = importlib.import_module("solstone.think.supervisor") + queue = mod.TaskQueue(on_queue_change=None) + monkeypatch.setattr( + mod.scheduler, + "collect_runtime_caps", + lambda: [(["journal", "think", "--weekly", "-v"], 60.0)], + ) + + for cmd, seconds in mod.scheduler.collect_runtime_caps(): + queue.set_cap(mod.TaskQueue.get_command_name(cmd), seconds) + + assert queue._caps == {"weekly": 60.0} + + +def test_queue_event_carries_mode_partition_name(monkeypatch): + mod = importlib.import_module("solstone.think.supervisor") + events = [] + queue = _fresh_task_queue( + mod, + on_queue_change=lambda command, running, queued: events.append( + (command, running, queued) + ), + ) + _capture_thread_starts(monkeypatch, mod) + + queue.submit( + ["journal", "think", "--day", "20260527", "--segment", "120000_300"], + ref="first-ref", + ) + queue.submit( + ["journal", "think", "--day", "20260527", "--segment", "120500_300"], + ref="second-ref", + ) + + assert events[-1][0] == "segment" + + +def test_no_literal_think_queue_keys_in_source(): + mod = importlib.import_module("solstone.think.supervisor") + + assert mod.TaskQueue.get_command_name( + ["journal", "think", "--day", "20260527"] + ) != ("think") + + def test_task_queue_same_command_queued(monkeypatch): """Test that same command is queued when already running.""" mod = importlib.import_module("solstone.think.supervisor") -- 2.51.2