mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
Measure a stalled supervisor loop by phase, CPU and event lag
On 2026-09-19 the loop stalled 64 times (95-1091 s) and the journal could say nothing but how long: no phase, no end, no CPU-vs-wall split. Every candidate cause - the inline 600 s custody block, the usage-ledger lock in the heartbeat handler, GIL starvation by heavy direct-chat threads of the same process, a daemon/pin mismatch generation - stays unproven because the onset row carries no fact that separates them. The loop now stamps liveness once per coarse tick phase (events | maintenance | assign) and publishes with each stamp the facts only the loop thread can take honestly: its own time.thread_time() delta over the interval that just ended (a thread that BURNED the wall gap reads differently from one blocked or starved), the worst worker-stamped lag of the last drain (events without a worker ts are skipped, never invented), and the in-memory daemon-pin match. supervisor_loop_stall carries them beside the wall gap and one supervisor_loop_stall_end closes the episode when the loop ticks again - an onset without an end is a generation that never recovered. The watchdog thread only reads that list: no lock, no daemon call, no disk, because a watchdog that waits on the thread it watches reports nothing. Its monotonic contract is unchanged, and so are the owner chat notice and toast. No registered-project count rides along: that set exists only as a full custody-log replay, and a measurement may not pay disk on the thread it measures. server.py sits at its pinned 1700-line bound, so the four added lines are paid inside it: a duplicated message_bus import, a single-use state temporary and two intermediate locals in _get_owner_chat_id. tests/test_server_shutdown.py drives the real loop against a stubbed clock namespace; it gains thread_time, which the loop now samples per phase.
This commit is contained in:
parent
42013142f0
commit
269ad6ced5
6 changed files with 359 additions and 21 deletions
|
|
@ -49,7 +49,7 @@ A spawned or respawned slot is not assignable until its child's PID-bound `worke
|
|||
|
||||
Unexpected worker death reserves exact custody under the queue lock and enqueues `confirmed_dead_worker` on the reaper (`worker_health.recover_confirmed_dead_worker`). A saved terminal source wins even after signal death; unknown or incomplete file publication keeps the same job (`TerminalFileRecoveryPending`); only confirmed absence of one reaches the crash policy: a signal is an infrastructure failure, an otherwise eligible non-signal crash retries within `QUEUE_MAX_RETRIES`, preserving owner-wait replay restrictions and cost. A crash storm suppresses respawn while terminal sources settle, then its fence stops pooled admission; direct chat stays available. Startup runs the same terminal-file recovery in `_run_supervisor` after process custody and before `_startup_prune_sweeps` (the no-provider lifespan runs it too, spawning nothing); unknown or still-live ownership defers it rather than racing a writer, and any unresolved or protected source, or an ownership/read error, sets `preserve_task_sources`, skipping task-drive deletion for that pass. For older canonical scheduled rows, `_recover_terminal_task_files` restores that start binding only from a known non-direct child's positive running/started-at record when the existing fresh-queue and later-worker-boot checks prove it orphaned, with no pending queue owner or active cancel; the normal orphan reconciler and terminal guards retain authority, without resuming work. The recovery report includes `rebound`.
|
||||
|
||||
Startup and throttled maintenance reconcile three residue classes. Process custody checks strict PID, start-time, command fingerprint, owner task, session and generation evidence before reaping an owned process. Delegated-run reconciliation applies the same owner-gone reasoning to external harness rows (§6 Delegated subagents). Task, review and project reconciliation repair durable records whose producer no longer exists. None of these are command-line-class kill sweeps, and one development or runtime instance never reaps another. The dedicated watchdog separately observes supervisor-loop liveness and every registered native actor; it alerts and recommends `/restart` but cannot kill an in-process thread. Other owner conversations run on independent native actors, without a second scheduler.
|
||||
Startup and throttled maintenance reconcile three residue classes. Process custody checks strict PID, start-time, command fingerprint, owner task, session and generation evidence before reaping an owned process. Delegated-run reconciliation applies the same owner-gone reasoning to external harness rows (§6 Delegated subagents). Task, review and project reconciliation repair durable records whose producer no longer exists. None of these are command-line-class kill sweeps, and one development or runtime instance never reaps another. The dedicated watchdog separately observes phase-stamped loop liveness and every registered native actor; it alerts and recommends `/restart` but cannot kill an in-process thread. Other owner conversations run on independent native actors, without a second scheduler.
|
||||
|
||||
Cooperative project checkpointing has two equivalent quiescence triggers: a host-minted genesis or cooperative tree is checked when its root settles with no live descendants, and again when the last child settles beneath an already-terminal root — a root-scope budget stop terminalizes the root before its children, so a root-only trigger would see a live tree once and never return. The bounded git chain runs on a daemon thread, revalidates quiescence under the queue lock immediately before mutation, and replays a trigger that arrives during an in-flight check. Only host-minted project roots are eligible: owner-attached folders are never auto-committed, credential-shaped files stay excluded and disclosed, and every material success, skip or error receives a durable receipt.
|
||||
|
||||
|
|
|
|||
|
|
@ -21,6 +21,7 @@ This chapter is the short list of properties the rest of the book must not contr
|
|||
17. **Review spend has one ceiling.** Every paid review gate shares `OUROBOROS_REVIEW_MAX_CYCLES`; the per-gate meanings live once in §6 Review stack, the SSOT in `review_cycles.py`. This paid ceiling does not remove the last author reaction. Blocking correction or stop grants no approval; informed Advisory finish keeps current-author authority separate from the critic.
|
||||
18. **A typed permanent engine refusal discharges a custody duty once.** A custody duty the engine refuses with a typed permanent code is discharged once, durably, under the engine's own code — never retried on a timer and never recorded as a deletion. Owner: `delegate_custody._retire_project_locked` (§9 registration pass).
|
||||
19. **An interrupted parent leaves no orphan.** A child that cannot start because its parent was interrupted is cancelled with the parent's cause, whichever door the interruption came through — Restart, crash or window close, pooled or direct. The planned path settles it in `kill_workers(preserve_pending=True)` and snapshot restore marks the same shutdown custody on boot (`pending_parent_interrupted`); both write a ledger-reconstructed cost and publish the terminal event, so no queued row is ever revived under a parent that no longer exists. Owners: `supervisor/workers.py`, `supervisor/queue_snapshot.py` (§5 for the flow, §9 for the doors).
|
||||
20. **A stalled loop says where it went silent.** `supervisor_loop_stall` carries the tick phase (`events`/`maintenance`/`assign`), the loop thread's CPU against the wall gap, the worst worker-stamped event lag and `daemon_pin_matched`; `supervisor_loop_stall_end` closes every alerted stall. The loop publishes these facts; the watchdog only reads them — no lock, daemon call or disk scan on either side.
|
||||
|
||||
### 10.1 Continuity data-flow map
|
||||
|
||||
|
|
|
|||
|
|
@ -1,8 +1,9 @@
|
|||
"""Wedge detection for the supervisor generation.
|
||||
|
||||
The two silent-wedge predicates (a stalled supervisor loop, a heartbeat-silent
|
||||
in-process chat turn), the owner alert one of them raises, and the dedicated
|
||||
watchdog thread that evaluates both outside the loop it watches.
|
||||
in-process chat turn), the owner alert one of them raises, the measurements the
|
||||
loop publishes about itself, and the dedicated watchdog thread that evaluates
|
||||
both outside the loop it watches.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
|
@ -10,10 +11,20 @@ from __future__ import annotations
|
|||
import os
|
||||
import threading
|
||||
import time
|
||||
from typing import Any, Optional
|
||||
|
||||
from ouroboros.deadline_utils import parse_deadline_ts, utc_now
|
||||
from ouroboros.server_process import DATA_DIR, log, _restart_requested
|
||||
from ouroboros.utils import utc_now_iso
|
||||
|
||||
# The loop hands the watchdog ONE list: [0] the monotonic stamp the watchdog
|
||||
# triggers on, [1] the facts published with that stamp, [2] the loop thread's CPU
|
||||
# base and [3] the worst worker-event lag of the drain in progress. The LOOP
|
||||
# THREAD owns every write; the watchdog only reads them, so it never takes a
|
||||
# lock, calls the daemon or touches disk — a watchdog that waits on the thread it
|
||||
# watches reports nothing. Older/foreign callers may pass the stamp alone.
|
||||
_STAMP, _FACTS, _CPU, _LAG = 0, 1, 2, 3
|
||||
|
||||
|
||||
def _supervisor_loop_stalled(last_tick: float, now: float, deadline_sec: int) -> bool:
|
||||
"""True when the supervisor loop has not published a liveness tick within the
|
||||
|
|
@ -21,6 +32,82 @@ def _supervisor_loop_stalled(last_tick: float, now: float, deadline_sec: int) ->
|
|||
return deadline_sec > 0 and (now - last_tick) > deadline_sec
|
||||
|
||||
|
||||
def _daemon_pin_matched() -> Optional[bool]:
|
||||
"""Does the engine this process already PROVED serve the next-spawn pin?
|
||||
|
||||
A generation whose daemon lags the installed pin is one of the suspects behind
|
||||
a stalled loop (every ``ensure`` on such a generation pays a probe/install
|
||||
path a matched one skips), so the stall row carries the answer the process
|
||||
ALREADY holds in memory: the version proven by the last successful handshake
|
||||
(``owned_engine_version``, explicitly no new I/O) against the loaded pin.
|
||||
``None`` means unknown — no handshake has succeeded here yet, or this install
|
||||
carries no pin — never a guess.
|
||||
"""
|
||||
try:
|
||||
from ouroboros.claudexor_daemon import owned_engine_version
|
||||
from ouroboros.claudexor_runtime import get_runtime_manager
|
||||
|
||||
pin = get_runtime_manager().pin
|
||||
proven = owned_engine_version()
|
||||
return proven == pin.version if (proven and pin is not None) else None
|
||||
except Exception:
|
||||
log.debug("owned-daemon pin match is unknown", exc_info=True)
|
||||
return None
|
||||
|
||||
|
||||
def loop_phase_facts(liveness: list, phase: str, *, new_tick: bool = False) -> dict:
|
||||
"""The measurements the LOOP THREAD publishes with its own liveness stamp.
|
||||
|
||||
``phase`` is the coarse tick phase the loop is entering — ``events`` |
|
||||
``maintenance`` | ``assign``, one stamp per phase and never per sub-step — so a
|
||||
stall names where the thread went silent instead of only how long it was.
|
||||
``loop_thread_cpu_sec`` is the ``time.thread_time()`` delta over the interval
|
||||
that ENDS with this stamp, sampled on the loop thread itself: read beside the
|
||||
wall gap it separates a thread that BURNED that gap from one blocked on a lock
|
||||
or starved of the GIL. ``max_event_lag_sec`` is the worst worker-stamped lag of
|
||||
the most recently completed drain, absent when no drained event carried a
|
||||
worker stamp. ``new_tick`` opens a fresh drain maximum (the events phase opens
|
||||
the tick), so a lag can never outlive the tick that observed it. No
|
||||
registered-project count rides here: the set the retire sweep walks exists
|
||||
only as a full custody-log replay, and a measurement may not pay disk on the
|
||||
thread it measures — an absent key beats a cheap-looking wrong one.
|
||||
"""
|
||||
cpu = time.thread_time()
|
||||
facts = {
|
||||
"phase": phase,
|
||||
"loop_thread_cpu_sec": round(cpu - liveness[_CPU], 3),
|
||||
"daemon_pin_matched": _daemon_pin_matched(),
|
||||
}
|
||||
if liveness[_LAG] is not None:
|
||||
facts["max_event_lag_sec"] = round(liveness[_LAG], 1)
|
||||
liveness[_CPU] = cpu
|
||||
if new_tick:
|
||||
liveness[_LAG] = None
|
||||
return facts
|
||||
|
||||
|
||||
def observe_worker_event_lag(liveness: list, evt: Any) -> None:
|
||||
"""Record how far behind the loop is on the worker event it is draining now.
|
||||
|
||||
The worker stamps ``ts`` in its OWN process when it queues the event, so this
|
||||
is a cross-process wall-clock gap — a measurement, never the watchdog's
|
||||
monotonic trigger. An event without a parsable worker stamp is skipped: the
|
||||
host's own clock is not evidence about when a worker spoke.
|
||||
"""
|
||||
stamped = parse_deadline_ts(evt.get("ts")) if isinstance(evt, dict) else None
|
||||
if stamped is None:
|
||||
return
|
||||
lag = (utc_now() - stamped).total_seconds()
|
||||
if liveness[_LAG] is None or lag > liveness[_LAG]:
|
||||
liveness[_LAG] = lag
|
||||
|
||||
|
||||
def _published_loop_facts(liveness: list) -> dict:
|
||||
"""The facts the loop published with its last stamp ({} when it published none)."""
|
||||
facts = liveness[_FACTS] if len(liveness) > _FACTS else None
|
||||
return dict(facts) if isinstance(facts, dict) else {}
|
||||
|
||||
|
||||
def _chat_turn_wedged(busy: bool, last_activity_ts, now: float, deadline_sec: int) -> bool:
|
||||
"""True when an IN-PROCESS direct-chat turn is busy but its liveness tick has been
|
||||
silent past the deadline (WS3). ``last_activity_ts is None`` => the turn has not
|
||||
|
|
@ -69,7 +156,10 @@ def _start_supervisor_liveness_watchdog(liveness: list, stop_event=None) -> None
|
|||
loop stall (new-message intake starvation) and a heartbeat-silent in-process
|
||||
direct-chat turn — converting a multi-hour silent wedge into an immediate signal.
|
||||
It deliberately does NOT kill a hung thread; independent native actors keep
|
||||
the chat responsive meanwhile. ``stop_event`` is
|
||||
the chat responsive meanwhile. A loop stall journals ``supervisor_loop_stall``
|
||||
with the phase facts the loop published with its last stamp and, once the loop
|
||||
ticks again, one ``supervisor_loop_stall_end`` — onset without an end is a
|
||||
generation that never recovered. ``stop_event`` is
|
||||
a PER-GENERATION token: when the supervisor loop that owns ``liveness`` exits (incl.
|
||||
the crash-storm death path, which never sets the global restart flag), it is set so
|
||||
this watchdog stops watching a now-stale liveness list (no false post-revival alert)."""
|
||||
|
|
@ -83,6 +173,7 @@ def _start_supervisor_liveness_watchdog(liveness: list, stop_event=None) -> None
|
|||
from supervisor.state import append_jsonl, load_state
|
||||
interval = min(15, max(1, deadline // 3))
|
||||
loop_alerted = False
|
||||
stall_onset: tuple = () # (stalled stamp, phase) of the OPEN alerted stall
|
||||
wedged_tasks: set[str] = set()
|
||||
while not _restart_requested.is_set() and not (stop_event is not None and stop_event.is_set()):
|
||||
time.sleep(interval)
|
||||
|
|
@ -93,16 +184,21 @@ def _start_supervisor_liveness_watchdog(liveness: list, stop_event=None) -> None
|
|||
# nor mask a real one on either half.
|
||||
now = time.monotonic()
|
||||
# (1) Supervisor loop stall — new-message intake starvation.
|
||||
if _supervisor_loop_stalled(liveness[0], now, deadline):
|
||||
if _supervisor_loop_stalled(liveness[_STAMP], now, deadline):
|
||||
if not loop_alerted:
|
||||
gap = now - liveness[0]
|
||||
gap = now - liveness[_STAMP]
|
||||
facts = _published_loop_facts(liveness)
|
||||
log.error(
|
||||
"Supervisor loop STALLED ~%.0fs — new-message intake starved (native "
|
||||
"chat still answers); investigate a blocking step.", gap,
|
||||
)
|
||||
try:
|
||||
# The facts the loop published with the stamp it went silent
|
||||
# on: where it was, what its own thread burned, how far
|
||||
# behind the drained worker events already were.
|
||||
append_jsonl(DATA_DIR / "logs" / "supervisor.jsonl", {
|
||||
"ts": utc_now_iso(), "type": "supervisor_loop_stall", "stalled_sec": round(gap, 1),
|
||||
"ts": utc_now_iso(), "type": "supervisor_loop_stall",
|
||||
"stalled_sec": round(gap, 1), **facts,
|
||||
})
|
||||
except Exception:
|
||||
log.debug("loop-stall log failed", exc_info=True)
|
||||
|
|
@ -120,13 +216,28 @@ def _start_supervisor_liveness_watchdog(liveness: list, stop_event=None) -> None
|
|||
# pid disambiguates server GENERATIONS: the monotonic stamp alone can
|
||||
# repeat at a similar uptime offset across restarts, and the browser's
|
||||
# toast-dedupe set outlives this process while the page stays open.
|
||||
"toast_once": f"supervisor-loop-stall:{os.getpid()}:{int(liveness[0])}",
|
||||
"toast_once": f"supervisor-loop-stall:{os.getpid()}:{int(liveness[_STAMP])}",
|
||||
},
|
||||
role="system", system_type="runtime_liveness_notice")
|
||||
except Exception:
|
||||
log.debug("loop-stall owner alert failed", exc_info=True)
|
||||
loop_alerted = True
|
||||
stall_onset = (liveness[_STAMP], facts.get("phase"))
|
||||
else:
|
||||
if loop_alerted:
|
||||
# The loop ticked again: close the episode ONCE, and only one
|
||||
# that was alerted. Both ends are stamps the LOOP published on
|
||||
# the monotonic clock, so the duration survives a wall-clock
|
||||
# jump; it rounds up by at most one watchdog interval, the
|
||||
# resolution at which recovery is observed at all.
|
||||
try:
|
||||
append_jsonl(DATA_DIR / "logs" / "supervisor.jsonl", {
|
||||
"ts": utc_now_iso(), "type": "supervisor_loop_stall_end",
|
||||
"stalled_sec": round(liveness[_STAMP] - stall_onset[0], 1),
|
||||
"phase": stall_onset[1],
|
||||
})
|
||||
except Exception:
|
||||
log.debug("loop-stall-end log failed", exc_info=True)
|
||||
loop_alerted = False
|
||||
# (2) Each native actor has its own liveness and alert identity.
|
||||
try:
|
||||
|
|
|
|||
24
server.py
24
server.py
|
|
@ -593,8 +593,7 @@ def _run_supervisor(settings: dict) -> None:
|
|||
try:
|
||||
ensure_legacy_imported(pathlib.Path(DATA_DIR))
|
||||
|
||||
from supervisor.message_bus import init as bus_init
|
||||
from supervisor.message_bus import LocalChatBridge
|
||||
from supervisor.message_bus import LocalChatBridge, init as bus_init
|
||||
|
||||
bridge = LocalChatBridge(settings)
|
||||
bridge._broadcast_fn = broadcast_ws_sync
|
||||
|
|
@ -699,9 +698,7 @@ def _run_supervisor(settings: dict) -> None:
|
|||
|
||||
def _get_owner_chat_id() -> Optional[int]:
|
||||
try:
|
||||
st = load_state()
|
||||
cid = st.get("owner_chat_id")
|
||||
return int(cid) if cid else None
|
||||
return int((load_state() or {}).get("owner_chat_id") or 0) or None
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
|
@ -709,8 +706,7 @@ def _run_supervisor(settings: dict) -> None:
|
|||
drive_root=DATA_DIR, repo_dir=REPO_DIR, owner_chat_id_fn=_get_owner_chat_id,
|
||||
routing_metadata_fn=lambda cid: main_lane_routing_metadata(_event_ctx, cid)) # _event_ctx is built below
|
||||
|
||||
_bg_st = load_state()
|
||||
if _bg_st.get("bg_consciousness_enabled"):
|
||||
if load_state().get("bg_consciousness_enabled"):
|
||||
_consciousness.start()
|
||||
log.info("Background consciousness auto-restored from saved state.")
|
||||
|
||||
|
|
@ -763,15 +759,16 @@ def _run_supervisor(settings: dict) -> None:
|
|||
_last_review_job_reconcile = [time.time()]
|
||||
# WS3: a dedicated watchdog thread (outside this loop, so it fires even if the
|
||||
# loop stalls) surfaces a wedge as an observable signal + owner alert instead
|
||||
# of silent hours; the loop publishes a liveness tick each iteration. The tick
|
||||
# is MONOTONIC: it is only ever read as an elapsed gap, so a wall-clock jump
|
||||
# must not turn a healthy loop into a phantom stall (nor hide a real one).
|
||||
_loop_liveness = [time.monotonic()]
|
||||
# of silent hours; the loop publishes a liveness tick at each tick PHASE. The
|
||||
# tick is MONOTONIC: it is only ever read as an elapsed gap, so a wall-clock
|
||||
# jump must not turn a healthy loop into a phantom stall (nor hide a real one).
|
||||
from ouroboros.server_liveness import loop_phase_facts, observe_worker_event_lag
|
||||
_loop_liveness = [time.monotonic(), {}, time.thread_time(), None] # slots: server_liveness.py
|
||||
_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() and not _supervisor_stop.is_set():
|
||||
try:
|
||||
_loop_liveness[0] = time.monotonic()
|
||||
_loop_liveness[1], _loop_liveness[0] = loop_phase_facts(_loop_liveness, "events", new_tick=True), time.monotonic()
|
||||
rotate_chat_log_if_needed(DATA_DIR)
|
||||
# progress.jsonl rotates on the same supervisor tick (v6.90.x P2); its
|
||||
# readers (history backfill, SSE replay, api_logs_tail, TB ATIF) are
|
||||
|
|
@ -800,6 +797,7 @@ def _run_supervisor(settings: dict) -> None:
|
|||
if evt.get("type") == "restart_request":
|
||||
_handle_restart_in_supervisor(evt, _event_ctx)
|
||||
continue
|
||||
observe_worker_event_lag(_loop_liveness, evt)
|
||||
dispatch_event(evt, _event_ctx)
|
||||
|
||||
if _restart_requested.is_set():
|
||||
|
|
@ -811,6 +809,7 @@ def _run_supervisor(settings: dict) -> None:
|
|||
# where no task_received fired for hours until a full restart).
|
||||
offset = _process_bridge_updates(bridge, offset, _event_ctx)
|
||||
|
||||
_loop_liveness[1], _loop_liveness[0] = loop_phase_facts(_loop_liveness, "maintenance"), time.monotonic()
|
||||
enforce_task_timeouts()
|
||||
try:
|
||||
from supervisor.queue import check_scheduled_tasks
|
||||
|
|
@ -821,6 +820,7 @@ def _run_supervisor(settings: dict) -> None:
|
|||
_last_custody_reap, _last_review_job_reconcile,
|
||||
on_orphans_healed=lambda count: _consciousness and _consciousness.notify(f"orphans_healed:{count}"),
|
||||
)
|
||||
_loop_liveness[1], _loop_liveness[0] = loop_phase_facts(_loop_liveness, "assign"), time.monotonic()
|
||||
# Loop-tick restart drain (no sleep, events keep flowing): while
|
||||
# draining a deferred restart, skip starting new work the restart
|
||||
# deadline would immediately chop (evolution / pending project tasks).
|
||||
|
|
|
|||
|
|
@ -648,7 +648,9 @@ def _supervisor_harness(monkeypatch, tmp_path, steps):
|
|||
noop = lambda *_a, **_k: None # noqa: E731
|
||||
monkeypatch.setattr(server, "DATA_DIR", tmp_path)
|
||||
# Patch the module's bound name, not the process-wide time.sleep.
|
||||
monkeypatch.setattr(server, "time", SimpleNamespace(sleep=noop, monotonic=time_mod.monotonic, time=time_mod.time))
|
||||
monkeypatch.setattr(server, "time", SimpleNamespace(
|
||||
sleep=noop, monotonic=time_mod.monotonic, time=time_mod.time,
|
||||
thread_time=time_mod.thread_time)) # the loop samples its OWN thread's CPU per phase stamp
|
||||
monkeypatch.setattr(server, "_supervisor_stop", rec.stop)
|
||||
monkeypatch.setattr(server, "_restart_requested", rec.restart)
|
||||
monkeypatch.setattr(server, "_supervisor_ready", rec.ready)
|
||||
|
|
|
|||
224
tests/test_supervisor_loop_measurements.py
Normal file
224
tests/test_supervisor_loop_measurements.py
Normal file
|
|
@ -0,0 +1,224 @@
|
|||
"""Honest measurements of a stalled supervisor loop (C1).
|
||||
|
||||
The watchdog used to record only the ONSET of a stall and nothing about it: no
|
||||
phase, no end, no CPU-vs-wall split. A 64-stall night could therefore not tell a
|
||||
loop thread BURNING its wall gap from one blocked on a lock or starved of the
|
||||
GIL, nor say which coarse phase it went silent in. These pins cover the
|
||||
measurement contract in both directions: the facts are published BY THE LOOP
|
||||
THREAD with its liveness stamp, the watchdog thread only reads them, and nothing
|
||||
is invented for an event a worker never stamped.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import inspect
|
||||
import re
|
||||
import threading
|
||||
import time
|
||||
|
||||
import pytest
|
||||
|
||||
|
||||
class _Clock:
|
||||
"""Controllable stand-in for the ``time`` module the watchdog reads.
|
||||
|
||||
Only ``monotonic``/``sleep`` are driven; the harness keeps the real clock, so
|
||||
a simulated stall cannot make the test's own timeouts lie.
|
||||
"""
|
||||
|
||||
def __init__(self, *, mono: float) -> None:
|
||||
self.mono = mono
|
||||
self.ticks = 0
|
||||
|
||||
def __getattr__(self, name):
|
||||
return getattr(time, name)
|
||||
|
||||
def monotonic(self) -> float:
|
||||
return self.mono
|
||||
|
||||
def sleep(self, _seconds: float) -> None:
|
||||
self.ticks += 1
|
||||
time.sleep(0.01) # a REAL yield; the fake clock moves only when a test says so
|
||||
|
||||
|
||||
def _wait_until(predicate, budget: float = 6.0) -> None:
|
||||
end = time.time() + budget # real clock: the harness never rides the fake one
|
||||
while not predicate() and time.time() < end:
|
||||
time.sleep(0.01)
|
||||
|
||||
|
||||
def _stop_watchdog(stop: threading.Event) -> None:
|
||||
stop.set()
|
||||
for thread in threading.enumerate():
|
||||
if thread.name == "supervisor-liveness-watchdog":
|
||||
thread.join(timeout=5)
|
||||
|
||||
|
||||
@pytest.fixture()
|
||||
def journal(monkeypatch):
|
||||
"""Collect the durable supervisor rows without touching a live data root."""
|
||||
rows: list = []
|
||||
|
||||
def _append(path, obj):
|
||||
rows.append((str(path), dict(obj)))
|
||||
return True
|
||||
|
||||
monkeypatch.setattr("supervisor.state.append_jsonl", _append)
|
||||
monkeypatch.setattr("supervisor.state.load_state", lambda: {}) # no owner chat: no alert send
|
||||
from supervisor.active_activity import get_direct_activity_registry
|
||||
|
||||
get_direct_activity_registry().clear() # isolate the loop-stall half
|
||||
return rows
|
||||
|
||||
|
||||
def _live_liveness(phase: str = "maintenance") -> list:
|
||||
"""A liveness list exactly as ``server.py::_run_supervisor`` publishes it."""
|
||||
from ouroboros import server_liveness
|
||||
|
||||
liveness = [time.monotonic(), {}, time.thread_time(), None]
|
||||
server_liveness.observe_worker_event_lag(liveness, {
|
||||
"type": "task_heartbeat", "task_id": "t-1", "ts": _iso_ago(4.0),
|
||||
})
|
||||
liveness[1] = server_liveness.loop_phase_facts(liveness, phase)
|
||||
return liveness
|
||||
|
||||
|
||||
def _iso_ago(seconds: float) -> str:
|
||||
import datetime as _dt
|
||||
|
||||
return (_dt.datetime.now(tz=_dt.timezone.utc) - _dt.timedelta(seconds=seconds)).isoformat()
|
||||
|
||||
|
||||
def test_stall_row_carries_phase_cpu_and_event_lag(monkeypatch, journal):
|
||||
"""The onset row names WHERE the loop went silent and what the thread was doing."""
|
||||
import server
|
||||
from ouroboros import server_liveness
|
||||
|
||||
monkeypatch.setenv("OUROBOROS_SUPERVISOR_LIVENESS_DEADLINE_SEC", "1")
|
||||
liveness = _live_liveness("maintenance")
|
||||
clock = _Clock(mono=liveness[0] + 100.0) # 100s of MONOTONIC silence
|
||||
monkeypatch.setattr(server_liveness, "time", clock)
|
||||
stop = threading.Event()
|
||||
try:
|
||||
server._start_supervisor_liveness_watchdog(liveness, stop)
|
||||
_wait_until(lambda: journal)
|
||||
finally:
|
||||
_stop_watchdog(stop)
|
||||
stalls = [row for _path, row in journal if row["type"] == "supervisor_loop_stall"]
|
||||
assert len(stalls) == 1, journal
|
||||
assert journal[0][0].endswith("logs/supervisor.jsonl")
|
||||
row = stalls[0]
|
||||
assert row["stalled_sec"] == pytest.approx(100.0, abs=1.0)
|
||||
assert row["phase"] == "maintenance"
|
||||
assert isinstance(row["loop_thread_cpu_sec"], float)
|
||||
assert row["max_event_lag_sec"] == pytest.approx(4.0, abs=2.0)
|
||||
assert "daemon_pin_matched" in row and row["daemon_pin_matched"] in (True, False, None)
|
||||
|
||||
|
||||
def test_stall_end_is_written_once_per_alerted_stall(monkeypatch, journal):
|
||||
"""Closing row: exactly one per alerted stall, carrying the STALLED phase."""
|
||||
import server
|
||||
from ouroboros import server_liveness
|
||||
|
||||
monkeypatch.setenv("OUROBOROS_SUPERVISOR_LIVENESS_DEADLINE_SEC", "1")
|
||||
liveness = _live_liveness("maintenance")
|
||||
clock = _Clock(mono=liveness[0] + 100.0)
|
||||
monkeypatch.setattr(server_liveness, "time", clock)
|
||||
stop = threading.Event()
|
||||
try:
|
||||
server._start_supervisor_liveness_watchdog(liveness, stop)
|
||||
_wait_until(lambda: journal)
|
||||
# The loop ticks again, in its NEXT phase: the episode closes exactly once.
|
||||
liveness[1] = server_liveness.loop_phase_facts(liveness, "assign")
|
||||
liveness[0] = clock.mono
|
||||
_wait_until(lambda: any(row["type"] == "supervisor_loop_stall_end" for _p, row in journal))
|
||||
ticks = clock.ticks
|
||||
_wait_until(lambda: clock.ticks > ticks + 2) # keep polling a healthy loop
|
||||
finally:
|
||||
_stop_watchdog(stop)
|
||||
kinds = [row["type"] for _path, row in journal]
|
||||
assert kinds.count("supervisor_loop_stall") == 1, journal
|
||||
assert kinds.count("supervisor_loop_stall_end") == 1, journal
|
||||
end = next(row for _p, row in journal if row["type"] == "supervisor_loop_stall_end")
|
||||
assert end["stalled_sec"] == pytest.approx(100.0, abs=1.0)
|
||||
assert end["phase"] == "maintenance" # where it was stuck, not where it resumed
|
||||
|
||||
|
||||
def test_a_healthy_loop_journals_neither_row(monkeypatch, journal):
|
||||
"""The quiet direction: a ticking loop writes no stall and no end row."""
|
||||
import server
|
||||
from ouroboros import server_liveness
|
||||
|
||||
monkeypatch.setenv("OUROBOROS_SUPERVISOR_LIVENESS_DEADLINE_SEC", "1")
|
||||
liveness = _live_liveness("events")
|
||||
clock = _Clock(mono=liveness[0])
|
||||
monkeypatch.setattr(server_liveness, "time", clock)
|
||||
stop = threading.Event()
|
||||
try:
|
||||
server._start_supervisor_liveness_watchdog(liveness, stop)
|
||||
_wait_until(lambda: clock.ticks >= 3)
|
||||
finally:
|
||||
_stop_watchdog(stop)
|
||||
assert journal == []
|
||||
|
||||
|
||||
def test_an_event_without_a_worker_stamp_never_produces_a_lag(monkeypatch):
|
||||
"""Events carry the worker's own ``ts``; an unstamped one is skipped, not invented."""
|
||||
from ouroboros import server_liveness
|
||||
|
||||
liveness = [time.monotonic(), {}, time.thread_time(), None]
|
||||
for unstamped in ({"type": "task_message_injected"}, {"type": "x", "ts": ""},
|
||||
{"type": "x", "ts": "not-a-timestamp"}, {"type": "x", "ts": None}):
|
||||
server_liveness.observe_worker_event_lag(liveness, unstamped)
|
||||
assert "max_event_lag_sec" not in server_liveness.loop_phase_facts(liveness, "maintenance")
|
||||
# The other direction: a worker-stamped event IS measured.
|
||||
server_liveness.observe_worker_event_lag(liveness, {"type": "task_heartbeat", "ts": _iso_ago(7.0)})
|
||||
facts = server_liveness.loop_phase_facts(liveness, "maintenance")
|
||||
assert facts["max_event_lag_sec"] == pytest.approx(7.0, abs=2.0)
|
||||
|
||||
|
||||
def test_the_drain_maximum_is_scoped_to_one_tick(monkeypatch):
|
||||
"""``new_tick`` opens a fresh drain maximum, so a stale lag cannot ride forever."""
|
||||
from ouroboros import server_liveness
|
||||
|
||||
liveness = [time.monotonic(), {}, time.thread_time(), None]
|
||||
server_liveness.observe_worker_event_lag(liveness, {"type": "task_heartbeat", "ts": _iso_ago(5.0)})
|
||||
opening = server_liveness.loop_phase_facts(liveness, "events", new_tick=True)
|
||||
assert opening["max_event_lag_sec"] == pytest.approx(5.0, abs=2.0)
|
||||
assert "max_event_lag_sec" not in server_liveness.loop_phase_facts(liveness, "maintenance")
|
||||
|
||||
|
||||
def test_loop_thread_cpu_separates_a_burning_thread_from_a_blocked_one(monkeypatch):
|
||||
"""The CPU-vs-wall split: the delta is the LOOP THREAD's own processor time."""
|
||||
from ouroboros import server_liveness
|
||||
|
||||
liveness = [time.monotonic(), {}, time.thread_time(), None]
|
||||
time.sleep(0.05) # wall time passes, this thread burns nothing
|
||||
idle = server_liveness.loop_phase_facts(liveness, "events")["loop_thread_cpu_sec"]
|
||||
burn_until = time.thread_time() + 0.05
|
||||
while time.thread_time() < burn_until:
|
||||
pass
|
||||
busy = server_liveness.loop_phase_facts(liveness, "maintenance")["loop_thread_cpu_sec"]
|
||||
assert idle < 0.03, idle
|
||||
assert busy >= 0.04, busy
|
||||
|
||||
|
||||
def test_the_loop_publishes_one_monotonic_stamp_per_tick_phase():
|
||||
"""Source pin of the PRODUCER (the behavioural tests feed hand-built stamps).
|
||||
|
||||
Three coarse phases, one stamp each, every stamp taken on ``time.monotonic()``
|
||||
(OB-03: a wall-clock jump must never fabricate or mask a stall), and the drain
|
||||
itself observes the worker-event lag it reports.
|
||||
"""
|
||||
import server
|
||||
|
||||
source = inspect.getsource(server._run_supervisor)
|
||||
stamps = re.findall(
|
||||
r'_loop_liveness\[1\], _loop_liveness\[0\] = loop_phase_facts\(\s*_loop_liveness, "(\w+)"'
|
||||
r'[^)]*\), time\.(\w+)\(\)',
|
||||
source,
|
||||
)
|
||||
assert [phase for phase, _clock in stamps] == ["events", "maintenance", "assign"], stamps
|
||||
assert {clock for _phase, clock in stamps} == {"monotonic"}, stamps
|
||||
assert "observe_worker_event_lag(_loop_liveness, evt)" in source
|
||||
|
||||
Loading…
Add table
Add a link
Reference in a new issue