Skip to content
Merged
1 change: 1 addition & 0 deletions contributors/emails/hermes@dasg.ltd
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
dasgltd
32 changes: 21 additions & 11 deletions gateway/stream_consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -574,13 +574,23 @@ async def run(self) -> None:
if self._should_edit(tick) and (
self._accumulated or (self._use_native_streaming and self._tool_progress_active)
):
# Seal first: it clears the message id, so a remainder still over the limit
# is split again below. A plain first send would let the adapter split it and
# adopt only the LAST chunk as the preview; the next seal then overwrites that
# chunk with the head of the whole remainder (duplicated + lost text, #25349).
await self._seal_overflow_heads()
# Overflow split. Native streaming bypasses this: the adapter
# truncates against the stream protocol's own limit.
if not self._use_native_streaming and self._first_send_overflows():
if await self._split_first_send(tick):
return
continue
await self._seal_overflow_heads()
if self._first_send_overflows():
# A head send failed: keep the full text for the fallback final, and
# skip the boundary reset below that would clear it.
self._signal_flush(tick.flush_event)
continue
# The split tail goes out now, so a commentary or tool boundary drained in
# this tick still lands after it instead of being dropped.
await self._push_update(tick)

if tick.got_done:
Expand Down Expand Up @@ -755,7 +765,11 @@ async def _split_first_send(self, tick: "_Tick") -> bool:
reply_to = new_id

if heads_delivered:
self._accumulated = chunks[-1]
# truncate_message suffixes multi-chunk output with " (n/n)"; the tail is the LIVE
# preview later deltas extend, so a kept indicator ends up embedded mid-reply.
tail = chunks[-1]
indicator = f" ({len(chunks)}/{len(chunks)})"
self._accumulated = tail[: -len(indicator)] if tail.endswith(indicator) else tail
# Flag BEFORE the tail send: fresh-final replaces every tracked preview
# with one message, which is only valid while the active message holds
# the whole answer — deleting sealed heads drops delivered text.
Expand All @@ -778,11 +792,6 @@ async def _split_first_send(self, tick: "_Tick") -> bool:
if tick.got_segment_break:
self._fallback_final_send = False
self._fallback_prefix = ""
if not self._accumulated:
return False
# Early `continue` skips the bottom-of-loop flush signal.
if tick.got_flush:
self._signal_flush(tick.flush_event)
return False

