Skip to content

fix(gateway): keep in-flight final delivery single-writer - #88166

Open
michaelversluis wants to merge 2 commits into
NousResearch:mainfrom
michaelversluis:fix/gateway-inflight-final-no-ceiling
Open

michaelversluis wants to merge 2 commits into
NousResearch:mainfrom
michaelversluis:fix/gateway-inflight-final-no-ceiling

Conversation

@michaelversluis

Copy link
Copy Markdown

Summary

Prevent duplicate gateway replies when a streamed final delivery is still in flight under a platform RetryAfter longer than the gateway's own wait ceiling.

This preserves and builds on the implementation from #79592 (original commit/authorship retained), then closes the remaining long-RetryAfter gap found in a live Telegram incident.

Root cause

The stream consumer and the normal gateway final-send path can both become writers for the same final response:

  1. the stream consumer starts its final send/edit;
  2. Telegram returns RetryAfter and the adapter keeps that attempt alive;
  3. the gateway waits for a fixed total duration;
  4. if the platform wait outlives that duration, the gateway cancels/abandons the consumer and starts its own final send;
  5. the already-issued first request can still land, producing two separate messages.

#79592 correctly introduced final_delivery_in_progress and waits past the original five seconds, but capped the total at 90 seconds. A production case on v0.20.2 returned 202–212 second waits, reproducing the duplicate beyond that cap:

06:11:29 one inbound message
06:11:34 streamed final path: RetryAfter 212s
06:11:44 one Codex call completes; normal final-send path starts
06:11:44 normal final path: RetryAfter 202s
06:15:06 both separately queued messages land

Changes

  • Preserve fix(gateway): await in-flight streamed final send before duplicate-resend decision #79592's explicit final_delivery_in_progress state and all three gateway call-site conversions.
  • Make an active final-delivery attempt the single writer until its platform adapter settles.
  • Remove the arbitrary gateway total ceiling; adapters retain ownership of retry counts and request timeouts.
  • Keep the short cancellation path unchanged when no final-delivery attempt is active.
  • Add deterministic coverage for:
    • an in-flight final that resolves after the former total ceiling;
    • an in-flight final that fails, followed by exactly one normal fallback;
    • a stalled non-final consumer that is still cancelled normally.

TDD evidence

The beyond-ceiling regression was run first against the #79592 base and failed with two visible messages:

FAILED test_retry_after_longer_than_wait_ceiling_still_delivers_once
AssertionError: final response visible in 2 separate messages
all sends=[... 'm-1', ... 'm-2']

After the single-writer change:

HERMES_TEST_FILE_RETRIES=0 scripts/run_tests.sh \
  tests/gateway/test_final_send_flood_race.py \
  tests/gateway/test_stream_consumer.py \
  tests/gateway/test_run_cleanup_progress.py \
  tests/gateway/test_run_progress_interrupt.py \
  tests/gateway/test_run_progress_topics.py -q

102 passed, 1 skipped

Additional checks:

ruff check gateway/run.py gateway/stream_consumer.py tests/gateway/test_final_send_flood_race.py
All checks passed!

git diff --check
clean

A local full tests/gateway/ run was attempted but the minimal dev environment lacks optional platform extras (aiohttp, Matrix, etc.). uv sync --all-extras is blocked on macOS by the existing python-olm==3.2.16 native build (cmake/gmake plus a const-qualified C++ compile error). The affected production and regression surfaces above are green; upstream CI should exercise the complete platform matrix.

Scope and attribution

Related: #74248, #76494, #79592.


Nederlandse samenvatting: bij Telegram-wachttijden boven 90 seconden kon de bestaande fix alsnog twee eindberichten versturen. Deze versie houdt één actieve eindlevering als enige schrijver, maar laat na een echte mislukking nog steeds precies één normale fallback toe.

JARVIS and others added 2 commits August 17, 2026 07:16
…send decision

Incident (2026-08-05 15:15 ET, Telegram group, 4070-char final response
delivered twice):

    15:15:10 MarkdownV2 edit failed, falling back to plain text: Flood
             control exceeded. Retry in 37 seconds
    15:15:10 Telegram flood control, waiting 37.0s
    15:15:10 Telegram flood control on send (attempt 1/3), retrying in
             37.0s: Flood control exceeded
    15:15:15 gateway.platforms.base: Sending response (4070 chars) to
             -1003725014629
    15:15:15 Telegram flood control on send (attempt 1/3), retrying in
             32.0s

gateway/run.py waited a flat 5 seconds (asyncio.wait_for(stream_task,
timeout=5.0)) for the stream consumer's finalize/fallback send before the
duplicate-send decision, then cancelled and moved on. Telegram's own
retry_after can mandate waits far longer than 5s, so the gateway routinely
gave up while the send was still in flight: final_response_sent /
final_content_delivered were still False (the attempt hadn't resolved
yet), suppression didn't trigger, and the gateway sent its own "final"
response. If the abandoned send still reached the platform after being
locally cancelled, the user got the answer twice.

Fix: GatewayStreamConsumer now tracks final_delivery_in_progress for the
duration of the got_done finalize/fallback sequence. gateway/run.py's new
_await_stream_task_before_final_decision waits a short base timeout first
(unchanged latency for the common case), and only when the consumer
reports an attempt still in flight extends the wait up to 90s before
giving up. Uses asyncio.wait (not wait_for) so checking the flag doesn't
race the task's own cancellation-triggered unwind. Applied at all three
call sites that raced this decision, including the proxy-relay path.

Regression test (tests/gateway/test_final_send_flood_race.py) drives the
real GatewayStreamConsumer and gateway helper directly, coordinating the
mocked flood-controlled send via asyncio.Event rather than sleep-timing
comparisons, so the race is deterministic in both directions. Verified
failing on unpatched main (2 duplicate messages, 3/3 runs) and passing
with the fix (1 message, 5/5 runs).

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
@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 area/streaming Streaming responses: gateway delivery, provider wire sweeper:risk-message-delivery Sweeper risk: may drop, duplicate, misroute, or suppress messages labels Aug 17, 2026
@Enough1122

Copy link
Copy Markdown

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

Reviewed the diff. Correct fix for a duplicate-send race with a genuinely subtle async detail handled properly: the old flat 5s wait_for cancelled the stream task on timeout and waited through its unwinding before raising, so by the time the caller checked final_delivery_in_progress the in-flight Telegram flood-control send had already been torn down (guaranteed False) and the gateway sent its own duplicate "final" response alongside the abandoned-but-delivered one. Switching to asyncio.wait (reports pending without touching the task), treating an active final-delivery attempt as the single writer awaited to completion under the adapter's own retry policy, and never reading "in flight" as "delivered" is exactly right. All three legacy wait sites are migrated, and the docstring explains the wait-vs-wait_for choice so nobody simplifies it back.

Nit: the no-timeout branch means an adapter with an unbounded retry loop would hold the turn forever; the comment leans on "adapters own bounded retries" - worth one assert/comment at the adapter layer stating retries are bounded by policy.

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 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