mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 12:18:39 +00:00
server: make the supervisor crash counter shutdown-aware
A graceful shutdown (window close, SIGTERM) tore the event bus/Manager down under the running supervisor tick, so the loop met BrokenPipe/EOF three times in a row, declared "Supervisor loop died after 3 consecutive crashes" and posted an owner alarm while the process was exiting on purpose. The lifespan teardown now sets a process-local stop event as its very first statement and joins the loop (bounded, 5 s) before the bridge and event bus shut down; the loop checks the event in its while condition and in its crash handler (an exception raised while a stop or restart is in progress exits at info level without counting, erroring or alarming), waits its crash backoff on that event instead of time.sleep so a shutdown is prompt, and reaches the one shared exit (watchdog stop, thread cleared) on every path. Three genuine consecutive crashes keep their visible death.
This commit is contained in:
parent
60a17a0a53
commit
f059e25fb6
3 changed files with 267 additions and 6 deletions
|
|
@ -886,7 +886,7 @@ The browser reconnects with bounded exponential delay, shows the reconnect overl
|
|||
Each Chat instance handles `open` by resynchronizing archive-aware durable history and `close` by withdrawing online/accounting presentation; reconnect deduplication covers overlap between live frames and REST replay, Logs merges the same way, and large history parsing runs off the server event loop. Delivery is live plus replay, not a promise that every transient frame is persisted: durable chat rows, task results, queue snapshots, Project revisions, review ledgers, cost ledgers, and lifecycle state remain the recovery authorities.
|
||||
## 5. Supervisor Loop
|
||||
|
||||
`server.py::_run_supervisor()` is the single scheduler for pooled tasks. A healthy tick publishes liveness, rotates the paired chat and progress logs, checks worker health, drains worker, direct-chat, and consciousness events, accepts owner bridge input, enforces deadlines and schedules, runs throttled reconciliation and evolution admission, assigns eligible work, and persists `state/queue_snapshot.json`. Bridge intake precedes timeout, maintenance, evolution, and assignment work so a slow control-plane step cannot make a new owner message invisible. Three consecutive loop failures clear supervisor readiness, stop its watchdog generation, and notify the owner instead of leaving a healthy-looking server that no longer assigns work.
|
||||
`server.py::_run_supervisor()` is the single scheduler for pooled tasks. A healthy tick publishes liveness, rotates the paired chat and progress logs, checks worker health, drains worker, direct-chat, and consciousness events, accepts owner bridge input, enforces deadlines and schedules, runs throttled reconciliation and evolution admission, assigns eligible work, and persists `state/queue_snapshot.json`. Bridge intake precedes timeout, maintenance, evolution, and assignment work so a slow control-plane step cannot make a new owner message invisible. Three consecutive loop failures clear supervisor readiness, stop its watchdog generation, and notify the owner instead of leaving a healthy-looking server that no longer assigns work; a failure raised while a shutdown or restart is already in progress (the lifespan teardown sets a process-local stop event first and joins the loop before the bridge and event bus go down) is not a crash — the loop exits quietly, without the counter, the error, or the alarm — and the crash backoff waits on that stop event so a shutdown is never held by it.
|
||||
|
||||
`PENDING` and `RUNNING`, guarded by `supervisor.queue._queue_lock`, are the live task-lifecycle authority. Admission reserves identity before project, workspace, attachment, or routing side effects can create a duplicate; refuses a disabled pool, duplicate task, project deletion, accepted or sealed root, or exhausted root budget; attaches the task contract; and preserves stable priority order. Assignment runs against the same locked state and skips reaping slots, budget-paused work, closed project roots, conflicting project writers, and tasks exceeding the root's subagent capacity or depth reservation; evolution tasks are dropped there when `evolution_block_reason()` is set (Light runtime mode, `supervisor/workers.py`). That is the last of three evolution-only runtime-mode fences: owner and post-task entry points refuse a campaign start, `enqueue_evolution_task_if_needed()` independently pauses and disables a carried campaign before queueing it, and assignment drops what still slipped through (`supervisor/evolution_lifecycle.py`). Generic `supervisor.queue.enqueue_task()` has no runtime-mode predicate at all. Configured worker count is therefore not available capacity: the truthful value is the currently assignable idle count after custody, reaping, and admission fences.
|
||||
|
||||
|
|
|
|||
27
server.py
27
server.py
|
|
@ -88,6 +88,11 @@ log = logging.getLogger("server")
|
|||
RESTART_EXIT_CODE = 42
|
||||
PANIC_EXIT_CODE = 99
|
||||
_restart_requested = threading.Event()
|
||||
# Set FIRST in the lifespan teardown: the supervisor loop reads it in its
|
||||
# ``while`` and in its crash handler, so the bus/Manager being torn down by
|
||||
# the shutdown itself never counts as a loop crash (no false "died after 3
|
||||
# consecutive crashes" alarm on a graceful window close / SIGTERM).
|
||||
_supervisor_stop = threading.Event()
|
||||
# Set only when the OWNER asked for the restart (the chat Restart button, and the
|
||||
# control endpoints that restart on the owner's behalf). The single fact the
|
||||
# re-exec needs to decide whether the runtime-mode ratchet pin rides along.
|
||||
|
|
@ -213,6 +218,7 @@ def _start_supervisor_if_needed(settings: dict) -> bool:
|
|||
if _supervisor_thread and _supervisor_thread.is_alive():
|
||||
return False
|
||||
_supervisor_error = None
|
||||
_supervisor_stop.clear() # in-process revival after a teardown-stopped generation
|
||||
_supervisor_thread = threading.Thread(
|
||||
target=_run_supervisor,
|
||||
args=(settings,),
|
||||
|
|
@ -2251,7 +2257,7 @@ def _run_supervisor(settings: dict) -> None:
|
|||
_loop_liveness = [time.monotonic()]
|
||||
_watchdog_stop = threading.Event() # per-generation: stops the watchdog when THIS loop exits
|
||||
_start_supervisor_liveness_watchdog(_loop_liveness, _watchdog_stop)
|
||||
while not _restart_requested.is_set():
|
||||
while not _restart_requested.is_set() and not _supervisor_stop.is_set():
|
||||
try:
|
||||
_loop_liveness[0] = time.monotonic()
|
||||
rotate_chat_log_if_needed(DATA_DIR)
|
||||
|
|
@ -2309,6 +2315,11 @@ def _run_supervisor(settings: dict) -> None:
|
|||
time.sleep(0.5)
|
||||
|
||||
except Exception as exc:
|
||||
if _supervisor_stop.is_set() or _restart_requested.is_set():
|
||||
# The shutdown/restart tore the bus down under this tick (a
|
||||
# Manager proxy raising BrokenPipe/EOF): not a crash, no alarm.
|
||||
log.info("Supervisor loop exiting on shutdown: %s", exc)
|
||||
break
|
||||
crash_count += 1
|
||||
log.error("Supervisor loop crash #%d: %s", crash_count, exc, exc_info=True)
|
||||
if crash_count >= 3:
|
||||
|
|
@ -2330,10 +2341,10 @@ def _run_supervisor(settings: dict) -> None:
|
|||
)
|
||||
except Exception:
|
||||
log.debug("Failed to notify owner about supervisor death", exc_info=True)
|
||||
_watchdog_stop.set() # this generation is dead — stop its liveness watchdog
|
||||
return
|
||||
time.sleep(min(30, 2 ** crash_count))
|
||||
_watchdog_stop.set() # loop exited (restart) — stop this generation's watchdog
|
||||
break # this generation is dead: the shared exit below stops its watchdog
|
||||
# Backoff on the stop event, not time.sleep, so a shutdown is prompt.
|
||||
_supervisor_stop.wait(min(30, 2 ** crash_count))
|
||||
_watchdog_stop.set() # every exit (restart, shutdown, crash death) stops this generation's watchdog
|
||||
_supervisor_thread = None
|
||||
|
||||
|
||||
|
|
@ -2955,6 +2966,7 @@ async def lifespan(app):
|
|||
try:
|
||||
yield
|
||||
finally:
|
||||
_supervisor_stop.set() # first: the loop must know a teardown owns what follows
|
||||
if extension_reconcile_task is not None:
|
||||
extension_reconcile_task.cancel()
|
||||
with suppress(asyncio.CancelledError, asyncio.TimeoutError):
|
||||
|
|
@ -3030,6 +3042,11 @@ async def lifespan(app):
|
|||
)
|
||||
except Exception:
|
||||
pass
|
||||
# Let the loop leave its current tick BEFORE the bridge/Manager go
|
||||
# down: a tick still running would otherwise meet BrokenPipe/EOF.
|
||||
supervisor_thread = _supervisor_thread
|
||||
if supervisor_thread is not None and supervisor_thread.is_alive():
|
||||
supervisor_thread.join(timeout=5)
|
||||
try:
|
||||
from supervisor.message_bus import get_bridge
|
||||
get_bridge().shutdown()
|
||||
|
|
|
|||
|
|
@ -369,3 +369,247 @@ def test_panic_stop_kills_services_without_log_finalization(monkeypatch, tmp_pat
|
|||
"force": True, "archive_service_logs": False,
|
||||
"reconcile_delegate_custody": False,
|
||||
}]
|
||||
|
||||
|
||||
# ---------------------------------------------------- shutdown-aware supervisor loop
|
||||
|
||||
class _FakeStopEvent:
|
||||
"""A stop event whose backoff wait returns at once and records its timeout."""
|
||||
|
||||
def __init__(self):
|
||||
self.flag = False
|
||||
self.waits = []
|
||||
|
||||
def is_set(self):
|
||||
return self.flag
|
||||
|
||||
def set(self):
|
||||
self.flag = True
|
||||
|
||||
def clear(self):
|
||||
self.flag = False
|
||||
|
||||
def wait(self, timeout=None):
|
||||
self.waits.append(timeout)
|
||||
return self.flag
|
||||
|
||||
|
||||
class _Recorder:
|
||||
def __init__(self):
|
||||
self.alerts = []
|
||||
self.watchdog_stops = []
|
||||
self.steps = []
|
||||
self.stop = _FakeStopEvent()
|
||||
self.restart = None
|
||||
self.ready = None
|
||||
|
||||
|
||||
def _supervisor_harness(monkeypatch, tmp_path, steps):
|
||||
"""Drive the REAL server._run_supervisor with every init/tick collaborator
|
||||
stubbed (no processes, ports, Manager or live data root). The scripted
|
||||
``steps`` fire from the first call of each tick: ``ok`` = healthy tick,
|
||||
``raise`` = a crash, ``stop``/``restart`` = set the flag (the loop exits at
|
||||
its next ``while`` check), ``raise_after_stop``/``raise_after_restart`` =
|
||||
the flag is set and the same tick then crashes (the shutdown race)."""
|
||||
import threading
|
||||
import queue as queue_mod
|
||||
|
||||
import server
|
||||
import supervisor.events as events_mod
|
||||
import supervisor.message_bus as bus_mod
|
||||
import supervisor.queue as queue_pkg
|
||||
import supervisor.state as state_mod
|
||||
import supervisor.workers as workers_mod
|
||||
|
||||
rec = _Recorder()
|
||||
rec.steps = list(steps)
|
||||
rec.restart = threading.Event()
|
||||
rec.ready = threading.Event()
|
||||
|
||||
def _tick_head(_data_dir):
|
||||
step = rec.steps.pop(0) if rec.steps else "stop"
|
||||
if step in ("stop", "raise_after_stop"):
|
||||
rec.stop.set()
|
||||
if step in ("restart", "raise_after_restart"):
|
||||
rec.restart.set()
|
||||
if step.startswith("raise"):
|
||||
raise BrokenPipeError(32, "Broken pipe")
|
||||
|
||||
class _Bridge:
|
||||
def __init__(self, _settings):
|
||||
self._broadcast_fn = None
|
||||
|
||||
class _Consciousness:
|
||||
def __init__(self, **_kwargs):
|
||||
pass
|
||||
|
||||
def start(self):
|
||||
pass
|
||||
|
||||
def stop(self):
|
||||
pass
|
||||
|
||||
noop = lambda *_a, **_k: None # noqa: E731
|
||||
monkeypatch.setattr(server, "DATA_DIR", tmp_path)
|
||||
monkeypatch.setattr(server, "_supervisor_stop", rec.stop)
|
||||
monkeypatch.setattr(server, "_restart_requested", rec.restart)
|
||||
monkeypatch.setattr(server, "_supervisor_ready", rec.ready)
|
||||
monkeypatch.setattr(server, "_supervisor_error", None)
|
||||
monkeypatch.setattr(server, "_supervisor_thread", None)
|
||||
monkeypatch.setattr(server, "_consciousness", None)
|
||||
monkeypatch.setattr(server, "_apply_settings_to_env", noop)
|
||||
monkeypatch.setattr(server, "ensure_legacy_imported", noop)
|
||||
monkeypatch.setattr(server, "_bootstrap_supervisor_repo", lambda _s: (True, "ok"))
|
||||
monkeypatch.setattr(server, "_runtime_branch_defaults", lambda: ("dev", "stable"))
|
||||
for name in (
|
||||
"_resume_interrupted_project_deletions", "_startup_prune_sweeps", "_startup_custody_sweep",
|
||||
"_startup_worktree_prune", "_prune_delegated_snapshots", "_periodic_supervisor_maintenance",
|
||||
):
|
||||
monkeypatch.setattr(server, name, noop)
|
||||
monkeypatch.setattr(server, "_start_supervisor_liveness_watchdog",
|
||||
lambda _liveness, stop_event=None: rec.watchdog_stops.append(stop_event))
|
||||
monkeypatch.setattr(server, "_process_bridge_updates", lambda _bridge, offset, _ctx: offset)
|
||||
monkeypatch.setattr(server, "_check_pending_restart_drain", lambda _ctx: True)
|
||||
monkeypatch.setattr(server.time, "sleep", noop)
|
||||
monkeypatch.setattr(bus_mod, "init", noop)
|
||||
monkeypatch.setattr(bus_mod, "LocalChatBridge", _Bridge)
|
||||
monkeypatch.setattr(bus_mod, "send_with_budget", lambda chat_id, text: rec.alerts.append((chat_id, text)))
|
||||
monkeypatch.setattr("ouroboros.utils.set_log_sink", noop)
|
||||
monkeypatch.setattr(events_mod, "make_server_log_sink", lambda *_a, **_k: None)
|
||||
monkeypatch.setattr(events_mod, "dispatch_event", noop)
|
||||
monkeypatch.setattr(state_mod, "init", noop)
|
||||
monkeypatch.setattr(state_mod, "init_state", noop)
|
||||
monkeypatch.setattr(state_mod, "load_state", lambda: {"owner_chat_id": 7})
|
||||
for name in ("save_state", "update_state", "append_jsonl", "update_budget_from_usage",
|
||||
"rotate_jsonl_log_if_needed"):
|
||||
monkeypatch.setattr(state_mod, name, noop)
|
||||
monkeypatch.setattr(state_mod, "rotate_chat_log_if_needed", _tick_head)
|
||||
for name in ("enqueue_task", "enforce_task_timeouts", "enqueue_evolution_task_if_needed",
|
||||
"persist_queue_snapshot", "cancel_task_by_id", "queue_deep_self_review_task",
|
||||
"sort_pending", "check_scheduled_tasks"):
|
||||
monkeypatch.setattr(queue_pkg, name, noop)
|
||||
monkeypatch.setattr(queue_pkg, "restore_pending_from_snapshot", lambda: 0)
|
||||
for name in ("init", "spawn_workers", "kill_workers", "assign_tasks", "ensure_workers_healthy",
|
||||
"auto_resume_after_restart"):
|
||||
monkeypatch.setattr(workers_mod, name, noop)
|
||||
monkeypatch.setattr(workers_mod, "get_event_q", lambda: queue_mod.Queue())
|
||||
monkeypatch.setattr("ouroboros.delegate_recovery.pre_adopt_planned_handoffs", noop)
|
||||
monkeypatch.setattr("ouroboros.observability.prune_observability_blobs", lambda _root: {})
|
||||
monkeypatch.setattr("ouroboros.tools.services.prune_service_logs", lambda _root: {})
|
||||
monkeypatch.setattr("ouroboros.consciousness.BackgroundConsciousness", _Consciousness)
|
||||
return rec
|
||||
|
||||
|
||||
def _run(rec, server):
|
||||
server._run_supervisor({})
|
||||
assert rec.ready.is_set() or server._supervisor_error # init reached the loop
|
||||
return rec
|
||||
|
||||
|
||||
def test_exception_while_stopping_exits_quietly_without_counting_a_crash(monkeypatch, tmp_path, caplog):
|
||||
"""The graceful-shutdown race: the teardown sets the stop flag, then the tick
|
||||
meets the torn-down bus (BrokenPipe). That is not a crash: no owner alarm,
|
||||
readiness untouched, no supervisor error — and the watchdog generation stops."""
|
||||
import logging
|
||||
|
||||
import server
|
||||
|
||||
rec = _supervisor_harness(monkeypatch, tmp_path, ["ok", "raise_after_stop"])
|
||||
with caplog.at_level(logging.INFO, logger="server"):
|
||||
_run(rec, server)
|
||||
|
||||
assert rec.alerts == []
|
||||
assert rec.ready.is_set() is True
|
||||
assert server._supervisor_error is None
|
||||
assert rec.watchdog_stops and rec.watchdog_stops[0].is_set()
|
||||
assert server._supervisor_thread is None
|
||||
assert any("exiting on shutdown" in record.getMessage() for record in caplog.records)
|
||||
assert not any(record.levelno >= logging.ERROR for record in caplog.records)
|
||||
|
||||
|
||||
def test_exception_while_restarting_exits_quietly_too(monkeypatch, tmp_path):
|
||||
import server
|
||||
|
||||
rec = _supervisor_harness(monkeypatch, tmp_path, ["raise_after_restart"])
|
||||
_run(rec, server)
|
||||
|
||||
assert rec.alerts == []
|
||||
assert rec.ready.is_set() is True
|
||||
assert server._supervisor_error is None
|
||||
assert rec.watchdog_stops[0].is_set()
|
||||
|
||||
|
||||
def test_three_genuine_consecutive_crashes_still_die_visibly(monkeypatch, tmp_path):
|
||||
"""Without a shutdown the contract stays: the third consecutive crash records
|
||||
the error, clears readiness, alerts the owner exactly once, stops the
|
||||
watchdog generation — and the backoff between crashes waits on the stop
|
||||
event (prompt shutdown), never on time.sleep."""
|
||||
import server
|
||||
|
||||
rec = _supervisor_harness(monkeypatch, tmp_path, ["raise", "raise", "raise", "ok"])
|
||||
_run(rec, server)
|
||||
|
||||
assert len(rec.alerts) == 1
|
||||
assert rec.alerts[0][0] == 7
|
||||
assert "died after repeated crashes" in rec.alerts[0][1]
|
||||
assert rec.ready.is_set() is False
|
||||
assert "3 consecutive crashes" in str(server._supervisor_error)
|
||||
assert rec.watchdog_stops[0].is_set()
|
||||
assert server._supervisor_thread is None
|
||||
assert rec.stop.waits == [2, 4]
|
||||
assert rec.steps == ["ok"] # the loop is dead: the next tick never ran
|
||||
|
||||
|
||||
def test_healthy_tick_between_crashes_resets_the_count(monkeypatch, tmp_path):
|
||||
import server
|
||||
|
||||
rec = _supervisor_harness(monkeypatch, tmp_path, ["raise", "raise", "ok", "raise", "raise", "stop"])
|
||||
_run(rec, server)
|
||||
|
||||
assert rec.alerts == []
|
||||
assert rec.ready.is_set() is True
|
||||
assert server._supervisor_error is None
|
||||
assert rec.steps == []
|
||||
assert rec.stop.waits == [2, 4, 2, 4]
|
||||
|
||||
|
||||
def test_lifespan_teardown_stops_and_joins_the_loop_before_the_bus_goes_down():
|
||||
"""Source-order pin (the file's style for lifespan ordering): the stop flag is
|
||||
the FIRST teardown statement, and the bounded join precedes both the bridge
|
||||
shutdown and the event-bus shutdown."""
|
||||
import inspect
|
||||
import server
|
||||
|
||||
source = inspect.getsource(server.lifespan)
|
||||
finally_idx = source.index(" finally:\n")
|
||||
stop_idx = source.index("_supervisor_stop.set()")
|
||||
join_idx = source.index("supervisor_thread.join(timeout=5)")
|
||||
bridge_idx = source.index("get_bridge().shutdown()")
|
||||
bus_idx = source.index("_shutdown_supervisor_event_bus()")
|
||||
assert finally_idx < stop_idx < join_idx < bridge_idx < bus_idx
|
||||
# Nothing between `finally:` and the stop flag but whitespace.
|
||||
assert source[finally_idx + len(" finally:\n"):stop_idx].strip() == ""
|
||||
|
||||
|
||||
def test_supervisor_revival_clears_a_stale_stop_flag(monkeypatch):
|
||||
import server
|
||||
|
||||
started = []
|
||||
|
||||
class _Thread:
|
||||
def __init__(self, **kwargs):
|
||||
self.kwargs = kwargs
|
||||
|
||||
def start(self):
|
||||
started.append(self.kwargs["target"])
|
||||
|
||||
monkeypatch.setattr(server, "has_startup_ready_provider", lambda _s: True)
|
||||
monkeypatch.setattr(server, "_supervisor_thread", None)
|
||||
monkeypatch.setattr(server.threading, "Thread", _Thread)
|
||||
server._supervisor_stop.set()
|
||||
try:
|
||||
assert server._start_supervisor_if_needed({}) is True
|
||||
assert server._supervisor_stop.is_set() is False
|
||||
assert started == [server._run_supervisor]
|
||||
finally:
|
||||
server._supervisor_stop.clear()
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue