mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 12:18:39 +00:00
P2-fix1d: Stamp the daemon-outage owner line per episode, not per task
S1(c) promised ONE owner line per unreachable-daemon EPISODE plus one recovery line, with the client's toast_once providing the dedup. _owner_line built a constant key and a constant text per task, so episodes 2..N were dropped by both surfaces: chat_media keeps a module-level set of shown toast_once keys, and the chat timeline dedupes a live card on phase|headline|body. In the very wave this step exists to close, 15:05-16:25, that is loss, not the noise the risk note predicted. The open episode now carries its own identity, following the loop_transport.incident precedent that stamps an episode's millisecond entry for exactly this pair: the boolean latch becomes the episode's UTC start plus its millisecond stamp, the stamp discriminates both toast keys, and each line names the episode it belongs to, so the text-keyed timeline dedup sees a second episode too. The outage and its recovery share one stamp, so a reader can tell which outage a recovery closes. No durable state is written and the beat is untouched. Residual: two episodes inside the same second render the same text and can still collapse on the timeline; their toast keys never collide, so the fact reaches the owner either way. Tests: two outage and recovery episodes in one wait now produce four distinct toast keys (two before), each pair sharing its own stamp, and the beat stays one tick per unreachable observation.
This commit is contained in:
parent
4805ec6067
commit
5ed3ba6762
2 changed files with 39 additions and 18 deletions
|
|
@ -879,7 +879,11 @@ def supervised_wait(
|
|||
"checkpoint_scheduled": bool(checkpoint_after_sec is not None),
|
||||
})
|
||||
|
||||
unobserved = False # an unreachable-daemon episode is open (one owner line each way)
|
||||
# The OPEN unreachable-daemon episode: its UTC start (empty when none) plus the
|
||||
# millisecond stamp that keeps two episodes of one wait distinct on both client
|
||||
# dedup surfaces (the loop_transport.incident precedent).
|
||||
outage_since = ""
|
||||
outage_stamp = ""
|
||||
gateway = None # the loop's own transport, only when it built the observing wait
|
||||
try:
|
||||
while True:
|
||||
|
|
@ -899,11 +903,14 @@ def supervised_wait(
|
|||
payload.get("status") == "observation_pending"
|
||||
and payload.get("reason") == _DAEMON_UNREACHABLE
|
||||
)
|
||||
if observed and unobserved and not unreachable:
|
||||
if observed and outage_since and not unreachable:
|
||||
# The first read the daemon answered again closes the episode.
|
||||
unobserved = False
|
||||
_owner_line(ctx, "Delegation daemon reachable again; delegated runs are "
|
||||
"being observed again.", "delegation_daemon_recovered", "ok")
|
||||
_owner_line(
|
||||
ctx,
|
||||
"Delegation daemon reachable again; the outage that began at "
|
||||
f"{outage_since} is over and delegated runs are being observed again.",
|
||||
f"delegation_daemon_recovered:{outage_stamp}", "ok")
|
||||
outage_since, outage_stamp = "", ""
|
||||
cursor = payload.get("last_seq")
|
||||
if isinstance(cursor, int):
|
||||
state["journal_cursor"] = max(int(state.get("journal_cursor") or 0), cursor)
|
||||
|
|
@ -974,16 +981,21 @@ def supervised_wait(
|
|||
"run_id": str(run_id), "reason": payload.get("reason"),
|
||||
"waited_sec": payload.get("waited_sec"),
|
||||
})
|
||||
if unreachable and not unobserved:
|
||||
if unreachable and not outage_since:
|
||||
# The class the model can do nothing about: a dead socket is a quiet
|
||||
# renewal on the same 3 s beat (no backoff, no durable counter), and
|
||||
# the owner hears about it exactly once per episode. Deadline, ceiling,
|
||||
# budget and cancel stay the outer bounds that cut a long unobserved
|
||||
# stretch.
|
||||
unobserved = True
|
||||
_owner_line(ctx, "Delegation daemon unreachable; delegated runs are not "
|
||||
"being observed, the runs themselves keep going.",
|
||||
"delegation_daemon_unreachable", "warn")
|
||||
# stretch. The episode stamps its own key and names its start in the
|
||||
# text, so a SECOND outage in the same wait is a new line on both
|
||||
# client surfaces instead of a duplicate the toast set already holds.
|
||||
opened = time.time()
|
||||
outage_since = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime(opened))
|
||||
outage_stamp = str(int(opened * 1000))
|
||||
_owner_line(ctx, f"Delegation daemon unreachable since {outage_since}; "
|
||||
"delegated runs are not being observed, the runs themselves "
|
||||
"keep going.",
|
||||
f"delegation_daemon_unreachable:{outage_stamp}", "warn")
|
||||
# A failed read is not a completed quiet window; retain the cursor
|
||||
# and avoid a busy loop if a transport fails before its read bound.
|
||||
time.sleep(_TICK_SEC)
|
||||
|
|
|
|||
|
|
@ -224,13 +224,22 @@ def test_unreachable_daemon_episode_is_one_owner_line_each_way(tmp_path, monkeyp
|
|||
assert [text.startswith("Delegation daemon unreachable") for text, _ in notes] == [
|
||||
True, False, True, False]
|
||||
outage, recovered = notes[0][1], notes[1][1]
|
||||
assert outage == {"task_incident": "delegation_daemon_unreachable",
|
||||
"toast_once": f"{ctx.task_id}:delegation_daemon_unreachable",
|
||||
"toast_tone": "warn"}
|
||||
assert recovered == {"task_incident": "delegation_daemon_unreachable",
|
||||
"toast_once": f"{ctx.task_id}:delegation_daemon_recovered",
|
||||
"toast_tone": "ok"}
|
||||
assert notes[2][1] == outage and notes[3][1] == recovered
|
||||
assert outage["task_incident"] == recovered["task_incident"] == "delegation_daemon_unreachable"
|
||||
assert outage["toast_tone"] == "warn" and recovered["toast_tone"] == "ok"
|
||||
assert outage["toast_once"].startswith(f"{ctx.task_id}:delegation_daemon_unreachable:")
|
||||
assert recovered["toast_once"].startswith(f"{ctx.task_id}:delegation_daemon_recovered:")
|
||||
# The SECOND episode is a second line on both client surfaces: the toast set
|
||||
# dedupes on the key and the timeline on the rendered text, so an episode
|
||||
# discriminator has to reach both or outages 2..N are dropped, not repeated.
|
||||
keys = [incident["toast_once"] for _text, incident in notes]
|
||||
assert len(set(keys)) == 4, keys
|
||||
# The text names the episode it belongs to, so the timeline's text-keyed
|
||||
# dedup sees a second episode too (two episodes inside one second would
|
||||
# still collapse there; the toast key above never does).
|
||||
assert all(text.count(":") >= 2 for text, _ in notes), notes
|
||||
# One episode's own pair shares its stamp: the recovery names the outage it closes.
|
||||
assert outage["toast_once"].rsplit(":", 1)[1] == recovered["toast_once"].rsplit(":", 1)[1]
|
||||
assert notes[2][1]["toast_once"].rsplit(":", 1)[1] == notes[3][1]["toast_once"].rsplit(":", 1)[1]
|
||||
assert sleeps == [delegate_supervision._TICK_SEC] * 4
|
||||
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue