Skip to content

fix(gateway): delete the frozen streamed bubble once the ledger redelivers the refused final - #106370

Open
AlexxRussell wants to merge 2 commits into
NousResearch:mainfrom
AlexxRussell:fix/redelivered-final-preview-cleanup
Open

AlexxRussell wants to merge 2 commits into
NousResearch:mainfrom
AlexxRussell:fix/redelivered-final-preview-cleanup

Conversation

@AlexxRussell

Copy link
Copy Markdown

Stacked on #105200 (its commit is the first one here; this PR is the second commit only). It closes the two cases that one does not reach.

What the reader sees

Telegram, a long reply streamed by editing one preview bubble. Twice in two days (8 Sep 12:05 UTC, 9 Sep 07:26 UTC) the chat ended up with a frozen partial bubble above the complete reply:

  1. Stale finalize. The consumer's finalize edit landed a stale snapshot, so _run_agent_mark_streamed_delivery tried the reconciliation edit; a flood window refused it (Stale-finalize reconciliation edit failed ... (flood_control:269.0)), the normal final send delivered the complete reply below, and that branch registers no cleanup at all.
  2. Ledger redelivery. In both incidents the gateway's own final was refused as well (retry after 102s and 269s, both over the 60 s inline cap), so _finalize_delivery_obligation marked the row failed and the flood timer redelivered it minutes later with the flood marker (Redelivered recovered final response ... attempt 1). The cleanup fix(gateway): delete the streaming preview the consumer abandoned once the final lands #105200 registers had already fired from the post-delivery hook with _hermes_final_delivered=False and stood down, correctly: at that moment the frozen bubble was the reader's only copy. Nothing fires it after the redelivery, because the redelivery runs outside the turn.

So the reader keeps a truncated copy with raw MarkdownV2 markers and the cursor, and directly under it a complete copy prefixed with "part of it may already have arrived above".

The change

  • run_turn.py: the stale-finalize branch schedules the same cleanup when its reconciliation edit is refused (only in the arm that actually attempted an edit; the split-chain and no-editable-message arms are untouched). The callback no longer simply stands down on a refused send: if the turn's own final was ledgered, it hands the delete to that row's redelivery.
  • platforms/base.py: the row id comes from the turn's own final send and nowhere else. _send_final_text (normal lane) now returns it, through a private _send_final_ledgered that the public send_final_ledgered wraps unchanged, and _process_message_background passes it to _fire_post_delivery_callback, which stamps it on the session event (_hermes_final_obligation_id) right before the callback fires and clears it afterwards. A queued-lane send_final_ledgered inside the handler (an earlier message's reply), a turn whose final was suppressed or never produced, and a stale value from an earlier turn on the shared event therefore all read as "no row".
  • run_startup.py: a small process-local registry on the runner, _register_redelivery_followup / _fire_redelivery_followups, keyed by obligation id (24 h cutoff to match the ledger's stale window, enforced on registration and by a timer armed with each new row; cap 64, oldest first). _redeliver_claimed_obligations fires the row's follow-ups after mark_delivered, only when the redelivered send was not truncated. Rows this process already redelivered are remembered, so a follow-up that registers late (the reconnect sweep redelivers inside the refused send's own finalization) runs at once instead of waiting forever.

Guards

  • The delete still runs only after a send that carried the complete reply landed: the inline final (existing stamp) or the ledger redelivery of the same obligation (new). A truncated redelivery does not fire it.
  • A follow-up fires only for its own obligation id, once. The row id is bound to the turn's own send, so the reply to message A can never be the trigger for deleting message B's only preview.
  • Process-local on purpose: a redelivery after a restart finds an empty registry and leaves the bubble. The consumer it would act for died with the old process, and that path already carries the restart marker.
  • Nothing changes for the boot sweep or the reconnect sweep beyond the one extra call in the success branch, which swallows its own exceptions. The queued lane's own frozen previews are out of scope here, as before: its callbacks run outside the hook and see no row.

Tests

tests/gateway/test_redelivered_final_preview_cleanup.py, 27 cases. 13 of the original 21 fail on the parent commit (stale-finalize registration, deferral, own-row-only firing, registry bounds, the redelivery loop's landed / refused / truncated outcomes, the base.py stamp lifecycle including a stale value from an earlier turn, and an end-to-end run on the real test-isolated ledger: refused final, row failed, redelivery, delete). The rest pass there as no-op guards. The six added after review cover the queued-lane-inside-the-handler ordering, a cancelled hook still clearing the stamp, the public bracket's unchanged return, a follow-up registering after its row already landed, and expiry without another registration. #105200's 24 still pass. The delivery ledger, restart-resume, post-delivery-callback, queued-final-ledger and platform base suites are unchanged.

Relation to #105340

#105340 stops the refused preview edits from escalating the flood penalty, which is why the gateway's final ran into a 102 s and a 269 s window at all. This PR cleans up after a refusal that still happens. Independent of each other, both wanted.

Follow-up, deliberately not here

Once the redelivery can prove that no partial copy remains above it, FLOOD_MARKER's "part of it may already have arrived above" can become evidence-based instead of unconditional. Separate change.

@AlexxRussell
AlexxRussell force-pushed the fix/redelivered-final-preview-cleanup branch from eac682d to 577c492 Compare September 9, 2026 08:20
@alt-glitch alt-glitch added type/bug Something isn't working P2 Medium — degraded but workaround exists comp/gateway Gateway runner, session dispatch, delivery platform/telegram Telegram bot adapter sweeper:risk-message-delivery Sweeper risk: may drop, duplicate, misroute, or suppress messages labels Sep 9, 2026
@Enough1122

Copy link
Copy Markdown

AI code review — automated review for reference; please use your judgment.

PR #106370 — fix(gateway): delete the frozen streamed bubble once the ledger redelivers the refused final

Summary: Large change across gateway/platforms/base.py, run_startup.py, run_turn.py, stream_consumer_transport.py: (1) tracks whether the COMPLETE final text actually reached the reader (SendResult.truncated, TTS-caption signal, text_delivered kept separate from aggregate delivery success); (2) schedules abandoned-preview cleanup via the post-delivery callback with a delivered stamp; (3) defers the delete to the ledger redelivery (_register_redelivery_followup, bounded 64 entries + 24h TTL) when the final was refused but ledgered.

Findings (all Non-blocking):

  • The three guards against deleting the reader's only copy are well-designed: skip when content delivered, segment-only single-bubble seam (abandoned_preview_ids returns empty on split delivery or ≠1 bubble, since capped multi-chunk sends still report success), and stand-down on failed send with ledger-deferral fallback. This is the safety-critical core and it reads correctly.
  • Stamping _hermes_final_delivered / _hermes_final_obligation_id as ad-hoc attributes on asyncio.Event, written in exactly one place and cleared in finally (cancel-safe), is hacky but contained and documented. Race note: the pre-drain firing (base.py:196) before handing the session to the drain task correctly avoids popping the next turn's callback with this turn's outcome.
  • Edge: _redelivered_obligation_ids deque has maxlen=64 — a landed id evicted before a late registration means that follow-up waits until TTL prune without ever firing. Bounded leak (closure held ≤24h), acceptable, but worth knowing.
  • Redelivery only fires follow-ups on untruncated sends — a truncated redelivery correctly does not authorize the delete.
  • ~1230 lines of new tests for the two cleanup branches is proportionate to the risk.

Verdict: Careful lifecycle work; guards read correctly. Safe to merge.

@AlexxRussell
AlexxRussell force-pushed the fix/redelivered-final-preview-cleanup branch from 577c492 to 40a125e Compare September 19, 2026 23:19
@AlexxRussell

Copy link
Copy Markdown
Author

Rebased onto current main (8a92051f20), replayed as one commit
(40a125ed14) on top of the rebased #105200 (6d5e879b95).

No conflicts of its own: everything upstream moved underneath it was resolved in
#105200, and that note is worth reading with this one. In particular
_adapter_for_source is gone from run_turn as of upstream 4c132febc9
(19 Sep), so the cleanup lookup is now _delivery_adapter_for; the two test
mocks in tests/gateway/test_redelivered_final_preview_cleanup.py moved with it.

One upstream test needed retargeting, and it is the part of this rebase most
worth a second opinion.
cd3de040ab added
tests/gateway/test_diagnostic_wake_presentation.py, which mocks the public
send_final_ledgered and asserts its await count. This PR splits that method
into a public wrapper plus a private _send_final_ledgered that also returns
the ledger row id, and the normal lane calls the private one, so the mock never
fires and the count is 0. I retargeted the mock to _send_final_ledgered
(3-tuple); the test's intent, that a suppressed diagnostic wake sends no final,
is unchanged. Public callers still get the 2-tuple and there are no in-tree
overrides of the method. If you would rather no upstream test moved, the
alternative is to carry the obligation id on SendResult and leave delivery
dispatching through the public method; say the word and I will reshape it that
way.

The one dash in this diff is upstream's own send_final_ledgered docstring,
carried verbatim onto the wrapper when the method was split; nothing in the
added text introduces one.

tests/gateway/test_redelivered_final_preview_cleanup.py with
test_abandoned_preview_cleanup.py is 51 passed. Full tests/gateway: 9438
passed with 31 failed and 57 skipped, against 9386 passed with 32 failed and 57
skipped on clean main at the same commit. No failure is unique to this branch:
the sets differ by one test that fails on main and passes here
(test_api_server_runs.py::TestStartRun::test_events_stream_forwards_interim_commentary),
which is order-dependent rather than anything this PR touches. The remaining 31
are a local macOS and Python 3.14 environment set.

@AlexxRussell
AlexxRussell force-pushed the fix/redelivered-final-preview-cleanup branch from 40a125e to 8420dc6 Compare September 28, 2026 07:47
@AlexxRussell

Copy link
Copy Markdown
Author

Rebased onto current main (together with #105200, which this stacks on). One conflict is worth pointing out because it only shows up as context in the diff: main now releases the turn marker inside send_final_ledgered once the obligation is recorded (360b969). This PR splits that method into a private _send_final_ledgered, which the queued lane's _send_final_text also uses, and a thin public wrapper. The release now sits in the private method, so both the normal and the queued lane keep it, in the same order as on main.

…e the final lands

The stream consumer deletes the previews it replaces, but only on paths where it sends the
replacement itself. When its edits stop working entirely (a Telegram flood window longer than the
5s inline cap) control returns to run_turn.py and the gateway sends the final instead. That branch
only logged, so the frozen preview stayed in the chat above the complete reply, showing raw
MarkdownV2 markers and the streaming cursor.

The cleanup registers the consumer's own preview deletion through the post-delivery callback, with
three guards so it can never remove the reader's only copy of the answer:

- skipped when the stream did deliver the content (those previews are the reply);
- a new public seam on the transport mixin that hands over only the single frozen bubble of a
  non-split segment: finalized earlier segments survive a tool boundary, and one bubble holds at
  most one platform message of the reply's prefix, which any successful final send covers even
  when an adapter capped a long reply by message count or length and still reported success;
- gated on a new delivery stamp base.py writes on the session event before firing the hook. The
  stamp means the COMPLETE final text reached the user: its own send, or a TTS caption carrying the
  whole text. Bare audio does not count, and neither does the plain-text fallback when it truncated
  the reply (new SendResult.truncated). On the pending-drain path the hook now fires before the
  hand-off, so a finishing turn can no longer pop the next turn's callback with its own outcome.

The post-delivery hook itself stays unconditional for its lifecycle consumers.

tests/gateway/test_abandoned_preview_cleanup.py: 24 cases, 18 fail on unmodified main.
…ivers the refused final

The abandoned-preview cleanup (parent commit) reached one branch and one outcome: the consumer
gave up and the gateway's own final landed inline. Two sibling cases kept leaving the frozen
bubble in the chat, seen on Telegram on 8 and 9 Sep 2026:

* Stale finalize. The consumer's finalize edit had landed a stale snapshot, so run_turn tried a
  reconciliation edit; a flood window refused it and the normal final send put the complete reply
  below the frozen bubble. That branch registered no cleanup at all.
* Ledger redelivery. The gateway's final was refused too (a 102 s and a 269 s penalty, both over
  the 60 s inline cap), the ledger redelivered it minutes later with the flood marker, and the
  registered cleanup had already stood down because the post-delivery stamp said the text did
  not land. The redelivery runs outside the turn, so nothing fired it.

Now the stale-finalize branch registers the same cleanup when its edit is refused, and a cleanup
that stands down on a refused send hands its delete to the ledger row carrying that final. The
row id comes from the turn's own final send: _send_final_text returns it and the post-delivery
hook stamps it on the session event beside the delivery stamp, then clears it (in finally, so a
cancelled callback cannot leave it behind), so neither a queued-lane send inside the handler
nor an earlier turn on the shared event can leave a row behind for a turn whose own final never
went out. The runner keeps a small process-local registry keyed by obligation id (24 h cutoff
enforced on registration and by a timer, cap 64, oldest first), and
_redeliver_claimed_obligations fires the registered follow-ups once the redelivered send has
landed untruncated. A row this process already redelivered runs a late follow-up at once (the
reconnect sweep redelivers inside the refused send's own finalization). A follow-up fires only
for its own row; with no row nothing will replace the bubble and it stays; a redelivery after a
restart finds an empty registry.

Tests: 27 in tests/gateway/test_redelivered_final_preview_cleanup.py; 13 of the original 21
fail on the parent commit, the rest pass there as no-op guards. The parent's 24 still pass.
@AlexxRussell
AlexxRussell force-pushed the fix/redelivered-final-preview-cleanup branch from 8420dc6 to f2a6e8b Compare October 2, 2026 14:06

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

comp/gateway Gateway runner, session dispatch, delivery P2 Medium — degraded but workaround exists platform/telegram Telegram bot adapter sweeper:risk-message-delivery Sweeper risk: may drop, duplicate, misroute, or suppress messages type/bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants