Skip to content

fix(gateway): await in-flight streamed final send before duplicate-resend decision - #79592

Open
0xjarvisagent wants to merge 1 commit into
NousResearch:mainfrom
0xjarvisagent:fix/gateway-final-send-flood-race
Open

0xjarvisagent wants to merge 1 commit into
NousResearch:mainfrom
0xjarvisagent:fix/gateway-final-send-flood-race

Conversation

@0xjarvisagent

Copy link
Copy Markdown

Incident

2026-08-05 15:15 ET, Telegram group — a 4070-char final response was 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

Mechanism

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 deciding whether to send its own "final" response, then cancelled the consumer's task on timeout and moved on. Telegram's own retry_after can mandate waits far longer than 5s (37s here). 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 never triggered, 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 (gateway/stream_consumer.py) now tracks final_delivery_in_progress for the duration of the got_done finalize/fallback sequence (cleared in the outer finally, so every exit path — normal return, cancellation, exception — resets it).
  • gateway/run.py's new _await_stream_task_before_final_decision waits a short base timeout first (unchanged latency for the common, non-flood case), and only when the consumer reports an attempt still in flight extends the wait up to 90s before giving up. It uses asyncio.wait rather than asyncio.wait_for — wait_for cancels the awaited task on timeout and waits for it to fully unwind before raising TimeoutError, so by the time a caller's except clause runs, the consumer's own cancellation handling has already reset the in-flight flag; checking it there would always see False. asyncio.wait just reports done/pending on timeout without touching the task, so the flag reflects genuinely live state.
  • Never suppresses on "in flight" alone — only a completed send sets the flags the gateway checks. If the in-flight send ultimately fails, the gateway resend still happens (lost-message safety over duplicate-message risk).
  • Applied at all three call sites that raced this decision (the two in the main agent-run path, plus the proxy-relay path _run_agent_via_proxy).

Sibling flood paths checked

  • gateway/stream_consumer.py's _fallback_flood_retry_delay users (_send_fallback_final, _send_empty_fallback_final) are unaffected — they already bound their own retry sleep to _max_fallback_flood_retry_seconds (5s) and bail out early on long waits, deliberately leaving final delivery to the gateway. That's the correct pattern; this fix makes the gateway side honor it properly.
  • gateway/run.py's progress-edit flood path (~line 4085-4106, inside send_progress_messages) is unaffected — it's cancelled harmlessly in the finally block alongside progress_task; it only ever sends non-critical progress/status messages, never the final answer, so there's no duplicate-content risk there.

Tests

New: tests/gateway/test_final_send_flood_race.py::test_flood_controlled_finalize_delivers_final_response_once

Drives the real GatewayStreamConsumer and the real gateway helper directly. The mocked flood-controlled send() blocks on an asyncio.Event (send_release) that only the test controls, and signals send_started the instant it begins blocking — this replaced an earlier sleep-duration-vs-scaled-timeout version that was flaky under scheduling jitter. The test awaits send_started (no timing guess), lets the gateway's base timeout genuinely elapse via a one-sided sleep (nothing races it, since the send cannot complete on its own), then releases the send — resolution after that is driven by the event loop waking the extended wait on task completion, not more sleeping.

Before fix (git stash the fix, keep the test) — 3/3 runs, always fails with the duplicate:

AssertionError: final response visible in 2 separate messages (expected exactly 1): [...{'message_id': 'm-1'}, {'message_id': 'm-2'}]
assert 2 == 1
1 failed in 0.91s-1.27s

After fix — 5/5 runs, passes:

tests/gateway/test_final_send_flood_race.py::test_flood_controlled_finalize_delivers_final_response_once PASSED
1 passed in 0.83s-1.03s

Full suite run (/home/deploy/.hermes/hermes-agent/venv/bin/python -m pytest):

tests/gateway/test_final_send_flood_race.py .
tests/gateway/test_telegram_final_delivery.py ....
tests/gateway/test_duplicate_reply_suppression.py .........
tests/gateway/test_stale_finalize_suppression.py .........
tests/gateway/test_send_retry.py .............
35 passed in 2.90s

Also ran tests/gateway/test_stream_consumer.py (66 passed) to confirm no regression in the consumer's broader finalize/fallback logic.

🤖 Generated with Claude Code

…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 sweeper:risk-message-delivery Sweeper risk: may drop, duplicate, misroute, or suppress messages labels Aug 5, 2026
@michaelversluis

Copy link
Copy Markdown

I reproduced a production case on Hermes Agent v0.20.2 that outlives this PR's fixed 90-second maximum, so the current patch still duplicates the final response for the reported class.

Sanitized Telegram gateway timeline (2026-08-17, one inbound, one Codex call):

06:11:29.511 inbound message
06:11:34.728 Telegram flood control, waiting 212.0s
06:11:39.671 live/interim send: RetryAfter 207.0s
06:11:44.625 response ready, api_calls=1, response=1304 chars
06:11:44.641 normal gateway final-send path starts
06:11:44.701 normal final send: RetryAfter 202.0s
06:15:06.777 normal final delivery obligation marked delivered

Telegram showed two separate messages: the first was the partial/live delivery, the second the full normal final response. The screenshot and ledger confirm separate sends rather than one edited message.

I also tested this PR's own regression harness at head a6cf7ef502a2a31d4ebc352fec1dabfbe11ea5ef. I extended _run_flood_race() only in the test so the mocked in-flight delivery resolves after _STREAM_FINAL_WAIT_MAX_TIMEOUT_SECONDS, then asserted one visible final message:

adapter, already_confirmed = await _run_flood_race(
    monkeypatch,
    base_timeout=0.05,
    max_timeout=0.10,
    release_delay=0.15,
)
assert len(adapter.visible_full_response_messages()) == 1
assert already_confirmed is True

Canonical run:

HERMES_TEST_FILE_RETRIES=0 scripts/run_tests.sh \
  tests/gateway/test_final_send_flood_race.py \
  -k retry_after_longer_than_wait_ceiling -v --tb=short

FAILED: final response visible in 2 separate messages after the wait ceiling expired
all sends=[... 'm-1', ... 'm-2']

This is the same deterministic mechanism as the PR's existing test; the only changed variable is that delivery resolves after the total wait ceiling. In production the observed RetryAfter was 202–212 seconds, while the patch stops waiting after 90 seconds, cancels the consumer, and lets the outer path send again.

Requested adjustment before merge:

  • Do not let a fixed total ceiling turn a still-confirmed final_delivery_in_progress attempt back into a duplicate-send decision.
  • Either derive the wait from the adapter's active retry deadline/contract, or keep waiting in bounded slices while the final attempt remains in flight and let the adapter's own bounded retry/timeout settle it.
  • Add the beyond-ceiling regression above. Merely increasing 90 to another constant leaves the class open (production reports already include much larger Telegram RetryAfter values).

The PR's in-flight state tracking and asyncio.wait direction are good; the remaining issue is specifically the fixed total ceiling.

@michaelversluis

Copy link
Copy Markdown

I opened #88166 as a strict superset after reproducing a production RetryAfter of 202–212 seconds that outlives this PR’s 90-second total ceiling. Your original commit is preserved verbatim as the first commit, with authorship intact. The follow-up commit adds the beyond-ceiling RED test, removes the arbitrary total ceiling once final delivery is actively in flight, and verifies failure still produces exactly one normal fallback plus non-final stalls still cancel. Focused result: 102 passed, 1 skipped; two independent reviews found no blockers. Thank you for the original in-flight state design.

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

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