diff --git a/ouroboros/delegate_supervision.py b/ouroboros/delegate_supervision.py index da894b3c5..cc00bd792 100644 --- a/ouroboros/delegate_supervision.py +++ b/ouroboros/delegate_supervision.py @@ -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) diff --git a/tests/test_delegate_observation_transport.py b/tests/test_delegate_observation_transport.py index 4220df2b6..95fdd9d47 100644 --- a/tests/test_delegate_observation_transport.py +++ b/tests/test_delegate_observation_transport.py @@ -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