Streamed replies that overflow the platform limit are no longer duplicated or left as a stale partial (#25349, salvage #114151) - #120315
Merged
Conversation
…inal lane _seal_overflow_heads clears the edit target mid-iteration, so the overflow gate evaluated at the top of run() is already stale when the leftover is pushed. The seal loop exits after ONE successful seal (its condition includes _message_id is not None), so a buffer far past the limit leaves most of itself behind. That leftover then reaches _push_update -> _first_send with no message to edit, on a NON-final tick, and adapter.send() chunks and caps it itself, numbering the pieces (i/n). The turn-final lane later publishes the same text again with its own denominator, because the two lanes use different budgets and different payload shapes. Observed in production on Discord: one inbound message, one API call, no tool turns, a 36642-char answer, and two interleaved sequences on screen - (i/10) and (i/9) - whose first chunks were byte-identical once the indicator was stripped and whose concatenations shared a prefix. The gateway's normal final send was correctly suppressed (streamed=True, content_delivered=True), which places the duplicate inside the consumer rather than in the gateway's delivery ledger. Fix: re-check the overflow gate AFTER the seal and hand a still-oversized leftover to _split_first_send, which is the consumer's own cap-aware, ledger-aware splitter and already owns sealing. This removes the same-iteration invalidation rather than compensating for it downstream, and it does not add a fourth delivery flag: the existing flags behaved correctly here. Negative control on this branch: with the change reverted, both new tests fail, including the one asserting that no line is published twice.
Two follow-ups to the re-split after _seal_overflow_heads (previous commit), found by the exactly-once E2E suite driving the real GatewayRunner on platforms whose limit a streamed reply overflows (Discord 2000): 1. _split_first_send kept the adapter's " (n/n)" chunk indicator on the tail chunk it reuses as the live preview, so later deltas were appended after it: the user saw "...wo (3/3)rd0575..." embedded mid-reply and the visible text no longer matched the transcript. Strip the indicator from the kept tail. 2. The re-split after a seal now `continue`s like the gate above it when the turn is not finished: the tail is still unsent, and falling through into the segment-break reset would clear it. tests/gateway/test_stream_consumer.py::...fence_aware_split asserted the old behaviour (the tail starts with the full indicator-suffixed chunk); it now asserts the indicator is dropped from the tail and never appears in an edit.
When the stream consumer's first send failed (e.g. a Telegram timeout that never reached the platform), _first_send disabled edits but left the message id unset, so the next tick sent another first send: a partial preview that could never be edited. The final reply then went out as a new message and the truncated preview stayed on screen next to it (on every first-send timeout, not only under load). With edits disabled, skip non-final first sends; the final is delivered once, complete. Found by the C12 exactly-once suite (fault_matrix[stream_timeout_first_send- fk_tg|fk_dc]).
- test_stream_consumer_oversized_leftover.py: keep the user-visible invariant (no line of the answer reaches the channel twice as a new message) and drop its lane-size twin, which pins the same re-split. - test_stream_consumer.py: a failed first send leaves no uneditable partial preview; the complete reply is the only message that reaches the chat. Both are red on origin/main and green with the fixes.
૮ >ﻌ< ა ci reviewran on 9b560ce — chore: retrigger CI (zero-job dispatch failure, auto-heal) debug infoCI timingsCI timings · View report · View jobWall time 6m40s vs 6m23s (+4.4%). 9 job(s) slower, 3 faster,
|
teknium1
added a commit
that referenced
this pull request
Sep 23, 2026
A child-process GatewayRunner (own HERMES_HOME, SIGKILL + restart) runs the real agent against the scripted fake LLM provider and three instances of a FakePlatformAdapter that implements the gateway/platforms/base.py contract (4096/2000 limits, edits vs no edits, threads, streaming on/off). Only the transport is fake; its fsynced op journal is ground truth. Faults per op: timeout, ack lost, 429 retry_after, connection reset, too long, and park points for SIGKILL before/after the platform applies an op or between provider completion and the delivery ledger. Oracle per inbound: exactly one complete unmarked visible reply (extra copies only with the ledger's duplicate marker), visible text == persisted assistant text, user row persisted once, model turn run once, no stray or stale partial reply; plus a whole-run audit across chats. Scenarios: 13 faults x 3 platforms, /queue and interrupt follow-ups on a slow turn, parallel threads, re-delivered inbound ids (during/after a turn, after a runner reconnect), five crash points with restart catch-up, and an unclean restart with nothing in flight. Red-proven against the reverted stream fixes (#120315), a hand-revert of c961e5b (#91653), retrying timed-out sends, boot redelivery without the marker, double-sent finals, a dropped queued follow-up and disabled inbound dedup. Five strict xfails pin four live gaps (streaming ack-lost duplicate, dedup lost on reconnect, streaming crash after accept, and the unclean-restart recency fallback that re-answers every session active in the last 120 s). Eight cells stay red until #120315 lands (every streamed reply over the limit, and the first-send timeout) and carry xfail 'fixed by #120315': strict where the bug fires every run, non-strict where it depends on the stream/edit tick. The Director counts a resume note merged into the original tagged user message as a resume, not a second run of that inbound.
…ill applies Review finding on the seal re-split: `continue` after `_split_first_send` skipped the rest of the tick, so a commentary drained with the oversized burst was dropped and a tool boundary was lost (post-tool text then extended the pre-tool preview). The initial-overflow gate had the same `continue`. Seal first (a no-op without a message id), split once, then push the tail through `_push_update` in the same tick so commentary, segment end and the flush barrier run normally after it. `continue` is kept only when a head send failed and the full text must stay for the fallback final.
Collaborator
Author
|
Review finding addressed in 7d5e9ae: the
|
teknium1
added a commit
that referenced
this pull request
Sep 23, 2026
A child-process GatewayRunner (own HERMES_HOME, SIGKILL + restart) runs the real agent against the scripted fake LLM provider and three instances of a FakePlatformAdapter that implements the gateway/platforms/base.py contract (4096/2000 limits, edits vs no edits, threads, streaming on/off). Only the transport is fake; its fsynced op journal is ground truth. Faults per op: timeout, ack lost, 429 retry_after, connection reset, too long, and park points for SIGKILL before/after the platform applies an op or between provider completion and the delivery ledger. Oracle per inbound: exactly one complete unmarked visible reply (extra copies only with the ledger's duplicate marker), visible text == persisted assistant text, user row persisted once, model turn run once, no stray or stale partial reply; plus a whole-run audit across chats. Scenarios: 13 faults x 3 platforms, /queue and interrupt follow-ups on a slow turn, parallel threads, re-delivered inbound ids (during/after a turn, after a runner reconnect), five crash points with restart catch-up, and an unclean restart with nothing in flight. Red-proven against the reverted stream fixes (#120315), a hand-revert of c961e5b (#91653), retrying timed-out sends, boot redelivery without the marker, double-sent finals, a dropped queued follow-up and disabled inbound dedup. Five strict xfails pin four live gaps (streaming ack-lost duplicate, dedup lost on reconnect, streaming crash after accept, and the unclean-restart recency fallback that re-answers every session active in the last 120 s). Eight cells stay red until #120315 lands (every streamed reply over the limit, and the first-send timeout) and carry xfail 'fixed by #120315': strict where the bug fires every run, non-strict where it depends on the stream/edit tick. The Director counts a resume note merged into the original tagged user message as a resume, not a second run of that inbound.
teknium1
added a commit
that referenced
this pull request
Sep 23, 2026
…y cells A strict xfail on a gap whose fix is an open PR turns main red the moment that fix merges (XPASS), and the non-strict ones guarded nothing. Each gap now has a probe (tests/e2e/core/delivery/_pending_fixes.py) that reproduces the defect's mechanism on the tree under test in a throwaway interpreter; expect_gap() applies the strict xfail only while the probe still reproduces it, so the cell becomes a plain test once the fix is in the tree, whatever the merge order. Covered: #120314, #119970 (soak), #120315, #120377, #120444 (C12) and #120450 (cell 5). Each probe was checked against every fix head: it flips on its own PR and on no other. C12 cells made deterministic (identical outcome on every run): - long_split streams the whole reply as one chunk; long_streamed and stream_timeout_first_send pace chunks so each lands in its own consumer tick. The five former coin-flip xfails are now two plain cells and three strict #120315 gap cells. - sent_ack_lost waits until the answer is persisted before the kill, so it pins the #120377 recovery; the streamed-before-persisted order is its own cell (stream_accepted_unpersisted, a strict live gap with a stalled provider stream). - zzz_unclean_restart compares director.resumes against a snapshot taken before its kill instead of requiring it empty: crash cells on their own homes may legitimately resume. - the whole-run audit skips the reconnect replay only while #120444's gap is open.
teknium1
added a commit
that referenced
this pull request
Sep 23, 2026
A child-process GatewayRunner (own HERMES_HOME, SIGKILL + restart) runs the real agent against the scripted fake LLM provider and three instances of a FakePlatformAdapter that implements the gateway/platforms/base.py contract (4096/2000 limits, edits vs no edits, threads, streaming on/off). Only the transport is fake; its fsynced op journal is ground truth. Faults per op: timeout, ack lost, 429 retry_after, connection reset, too long, and park points for SIGKILL before/after the platform applies an op or between provider completion and the delivery ledger. Oracle per inbound: exactly one complete unmarked visible reply (extra copies only with the ledger's duplicate marker), visible text == persisted assistant text, user row persisted once, model turn run once, no stray or stale partial reply; plus a whole-run audit across chats. Scenarios: 13 faults x 3 platforms, /queue and interrupt follow-ups on a slow turn, parallel threads, re-delivered inbound ids (during/after a turn, after a runner reconnect), five crash points with restart catch-up, and an unclean restart with nothing in flight. Red-proven against the reverted stream fixes (#120315), a hand-revert of c961e5b (#91653), retrying timed-out sends, boot redelivery without the marker, double-sent finals, a dropped queued follow-up and disabled inbound dedup. Five strict xfails pin four live gaps (streaming ack-lost duplicate, dedup lost on reconnect, streaming crash after accept, and the unclean-restart recency fallback that re-answers every session active in the last 120 s). Eight cells stay red until #120315 lands (every streamed reply over the limit, and the first-send timeout) and carry xfail 'fixed by #120315': strict where the bug fires every run, non-strict where it depends on the stream/edit tick. The Director counts a resume note merged into the original tagged user message as a resume, not a second run of that inbound.
teknium1
added a commit
that referenced
this pull request
Sep 23, 2026
…y cells A strict xfail on a gap whose fix is an open PR turns main red the moment that fix merges (XPASS), and the non-strict ones guarded nothing. Each gap now has a probe (tests/e2e/core/delivery/_pending_fixes.py) that reproduces the defect's mechanism on the tree under test in a throwaway interpreter; expect_gap() applies the strict xfail only while the probe still reproduces it, so the cell becomes a plain test once the fix is in the tree, whatever the merge order. Covered: #120314, #119970 (soak), #120315, #120377, #120444 (C12) and #120450 (cell 5). Each probe was checked against every fix head: it flips on its own PR and on no other. C12 cells made deterministic (identical outcome on every run): - long_split streams the whole reply as one chunk; long_streamed and stream_timeout_first_send pace chunks so each lands in its own consumer tick. The five former coin-flip xfails are now two plain cells and three strict #120315 gap cells. - sent_ack_lost waits until the answer is persisted before the kill, so it pins the #120377 recovery; the streamed-before-persisted order is its own cell (stream_accepted_unpersisted, a strict live gap with a stalled provider stream). - zzz_unclean_restart compares director.resumes against a snapshot taken before its kill instead of requiring it empty: crash cells on their own homes may legitimately resume. - the whole-run audit skips the reconnect replay only while #120444's gap is open.
teknium1
added a commit
that referenced
this pull request
Sep 23, 2026
A child-process GatewayRunner (own HERMES_HOME, SIGKILL + restart) runs the real agent against the scripted fake LLM provider and three instances of a FakePlatformAdapter that implements the gateway/platforms/base.py contract (4096/2000 limits, edits vs no edits, threads, streaming on/off). Only the transport is fake; its fsynced op journal is ground truth. Faults per op: timeout, ack lost, 429 retry_after, connection reset, too long, and park points for SIGKILL before/after the platform applies an op or between provider completion and the delivery ledger. Oracle per inbound: exactly one complete unmarked visible reply (extra copies only with the ledger's duplicate marker), visible text == persisted assistant text, user row persisted once, model turn run once, no stray or stale partial reply; plus a whole-run audit across chats. Scenarios: 13 faults x 3 platforms, /queue and interrupt follow-ups on a slow turn, parallel threads, re-delivered inbound ids (during/after a turn, after a runner reconnect), five crash points with restart catch-up, and an unclean restart with nothing in flight. Red-proven against the reverted stream fixes (#120315), a hand-revert of c961e5b (#91653), retrying timed-out sends, boot redelivery without the marker, double-sent finals, a dropped queued follow-up and disabled inbound dedup. Five strict xfails pin four live gaps (streaming ack-lost duplicate, dedup lost on reconnect, streaming crash after accept, and the unclean-restart recency fallback that re-answers every session active in the last 120 s). Eight cells stay red until #120315 lands (every streamed reply over the limit, and the first-send timeout) and carry xfail 'fixed by #120315': strict where the bug fires every run, non-strict where it depends on the stream/edit tick. The Director counts a resume note merged into the original tagged user message as a resume, not a second run of that inbound.
teknium1
added a commit
that referenced
this pull request
Sep 23, 2026
…y cells A strict xfail on a gap whose fix is an open PR turns main red the moment that fix merges (XPASS), and the non-strict ones guarded nothing. Each gap now has a probe (tests/e2e/core/delivery/_pending_fixes.py) that reproduces the defect's mechanism on the tree under test in a throwaway interpreter; expect_gap() applies the strict xfail only while the probe still reproduces it, so the cell becomes a plain test once the fix is in the tree, whatever the merge order. Covered: #120314, #119970 (soak), #120315, #120377, #120444 (C12) and #120450 (cell 5). Each probe was checked against every fix head: it flips on its own PR and on no other. C12 cells made deterministic (identical outcome on every run): - long_split streams the whole reply as one chunk; long_streamed and stream_timeout_first_send pace chunks so each lands in its own consumer tick. The five former coin-flip xfails are now two plain cells and three strict #120315 gap cells. - sent_ack_lost waits until the answer is persisted before the kill, so it pins the #120377 recovery; the streamed-before-persisted order is its own cell (stream_accepted_unpersisted, a strict live gap with a stalled provider stream). - zzz_unclean_restart compares director.resumes against a snapshot taken before its kill instead of requiring it empty: crash cells on their own homes may legitimately resume. - the whole-run audit skips the reconnect replay only while #120444's gap is open.
…loses the tail After an overflow split the heads are sealed and the message id is cleared, so the tail goes out as a first send in the tick that carries the tool boundary. When that send failed, _end_segment skipped the unseen-tail flush (it required a truthy _message_id) and _reset_segment_state cleared _accumulated: ~1.1k chars vanished. The same loss existed for a plain failed first send at a segment break. The flush now runs whenever the boundary update did not land, except for the __no_edit__ sentinel (its continuation still goes out once via the fallback final). With a real message id the condition is unchanged, so no new send there. The flush also honours the egress guard like the other new-message fallbacks, since it can now run on the path where a draft was declined. Tests: the oversized-leftover file now pins one invariant, parametrized over none / commentary / segment_break / tail_send_fails: every tail line reaches the channel, never twice as a new message, and post-boundary text is not glued onto the pre-boundary preview. tail_send_fails is red on f361b39 (24 lines lost), green here.
teknium1
added a commit
that referenced
this pull request
Sep 23, 2026
A child-process GatewayRunner (own HERMES_HOME, SIGKILL + restart) runs the real agent against the scripted fake LLM provider and three instances of a FakePlatformAdapter that implements the gateway/platforms/base.py contract (4096/2000 limits, edits vs no edits, threads, streaming on/off). Only the transport is fake; its fsynced op journal is ground truth. Faults per op: timeout, ack lost, 429 retry_after, connection reset, too long, and park points for SIGKILL before/after the platform applies an op or between provider completion and the delivery ledger. Oracle per inbound: exactly one complete unmarked visible reply (extra copies only with the ledger's duplicate marker), visible text == persisted assistant text, user row persisted once, model turn run once, no stray or stale partial reply; plus a whole-run audit across chats. Scenarios: 13 faults x 3 platforms, /queue and interrupt follow-ups on a slow turn, parallel threads, re-delivered inbound ids (during/after a turn, after a runner reconnect), five crash points with restart catch-up, and an unclean restart with nothing in flight. Red-proven against the reverted stream fixes (#120315), a hand-revert of c961e5b (#91653), retrying timed-out sends, boot redelivery without the marker, double-sent finals, a dropped queued follow-up and disabled inbound dedup. Five strict xfails pin four live gaps (streaming ack-lost duplicate, dedup lost on reconnect, streaming crash after accept, and the unclean-restart recency fallback that re-answers every session active in the last 120 s). Eight cells stay red until #120315 lands (every streamed reply over the limit, and the first-send timeout) and carry xfail 'fixed by #120315': strict where the bug fires every run, non-strict where it depends on the stream/edit tick. The Director counts a resume note merged into the original tagged user message as a resume, not a second run of that inbound.
teknium1
added a commit
that referenced
this pull request
Sep 23, 2026
…y cells A strict xfail on a gap whose fix is an open PR turns main red the moment that fix merges (XPASS), and the non-strict ones guarded nothing. Each gap now has a probe (tests/e2e/core/delivery/_pending_fixes.py) that reproduces the defect's mechanism on the tree under test in a throwaway interpreter; expect_gap() applies the strict xfail only while the probe still reproduces it, so the cell becomes a plain test once the fix is in the tree, whatever the merge order. Covered: #120314, #119970 (soak), #120315, #120377, #120444 (C12) and #120450 (cell 5). Each probe was checked against every fix head: it flips on its own PR and on no other. C12 cells made deterministic (identical outcome on every run): - long_split streams the whole reply as one chunk; long_streamed and stream_timeout_first_send pace chunks so each lands in its own consumer tick. The five former coin-flip xfails are now two plain cells and three strict #120315 gap cells. - sent_ack_lost waits until the answer is persisted before the kill, so it pins the #120377 recovery; the streamed-before-persisted order is its own cell (stream_accepted_unpersisted, a strict live gap with a stalled provider stream). - zzz_unclean_restart compares director.resumes against a snapshot taken before its kill instead of requiring it empty: crash cells on their own homes may legitimately resume. - the whole-run audit skips the reconnect replay only while #120444's gap is open.
teknium1
added a commit
that referenced
this pull request
Sep 23, 2026
A child-process GatewayRunner (own HERMES_HOME, SIGKILL + restart) runs the real agent against the scripted fake LLM provider and three instances of a FakePlatformAdapter that implements the gateway/platforms/base.py contract (4096/2000 limits, edits vs no edits, threads, streaming on/off). Only the transport is fake; its fsynced op journal is ground truth. Faults per op: timeout, ack lost, 429 retry_after, connection reset, too long, and park points for SIGKILL before/after the platform applies an op or between provider completion and the delivery ledger. Oracle per inbound: exactly one complete unmarked visible reply (extra copies only with the ledger's duplicate marker), visible text == persisted assistant text, user row persisted once, model turn run once, no stray or stale partial reply; plus a whole-run audit across chats. Scenarios: 13 faults x 3 platforms, /queue and interrupt follow-ups on a slow turn, parallel threads, re-delivered inbound ids (during/after a turn, after a runner reconnect), five crash points with restart catch-up, and an unclean restart with nothing in flight. Red-proven against the reverted stream fixes (#120315), a hand-revert of c961e5b (#91653), retrying timed-out sends, boot redelivery without the marker, double-sent finals, a dropped queued follow-up and disabled inbound dedup. Five strict xfails pin four live gaps (streaming ack-lost duplicate, dedup lost on reconnect, streaming crash after accept, and the unclean-restart recency fallback that re-answers every session active in the last 120 s). Eight cells stay red until #120315 lands (every streamed reply over the limit, and the first-send timeout) and carry xfail 'fixed by #120315': strict where the bug fires every run, non-strict where it depends on the stream/edit tick. The Director counts a resume note merged into the original tagged user message as a resume, not a second run of that inbound.
teknium1
added a commit
that referenced
this pull request
Sep 23, 2026
…y cells A strict xfail on a gap whose fix is an open PR turns main red the moment that fix merges (XPASS), and the non-strict ones guarded nothing. Each gap now has a probe (tests/e2e/core/delivery/_pending_fixes.py) that reproduces the defect's mechanism on the tree under test in a throwaway interpreter; expect_gap() applies the strict xfail only while the probe still reproduces it, so the cell becomes a plain test once the fix is in the tree, whatever the merge order. Covered: #120314, #119970 (soak), #120315, #120377, #120444 (C12) and #120450 (cell 5). Each probe was checked against every fix head: it flips on its own PR and on no other. C12 cells made deterministic (identical outcome on every run): - long_split streams the whole reply as one chunk; long_streamed and stream_timeout_first_send pace chunks so each lands in its own consumer tick. The five former coin-flip xfails are now two plain cells and three strict #120315 gap cells. - sent_ack_lost waits until the answer is persisted before the kill, so it pins the #120377 recovery; the streamed-before-persisted order is its own cell (stream_accepted_unpersisted, a strict live gap with a stalled provider stream). - zzz_unclean_restart compares director.resumes against a snapshot taken before its kill instead of requiring it empty: crash cells on their own homes may legitimately resume. - the whole-run audit skips the reconnect replay only while #120444's gap is open.
teknium1
added a commit
that referenced
this pull request
Sep 23, 2026
A child-process GatewayRunner (own HERMES_HOME, SIGKILL + restart) runs the real agent against the scripted fake LLM provider and three instances of a FakePlatformAdapter that implements the gateway/platforms/base.py contract (4096/2000 limits, edits vs no edits, threads, streaming on/off). Only the transport is fake; its fsynced op journal is ground truth. Faults per op: timeout, ack lost, 429 retry_after, connection reset, too long, and park points for SIGKILL before/after the platform applies an op or between provider completion and the delivery ledger. Oracle per inbound: exactly one complete unmarked visible reply (extra copies only with the ledger's duplicate marker), visible text == persisted assistant text, user row persisted once, model turn run once, no stray or stale partial reply; plus a whole-run audit across chats. Scenarios: 13 faults x 3 platforms, /queue and interrupt follow-ups on a slow turn, parallel threads, re-delivered inbound ids (during/after a turn, after a runner reconnect), five crash points with restart catch-up, and an unclean restart with nothing in flight. Red-proven against the reverted stream fixes (#120315), a hand-revert of c961e5b (#91653), retrying timed-out sends, boot redelivery without the marker, double-sent finals, a dropped queued follow-up and disabled inbound dedup. Five strict xfails pin four live gaps (streaming ack-lost duplicate, dedup lost on reconnect, streaming crash after accept, and the unclean-restart recency fallback that re-answers every session active in the last 120 s). Eight cells stay red until #120315 lands (every streamed reply over the limit, and the first-send timeout) and carry xfail 'fixed by #120315': strict where the bug fires every run, non-strict where it depends on the stream/edit tick. The Director counts a resume note merged into the original tagged user message as a resume, not a second run of that inbound.
teknium1
added a commit
that referenced
this pull request
Sep 23, 2026
…y cells A strict xfail on a gap whose fix is an open PR turns main red the moment that fix merges (XPASS), and the non-strict ones guarded nothing. Each gap now has a probe (tests/e2e/core/delivery/_pending_fixes.py) that reproduces the defect's mechanism on the tree under test in a throwaway interpreter; expect_gap() applies the strict xfail only while the probe still reproduces it, so the cell becomes a plain test once the fix is in the tree, whatever the merge order. Covered: #120314, #119970 (soak), #120315, #120377, #120444 (C12) and #120450 (cell 5). Each probe was checked against every fix head: it flips on its own PR and on no other. C12 cells made deterministic (identical outcome on every run): - long_split streams the whole reply as one chunk; long_streamed and stream_timeout_first_send pace chunks so each lands in its own consumer tick. The five former coin-flip xfails are now two plain cells and three strict #120315 gap cells. - sent_ack_lost waits until the answer is persisted before the kill, so it pins the #120377 recovery; the streamed-before-persisted order is its own cell (stream_accepted_unpersisted, a strict live gap with a stalled provider stream). - zzz_unclean_restart compares director.resumes against a snapshot taken before its kill instead of requiring it empty: crash cells on their own homes may legitimately resume. - the whole-run audit skips the reconnect replay only while #120444's gap is open.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Streamed replies that overflow the platform's message limit (Discord 2000, Telegram 4096) now reach the chat exactly once and complete: no duplicated block of text, no stray
(3/3)embedded mid-reply, and no truncated preview left above the final after a failed first send._seal_overflow_headsclears the message id, so a remainder still over the limit went out as a plain first send; the adapter split it and returned the LAST chunk's id, the consumer adopted that chunk as a preview holding the whole remainder, and the next seal overwrote it with the remainder's head and re-sent the rest (~4k chars duplicated on Discord). The remainder now goes back through_split_first_send.continue-ing, so a commentary or tool boundary drained together with the oversized burst still applies (thecontinuedropped the commentary and lost the boundary, gluing post-tool text onto the pre-tool preview). The initial-overflow gate had the samecontinueand now shares the path;continueremains only when a head send failed and the full text must stay for the fallback final.gateway/stream_consumer.py::_end_segment): after a split the message id is cleared, so the tail goes out as a first send in the tick that carries the tool boundary. When that send failed, the unseen-tail flush was skipped (it required a message id) and the segment reset cleared the buffer, losing ~1.1k chars. The flush now runs whenever the boundary update did not land, except for the__no_edit__sentinel; with a real message id the condition is unchanged, so no extra send. The flush also respects the egress guard, like the other new-message fallbacks. The same loss existed onorigin/mainfor a plain failed first send at a segment break.(n/n)in the live preview (ours,gateway/stream_consumer.py::_split_first_send): the tail chunk kept as the live preview still had the adapter's(n/n)suffix, so later deltas were appended after it and the user saw…word0574 wo (3/3)rd0575…. The indicator is now stripped from the kept tail.gateway/stream_consumer_transport.py::_send_or_edit): a failed first send (for example a Telegram timeout) disabled edits but left the message id unset, so the next tick sent another first send. That preview could never be edited, the final went out as a new message, and the truncated preview stayed on screen. With edits disabled, non-final first sends are now skipped, so only the complete final is delivered.Root cause: after a split or a failed send, the overflow and first-send paths lost track of which message is the live preview and kept publishing text that nothing would ever edit or seal again.
Live repro (exactly-once E2E through the real
GatewayRunnerwith fake Telegram/Discord platforms and a fake LLM provider:tests/e2e/core/delivery/test_messaging_exactly_once.py -k 'burst_overflow or long_split or stream_timeout_first_send', from the separate delivery E2E-suite PR and copied in temporarily, not part of this PR):origin/main8fb0fc6burst_overflow/long_spliton fk_tg + fk_dc givevisible reply != persisted assistant text … first diff at 4002/7182;stream_timeout_first_sendon fk_tg + fk_dc givea truncated/stale partial rendition … stayed visible: [Copy(text='<<A-fk_dc.stream_timeout_first', complete=False …)]long_split-fk_dcturns green (that cell is timing-dependent);burst_overflowstill fails at 4002 (the(n/n)tail) andstream_timeout_first_sendstill leaves the partial previewLive repro of the review finding (real
GatewayRunner, fake Telegram/Discord, fake LLM streaming a short lead, then an ~11k burst followed by aterminaltool call; the lead's first send is held 1.5 s so the burst and the tool boundary land in ONE consumer tick; scratch probe, not committed): on 075b22a a trace showsRESPLIT seg=Truefollowed by no segment end for that tick, so the boundary was dropped (the visible glue is masked on this path only because the tool round emits a secondstream_delta_callback(None)); on this head the same tick logsRESPLIT seg=True→END_SEGMENT acc=3255 visible=True(tg) /acc=1807(dc), and 6/6 runs show the post-tool answer as its own message with no pre-tool word repeated. Commentary has no second emit, so there the drop was user-visible (unit test above). Delivery E2E matrix re-run on this head (burst_overflow,long_split,long_streamed,stream_timeout_first_send,tool_then_answer,clean,slow_stream× 3 platforms): 21 passed.Tests (2 new invariants):
tests/gateway/test_stream_consumer_oversized_leftover.py::test_every_tail_line_reaches_the_channel_exactly_once[none|commentary|segment_break|tail_send_fails]: every tail line reaches the channel, never twice as a new message, and post-boundary text is never glued onto the pre-boundary preview. Merges fix(stream): a seal must not leave an oversized leftover on the non-final lane #114151'stest_no_tail_line_is_published_twice_across_new_messages(thenonecase; I dropped its lane-size twin, which pins the same re-split) and the earlier boundary test.noneandtail_send_failsare red on the merge-base;commentary/segment_breakwere red on 075b22a;tail_send_failswas red on f361b39 (24 tail lines lost). All 4 are green here.tests/gateway/test_stream_consumer.py::TestFinalResponseDeliveryGuard::test_failed_first_send_leaves_no_uneditable_partial_preview(red onorigin/main)Changed assertion in an existing test:
TestInitialOverflowRollingEdit::test_initial_overflow_uses_adapter_fence_aware_splitassertedsent_texts[-1].startswith(expected_chunks[-1]), meaning the live tail preview began with the full indicator-suffixed chunk. That is the stray-indicator bug above, written into the test. It now asserts that the tail starts with the chunk minus(n/n)and that no edit ever contains the indicator. The other assertions in that test are unchanged: fence balance, chunk count, sealed heads byte-identical, limit respected.Gates on 443bd7d:
scripts/run_tests.sh tests/gateway/test_stream_consumer*.pygave 9 files, 132 passed, 0 failed; the 24 other test files touching the consumer gave 522 passed, 1 failed (test_compression_failure_session_sync, which times out the same way on f361b39 without this change at load ~127). Earlier run ofscripts/run_tests.sh tests/gateway/gave 951 files, 8682 passed, 4 failed. All 4 are outside the stream consumer (approval timeout, compression session sync, API drain, turn-hold adoption), and those 4 files passed 43/43 when re-run alone, so the failures came from host load (~200).ruff checkon the touched files,scripts/check_no_tmp_literals.py gateway tests/gateway,git diff --checkandscripts/audit_pr_attribution.pyare all clean.Credit: @dasgltd (#114151, earliest fix for the seal re-split, salvaged here). Supersedes #114151.
Related, not superseded: #80990 (flushes the unsent tail at a tool boundary after a failed first send; a different path), #105200 (deletes an abandoned preview after the gateway sends the final; cleanup rather than prevention, on the flood-control path).
Refs #25349.
Infographic