Skip to content

fix(gateway): delete the streaming preview the consumer abandoned once the final lands - #105200

Open
AlexxRussell wants to merge 1 commit into
NousResearch:mainfrom
AlexxRussell:fix/abandoned-preview-cleanup
Open

AlexxRussell wants to merge 1 commit into
NousResearch:mainfrom
AlexxRussell:fix/abandoned-preview-cleanup

Conversation

@AlexxRussell

Copy link
Copy Markdown

A streaming preview that the consumer abandons is left in the chat when the gateway sends the final reply itself, so the reader sees a truncated message with raw MarkdownV2 markers and the streaming cursor sitting above the complete answer.

What happens

The stream consumer deletes the previews it replaces. _delete_previews has four call sites and every one of them is inside the consumer, on a path where the consumer itself sends the replacement: the fresh-final path in gateway/stream_consumer_transport.py, and the partial, empty and silence-marker paths in gateway/stream_consumer_fallback.py.

There is a fifth case the consumer does not own. When its edits stop working entirely, control returns to gateway/run_turn.py and the gateway sends the final instead. That branch, the elif _sc is not None: arm of _run_agent_mark_streamed_delivery, is labelled a duplicate-risk diagnostic and only logs. Nothing cleans up behind it, so the abandoned preview stays.

Observed

On a 1 GB deployment running the Telegram adapter, 6 September 2026:

16:02:10  [Telegram] Telegram flood control, waiting 9.0s      (x4)
16:02:16  Normal final-send NOT suppressed despite active stream consumer
          for session ...: streamed=False previewed=False content_delivered=...
16:02:16  [Telegram] Sending response (2888 chars) to ...

The four preview edits asked for a 9 second wait, which exceeds the adapter's 5 second inline cap, so each failed closed rather than sleeping. The preview froze mid-sentence. Six seconds later the complete 2888 character reply arrived as a separate message. Nothing was lost, but the frozen preview remained, and because a partial preview is rendered without MarkdownV2 its asterisks and heading markers were literal.

This pattern appears 6 times in that deployment's logs since 25 May, 4 of them in the last 8 days as traffic grew. It is cosmetic, never a lost reply.

_delete_previews already anticipates the neighbouring case: its docstring notes that "the same flood window that broke the finalize edit can reject this delete too, leaving the preview bubble next to the fresh final (#71047 Problem B)". This is the same class of problem on the one path that had no cleanup at all.

The change

Three files, one small piece each.

gateway/stream_consumer_transport.py: two public methods on the transport mixin, abandoned_preview_ids() and delete_abandoned_previews(), so the gateway never reaches into consumer privates. The ids are the current segment's only. A tool boundary finalizes the segment before it and resets the per-segment state, but the turn-wide preview set keeps those ids; handing that set over would delete delivered preambles whose text is absent from the final answer. And only a single frozen bubble of a non-split segment is handed over. Adapters cap a long final by message count or length and still report success (Discord keeps the first N-1 chunks and replaces the rest with a notice, WeCom and others cut at their limit, none of them marks the result). One bubble holds at most one platform message of the reply's prefix, which any successful final send covers; a split chain or several bubbles could hold more than a capped final delivered, so those stay. The delete goes through the existing _delete_previews with retry_on_false=True, because Telegram reports a refused delete by returning False and the flood window that stranded the preview can refuse the delete too.

gateway/platforms/base.py: _fire_post_delivery_callback takes delivered and stamps it on the session's interrupt event before popping the callback, the same way the run generation already travels on that event. The value is a new text_delivered flag, kept apart from the aggregate delivery_succeeded: it is set only by the final text send itself, or by a TTS caption that carried the complete text. A voice reply can land its audio while the text send is refused, and that must not read as "text delivered". Nor may a truncated fallback: after a formatting failure _send_with_retry sends the first 3500 characters as plain text and returns that success, so SendResult gains a truncated flag the fallback sets when it cut the text, and the recorder ignores such results. The hook itself stays unconditional: goal continuation and the background-review release depend on it firing either way.

One ordering change on the pending-drain path. A follow-up queued mid-turn is drained by a new task that shares the session's interrupt event and re-stamps its run generation as soon as it starts. Popped from finally after the hand-off, the callback could be the next turn's, fired with this turn's outcome while the next final is still in flight. The hook is therefore fired before _spawn_drain_task on that branch, and finally skips it. Everywhere else the timing is unchanged.

gateway/run_turn.py: _run_agent_schedule_abandoned_preview_cleanup, a sibling of _run_agent_schedule_bubble_cleanup using the same post-delivery registration, called as the last statement of that elif arm.

Three guards keep this from ever removing the reader's only copy of the answer:

  • It is skipped when the stream delivered the content. Those previews hold the answer.
  • The ids are segment-only, as above.
  • The delete is deferred to the post-delivery callback and gated on the stamp. The branch executes before the normal final send, and the post-delivery hook fires whether or not that send succeeded. Without the gate, a failed final send would have deleted the frozen preview while no replacement existed. The callback reads the stamp off the active-session event and stands down when it is missing or False.

A consumer without the transport mixin, an adapter without the post-delivery hook, an empty stale set, a failed registration and a refused delete are all no-ops.

Tests

tests/gateway/test_abandoned_preview_cleanup.py, 24 cases. 18 fail without the change and the 6 pure guard tests pass there as no-ops.

  • The branch registers the cleanup and nothing is deleted inline.
  • A landed final lets the callback delete the stale previews.
  • A final that did not land keeps the preview, parametrized over a False stamp, a missing stamp and a missing active session.
  • Seam tests on the real GatewayStreamConsumer: after a segment break, abandoned_preview_ids() returns the new segment's bubble only; a split chain and a multi-bubble segment return nothing; the __no_edit__ sentinel is never included.
  • Two end-to-end cases drive the real _process_message_background with a BasePlatformAdapter subclass and a mocked _send_with_retry: success deletes, a refused send keeps the preview.
  • Two voice cases through the same processor: audio that landed bare while the text send was refused keeps the preview; a delivered caption counts as the text landing and the preview goes.
  • A two-turn drain case queues a follow-up mid-turn, blocks and then refuses the second turn's final send, and checks that the first turn fired only its own cleanup and the second turn's preview stayed.
  • Two cases drive the real _send_with_retry: a formatting failure whose plain-text fallback truncates a long reply keeps the preview, and a complete fallback deletes it.
  • The remaining cases pin that the guards stay no-ops and that a confirmed streamed delivery still suppresses the normal send.

Scope: this is the abandoned-preview path only. It does not change when a preview is created, when suppression fires, or the adapter's flood handling. Adapters that cap a send could set the new SendResult.truncated themselves in a follow-up, which would let the cleanup cover split chains too; this PR does not depend on that. One known cosmetic edge remains: Slack's rich-blocks final clips a heading over 150 characters in the rendered message while the fallback text keeps it, so deleting a preview that showed the raw heading loses that tail from the rendered view. The consumer's own fresh-final path has the same property, so it is left alone here.

@alt-glitch alt-glitch added type/bug Something isn't working P3 Low — cosmetic, nice to have comp/gateway Gateway runner, session dispatch, delivery platform/telegram Telegram bot adapter area/streaming Streaming responses: gateway delivery, provider wire sweeper:risk-message-delivery Sweeper risk: may drop, duplicate, misroute, or suppress messages labels Sep 7, 2026
@Enough1122

Copy link
Copy Markdown

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

Summary

Deletes the streaming preview the consumer abandoned once the gateway-sent final lands: SendResult.truncated distinguishes prefix-only fallback sends, base.py stamps _hermes_final_delivered (complete text via own send or TTS caption only) before firing the post-delivery hook, and the cleanup runs from that hook — standing down when the final never landed. Segment-only/single-bubble seam keeps finalized preambles and split-chain heads safe.

Findings

  • Non-blocking — gateway/platforms/base.py:50: stamping _hermes_final_delivered on the shared interrupt event before popping the generation-checked callback is the load-bearing ordering, and firing the callback before the follow-up drain hand-off (base.py:102) closes the exact race where a finishing turn could fire the next turn's cleanup. Both races are pinned by tests.
  • Non-blocking — gateway/stream_consumer_transport.py:201: the single-frozen-bubble rule is honestly conservative — adapters that cap long finals (Discord N-1, WeCom length cut) still report success, so multi-bubble segments are never handed over. Correct to under-delete here.
  • Non-blocking — gateway/run_turn.py:164: the loop snapshot + suppress(Exception) around scheduling means a closing loop silently drops the cleanup, leaving a stale preview. Acceptable (cosmetic leftover, retry path exists via bounded retry on the delete itself), just noting the silent path.

Verified

  • 604-line test module: streamed-delivery suppression path untouched, segment/split/multi-bubble/__no_edit__ exclusions, bounded-retry delete, end-to-end landed-vs-refused, voice (audio≠text, caption=text), and the queued-follow-up generation race.

Verdict

Looks good. Three independent guards with the failure modes tested individually.

@AlexxRussell

Copy link
Copy Markdown
Author

Note on how this relates to #105340, since the two came out of the same incident and could look overlapping.

#105340 fixes the cause. A flood-refused preview edit ignored the retry_after the adapter had already attached to the result, so from a sub-second edit interval all three strikes were spent in under half a second and the stream was abandoned. Each refused edit was itself another request against the same limit, which is how a 9 second penalty became 271 seconds on the final send.

#105200 fixes the leftover. On the one path the consumer does not own, where the gateway sends the final itself, nothing deleted the frozen preview, so it stayed above the real answer showing raw MarkdownV2 and the streaming cursor.

They are complementary rather than alternatives, and neither depends on the other. #105340 makes the abandoned-preview path much rarer without removing it: a penalty above the cap, an adapter that reports no retry_after, and continued refusal after compliance all still fall back to the gateway send. #105200 is also the only one of the two that helps a deployment already running the current consumer.

Two files appear in both, gateway/run_turn.py and gateway/stream_consumer_transport.py, but no function is touched by both, so either can land first. I checked rather than assumed: both commits cherry-pick onto the current main (fef0e16fe1) in sequence with no conflict, and on that combined tree the two new test modules plus the surrounding streaming, flood, post-delivery and final-contract suites give 238 passed, 1 skipped.

Happy to combine them into one PR if that is easier to review, or to rebase either onto a newer main on request.

@AlexxRussell

Copy link
Copy Markdown
Author

Rebased onto current main (8a92051f20), replayed as one commit
(6d5e879b95).

Three things changed in the rebase. Two of them are load-bearing, so they are
worth reading before the diff.

1. _adapter_for_source no longer exists. Upstream 4c132febc9 (19 Sep)
moved run_turn onto the intake/delivery seams. The cleanup scheduler here
called self._adapter_for_source(source), which was correct at this branch's
base and is an AttributeError on today's main, in production and not only
under test. It now calls self._delivery_adapter_for(source), the same lookup
run_turn uses for its own progress-bubble cleanup adapter
(gateway/run_turn.py:3042), and the test's mock moved with it. A cleanup that
deletes a message in the chat it answered wants the delivery seam, so this is
the right successor rather than the nearest one.

2. The truncation flag moved to where the truncation happens. Upstream
extracted the plain-text fallback into _send_plain_fallback, which is an
override point (Photon drops rich links instead of adding the banner). Stamping
SendResult.truncated at the old call site would have hardcoded the 3500
character cap outside the method that now applies it, and would have mis-stamped
any override that does not truncate at all. The stamp lives in the default
_send_plain_fallback instead, so an override opts in by truncating, and
_send_with_retry carries no change from this PR at all.

3. Guard order. abandoned_preview_ids() is now read before the adapter
lookup, so a turn with nothing abandoned does no adapter work and registers no
callback. That also keeps
tests/gateway/test_stream_consumer.py::TestConfirmedFinalDeliveryConsultsCommentaryRecord::test_draft_streamed_text_with_failed_finalize_send_keeps_fallback_final_send
from reaching a sibling mixin's attribute on a bare GatewayTurnMixin(), which
is how the stale lookup in point 1 surfaced.

The change itself is unchanged: the previews the consumer abandoned are deleted
only when the gateway sent the replacement, only from the segment-only seam, and
only after the post-delivery callback confirms the replacement landed.

tests/gateway/test_abandoned_preview_cleanup.py is 24 passed.
Full tests/gateway: 9410 passed with 32 failed and 57 skipped, against 9386
passed with 32 failed and 57 skipped on clean main at the same commit. The two
failure sets are identical name for name (a local macOS and Python 3.14
environment set, nothing this branch touches), and the 24 added passes are this
PR's own file.

…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.
@AlexxRussell
AlexxRussell force-pushed the fix/abandoned-preview-cleanup branch from 6a2bf8e to 369379d Compare October 2, 2026 14:05

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

area/streaming Streaming responses: gateway delivery, provider wire comp/gateway Gateway runner, session dispatch, delivery P3 Low — cosmetic, nice to have 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