def _overflows(self) -> bool:
Expand Down Expand Up @@ -917,9 +926,10 @@ async def _end_segment(self, tick: "_Tick") -> None:
continuation goes out once via _send_fallback_final."""
if self._cumulative_transport():
return
# If the segment-break edit didn't land (flood control / fallback mode),
# _accumulated holds unseen pre-boundary text — flush it before the reset.
if (self._accumulated and not tick.update_visible and self._message_id
# If the segment-break edit or send didn't land (flood control / fallback mode, or a
# failed first send of a split tail with no message id yet), _accumulated holds unseen
# pre-boundary text — flush it before the reset clears it.
if (self._accumulated and not tick.update_visible
and self._message_id != "__no_edit__"):
await self._flush_segment_tail_on_edit_failure()
self._reset_segment_state(preserve_no_edit=True)
Expand Down
2 changes: 2 additions & 0 deletions gateway/stream_consumer_fallback.py
Original file line number Diff line number Diff line change
Expand Up @@ -305,6 +305,8 @@ def _is_flood_error(self, result) -> bool:
async def _flush_segment_tail_on_edit_failure(self) -> None:
"""Before a segment reset, send the unseen tail as a new message (and best-effort
strip the stuck cursor from the partial)."""
if getattr(self, "_egress_declined", False):
return # a new message is exactly the re-addressing the egress guard refused
if not self._fallback_final_send:
await self._try_strip_cursor()
visible = self._fallback_prefix or self._visible_prefix()
Expand Down
4 changes: 4 additions & 0 deletions gateway/stream_consumer_transport.py
Original file line number Diff line number Diff line change
Expand Up @@ -349,6 +349,10 @@ async def _send_or_edit(
self._last_edit_overflowed = False
try:
if self._message_id is None:
if not self._edit_supported and not finalize:
# A failed send disabled edits: a preview sent now could never be updated and
# would stay on screen truncated next to the final reply. Send only the final.
return False
return await self._first_send(text, finalize=finalize)
if not self._edit_supported:
return False # edits unsupported; fallback path sends the final
Expand Down
39 changes: 38 additions & 1 deletion tests/gateway/test_stream_consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -581,6 +581,38 @@ async def test_split_overflow_failed_send_does_not_mark_final_sent(self):
)


@pytest.mark.asyncio
async def test_failed_first_send_leaves_no_uneditable_partial_preview(self):
"""A failed first send disables edits: a preview sent after it could never be updated and
would stay on screen truncated next to the final. Only the complete reply reaches the chat."""
delivered = []
results = iter([SimpleNamespace(success=False, error="timeout")])

async def send(**kw):
result = next(results, None) or SimpleNamespace(success=True, message_id=f"m{len(delivered) + 1}")
if result.success:
delivered.append(kw["content"])
return result

adapter = MagicMock()
adapter.send = AsyncMock(side_effect=send)
adapter.edit_message = AsyncMock(return_value=SimpleNamespace(success=True))
adapter.MAX_MESSAGE_LENGTH = 4096
config = StreamConsumerConfig(edit_interval=0.01, buffer_threshold=5, cursor="")
consumer = GatewayStreamConsumer(adapter, "chat_123", config)

consumer.on_delta("preview never landed ")
task = asyncio.create_task(consumer.run())
await asyncio.sleep(0.08)
consumer.on_delta("and more streamed text ")
await asyncio.sleep(0.08)
consumer.on_delta("then the end.")
consumer.finish()
await asyncio.wait_for(task, timeout=10)

assert delivered == ["preview never landed and more streamed text then the end."]


class TestFinalContentDeliveredGuard:
"""Regression coverage for #25010 — _final_content_delivered must only be
set when the final response is actually confirmed delivered to the user,
Expand Down Expand Up @@ -736,7 +768,12 @@ async def test_initial_overflow_uses_adapter_fence_aware_split(self):
assert all(text.count("```") % 2 == 0 for text in sent_texts + edited_texts)
assert len(sent_texts) == len(expected_chunks)
assert sent_texts[:-1] == expected_chunks[:-1]
assert sent_texts[-1].startswith(expected_chunks[-1])
# The tail is the live preview later deltas extend: it drops the " (n/n)" indicator, which
# would otherwise end up embedded mid-reply ("``` (2/2)\nTail after ...").
indicator = f" ({len(expected_chunks)}/{len(expected_chunks)})"
assert expected_chunks[-1].endswith(indicator)
assert sent_texts[-1].startswith(expected_chunks[-1][: -len(indicator)])
assert not any(indicator in text for text in edited_texts)
assert any("Tail after the fenced stream." in text for text in edited_texts)
assert all(utf16_len(text) <= safe_limit for text in sent_texts)

Expand Down
129 changes: 129 additions & 0 deletions tests/gateway/test_stream_consumer_oversized_leftover.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,129 @@
"""A seal must not leave an oversized leftover on the non-final send lane.

``_seal_overflow_heads`` clears the edit target MID-ITERATION, so the overflow
gate evaluated at the top of ``run()`` is already stale by the time the leftover
is pushed. When that leftover still exceeds one message, it reaches
``_push_update`` -> ``_first_send`` with no message to edit, on a NON-final tick.
``adapter.send()`` then chunks and caps the payload itself, numbering the pieces
``(i/n)``; the turn-final lane afterwards publishes the SAME text with its own
denominator.

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, which is what places the duplicate inside the
consumer rather than in the gateway's delivery ledger.

Contract pinned here: on a non-final lane the CONSUMER owns chunking, so no line
of the answer reaches the channel twice as a new message, and none is lost when a
boundary or a failed tail send lands in the sealing tick.
"""

import asyncio
from types import SimpleNamespace
from unittest.mock import AsyncMock

import pytest

from gateway.stream_consumer import GatewayStreamConsumer, StreamConsumerConfig

# The consumer floors its own length budget at 500 chars, so a toy adapter limit
# below that floor would make it legitimately emit chunks larger than the limit -
# a measurement artifact, not a leak. Use a realistic limit instead.
LIMIT = 2000
HEAD = "head line\n" * 60
TAIL = "".join(f"tail line {i:03d} filler filler filler filler filler\n" for i in range(240))


def _make_plain_adapter():
"""No rich tail at all: ``preferred_final_chunks`` declines, so an oversized
payload has no atomic-chunk excuse and must be split by the consumer."""
from gateway.platforms.base import BasePlatformAdapter

cls = type(
"PlainAdapter",
(BasePlatformAdapter,),
{
"MAX_MESSAGE_LENGTH": LIMIT,
"preferred_final_chunks": lambda self, text, budget: None,
"preferred_final_split_index": lambda self, text, budget: None,
},
)
cls.__abstractmethods__ = frozenset()
adapter = cls.__new__(cls)
adapter._typing_paused = set()
adapter._fatal_error_message = None
return adapter


def _wire(adapter, fail_first_send_of=None):
"""Record delivered sends/edits. ``fail_first_send_of``: the first send whose content
contains that marker fails (success=False) and is not delivered."""
sends, edits = [], []
failed = []

async def fake_send(**kw):
content = kw.get("content", "")
if fail_first_send_of and not failed and fail_first_send_of in content:
failed.append(content)
return SimpleNamespace(success=False, error="timeout")
sends.append((content, (kw.get("metadata") or {}).get("notify")))
return SimpleNamespace(success=True, message_id=f"m{len(sends)}")

async def fake_edit(**kw):
edits.append(kw.get("content", ""))
return SimpleNamespace(success=True, message_id=kw.get("message_id", "m1"))

adapter.send = AsyncMock(side_effect=fake_send)
adapter.edit_message = AsyncMock(side_effect=fake_edit)
return sends, edits, failed


@pytest.mark.asyncio
@pytest.mark.parametrize("boundary", ["none", "commentary", "segment_break", "tail_send_fails"])
async def test_every_tail_line_reaches_the_channel_exactly_once(boundary):
"""Live shape: a streamed head is already on screen when a long tail pushes the buffer far
past the limit, so the seal leaves a leftover that STILL overflows. A commentary or tool
boundary may be drained in that same tick, and the tail send carrying the boundary may fail.
Every tail line must reach the channel, never twice as a new message, and post-boundary
text must start a new message instead of being glued onto the pre-boundary preview."""
adapter = _make_plain_adapter()
sends, edits, failed = _wire(
adapter, fail_first_send_of="tail line 239" if boundary == "tail_send_fails" else None)
config = StreamConsumerConfig(edit_interval=0.01, buffer_threshold=5, cursor="")
consumer = GatewayStreamConsumer(adapter, "chat_plain", config)
consumer.on_delta(HEAD)
task = asyncio.create_task(consumer.run())
await asyncio.sleep(0.06)
# Enqueued synchronously, so one drain folds TAIL and stops at the boundary.
consumer.on_delta(TAIL)
if boundary == "commentary":
consumer.on_commentary("COMMENTARY-MARKER")
if boundary != "none":
consumer.on_segment_break()
consumer.on_delta("POST-TOOL-MARKER")
await asyncio.sleep(0.12)
consumer.finish()
# asyncio.wait_for, never a bare await: a hold branch that fails to yield
# spins hot and would hang the suite instead of failing it.
await asyncio.wait_for(task, timeout=10)

if boundary == "tail_send_fails":
assert failed, "the harness never failed the tail send"
joined = "".join(c for c, _ in sends)
repeated = [f"tail line {i:03d}" for i in range(240) if joined.count(f"tail line {i:03d}") > 1]
assert not repeated, (
f"{len(repeated)} tail line(s) reached the channel more than once as NEW "
f"messages (first: {repeated[:3]}) - that is the interleaved duplicate"
)
texts = [c for c, _ in sends] + edits
lost = [f"tail line {i:03d}" for i in range(240)
if not any(f"tail line {i:03d}" in t for t in texts)]
assert not lost, f"{len(lost)} tail line(s) never reached the channel (first: {lost[:3]})"
if boundary == "commentary":
assert any("COMMENTARY-MARKER" in t for t in texts), "commentary was dropped"
if boundary != "none":
assert any("POST-TOOL-MARKER" in t for t in texts)
glued = [t for t in texts if "POST-TOOL-MARKER" in t and "tail line" in t]
assert not glued, "post-boundary text was glued onto the pre-boundary preview"
Loading