ouroboros/tests/test_transport_repeat_owner_stop.py
Ouroboros e33fd69115 Make the external-executor verbs answer with a native typed result
`delegate_shared._fail` rendered `{"status":"refused", ...}` as a plain string,
and the registry's legacy text adapter classified it as OK: it only understands
a top-level `ok:false` or a first-line `⚠️ IDENTIFIER` marker. So a refused
`delegate_wait`/`delegate_cancel` — daemon unreachable, run not owned, a
containment fault, a cancel the daemon refused — was recorded as a SUCCESSFUL
tool call on the outcome axis, in the acceptance packet and in the supervising
task's own reasoning.

Owner decision Q8A: the fix is a native structured result INSIDE the family,
not a repo-wide ABI migration. `_fail` now returns a `ToolResult` whose text is
the same JSON the callers emitted, plus two ADDITIVE envelope keys — `ok:false`
and `host_code` — written beside the domain payload. The domain `reason` is
never renamed into `ToolResult.code`, and the domain extras
(`definitely_unrun`, `pending_invocation_id`, `run_id`, `reset_at`, the custody
facts) keep their places. One exact, closed table maps a reason to its class:
substrate refusals (the daemon, the engine, custody or the run said no) are
`TOOL_REPORTED_FAILURE`, recorded and never degrading; malformed or
self-contradictory calls are `TOOL_ARG_ERROR`, which degrades and feeds
reflection. Neither is a timeout or the generic tool error, and an unclassified
reason defaults to the substrate class — the safe direction.

The second literal refusal author, `supervised_wait`'s checkpoint argument
check, folds into `_fail`. `_delegate_cancel`'s own outcomes join it: `failed`
and `containment_fault_run_may_still_be_live` report a run that may still be
live and mutating, so they publish as failures, while `confirmed` and
`requested` stay successful observations of the control surface.

Publication happens once, at the four REGISTERED entries, after every
decoration and immediately before the string is returned: `subagent_runtime`
mutates the start payload after `_delegate_start`, and an earlier publish fails
the registry's equality gate silently. `exact_start` now reads the native
payload, adds the actor identity and the work-order source, and returns
`_replace_tool_result`; its `json.loads → TypeError → return result` bypass,
which dropped that decoration without saying so, is deleted.
`_mark_actor_physical_start` still reads the DOMAIN `started` /
`started_uncustodied` status, never the host class.

The consumers migrate in the same change as the type they consume: the
configured-session bootstrap, the recovery handoff, the unknown-provider hold
and the pending-wake replay all read the producer's own payload rather than a
stringified result — without this, every leaf wake would have failed its
acknowledgement and taken the no-resend terminal. The wake envelope carries
`ok`/`host_code` through both fitted-spill shapes, so a refusal too large to
inline cannot read as a successful wait. `wait_once` deliberately keeps its
`str` tick contract, and `integrate_delegated_patch` keeps its own string ABI
at the one helper the two families share.

Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
2026-09-15 18:01:50 +03:00

213 lines
12 KiB
Python

"""A cooperative stop interrupts only the unsent paid transport-repeat grant."""
import threading
import time
from concurrent.futures import ThreadPoolExecutor
from types import SimpleNamespace
import httpx
import pytest
from ouroboros import cancel_intents, loop, loop_llm_call, loop_transport, owner_mailbox
from ouroboros import usage_accounting as accounting
from ouroboros.delegate_shared import delegate_result
from ouroboros.outcomes import REASON_OWNER_REQUESTED_FINALIZATION
from ouroboros.task_results import write_task_result
from supervisor.owner_stop import REASON_OWNER_STOPPED_DIRECT_TURN, owner_stop_control_id
from tests.test_transport_death_retry import (
EMPTY_RESPONSE, _LedgerLLM, _ScriptedLLM, _events, _ledger, _loop_kwargs,
_no_chain, _primary_call, _status_failure,
)
@pytest.mark.parametrize("after_attempt,owner_grace", [(1, True), (1, False), (2, False)])
def test_stop_arriving_in_real_backoff_keeps_only_actual_physical_attempts(
tmp_path, monkeypatch, after_attempt, owner_grace,
):
canonical = tmp_path / "canonical"
execution = tmp_path / "execution" if owner_grace else canonical
canonical.mkdir()
execution.mkdir(exist_ok=True)
write_task_result(canonical, "t-death", "running",
**({"child_drive_root": str(execution)} if owner_grace else {}))
monkeypatch.setenv("OUROBOROS_TASK_REVIEW_MODE", "off")
monkeypatch.delenv("USE_LOCAL_FALLBACK", raising=False)
monkeypatch.setattr(loop, "_run_cross_model_fallback_chain", _no_chain)
llm = _LedgerLLM(canonical, *([lambda: httpx.ReadError("controlled socket death")] * 3))
kwargs = _loop_kwargs(execution, llm, [])
kwargs["drive_logs"] = execution / "logs"
ctx = kwargs["tools"]._ctx
ctx.task_id, ctx.task_attempt, ctx.budget_drive_root = "t-death", 1, canonical
ctx.is_direct_chat = not owner_grace
ctx.task_metadata = {"root_task_id": "t-death"}
first_empty_check = threading.Event()
real_wait = loop_transport.interruptible_wait_sleep
def observed_wait(seconds, check):
def checked():
result = check()
if not result and llm.calls == after_attempt:
first_empty_check.set()
return result
return real_wait(seconds, checked)
monkeypatch.setattr(loop_transport, "interruptible_wait_sleep", observed_wait)
def write_stop():
assert first_empty_check.wait(10)
if owner_grace:
intent = cancel_intents.request_cancel(canonical, "t-death", source="isolated-test",
requested_stop_policy=cancel_intents.STOP_POLICY_FINALIZE)
text, mid = REASON_OWNER_REQUESTED_FINALIZATION, owner_stop_control_id(intent)
else:
text, mid = REASON_OWNER_STOPPED_DIRECT_TURN, "direct-stop"
assert owner_mailbox.write_owner_message(execution, text, "t-death", msg_id=mid,
kind=owner_mailbox.KIND_FINALIZE_NOW)
return mid
with accounting.usage_scope(accounting.UsageScope(drive_root=canonical, task_id="t-death",
root_task_id="t-death", global_limit_usd=100.0)), ThreadPoolExecutor(max_workers=1) as pool:
writing = pool.submit(write_stop)
_text, usage, trace = loop.run_llm_loop(**kwargs)
mid = writing.result(timeout=5)
assert llm.calls == after_attempt
rows = _ledger(canonical)
assert len({row["attempt_id"] for row in rows}) == after_attempt
assert [row["state"] for row in rows] == ["reserved", "dispatched", "unresolved"] * after_attempt
assert accounting.usage_projection(canonical)["unresolved_upper_bound_usd"] == float(after_attempt)
assert usage.get(loop_llm_call.TRANSPORT_DEATHS_KEY, {}).get("count", 0) == after_attempt - 1
assert usage["_last_llm_retry_same_request"] is False
assert usage["_last_llm_error_kind"] == "provider_outcome_unknown"
assert trace["forced_finalization"]["source"] == "provider_outcome_unknown_no_resend"
assert ("owner requested Wrap up" if owner_grace else "owner requested Stop") in usage["terminal_provider_notice"]
assert "no terminal provider outcome" in usage["terminal_provider_notice"]
assert [row["reason_code"] for row in _events(execution / "logs", "llm_not_dispatched")] == ([] if owner_grace else ["finalize_control_pending"])
assert not _events(execution / "logs", "llm_retry_deadline_exhausted")
if owner_grace:
# Managed unknown waits use the ordinary round-top owner drain; no paid
# repeat was granted, and the current finalize intent is consumed there.
assert cancel_intents.active_intent(canonical, "t-death").get("control_drained_at")
else:
assert mid not in ctx._loop_mailbox_seen_ids
assert mid not in owner_mailbox.acknowledged_task_message_ids(execution, "t-death", attempt_key=1)
assert getattr(ctx, "_skip_post_task_synthesis", False) # same Stop-now contract after loop exit
from ouroboros import agent_task_pipeline
post_task_calls = []
monkeypatch.setattr(agent_task_pipeline, "_run_post_task_processing_async",
lambda *_a, **_k: post_task_calls.append(True))
task = {"id": "t-death", "type": "task", "chat_id": 7, "_is_direct_chat": True}
agent_task_pipeline.emit_task_results(SimpleNamespace(drive_root=canonical, repo_dir=canonical),
None, None, [], task, _text, usage, trace, start_time=0.0, drive_logs=canonical / "logs", ctx=ctx)
assert task["_skip_post_task_synthesis"] is True
assert post_task_calls == []
@pytest.mark.parametrize("kind,revoked,seen,expected", [
("owner_text", False, False, False), ("hurry", False, False, False),
("finalize_now", True, False, False), ("finalize_now", False, True, False),
("finalize_now", False, False, True),
])
def test_only_current_unseen_finalize_controls_request_a_retry_stop(tmp_path, kind, revoked, seen, expected):
ctx = SimpleNamespace(drive_root=tmp_path, task_id="t-death", task_attempt=1,
_loop_mailbox_seen_ids={"control"} if seen else set())
owner_mailbox.write_owner_message(tmp_path, "deadline", "t-death", msg_id="control", kind=kind)
if revoked:
owner_mailbox.revoke_owner_control(tmp_path, "t-death", "control")
assert loop_transport.transport_repeat_stop_requested(ctx) is expected
assert ctx._loop_mailbox_seen_ids == ({"control"} if seen else set())
def test_stale_owner_grace_control_cannot_cancel_a_repeat(tmp_path):
ctx = SimpleNamespace(drive_root=tmp_path, task_id="t-death", task_attempt=1, _loop_mailbox_seen_ids=set())
owner_mailbox.write_owner_message(tmp_path, REASON_OWNER_REQUESTED_FINALIZATION,
"t-death", msg_id="ownerstop:old-request", kind=owner_mailbox.KIND_FINALIZE_NOW)
assert not loop_transport.transport_repeat_stop_requested(ctx)
def test_deadline_refusal_is_not_misattributed_to_an_unchecked_control(tmp_path):
def forbidden():
raise AssertionError("deadline already refused this wait")
usage = {"_last_llm_error_kind": "provider_outcome_unknown", "_last_llm_retry_same_request": True,
loop_llm_call.TRANSPORT_DEATHS_KEY: {"round_id": "r", "count": 1, "backoff_sec": 4.0}}
ctx = loop_llm_call._LlmErrorContext(task_id="t-death", task_type="task", execution_id="e",
round_id="r", llm_call_id="old-call", round_idx=1, attempt=0, model="fixture",
request_ref=None, drive_logs=tmp_path, event_queue=None, accumulated_usage=usage,
deadline_ts=time.time() + 1, stop_retry_check=forbidden)
assert loop_llm_call._stop_after_llm_error(ctx)
assert not _events(tmp_path, "llm_not_dispatched")
assert len(_events(tmp_path, "llm_retry_deadline_exhausted")) == 1
assert usage["_last_llm_error_kind"] == "provider_outcome_unknown"
assert loop_llm_call.TRANSPORT_DEATHS_KEY not in usage
@pytest.mark.parametrize("empty", [False, True])
def test_control_callback_does_not_change_other_transient_or_empty_retry_contracts(tmp_path, monkeypatch, empty):
def forbidden():
raise AssertionError("only a paid transport-death repeat opts into this stop check")
def sleep(_seconds, _deadline, **options):
assert options == {}
return True
monkeypatch.setattr(loop_llm_call, "_sleep_within_deadline", sleep)
llm = _ScriptedLLM(EMPTY_RESPONSE if empty else lambda: _status_failure(503))
message, _cost = _primary_call(llm, tmp_path, {}, stop_retry_check=forbidden)
assert message["content"] == "done" and llm.calls == 2
def test_paid_repeat_empty_peek_reuses_existing_wait_proof(tmp_path, monkeypatch):
ctx = SimpleNamespace(drive_root=tmp_path, task_id="t-death", task_attempt=1, _loop_mailbox_seen_ids={"old"})
owner_mailbox.write_owner_message(tmp_path, "seen text", "t-death", msg_id="old")
calls = []
original = owner_mailbox.drain_owner_entries
def read(*args, **kwargs):
calls.append(True)
return original(*args, **kwargs)
monkeypatch.setattr(owner_mailbox, "drain_owner_entries", read)
peek = owner_mailbox.OwnerMailboxPeek()
for _ in range(4):
assert not loop_transport.transport_repeat_stop_requested(ctx, mailbox_peek=peek)
assert len(calls) == 1
owner_mailbox.write_owner_message(tmp_path, REASON_OWNER_STOPPED_DIRECT_TURN, "t-death", msg_id="stop", kind=owner_mailbox.KIND_FINALIZE_NOW)
assert loop_transport.transport_repeat_stop_requested(ctx, mailbox_peek=peek)
assert ctx._transport_repeat_control_reason == REASON_OWNER_STOPPED_DIRECT_TURN
assert ctx._loop_mailbox_seen_ids == {"old"}
@pytest.mark.parametrize("with_leaf, during_hold", [(False, False), (True, False), (True, True)])
def test_wrapup_reason_survives_the_live_delegate_hold(tmp_path, monkeypatch, with_leaf, during_hold):
from ouroboros import claudexor_daemon, delegate_custody, delegate_progress, delegate_hold
from tests.test_delegate_hold import _configured_registry, _start_leaf, _loop_kwargs as hold_kwargs
monkeypatch.setenv("OUROBOROS_TASK_REVIEW_MODE", "off")
monkeypatch.delenv("USE_LOCAL_FALLBACK", raising=False)
registry = _configured_registry(tmp_path, task_id="t-death")
registry._ctx.budget_drive_root = tmp_path
registry._ctx.task_attempt = 1
if with_leaf:
_start_leaf(tmp_path, task_id="t-death", run_id="fixture-leaf")
def request_wrapup():
intent = cancel_intents.request_cancel(tmp_path, "t-death", requested_stop_policy=cancel_intents.STOP_POLICY_FINALIZE)
owner_mailbox.write_owner_message(tmp_path, REASON_OWNER_REQUESTED_FINALIZATION, "t-death",
msg_id=owner_stop_control_id(intent), kind=owner_mailbox.KIND_FINALIZE_NOW)
def death():
if not during_hold:
request_wrapup()
return httpx.ReadError("controlled post-dispatch failure")
def hold(*args):
request_wrapup()
return delegate_result({"status": "progress", "wake_events": [{"kind": "finalize_now"}]})
if during_hold:
monkeypatch.setattr(delegate_hold, "supervised_wait", hold)
llm = _LedgerLLM(tmp_path, death)
kwargs = hold_kwargs(tmp_path, registry, [])
kwargs["llm"] = llm
monkeypatch.setattr(claudexor_daemon, "ensure_owned_gateway", lambda **kw: SimpleNamespace(close=lambda: None))
monkeypatch.setattr(delegate_progress, "bounded_poll", lambda *a, **kw: {"summary": {"state": "running"}})
monkeypatch.setattr(delegate_custody, "release_task_runs", lambda *a, **kw: None)
with accounting.usage_scope(accounting.UsageScope(
drive_root=tmp_path, task_id="t-death", root_task_id="t-death", global_limit_usd=100.0)):
_text, usage, _trace = loop.run_llm_loop(**kwargs)
assert llm.calls == 1
assert accounting.usage_projection(tmp_path)["unresolved_upper_bound_usd"] == 1.0
assert "owner requested Wrap up" in usage["terminal_provider_notice"]