fix(gateway): gate sub-floor preamble finalizes on tool boundaries
A segment-break finalize never carries the streaming cursor (the cursor is only appended to mid-stream frames), so the sub-floor standalone-message guard's cursor-membership test let every 1-2 token preamble land as its own durable message at each tool boundary. Each landing fired on_new_message, resetting the gateway's tool-progress anchor and fragmenting accumulated tool progress into one persistent message per tool on draft-streaming platforms (#99026). Extend the guard to cover segment-break finalizes (finalize and not is_turn_final). Turn finals stay exempt so a complete short answer is never swallowed. Update the existing per-segment callback test to use segment text above the floor, and add draft-mode regression coverage: - short preamble rounds: no durable sends, no progress resets, drafts and the gateway final path unaffected - full preamble rounds: still one durable send + one reset per segment (#17280 chronological-order contract preserved) - turn-final short answer: always delivered (cherry picked from commit 5ac896e1f22e6b2164544287c342c7fc7ec2c373)
This commit is contained in:
@@ -331,8 +331,12 @@ class StreamTransportMixin:
|
||||
return True # cursor-only / whitespace-only update
|
||||
# Don't open a new message for 1-2 tokens + cursor (rapid tool-calling): if
|
||||
# the cursor-strip edit is then rate-limited, "X ▉" stays forever.
|
||||
if (self._message_id is None and self.cfg.cursor and self.cfg.cursor in text
|
||||
and len(visible_stripped) < self._MIN_NEW_MSG_CHARS):
|
||||
# A segment-break finalize never carries the cursor, so gate it too or a 1-2 token
|
||||
# preamble lands durably at every tool boundary and resets the progress anchor
|
||||
# (#99026). Turn finals are exempt: a short complete answer must be delivered.
|
||||
preamble_finalize = finalize and not is_turn_final
|
||||
if (self._message_id is None and len(visible_stripped) < self._MIN_NEW_MSG_CHARS
|
||||
and (preamble_finalize or (self.cfg.cursor and self.cfg.cursor in text))):
|
||||
return True # too short for a standalone message — accumulate more
|
||||
|
||||
# A failed native/draft transport disables itself and falls through so the
|
||||
|
||||
@@ -1145,15 +1145,17 @@ class TestOnNewMessageCallback:
|
||||
on_new_message=lambda: events.append("reset"),
|
||||
)
|
||||
|
||||
consumer.on_delta("A")
|
||||
consumer.on_delta("First preamble")
|
||||
consumer.on_delta(None)
|
||||
consumer.on_delta("B")
|
||||
consumer.on_delta("Second preamble")
|
||||
consumer.on_delta(None)
|
||||
consumer.on_delta("C")
|
||||
consumer.on_delta("Final answer")
|
||||
consumer.finish()
|
||||
await consumer.run()
|
||||
|
||||
# Three content bubbles ⇒ three reset notifications
|
||||
# Three content bubbles ⇒ three reset notifications. (Segment text
|
||||
# must clear the sub-floor standalone-message guard — a 1-2 token
|
||||
# preamble is intentionally held back now, see #99026.)
|
||||
assert events == ["reset", "reset", "reset"]
|
||||
|
||||
|
||||
|
||||
@@ -337,6 +337,103 @@ def _make_fresh_final_adapter():
|
||||
return adapter
|
||||
|
||||
|
||||
class TestSubFloorPreambleToolBoundaries:
|
||||
"""Regression for #99026: a 1-2 token preamble finalized at a tool
|
||||
boundary must NOT land as its own durable message.
|
||||
|
||||
The sub-floor standalone-message guard keyed on ``cursor in text``,
|
||||
but a segment-break finalize's text never carries the cursor (the
|
||||
cursor is only appended to mid-stream frames), so every short
|
||||
preamble became a real sendMessage, fired on_new_message, and reset
|
||||
the gateway's tool-progress anchor — fragmenting accumulated progress
|
||||
into one bubble per tool on draft-streaming platforms (Telegram).
|
||||
"""
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_short_preamble_rounds_keep_progress_anchor(self):
|
||||
adapter = _make_draft_capable_adapter()
|
||||
cfg = StreamConsumerConfig(
|
||||
transport="auto", chat_type="dm",
|
||||
edit_interval=0.01, buffer_threshold=5, cursor="▉",
|
||||
)
|
||||
consumer = GatewayStreamConsumer(adapter, "12345", cfg)
|
||||
resets: list[str] = []
|
||||
consumer._on_new_message = lambda: resets.append("reset")
|
||||
|
||||
task = asyncio.create_task(consumer.run())
|
||||
for _ in range(3):
|
||||
consumer.on_delta("Ok")
|
||||
await asyncio.sleep(0.05)
|
||||
consumer.on_delta(None) # tool boundary
|
||||
await asyncio.sleep(0.05)
|
||||
consumer.finish("Done.")
|
||||
await task
|
||||
|
||||
# No sub-floor preamble landed as a standalone durable message, so
|
||||
# the tool-progress anchor was never reset between tool rounds.
|
||||
assert resets == []
|
||||
sends = [c.kwargs.get("content") for c in adapter.send.await_args_list]
|
||||
assert sends == [], sends
|
||||
# The sub-floor preamble was held back at every stage: the mid-stream
|
||||
# frame guard suppressed the "Ok ▉" draft frame (legacy behavior),
|
||||
# and the segment-break finalize no longer turns it into a real
|
||||
# message either.
|
||||
assert all(len((c["content"] or "").replace("▉", "").strip()) >= 4
|
||||
for c in adapter.draft_calls) or not adapter.draft_calls
|
||||
# The consumer never claimed final delivery (drafts don't set
|
||||
# already_sent), so the gateway's final-send path owns "Done.".
|
||||
assert consumer.final_response_sent is False
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_full_preamble_rounds_still_reset_once_each(self):
|
||||
"""Counterpart: a preamble long enough to stand alone SHOULD land as
|
||||
a durable message and reset the progress anchor below it (the
|
||||
#17280 chronological-order contract)."""
|
||||
adapter = _make_draft_capable_adapter()
|
||||
cfg = StreamConsumerConfig(
|
||||
transport="auto", chat_type="dm",
|
||||
edit_interval=0.01, buffer_threshold=5, cursor="▉",
|
||||
)
|
||||
consumer = GatewayStreamConsumer(adapter, "12345", cfg)
|
||||
resets: list[str] = []
|
||||
consumer._on_new_message = lambda: resets.append("reset")
|
||||
|
||||
task = asyncio.create_task(consumer.run())
|
||||
for i in range(3):
|
||||
consumer.on_delta(f"Round {i} preamble text")
|
||||
await asyncio.sleep(0.05)
|
||||
consumer.on_delta(None) # tool boundary
|
||||
await asyncio.sleep(0.05)
|
||||
consumer.finish("Done.")
|
||||
await task
|
||||
|
||||
assert len(resets) == 3, resets
|
||||
sends = [c.kwargs.get("content") for c in adapter.send.await_args_list]
|
||||
assert sends == [
|
||||
"Round 0 preamble text",
|
||||
"Round 1 preamble text",
|
||||
"Round 2 preamble text",
|
||||
], sends
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_turn_final_short_answer_is_never_swallowed(self):
|
||||
"""The sub-floor guard must not eat a turn-final answer, or the
|
||||
gateway's final-send suppression would drop it entirely."""
|
||||
adapter = _make_draft_capable_adapter()
|
||||
cfg = StreamConsumerConfig(
|
||||
transport="auto", chat_type="dm",
|
||||
edit_interval=0.01, buffer_threshold=5, cursor="▉",
|
||||
)
|
||||
consumer = GatewayStreamConsumer(adapter, "12345", cfg)
|
||||
|
||||
consumer.on_delta("Ok")
|
||||
consumer.finish()
|
||||
await consumer.run()
|
||||
|
||||
sends = [c.kwargs.get("content") for c in adapter.send.await_args_list]
|
||||
assert any("Ok" in (s or "") for s in sends), sends
|
||||
|
||||
|
||||
class TestAdapterPrefersFreshFinal:
|
||||
"""An adapter whose send path is richer than its edit path (e.g. Telegram
|
||||
rich messages) finalizes a streamed reply by sending a fresh final message
|
||||
|
||||
Reference in New Issue
Block a user