mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-04 22:31:51 +00:00
Request delegated cancellation when settling a missing worker
Forward the impending cancelled verdict through the existing custody audit
only after child-result copyback confirms the owner has not settled. This
sends Stop and shutdown-restore cancellation requests before publication
while preserving healthy delegated runs behind provider failures.
Exercise real custody replay and result writes for both cancellation causes
and for prior or copied-back infrastructure failures.
(cherry picked from commit 3b9fa7fab7)
This commit is contained in:
parent
f5ca429be2
commit
4101c46d7a
2 changed files with 75 additions and 2 deletions
|
|
@ -693,9 +693,13 @@ def _finalize_cancel_intent_on_miss(
|
|||
# GR5-3: neither queued nor running — the worker is gone, but its
|
||||
# delegated runs may still be live; audit custody like the kill path
|
||||
# and thread the disclosure into every miss-lane delivery below.
|
||||
audit = _audit_delegated_runs_on_kill(q, task_id)
|
||||
unreconciled = list(audit.get("unreconciled") or [])
|
||||
# Read after child copyback: only an unsettled task will receive our
|
||||
# cancelled write. An existing terminal keeps its own custody verdict.
|
||||
settled = _settled_status(q.DRIVE_ROOT, task_id)
|
||||
audit = _audit_delegated_runs_on_kill(
|
||||
q, task_id, **({} if settled else {"deliberate_terminal": STATUS_CANCELLED}),
|
||||
)
|
||||
unreconciled = list(audit.get("unreconciled") or [])
|
||||
if settled:
|
||||
_recover_stranded_reaping_slot(q, task_id, active)
|
||||
# D1b (R4): this branch performs no terminal write of its own, so
|
||||
|
|
|
|||
|
|
@ -19,6 +19,7 @@ import pytest
|
|||
from ouroboros.config import CLAUDEXOR_DELEGATED_MARKER_MIN_VERSION
|
||||
from ouroboros.gateways import claudexor as cx
|
||||
|
||||
from tests._cancel_intents_shared import qenv as _qenv
|
||||
from tests._delegated_transport_shared import ( # noqa: F401 (autouse fixture applies on import)
|
||||
_LiveRunStub,
|
||||
_event_types,
|
||||
|
|
@ -28,6 +29,8 @@ from tests._delegated_transport_shared import ( # noqa: F401 (autouse fixture
|
|||
_write_attempt,
|
||||
)
|
||||
|
||||
qenv = _qenv
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def startup_owners(tmp_path, monkeypatch):
|
||||
|
|
@ -594,3 +597,69 @@ def test_a_settled_owner_lets_its_durable_outcome_decide_the_live_run(
|
|||
assert [row["action"] for row in outcomes] == [expected]
|
||||
assert transport.cancels == cancelled
|
||||
dc._CUSTODY.clear()
|
||||
|
||||
|
||||
@pytest.mark.parametrize("case", ["owner_stop", "snapshot_restore", "prior_infra", "child_infra"])
|
||||
def test_miss_lane_cancels_live_delegation_only_for_its_own_terminal(qenv, monkeypatch, case):
|
||||
"""A fresh Stop requests cancellation before publication, including after restore.
|
||||
|
||||
A prior provider failure, even one only found by child-result copyback, keeps
|
||||
its outcome and its healthy delegated run. Exercise the real custody/audit
|
||||
path; only the external transport and owner delivery are replaced.
|
||||
"""
|
||||
from ouroboros import cancel_intents, claudexor_daemon, delegate_custody as dc
|
||||
from ouroboros.outcomes import infra_failed_axes
|
||||
from ouroboros.task_results import load_task_result, write_task_result
|
||||
from supervisor import cancel_publication
|
||||
|
||||
root, task_id, run_id = qenv.drive, "owner", "run-owned"
|
||||
status_at_cancel = []
|
||||
|
||||
class Gateway(_LiveRunStub):
|
||||
def cancel_run(self, rid, reason=""):
|
||||
status_at_cancel.append(load_task_result(root, task_id)["status"])
|
||||
return super().cancel_run(rid, reason)
|
||||
|
||||
transport = Gateway(run_id=run_id)
|
||||
monkeypatch.setattr(claudexor_daemon, "ensure_owned_gateway", lambda: transport)
|
||||
monkeypatch.setattr(claudexor_daemon, "read_owned_gateway", lambda: transport)
|
||||
monkeypatch.setattr(cancel_publication, "_deliver_on_miss", lambda *a, **kw: True)
|
||||
monkeypatch.setattr(qenv.tl, "_deliver_on_miss", lambda *a, **kw: True)
|
||||
monkeypatch.setattr(qenv.q, "_emit_cancel_task_done", lambda *a, **kw: None)
|
||||
monkeypatch.setattr(dc, "_CUSTODY", {})
|
||||
dc.record_started(root, dc.RunCustody(
|
||||
run_id=run_id, task_id=task_id, route_id="r", model="m", project_id="p",
|
||||
project_owned=False, root_task_id=task_id, ledger_root=str(root),
|
||||
))
|
||||
dc._CUSTODY.clear() # Recover the actual durable custody, as after restart.
|
||||
|
||||
fields = {"root_task_id": task_id, "result": "working"}
|
||||
failure = {"result": "provider died", "outcome_axes": infra_failed_axes("provider_unavailable")}
|
||||
if case == "child_infra":
|
||||
child = root / "state" / "headless_tasks" / task_id / "data"
|
||||
write_task_result(child, task_id, "failed", root_task_id=task_id, **failure)
|
||||
fields.update(child_drive_root=str(child), delegation_role="subagent")
|
||||
if case == "prior_infra":
|
||||
fields.update(failure)
|
||||
write_task_result(root, task_id, "failed" if case == "prior_infra" else "running", **fields)
|
||||
restoring = case == "snapshot_restore"
|
||||
cancel_intents.request_cancel(
|
||||
root, task_id, source="snapshot_restore" if restoring else "owner",
|
||||
reason="server_shutdown" if restoring else "Stop",
|
||||
allow_settled_target=case == "prior_infra",
|
||||
)
|
||||
|
||||
outcome = qenv.tl.cancel_task_custody(task_id, deliver=False)
|
||||
stored = load_task_result(root, task_id)
|
||||
if case in {"owner_stop", "snapshot_restore"}:
|
||||
assert outcome == qenv.tl.CANCEL_CANCELLED
|
||||
assert stored["status"] == "cancelled"
|
||||
assert transport.cancels == [(run_id, "owner_task_gone")]
|
||||
assert status_at_cancel == ["running"], "request precedes the cancelled write"
|
||||
# The transport still reports running: a request cannot claim physical death.
|
||||
assert stored["delegated_runs_unreconciled"] == [run_id]
|
||||
else:
|
||||
assert outcome == qenv.tl.CANCEL_ALREADY_SETTLED
|
||||
assert stored["status"] == "failed" and stored["result"] == "provider died"
|
||||
assert stored["outcome_axes"] == failure["outcome_axes"]
|
||||
assert transport.cancels == [] and status_at_cancel == []
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue