From 5203d010bf024139976771e5285d455495df78cb Mon Sep 17 00:00:00 2001 From: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com> Date: Mon, 21 Sep 2026 05:04:30 +0300 Subject: [PATCH] Carry the fence answer typed through the loop: ok, refused or unknown `_begin_task_acceptance_fence` / `_end_task_acceptance_fence` collapsed three different facts - the supervisor said yes, said no, or never answered - into a bool with a DEBUG line, so a stalled supervisor read downstream as a refusal and left no trace in the task's Logs. Both now return a frozen `FenceOutcome` (truthy iff ok; refused carries the reason, unknown the seconds waited), expose it on the context for the call sites that still treat a failed begin as before, and every refusal or gap leaves ONE durable worker-side row `supervisor_ack_unavailable {task_id, root_task_id, op, outcome, reason, waited_sec}`. A healthy path writes none. A refused inspection drops the stale binding and begins afresh (the supervisor re-adopts or reopens); an unanswered one keeps the token and its known generation and stacks no second request on a silent supervisor. A failed end drops the binding. `end` always carries the known generation. Only `sealed` is a seal. A `released` answer to a terminal end with no known generation mismatch may be the echo of a lost `released + generation_mismatch` ack (the re-send then finds no row): owner mail is durably written before the generation moves, so the local mailbox decides and the caller revises instead of sealing blind (#406). `loop_acceptance.py` pays its cap by reusing `_resolve_ctx_lineage` for two duplicated lineage reads. Two tests pinned the bool return and are rewritten to the typed truth. --- ouroboros/loop_acceptance.py | 250 +++++++++++---------- ouroboros/loop_acceptance_review.py | 4 +- tests/test_acceptance_fence_outcome.py | 261 ++++++++++++++++++++++ tests/test_acceptance_source_ack_truth.py | 3 +- tests/test_v664_acceptance_planning.py | 6 +- 5 files changed, 402 insertions(+), 122 deletions(-) create mode 100644 tests/test_acceptance_fence_outcome.py diff --git a/ouroboros/loop_acceptance.py b/ouroboros/loop_acceptance.py index f6ff1eadd..f6a11bc63 100644 --- a/ouroboros/loop_acceptance.py +++ b/ouroboros/loop_acceptance.py @@ -7,8 +7,10 @@ from __future__ import annotations import json import logging import pathlib +import time +from dataclasses import dataclass -from typing import Any, Callable, Dict, List, Optional +from typing import TYPE_CHECKING, Any, Callable, Dict, List, Optional from ouroboros.acceptance_settlement import forced_rail_panel_verdict from ouroboros.review_cycles import REASON_REVIEW_CYCLES_EXHAUSTED from ouroboros.review_projection import publish_acceptance_checkpoint @@ -16,9 +18,6 @@ from ouroboros.outcomes import ACCEPTANCE_ACCEPTED, ACCEPTANCE_BYPASS_REASONS, A from ouroboros.tools.registry import ToolRegistry from ouroboros.utils import truncate_review_artifact - -from typing import TYPE_CHECKING - if TYPE_CHECKING: # annotation-only names; lazy under future annotations, never imported at runtime from ouroboros.loop_delivery import DeliveryCandidate from ouroboros.loop_round_limits import _RoundLimitContext @@ -93,7 +92,71 @@ from ouroboros.loop_messages import ( # noqa: F401 — shared owner-source surf ) -def _begin_task_acceptance_fence(ctx: Any, task_id: str) -> tuple[bool, Any]: +@dataclass(frozen=True) +class FenceOutcome: + """Typed answer to one queue-fence request; truthy iff the supervisor said yes. + + ``refused``: it answered no (``reason``). ``unknown``: no answer arrived within + ``waited_sec`` — a gap, never a refusal, a verdict or an owner message. + """ + + status: str = "ok" + op: str = "" + reason: str = "" + waited_sec: float = 0.0 + + def __bool__(self) -> bool: + return self.status == "ok" + + +def _root_task_id(ctx: Any, task_id: str) -> str: + meta = getattr(ctx, "task_metadata", {}) + meta = meta if isinstance(meta, dict) else {} + return str(meta.get("root_task_id") or getattr(ctx, "root_task_id", "") or task_id) + + +def _settle_fence_outcome(ctx: Any, outcome: FenceOutcome) -> FenceOutcome: + """Expose the outcome on ctx; a refusal or a gap leaves ONE durable worker-side row.""" + ctx._task_acceptance_fence_outcome = outcome + if not outcome: + from ouroboros import task_pacing + from ouroboros.utils import append_jsonl, utc_now_iso + + task_id = str(getattr(ctx, "task_id", "") or "") + try: + append_jsonl(task_pacing.acceptance_timing_events_path(ctx), { + "ts": utc_now_iso(), "type": "supervisor_ack_unavailable", "task_id": task_id, + "root_task_id": _root_task_id(ctx, task_id), "op": outcome.op, "outcome": outcome.status, + "reason": outcome.reason, "waited_sec": outcome.waited_sec, + }) + except Exception: + log.warning("supervisor_ack_unavailable row could not be written for %s", task_id, exc_info=True) + return outcome + + +def _fence_request(ctx: Any, op: str, callback: Callable[..., Any], **kwargs: Any) -> tuple[FenceOutcome, Any]: + """One request to the queue-owned fence: ok | refused(reason) | unknown(waited_sec).""" + started = time.monotonic() + try: + response = callback(**kwargs) + if not (isinstance(response, dict) and not response.get("ok", True)): + return _settle_fence_outcome(ctx, FenceOutcome(op=op)), response + outcome = FenceOutcome("refused", op, str(response.get("error") or response.get("status") or "")) + except RuntimeError as exc: + outcome = FenceOutcome("refused", op, str(exc)) + except Exception as exc: + outcome = FenceOutcome("unknown", op, type(exc).__name__, round(time.monotonic() - started, 3)) + return _settle_fence_outcome(ctx, outcome), None + + +def _drop_fence_binding(ctx: Any) -> None: + """The queue owns the fence; a token whose state is unproven is never retained.""" + ctx._task_acceptance_fence_token = None + ctx._task_acceptance_fence_generation = None + ctx._task_acceptance_queue_descendants = [] + + +def _begin_task_acceptance_fence(ctx: Any, task_id: str) -> tuple[FenceOutcome, Any]: """Optional seam implemented by the supervisor under its queue lock.""" admission_lock = getattr(ctx, "owner_message_admission_lock", None) admission_agent = getattr(ctx, "owner_message_admission_agent", None) @@ -101,57 +164,38 @@ def _begin_task_acceptance_fence(ctx: Any, task_id: str) -> tuple[bool, Any]: with admission_lock: ctx._task_acceptance_owner_generation = int(getattr(admission_agent, "_owner_message_generation", 0) or 0) existing = getattr(ctx, "_task_acceptance_fence_token", None) + inspect = getattr(ctx, "inspect_acceptance_fence", None) + if existing is not None and not callable(inspect): + return FenceOutcome(op="begin"), existing if existing is not None: - inspect = getattr(ctx, "inspect_acceptance_fence", None) - if callable(inspect): - try: - refreshed = inspect(token=str(existing)) - ctx._task_acceptance_queue_descendants = ( - list(refreshed.get("queue_descendants") or []) - if isinstance(refreshed, dict) else [] - ) - if isinstance(refreshed, dict): - ctx._task_acceptance_fence_generation = int( - refreshed.get("owner_message_generation") or 0 - ) - except Exception: - log.debug("Queue-owned acceptance fence inspection failed", exc_info=True) - return False, existing - return True, existing + outcome, refreshed = _fence_request(ctx, "inspect", inspect, token=str(existing)) + if outcome: + ctx._task_acceptance_queue_descendants = [] + if outcome and isinstance(refreshed, dict): + ctx._task_acceptance_queue_descendants = list(refreshed.get("queue_descendants") or []) + ctx._task_acceptance_fence_generation = int(refreshed.get("owner_message_generation") or 0) + if outcome or outcome.status == "unknown": + return outcome, existing # no answer is a gap: the binding and its known generation stay + _drop_fence_binding(ctx) # refused: the row is gone — rebind through the idempotent begin callback = getattr(ctx, "begin_acceptance_fence", None) if not callable(callback): - return True, None # one-minor/direct-context compatibility - try: - meta = getattr(ctx, "task_metadata", {}) - meta = meta if isinstance(meta, dict) else {} - response = callback( - root_task_id=str( - meta.get("root_task_id") or getattr(ctx, "root_task_id", "") or task_id - ), - task_id=str(task_id), - ) - except Exception: - log.debug("Queue-owned acceptance fence begin failed", exc_info=True) - return False, None - if isinstance(response, dict): - token = response.get("token") - ctx._task_acceptance_queue_descendants = list(response.get("queue_descendants") or []) - ctx._task_acceptance_fence_generation = int( - response.get("owner_message_generation") or 0 - ) - else: - token = response - ctx._task_acceptance_queue_descendants = [] - ctx._task_acceptance_fence_generation = None - if token in (None, False, ""): - return False, None - ctx._task_acceptance_fence_token = token - return True, token + return FenceOutcome(op="begin"), None # one-minor/direct-context compatibility + outcome, response = _fence_request( + ctx, "begin", callback, root_task_id=_root_task_id(ctx, task_id), task_id=str(task_id)) + answer = response if isinstance(response, dict) else {"token": response} + if outcome and answer.get("token") in (None, False, ""): + outcome = _settle_fence_outcome(ctx, FenceOutcome("refused", "begin", "no_token")) + if not outcome: + return outcome, None + ctx._task_acceptance_queue_descendants = list(answer.get("queue_descendants") or []) + ctx._task_acceptance_fence_generation = ( + int(answer.get("owner_message_generation") or 0) if isinstance(response, dict) else None + ) + ctx._task_acceptance_fence_token = answer["token"] + return outcome, answer["token"] -def _end_task_acceptance_fence( - ctx: Any, *, outcome: str, admission_locked: bool = False, -) -> bool: +def _end_task_acceptance_fence(ctx: Any, *, outcome: str, admission_locked: bool = False) -> FenceOutcome: if getattr(ctx, "_acceptance_review_only", False) and outcome != "revision": outcome = "revision" # Early feedback never closes the root's future work. token = getattr(ctx, "_task_acceptance_fence_token", None) @@ -169,57 +213,44 @@ def _end_task_acceptance_fence( from ouroboros.loop_messages import owner_source_sha256 from ouroboros.loop_transport import _owner_signal_pending - acknowledged_source = getattr(ctx, "_acceptance_ack_source_sha256", "") - direct_generation_mismatch = bool( - (acknowledged_source and ( - acknowledged_source != owner_source_sha256(ctx) - or _owner_signal_pending( - getattr(ctx, "_acceptance_observation_incoming", None), getattr(ctx, "drive_root", None), - str(getattr(ctx, "task_id", "") or ""), getattr(ctx, "_loop_mailbox_seen_ids", None), - getattr(ctx, "task_attempt", None) or 1, - owner_authority_only=True, - ) - )) or ( - expected_owner_generation is not None - and admission_agent is not None - and int(getattr(admission_agent, "_owner_message_generation", 0) or 0) - != int(expected_owner_generation)) - ) - effective_outcome = "revision" if direct_generation_mismatch else str(outcome) - if token is None or not callable(callback): - ctx._task_acceptance_fence_generation_mismatch = direct_generation_mismatch - return True - expected_queue_generation = getattr(ctx, "_task_acceptance_fence_generation", None) - if expected_queue_generation is None: - response = callback(token=token, outcome=effective_outcome) - else: - response = callback( - token=token, - outcome=effective_outcome, - expected_generation=int(expected_queue_generation), + def owner_mail_pending() -> bool: + return _owner_signal_pending( + getattr(ctx, "_acceptance_observation_incoming", None), getattr(ctx, "drive_root", None), + str(getattr(ctx, "task_id", "") or ""), getattr(ctx, "_loop_mailbox_seen_ids", None), + getattr(ctx, "task_attempt", None) or 1, owner_authority_only=True, ) - except Exception: - log.debug("Queue-owned acceptance fence transition failed", exc_info=True) - return False + + acknowledged_source = getattr(ctx, "_acceptance_ack_source_sha256", "") + generation_mismatch = bool( + (acknowledged_source and (acknowledged_source != owner_source_sha256(ctx) or owner_mail_pending())) + or (expected_owner_generation is not None and admission_agent is not None + and int(getattr(admission_agent, "_owner_message_generation", 0) or 0) != int(expected_owner_generation)) + ) + effective_outcome = "revision" if generation_mismatch else str(outcome) + if token is None or not callable(callback): + ctx._task_acceptance_fence_generation_mismatch = generation_mismatch + return FenceOutcome(op="end") + expected_queue_generation = getattr(ctx, "_task_acceptance_fence_generation", None) + result, response = _fence_request( + ctx, "end", callback, token=token, outcome=effective_outcome, + **({} if expected_queue_generation is None else {"expected_generation": int(expected_queue_generation)}), + ) + status = str(response.get("status") or "") if isinstance(response, dict) else "" + generation_mismatch = generation_mismatch or bool(isinstance(response, dict) and response.get("generation_mismatch")) + if result and status == "released" and effective_outcome != "revision" and not generation_mismatch: + # Only ``sealed`` is a seal. An unexplained release may echo a lost ``released + + # generation_mismatch`` answer: owner mail is durably written before the generation + # moves, so the local mailbox decides here, never a blind seal. + generation_mismatch = owner_mail_pending() finally: if acquired: admission_lock.release() - if isinstance(response, dict) and not bool(response.get("ok", True)): - return False - status = str((response or {}).get("status") or "") if isinstance(response, dict) else "" - generation_mismatch = bool( - direct_generation_mismatch - or (isinstance(response, dict) and response.get("generation_mismatch")) - ) - ctx._task_acceptance_fence_generation_mismatch = generation_mismatch - ctx._task_acceptance_fence_token = None - ctx._task_acceptance_fence_generation = None - ctx._task_acceptance_queue_descendants = [] - if status == "sealed" or (not status and effective_outcome != "revision"): - ctx._task_acceptance_sealed_fence_token = token - else: - ctx._task_acceptance_sealed_fence_token = None - return True + _drop_fence_binding(ctx) # also after a refusal or a gap: the next begin re-adopts or reopens + if result: + ctx._task_acceptance_fence_generation_mismatch = generation_mismatch + sealed = status == "sealed" or (not status and effective_outcome != "revision") + ctx._task_acceptance_sealed_fence_token = token if sealed else None + return result def _supersede_delivery_acceptance_binding( @@ -439,17 +470,15 @@ def _task_acceptance_subtree_snapshot( from ouroboros.tools.join_ledger import _child_result_sha256 meta = getattr(ctx, "task_metadata", {}) - meta = meta if isinstance(meta, dict) else {} - root_id = str(meta.get("root_task_id") or getattr(ctx, "root_task_id", "") or task_id) status_root = pathlib.Path(str( - meta.get("budget_drive_root") + (meta.get("budget_drive_root") if isinstance(meta, dict) else "") or getattr(ctx, "budget_drive_root", "") or drive_root )) rows = find_child_tasks( status_root, parent_task_id=str(task_id), - root_task_id=root_id, + root_task_id=_root_task_id(ctx, task_id), exclude_task_id=str(task_id), scope="subtree", ) @@ -497,21 +526,9 @@ def _mark_root_acceptance_checkpoint( ctx: Any, llm_trace: Dict[str, Any], *, status: str, pass_index: int = 0, ) -> None: """Minimal in-result phase checkpoint; no parallel acceptance journal.""" - from ouroboros.task_results import resolve_task_lineage + from ouroboros.loop_acceptance_review import _resolve_ctx_lineage - meta = getattr(ctx, "task_metadata", {}) - meta = meta if isinstance(meta, dict) else {} - task_id = str(getattr(ctx, "task_id", "") or "") - lineage = resolve_task_lineage( - task_id, - metadata=meta, - root_task_id=getattr(ctx, "root_task_id", None), - parent_task_id=getattr(ctx, "parent_task_id", None), - delegation_role=getattr(ctx, "delegation_role", None), - original_task_id=getattr(ctx, "original_task_id", None), - timeout_retry_from=getattr(ctx, "timeout_retry_from", None), - ) - if not lineage["is_root_task"]: + if not _resolve_ctx_lineage(ctx)["is_root_task"]: return llm_trace["root_phase_checkpoint"] = { "phase": "task_acceptance", @@ -677,7 +694,6 @@ def merge_agent_acceptance_stance(trace: Dict[str, Any], decision: dict, ctx: An trace["acceptance_decision"] = merged - def _collect_acceptance_obligations(llm_trace: Dict[str, Any], result: Any) -> None: """Typed PER-TASK obligations from critical contributing findings (v6.54.4). diff --git a/ouroboros/loop_acceptance_review.py b/ouroboros/loop_acceptance_review.py index 70e4f6bd2..0e6b3a2de 100644 --- a/ouroboros/loop_acceptance_review.py +++ b/ouroboros/loop_acceptance_review.py @@ -672,7 +672,7 @@ def _finish_cyber_acceptance(ctx: _TaskAcceptanceContext, result: Any) -> bool: enforcement="advisory", source="author_final_response", ) ctx.llm_trace["review_decision"].update(author_finish=True, review_pending=pending, - admission_released=released) + admission_released=bool(released)) _loop()._set_acceptance_decision(ctx.llm_trace, { "status": ACCEPTANCE_ACCEPTED if clean else ACCEPTANCE_FINALIZED_UNACCEPTED, "reason": "clean_pass" if clean else "author_finish", "source": "task_acceptance_review", @@ -1383,7 +1383,7 @@ def _run_task_acceptance_review_once( ) emit_progress("Task acceptance review waiting for recursive subtree quiescence.") return True - llm_trace["review_decision"].update(admission_fence_available=fence_ok, subtree_quiescent=quiescent) + llm_trace["review_decision"].update(admission_fence_available=bool(fence_ok), subtree_quiescent=quiescent) # One effective profile carries explicit author caps/Hurry to gates and display. budget_profile = effective_budget_profile( tools._ctx, task_pacing.resolve_budget_profile(tools._ctx), diff --git a/tests/test_acceptance_fence_outcome.py b/tests/test_acceptance_fence_outcome.py new file mode 100644 index 000000000..338c0a789 --- /dev/null +++ b/tests/test_acceptance_fence_outcome.py @@ -0,0 +1,261 @@ +"""The acceptance fence's TYPED OUTCOME through the loop's middle layer. + +``_begin_task_acceptance_fence`` / ``_end_task_acceptance_fence`` answer +ok | refused(reason) | unknown(waited_sec) instead of a bool with a DEBUG line: +an answer that has not arrived is a gap, every refusal or gap leaves ONE durable +``supervisor_ack_unavailable`` row, and only ``sealed`` is a seal — an unexplained +``released`` on a terminal outcome is decided by the local mailbox (#406). +""" + +from __future__ import annotations + +import json +import queue as stdqueue +from types import SimpleNamespace + +import pytest + +from tests.test_acceptance_fence import _isolated_queue +from tests.test_acceptance_fence_transport import WAIT_SEC, _Supervisor, _pooled_agent + + +@pytest.fixture +def short_wait(monkeypatch): + from ouroboros import runtime_limits + + monkeypatch.setattr(runtime_limits, "get_acceptance_fence_ack_wait_sec", lambda: WAIT_SEC) + + +def _loop_ctx(tmp_path, agent, task_id="root-1", **extra): + ctx = SimpleNamespace( + task_metadata={"root_task_id": task_id}, task_id=task_id, drive_root=tmp_path, + _task_acceptance_fence_token=None, _task_acceptance_sealed_fence_token=None, + _task_acceptance_fence_generation=None, _task_acceptance_queue_descendants=[], + **extra, + ) + if agent is not None: + ctx.begin_acceptance_fence = agent._begin_acceptance_fence + ctx.inspect_acceptance_fence = agent._inspect_acceptance_fence + ctx.end_acceptance_fence = agent._end_acceptance_fence + return ctx + + +def _unavailable_rows(tmp_path): + path = tmp_path / "logs" / "events.jsonl" + if not path.is_file(): + return [] + rows = [json.loads(line) for line in path.read_text(encoding="utf-8").splitlines() if line.strip()] + return [row for row in rows if row.get("type") == "supervisor_ack_unavailable"] + + +# --- the middle layer: ok | refused(reason) | unknown(waited_sec) --------------------------- + + +def test_one_lost_ack_and_resend_puts_the_fence_up(monkeypatch, tmp_path, short_wait): + from ouroboros.loop import _begin_task_acceptance_fence + + queue_mod, _pending = _isolated_queue(monkeypatch, tmp_path) + events: stdqueue.Queue = stdqueue.Queue() + supervisor = _Supervisor(events, tmp_path, lose_acks=1) + ctx = _loop_ctx(tmp_path, _pooled_agent(tmp_path, events)) + try: + fence_ok, token = _begin_task_acceptance_fence(ctx, "root-1") + finally: + supervisor.stop() + assert fence_ok and token == ctx._task_acceptance_fence_token + assert queue_mod.ACCEPTANCE_FENCES["root-1"]["token"] == token + assert ctx._task_acceptance_fence_outcome.status == "ok" + assert _unavailable_rows(tmp_path) == [] # the fence is up: no model round was owed + + +def test_never_acked_begin_is_typed_unknown_with_one_durable_row(monkeypatch, tmp_path, short_wait): + from ouroboros.loop import _begin_task_acceptance_fence + + _isolated_queue(monkeypatch, tmp_path) + ctx = _loop_ctx(tmp_path, _pooled_agent(tmp_path, stdqueue.Queue())) + outcome, token = _begin_task_acceptance_fence(ctx, "root-1") + assert not outcome and token is None + assert (outcome.status, outcome.op) == ("unknown", "begin") + assert outcome.waited_sec >= WAIT_SEC * 2 - 0.1 + assert ctx._task_acceptance_fence_outcome is outcome # the next package reads it here + rows = _unavailable_rows(tmp_path) + assert len(rows) == 1 + assert {key: rows[0][key] for key in ("task_id", "root_task_id", "op", "outcome")} == { + "task_id": "root-1", "root_task_id": "root-1", "op": "begin", "outcome": "unknown"} + assert rows[0]["waited_sec"] == outcome.waited_sec and "reason" in rows[0] + + +def test_healthy_supervisor_leaves_no_unavailable_row(monkeypatch, tmp_path, short_wait): + from ouroboros.loop import _begin_task_acceptance_fence, _end_task_acceptance_fence + + _isolated_queue(monkeypatch, tmp_path) + events: stdqueue.Queue = stdqueue.Queue() + supervisor = _Supervisor(events, tmp_path, delay=0.05) + ctx = _loop_ctx(tmp_path, _pooled_agent(tmp_path, events)) + try: + assert _begin_task_acceptance_fence(ctx, "root-1")[0] + assert _begin_task_acceptance_fence(ctx, "root-1")[0] # refresh through inspect + sealed = _end_task_acceptance_fence(ctx, outcome="terminal") + finally: + supervisor.stop() + assert sealed and sealed.status == "ok" + assert ctx._task_acceptance_sealed_fence_token + assert _unavailable_rows(tmp_path) == [] + + +def test_refused_begin_is_typed_refused_with_its_reason(monkeypatch, tmp_path, short_wait): + from ouroboros.loop import _begin_task_acceptance_fence + + queue_mod, _pending = _isolated_queue(monkeypatch, tmp_path) + queue_mod.transition_acceptance_fence( + action="begin", token="a" * 32, root_task_id="root-1", task_id="root-1") + queue_mod.transition_acceptance_fence(action="end", token="a" * 32, outcome="terminal") + events: stdqueue.Queue = stdqueue.Queue() + supervisor = _Supervisor(events, tmp_path) + ctx = _loop_ctx(tmp_path, _pooled_agent(tmp_path, events)) + try: + outcome, token = _begin_task_acceptance_fence(ctx, "root-1") + finally: + supervisor.stop() + assert not outcome and token is None + assert outcome.status == "refused" and "already sealed" in outcome.reason + rows = _unavailable_rows(tmp_path) + assert [(row["op"], row["outcome"]) for row in rows] == [("begin", "refused")] + assert "already sealed" in rows[0]["reason"] + + +def test_begin_with_stale_token_rebinds_through_fresh_begin(tmp_path): + """A lost end/inspect ack leaves a stale local token while the supervisor already + released the fence: a REFUSED inspection drops the binding and begins afresh.""" + from ouroboros.loop import _begin_task_acceptance_fence + + ctx = _loop_ctx(tmp_path, None) + ctx._task_acceptance_fence_token, ctx._task_acceptance_fence_generation = "stale-token", 0 + ctx.inspect_acceptance_fence = lambda **_kwargs: (_ for _ in ()).throw( + RuntimeError("acceptance fence inspect failed")) + ctx.begin_acceptance_fence = lambda **_kwargs: { + "ok": True, "status": "active", "token": "fresh-token", + "owner_message_generation": 1, "queue_descendants": [], + } + ok, token = _begin_task_acceptance_fence(ctx, "root-1") + assert ok and token == "fresh-token" + assert ctx._task_acceptance_fence_token == "fresh-token" + assert ctx._task_acceptance_fence_generation == 1 + + +def test_unanswered_inspection_keeps_the_binding_and_asks_nothing_more(tmp_path): + """The other direction: no answer is a gap, not a release — the token and the + known generation stay, and no second request is stacked on a silent supervisor.""" + from ouroboros.loop import _begin_task_acceptance_fence + + begins: list = [] + ctx = _loop_ctx(tmp_path, None) + ctx._task_acceptance_fence_token, ctx._task_acceptance_fence_generation = "held-token", 3 + ctx.inspect_acceptance_fence = lambda **_kwargs: (_ for _ in ()).throw(TimeoutError("no ack")) + ctx.begin_acceptance_fence = lambda **kwargs: begins.append(kwargs) or {"token": "never"} + outcome, token = _begin_task_acceptance_fence(ctx, "root-1") + assert not outcome and outcome.status == "unknown" and token == "held-token" + assert begins == [] + assert (ctx._task_acceptance_fence_token, ctx._task_acceptance_fence_generation) == ("held-token", 3) + + +def test_end_failure_drops_binding_so_next_begin_is_fresh(tmp_path): + from ouroboros.loop import _begin_task_acceptance_fence, _end_task_acceptance_fence + + ctx = _loop_ctx(tmp_path, None) + ctx._task_acceptance_fence_token, ctx._task_acceptance_fence_generation = "token-1", 0 + + def failing_end(**_kwargs): + raise TimeoutError("supervisor did not acknowledge acceptance fence token-1") + + ctx.end_acceptance_fence = failing_end + ended = _end_task_acceptance_fence(ctx, outcome="revision") + assert not ended and (ended.status, ended.op) == ("unknown", "end") + assert ctx._task_acceptance_fence_token is None + assert ctx._task_acceptance_fence_generation is None + assert [(row["op"], row["outcome"]) for row in _unavailable_rows(tmp_path)] == [("end", "unknown")] + + ctx.begin_acceptance_fence = lambda **_kwargs: { + "ok": True, "status": "active", "token": "token-2", + "owner_message_generation": 0, "queue_descendants": [], + } + ok, token = _begin_task_acceptance_fence(ctx, "root-1") + assert ok and token == "token-2" + + +def test_end_carries_the_known_generation(tmp_path): + """``end`` is never sent without ``expected_generation`` once a generation was known — + an unanswered refresh must not erase it (the compare-and-seal would vanish with it).""" + from ouroboros.loop import _begin_task_acceptance_fence, _end_task_acceptance_fence + + sent: list = [] + ctx = _loop_ctx(tmp_path, None) + ctx.begin_acceptance_fence = lambda **_kwargs: {"token": "t", "owner_message_generation": 3} + ctx.inspect_acceptance_fence = lambda **_kwargs: (_ for _ in ()).throw(TimeoutError("no ack")) + ctx.end_acceptance_fence = lambda **kwargs: sent.append(kwargs) or {"ok": True, "status": "sealed"} + assert _begin_task_acceptance_fence(ctx, "root-1")[0] + assert not _begin_task_acceptance_fence(ctx, "root-1")[0] # the refresh went unanswered + assert _end_task_acceptance_fence(ctx, outcome="terminal") + assert sent == [{"token": "t", "outcome": "terminal", "expected_generation": 3}] + + +# --- ``released`` is not a seal -------------------------------------------------------------- + + +def _owner_mail(tmp_path, task_id="root-1"): + from ouroboros.owner_mailbox import write_owner_message + + assert write_owner_message(tmp_path, "Use the blue variant", task_id, msg_id="owner-blue") + + +@pytest.mark.parametrize("mail,expected_mismatch", [(True, True), (False, False)]) +def test_released_is_not_a_seal_and_owner_mail_forces_revision(tmp_path, mail, expected_mismatch): + from ouroboros.loop import _end_task_acceptance_fence + + ctx = _loop_ctx(tmp_path, None) + ctx._task_acceptance_fence_token, ctx._task_acceptance_fence_generation = "token-1", 0 + ctx.end_acceptance_fence = lambda **_kwargs: {"ok": True, "status": "released", "row_absent": True} + if mail: + _owner_mail(tmp_path) + ended = _end_task_acceptance_fence(ctx, outcome="terminal") + assert ended # the supervisor answered; an absent row is not a refusal + assert ctx._task_acceptance_sealed_fence_token is None # and it is not a seal either + assert ctx._task_acceptance_fence_generation_mismatch is expected_mismatch + + +def test_a_sealed_answer_is_a_seal_without_consulting_the_mailbox(tmp_path): + from ouroboros.loop import _end_task_acceptance_fence + + ctx = _loop_ctx(tmp_path, None) + ctx._task_acceptance_fence_token, ctx._task_acceptance_fence_generation = "token-1", 0 + ctx.end_acceptance_fence = lambda **_kwargs: {"ok": True, "status": "sealed"} + _owner_mail(tmp_path) # queue authority already compared the generation + assert _end_task_acceptance_fence(ctx, outcome="terminal") + assert ctx._task_acceptance_sealed_fence_token == "token-1" + assert ctx._task_acceptance_fence_generation_mismatch is False + + +def test_lost_generation_mismatch_ack_cannot_produce_a_blind_seal(monkeypatch, tmp_path, short_wait): + """#406: the first ``end(terminal)`` is applied as ``released + generation_mismatch`` and + its ack is lost; the re-send finds no row. Owner mail is durably written before the + generation moves, so the local mailbox — not the bare ``released`` — decides.""" + from ouroboros.loop import _begin_task_acceptance_fence, _end_task_acceptance_fence + + queue_mod, _pending = _isolated_queue(monkeypatch, tmp_path) + events: stdqueue.Queue = stdqueue.Queue() + agent = _pooled_agent(tmp_path, events) + ctx = _loop_ctx(tmp_path, agent) + supervisor = _Supervisor(events, tmp_path) + try: + assert _begin_task_acceptance_fence(ctx, "root-1")[0] + with queue_mod._queue_lock: # steering: durable mail first, then the generation + _owner_mail(tmp_path) + queue_mod.ACCEPTANCE_FENCES["root-1"]["owner_message_generation"] += 1 + supervisor.lose_acks = 1 + ended = _end_task_acceptance_fence(ctx, outcome="terminal") + finally: + supervisor.stop() + assert [evt.get("expected_generation") for evt in supervisor.seen if evt["action"] == "end"] == [0, 0] + assert ended and queue_mod.ACCEPTANCE_FENCES == {} + assert ctx._task_acceptance_sealed_fence_token is None + assert ctx._task_acceptance_fence_generation_mismatch is True # the caller revises, never seals diff --git a/tests/test_acceptance_source_ack_truth.py b/tests/test_acceptance_source_ack_truth.py index c8c44bf80..a78aa5879 100644 --- a/tests/test_acceptance_source_ack_truth.py +++ b/tests/test_acceptance_source_ack_truth.py @@ -148,7 +148,8 @@ def test_unknown_at_capture_and_ack_still_fails_closed_at_the_end_seal(case): tool_ctx.end_acceptance_fence = end tool_ctx.inspect_acceptance_fence = lambda **_kw: (_ for _ in ()).throw(INSPECT_FAILURE) tool_ctx._execution_trace = trace - assert _begin_task_acceptance_fence(tool_ctx, "root") == (True, "final") + opened, token = _begin_task_acceptance_fence(tool_ctx, "root") + assert opened and token == "final" assert tool_ctx._task_acceptance_fence_generation == 0 observed = capture_acceptance_observation(tool_ctx, trace, ctx.incoming_messages) diff --git a/tests/test_v664_acceptance_planning.py b/tests/test_v664_acceptance_planning.py index 281b23810..d2aa3af75 100644 --- a/tests/test_v664_acceptance_planning.py +++ b/tests/test_v664_acceptance_planning.py @@ -570,8 +570,10 @@ def test_queue_owned_acceptance_fence_uses_only_optional_ctx_hooks(): begin_acceptance_fence=begin, end_acceptance_fence=end, ) - assert _begin_task_acceptance_fence(ctx, "root") == (True, "fence-1") - assert _end_task_acceptance_fence(ctx, outcome="revision") is True + opened, token = _begin_task_acceptance_fence(ctx, "root") + assert opened and opened.status == "ok" and token == "fence-1" + released = _end_task_acceptance_fence(ctx, outcome="revision") + assert released and released.status == "ok" assert calls == [ ("begin", {"root_task_id": "root", "task_id": "root"}), ("end", {"token": "fence-1", "outcome": "revision"}),