mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
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.
This commit is contained in:
parent
10a1d0ca43
commit
5203d010bf
5 changed files with 402 additions and 122 deletions
|
|
@ -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).
|
||||
|
||||
|
|
|
|||
|
|
@ -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),
|
||||
|
|
|
|||
261
tests/test_acceptance_fence_outcome.py
Normal file
261
tests/test_acceptance_fence_outcome.py
Normal file
|
|
@ -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
|
||||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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"}),
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue