mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
Never re-send a fence read
inspect is a read asked many times per turn and its loss is harmless (the last known generation stays, the end compare-and-seal decides). Re-sending it doubled the host blocking of every read under a stalled supervisor; only the transitions (begin, end) keep their single re-send.
This commit is contained in:
parent
0b9413ce4a
commit
101e5cd940
2 changed files with 16 additions and 6 deletions
|
|
@ -268,8 +268,9 @@ class OuroborosAgent:
|
|||
transition (``fence_transition``), it applies the fence in-process — no event,
|
||||
no ack file, no wait (lock order: admission lock, then ``_queue_lock``). A pooled
|
||||
worker cannot share ``_queue_lock``: it sends the event, polls its own one-shot
|
||||
ack and re-sends the SAME request once; no answer raises ``TimeoutError`` (a gap),
|
||||
a refusal raises ``RuntimeError``. The ack is transport, not a lifecycle authority.
|
||||
ack and re-sends the SAME transition once — a read (``inspect``, asked many times a
|
||||
turn) is never re-sent, its loss is harmless; no answer raises ``TimeoutError`` (a
|
||||
gap), a refusal raises ``RuntimeError``. The ack is transport, not an authority.
|
||||
"""
|
||||
request.setdefault("task_id", str(self._current_task_id or ""))
|
||||
transition = getattr(self, "fence_transition", None)
|
||||
|
|
@ -279,7 +280,7 @@ class OuroborosAgent:
|
|||
raise RuntimeError("acceptance fence requires a supervisor event queue")
|
||||
else:
|
||||
event = {"type": "acceptance_fence", "req": uuid.uuid4().hex, **request}
|
||||
ack = self._send_fence_event(event) or self._send_fence_event(event)
|
||||
ack = self._send_fence_event(event) or (request["action"] != "inspect" and self._send_fence_event(event)) or {}
|
||||
if not ack:
|
||||
raise TimeoutError(f"supervisor did not acknowledge acceptance fence {request['action']}")
|
||||
if not ack.get("ok", True) or str(ack.get("status") or "") not in accept:
|
||||
|
|
|
|||
|
|
@ -266,7 +266,7 @@ def test_late_ack_of_a_previous_request_is_read_by_nobody(monkeypatch, tmp_path,
|
|||
with pytest.raises(TimeoutError):
|
||||
agent._inspect_acceptance_fence(token=token)
|
||||
stalled.resume() # the loop catches up and answers the inspect LATE
|
||||
assert stalled.drained(2)
|
||||
assert stalled.drained(1) # a read is sent once, never re-sent
|
||||
late = _ack_names(tmp_path)
|
||||
assert late and all(name.startswith(f"{token}.") for name in late)
|
||||
|
||||
|
|
@ -297,15 +297,24 @@ def test_pooled_request_is_resent_once_with_the_same_identity(monkeypatch, tmp_p
|
|||
assert (first["token"], first["req"]) == (second["token"], second["req"]) == (ack["token"], first["req"])
|
||||
assert queue_mod.ACCEPTANCE_FENCES["root-1"]["token"] == ack["token"]
|
||||
|
||||
# Never answered: exactly ONE re-send, then a TimeoutError — no loop, no long wait.
|
||||
# A transition never answered: exactly ONE re-send, then a TimeoutError — no loop, no long wait.
|
||||
silent: stdqueue.Queue = stdqueue.Queue()
|
||||
started = time.monotonic()
|
||||
with pytest.raises(TimeoutError):
|
||||
_pooled_agent(tmp_path, silent)._inspect_acceptance_fence(token=ack["token"])
|
||||
_pooled_agent(tmp_path, silent)._end_acceptance_fence(token=ack["token"], outcome="revision")
|
||||
assert time.monotonic() - started < WAIT_SEC * 2 + 2.0
|
||||
sent = [silent.get_nowait() for _ in range(silent.qsize())]
|
||||
assert len(sent) == 2 and sent[0]["req"] == sent[1]["req"]
|
||||
|
||||
# A READ (inspect is asked many times a turn; its loss is harmless) is never re-sent:
|
||||
# one wait, one event — a stalled supervisor cannot multiply into minutes of host blocking.
|
||||
quiet: stdqueue.Queue = stdqueue.Queue()
|
||||
started = time.monotonic()
|
||||
with pytest.raises(TimeoutError):
|
||||
_pooled_agent(tmp_path, quiet)._inspect_acceptance_fence(token=ack["token"])
|
||||
assert time.monotonic() - started < WAIT_SEC + 2.0
|
||||
assert quiet.qsize() == 1
|
||||
|
||||
|
||||
def test_every_request_carries_a_fresh_req_and_one_token_per_logical_begin(monkeypatch, tmp_path, short_wait):
|
||||
_isolated_queue(monkeypatch, tmp_path)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue