diff --git a/solstone/think/doctor.py b/solstone/think/doctor.py index c4293722a..51843142a 100644 --- a/solstone/think/doctor.py +++ b/solstone/think/doctor.py @@ -108,11 +108,16 @@ _DEFAULT_STT_MODEL_FIX = ( "parakeet model is not downloaded — fetch it with: journal install-models" ) JOURNAL_CAUGHT_UP_CHECK = Check("journal_caught_up", "advisory", ("linux", "darwin")) +TASK_PACE_CHECK = Check("task_pace", "advisory", ("linux", "darwin")) _CAUGHT_UP_BACKLOG_FIX = ( "solstone catches up on its own; reprocess a day from the health surface " "to prioritize it" ) _CAUGHT_UP_CANT_TELL_FIX = "re-run journal doctor; check the health logs if it persists" +_TASK_PACE_FIX = ( + "a job is running long; it will be stopped automatically if it passes its cap " + "— no action needed unless it persists" +) def python_sanity_check(args: Args) -> CheckResult: @@ -600,6 +605,24 @@ def journal_caught_up_check(args: Args) -> CheckResult: return make_result(check, "warn", detail, _CAUGHT_UP_BACKLOG_FIX) +def task_pace_check(args: Args) -> CheckResult: + del args + check = TASK_PACE_CHECK + status = fetch_supervisor_status() + if status is None: + return make_result(check, "skip", "supervisor status unavailable") + tasks = status.get("tasks") or [] + slow = [t for t in tasks if t.get("slow") or t.get("stuck")] + if not slow: + return make_result(check, "ok", "tasks on pace") + names = ", ".join( + f"{t.get('name', '?')} " + f"({t.get('duration_seconds', 0)}s of {t.get('max_runtime_seconds', '?')}s cap)" + for t in slow + ) + return make_result(check, "warn", f"running long: {names}", _TASK_PACE_FIX) + + def _resolve_configured_backend() -> str | None: """Read transcribe.backend from an existing journal config without creating anything. @@ -694,6 +717,7 @@ JOURNAL_CHECKS: list[tuple[Check, Runner]] = [ (SERVICE_RUNNING_CHECK, service_running_check), (JOURNAL_SYNC_CHECK, journal_sync_check), (JOURNAL_CAUGHT_UP_CHECK, journal_caught_up_check), + (TASK_PACE_CHECK, task_pace_check), (STALE_ALIAS_CHECK, partial(stale_alias_symlink_check, binary="journal")), (LAUNCHD_STALE_PLIST_CHECK, launchd_stale_plist_check), (DEFAULT_STT_READY_CHECK, default_stt_ready_check), diff --git a/solstone/think/health_cli.py b/solstone/think/health_cli.py index 8289c8ed7..e9fa6bfa7 100644 --- a/solstone/think/health_cli.py +++ b/solstone/think/health_cli.py @@ -75,6 +75,8 @@ def print_status(status: dict[str, Any]) -> None: line = f" {name:16} {duration}s" if task.get("stuck"): line += f" STUCK (cap {task['max_runtime_seconds']}s)" + elif task.get("slow"): + line += f" SLOW (cap {task['max_runtime_seconds']}s)" print(line) for name, count in non_zero_queues: print(f" queued {name:9} {count}") diff --git a/solstone/think/supervisor.py b/solstone/think/supervisor.py index d89fecef0..c8f4918d2 100644 --- a/solstone/think/supervisor.py +++ b/solstone/think/supervisor.py @@ -48,6 +48,7 @@ from solstone.think.sync_check import ( from solstone.think.utils import ( EXIT_EMPTY, EXIT_TEMPFAIL, + SOFT_RUNTIME_FRACTION, day_path, find_available_port, get_journal, @@ -932,6 +933,7 @@ class TaskQueue: "name": cmd_name, "duration_seconds": duration, "max_runtime_seconds": cap, + "slow": duration >= cap * SOFT_RUNTIME_FRACTION, "stuck": duration > cap, } ) diff --git a/solstone/think/utils.py b/solstone/think/utils.py index 17ba32589..c6ec08823 100644 --- a/solstone/think/utils.py +++ b/solstone/think/utils.py @@ -36,6 +36,7 @@ CHRONICLE_DIR = "chronicle" DEFAULT_STREAM = "_default" EXIT_TEMPFAIL = 75 EXIT_EMPTY = 66 # EX_NOINPUT: a rollup ran over zero inputs (nothing to roll up) +SOFT_RUNTIME_FRACTION = 0.75 class SolstoneNotConfigured(RuntimeError): diff --git a/tests/test_health_cli.py b/tests/test_health_cli.py index 4b5c1f062..dfda3ac08 100644 --- a/tests/test_health_cli.py +++ b/tests/test_health_cli.py @@ -143,6 +143,50 @@ def test_health_check_renders_stuck_marker(capsys): assert " providers 313656s STUCK (cap 300s)" in output +def test_health_check_renders_slow_marker(capsys): + status = { + "services": [], + "tasks": [ + { + "name": "providers", + "duration_seconds": 240, + "max_runtime_seconds": 300, + "slow": True, + "stuck": False, + } + ], + "queues": {}, + } + + print_status(status) + + output = capsys.readouterr().out + assert " providers 240s SLOW (cap 300s)" in output + assert "STUCK" not in output + + +def test_health_check_renders_stuck_only_when_slow_and_stuck(capsys): + status = { + "services": [], + "tasks": [ + { + "name": "providers", + "duration_seconds": 313656, + "max_runtime_seconds": 300, + "slow": True, + "stuck": True, + } + ], + "queues": {}, + } + + print_status(status) + + output = capsys.readouterr().out + assert " providers 313656s STUCK (cap 300s)" in output + assert "SLOW" not in output + + def test_healthy_task_rendering_unchanged(capsys): status = { "services": [], diff --git a/tests/test_supervisor.py b/tests/test_supervisor.py index e4e02787d..2eb90d308 100644 --- a/tests/test_supervisor.py +++ b/tests/test_supervisor.py @@ -1726,6 +1726,7 @@ def test_collect_task_status_reports_default_cap(monkeypatch): "name": "providers", "duration_seconds": 12, "max_runtime_seconds": mod.DEFAULT_TASK_MAX_RUNTIME, + "slow": False, "stuck": False, } ] @@ -1742,6 +1743,22 @@ def test_collect_task_status_under_cap(monkeypatch): status = queue.collect_task_status() assert status[0]["max_runtime_seconds"] == 300 + assert status[0]["slow"] is False + assert status[0]["stuck"] is False + + +def test_collect_task_status_slow_under_cap(monkeypatch): + mod = importlib.import_module("solstone.think.supervisor") + queue = mod.TaskQueue(on_queue_change=None) + managed = _TaskManagedStub(cmd=["journal", "providers"], start_time=100.0) + queue._active["ref-1"] = managed + queue.set_cap("providers", 15) + monkeypatch.setattr(mod.time, "time", lambda: 112.0) + + status = queue.collect_task_status() + + assert status[0]["max_runtime_seconds"] == 15 + assert status[0]["slow"] is True assert status[0]["stuck"] is False @@ -1756,6 +1773,7 @@ def test_collect_task_status_over_cap(monkeypatch): status = queue.collect_task_status() assert status[0]["max_runtime_seconds"] == 5 + assert status[0]["slow"] is True assert status[0]["stuck"] is True diff --git a/tests/test_task_pace.py b/tests/test_task_pace.py new file mode 100644 index 000000000..5bf77a3ab --- /dev/null +++ b/tests/test_task_pace.py @@ -0,0 +1,106 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +import pytest + + +@pytest.fixture +def doctor(): + from solstone.think import doctor as doctor_module + + return doctor_module + + +def args(doctor): + return doctor.Args(verbose=False, json=False, jsonl=False, port=5015) + + +def test_task_pace_skips_when_supervisor_status_unavailable(doctor, monkeypatch): + monkeypatch.setattr(doctor, "fetch_supervisor_status", lambda: None) + + result = doctor.task_pace_check(args(doctor)) + + assert result.status == "skip" + + +def test_task_pace_ok_when_idle(doctor, monkeypatch): + monkeypatch.setattr(doctor, "fetch_supervisor_status", lambda: {"tasks": []}) + + result = doctor.task_pace_check(args(doctor)) + + assert result.status == "ok" + + +def test_task_pace_warns_on_slow_task(doctor, monkeypatch): + monkeypatch.setattr( + doctor, + "fetch_supervisor_status", + lambda: { + "tasks": [ + { + "name": "providers", + "duration_seconds": 80, + "max_runtime_seconds": 100, + "slow": True, + "stuck": False, + } + ] + }, + ) + + result = doctor.task_pace_check(args(doctor)) + + assert result.status == "warn" + assert "providers" in result.detail + + +def test_task_pace_ok_under_soft_fraction(doctor, monkeypatch): + monkeypatch.setattr( + doctor, + "fetch_supervisor_status", + lambda: { + "tasks": [ + { + "name": "providers", + "duration_seconds": 50, + "max_runtime_seconds": 100, + "slow": False, + "stuck": False, + } + ] + }, + ) + + result = doctor.task_pace_check(args(doctor)) + + assert result.status == "ok" + + +def test_task_pace_warns_on_stuck_task(doctor, monkeypatch): + monkeypatch.setattr( + doctor, + "fetch_supervisor_status", + lambda: { + "tasks": [ + { + "name": "providers", + "duration_seconds": 120, + "max_runtime_seconds": 100, + "slow": True, + "stuck": True, + } + ] + }, + ) + + result = doctor.task_pace_check(args(doctor)) + + assert result.status == "warn" + assert "providers" in result.detail + + +def test_task_pace_registered_in_journal_checks(doctor): + assert "task_pace" in {check.name for check, _ in doctor.JOURNAL_CHECKS} + assert doctor.CHECK_MAP["task_pace"].severity == "advisory"