mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
Bind lifecycle tests to their own event and maintenance completion
This commit is contained in:
parent
a3d53963f1
commit
71729d2812
3 changed files with 53 additions and 4 deletions
|
|
@ -378,7 +378,13 @@ def test_s6_subagent_tree_lineage_quiescence_and_child_result_handoff(
|
|||
assert S6_CHILD_MARKER in wait_blob, "child result text never reached the parent"
|
||||
|
||||
# Quiescence: the child's terminal task_done precedes the parent's.
|
||||
done_ids = [str(row.get("task_id") or "") for row in oracle.events("task_done")]
|
||||
def terminal_done_ids():
|
||||
rows = oracle.events("task_done")
|
||||
ids = [str(row.get("task_id") or "") for row in rows]
|
||||
return ids if child_id in ids and parent_id in ids else None
|
||||
|
||||
done_ids = wait_until(terminal_done_ids, 60)
|
||||
assert done_ids is not None, done_ids
|
||||
assert child_id in done_ids and parent_id in done_ids, done_ids
|
||||
assert done_ids.index(child_id) < done_ids.index(parent_id), done_ids
|
||||
|
||||
|
|
|
|||
|
|
@ -97,17 +97,36 @@ def test_startup_revival_keeps_actual_current_owners(tmp_path, monkeypatch, star
|
|||
assert {row.run_id for row in dc.open_runs(tmp_path)} == {f"run-{task_id}" for task_id in expected}
|
||||
|
||||
|
||||
def test_both_custody_surfaces_see_the_same_live_task_set(monkeypatch):
|
||||
def test_both_custody_surfaces_see_the_same_live_task_set(tmp_path, monkeypatch, startup_owners):
|
||||
"""The periodic sweep must hand the delegated reconciler the SAME live task set the
|
||||
process reaper gets. Two copies of "is the owner still running" is exactly how one
|
||||
custody surface ends up reaping while its twin does not."""
|
||||
import time
|
||||
import threading
|
||||
|
||||
import ouroboros.server_maintenance as sm
|
||||
import ouroboros.delegate_custody as dc
|
||||
import ouroboros.process_custody as pc
|
||||
import supervisor.queue as queue
|
||||
import supervisor.task_lifecycle as lifecycle
|
||||
|
||||
monkeypatch.setattr(queue, "DRIVE_ROOT", tmp_path)
|
||||
lock = threading.Lock()
|
||||
monkeypatch.setattr(sm, "_CANCEL_INTENT_SWEEP_LOCK", lock)
|
||||
monkeypatch.setattr(sm, "_LAST_CANCEL_INTENT_SWEEP", [0.0])
|
||||
entered, release = threading.Event(), threading.Event()
|
||||
threads = []
|
||||
def tracked_thread(**kwargs):
|
||||
thread = threading.Thread(**kwargs)
|
||||
threads.append(thread)
|
||||
return thread
|
||||
monkeypatch.setattr(sm, "threading", SimpleNamespace(Thread=tracked_thread))
|
||||
real_sweep = lifecycle.sweep_cancel_intents
|
||||
def delayed_sweep():
|
||||
entered.set()
|
||||
assert release.wait(5), "test must release its maintenance work"
|
||||
return real_sweep()
|
||||
monkeypatch.setattr(lifecycle, "sweep_cancel_intents", delayed_sweep)
|
||||
seen = {}
|
||||
monkeypatch.setattr(pc, "reap_orphaned_processes",
|
||||
lambda root, **kw: seen.__setitem__("processes", kw.get("running_task_ids")) or [])
|
||||
|
|
@ -121,8 +140,19 @@ def test_both_custody_surfaces_see_the_same_live_task_set(monkeypatch):
|
|||
from supervisor.active_activity import get_direct_activity_registry
|
||||
|
||||
get_direct_activity_registry().register("native-live", 1)
|
||||
sm._periodic_supervisor_maintenance([0.0], [time.time()])
|
||||
assert seen["processes"] == seen["delegated"] == {"t-live", "native-live"}, seen
|
||||
try:
|
||||
sm._periodic_supervisor_maintenance([0.0], [time.time()])
|
||||
assert seen["processes"] == seen["delegated"] == {"t-live", "native-live"}, seen
|
||||
assert entered.wait(2) and len(threads) == 1
|
||||
assert threads[0].name == "terminal-maintenance" and threads[0].is_alive()
|
||||
finally:
|
||||
# The actual maintenance owner finishes before monkeypatch restores its
|
||||
# root, lock and dependent functions, including when an assertion fails.
|
||||
release.set()
|
||||
for thread in threads:
|
||||
thread.join(timeout=5)
|
||||
assert all(not thread.is_alive() for thread in threads)
|
||||
assert not lock.locked(), "the real maintenance finally released its latch"
|
||||
|
||||
|
||||
def test_an_orphaned_delegated_run_is_reconciled_when_its_owner_is_gone(tmp_path, monkeypatch):
|
||||
|
|
|
|||
|
|
@ -5,6 +5,8 @@ import json
|
|||
import time
|
||||
from types import SimpleNamespace
|
||||
|
||||
import pytest
|
||||
|
||||
|
||||
def _patch_queue(queue_module, workers_module, monkeypatch, tmp_path, workers):
|
||||
monkeypatch.setattr(queue_module, "DRIVE_ROOT", tmp_path)
|
||||
|
|
@ -415,6 +417,13 @@ def test_reaper_finalizes_stuck_artifact_on_self_finalized_result(tmp_path, monk
|
|||
workers = {4: SimpleNamespace(busy_task_id=None, proc=_FakeProc(), reaping=True)}
|
||||
_patch_queue(q, w, monkeypatch, tmp_path, workers)
|
||||
monkeypatch.setattr(q, "_kept_service_pids", lambda: set(), raising=False)
|
||||
# A prior server lifespan can legitimately close the process-global bus.
|
||||
# This fixture owns its publication collector, not a new supervisor bus.
|
||||
monkeypatch.setattr(w, "_EVENT_Q_SHUTDOWN", True)
|
||||
with pytest.raises(RuntimeError, match="supervisor event bus is shutting down"):
|
||||
w.get_event_q()
|
||||
emitted = []
|
||||
monkeypatch.setattr(w, "get_event_q", lambda: SimpleNamespace(put=emitted.append))
|
||||
|
||||
calls = []
|
||||
def finalize(root, task):
|
||||
|
|
@ -448,6 +457,10 @@ def test_reaper_finalizes_stuck_artifact_on_self_finalized_result(tmp_path, monk
|
|||
_run({"id": "wt3", "type": "task", "delegation_role": "subagent",
|
||||
"task_constraint": {"mode": "local_readonly_subagent"}}, "finalizing")
|
||||
assert calls == [], "a readonly subagent has no durable artifacts to finalize"
|
||||
terminals = [event for event in emitted if event.get("type") == "task_done"]
|
||||
assert [event["task_id"] for event in terminals] == ["wt1", "wt2", "wt3"]
|
||||
assert all(event["status"] == "completed" and event["chat_id"] == 0 for event in terminals)
|
||||
assert all(event["_files_prepared_attempt"] == 1 and event["worker_id"] == 4 for event in terminals)
|
||||
|
||||
|
||||
def test_task_is_readonly_subagent_gate():
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue