Skip to content

Keep subscribe_with_handler running until the caller closes it - #1235

Draft
CasperGN wants to merge 5 commits into
dapr:mainfrom
CasperGN:fix/subscribe-with-handler-loop
Draft

CasperGN wants to merge 5 commits into
dapr:mainfrom
CasperGN:fix/subscribe-with-handler-loop

Conversation

@CasperGN

@CasperGN CasperGN commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor

Description

Based on #1231. Merge it after #1231. Until #1231 merges, this diff also shows #1231's commits; the changes of its own are the last two commits (0844c95, 942c58d). It reuses _closed, _close_stream(), the per-stream send queue and the fake sidecar hooks from #1231.

subscribe_with_handler could stop delivering messages without any error. This PR fixes it in both clients.

  • Sync, retry path: after a failed reconnect the loop slept 5 seconds and then exited. It now retries every SUBSCRIPTION_RECONNECT_BACKOFF_SECONDS (5) until the close function is called. The close function also ends a loop that is waiting to retry, so it returns right away.
  • Both, pause after a stream error: when next_message() gives up (after its own bounded reconnects from Return the next message after a streaming subscription reconnects #1231) or raises a non-retryable error, the loop now waits SUBSCRIPTION_RECONNECT_BACKOFF_SECONDS before reconnecting. Before, a stream that kept failing right after connecting was reconnected with no pause, forever. Closing ends the wait at once.
  • Both, handler errors: an exception raised by the handler was treated as a stream error. The stream was reconnected and the message was delivered again. Now the exception is logged with its traceback, the message gets a retry response, and the next message is read. The stream stays open.
  • Async, errors ended the task: any error other than StreamInactiveError ended the handler task, and it only showed up later as "Task exception was never retrieved". The task now follows the same rules as the sync loop. It reconnects on stream errors, retries after a failed reconnect, and stops only after the caller has closed the subscription. grpc.aio raises asyncio.CancelledError when reading a stream that close() cancelled. That counts as a normal exit only after a close, and is re-raised otherwise.
  • Async, task reference: the SDK now keeps a reference to the handler task, so it can't be garbage-collected while it runs.
  • Async, close waits: the close function now waits up to SUBSCRIPTION_CLOSE_TIMEOUT_SECONDS for the handler task to finish. It skips the wait when it is called from inside the handler.

Tests:

  • The async test_subscribe_topic_with_handler is restored.
  • New tests for both clients cover:
    • a reconnect that fails and then succeeds
    • a handler that raises once
    • a stream that fails with CANCELLED or NOT_FOUND
    • the pause after a stream error
    • closing while the loop is waiting
  • The async tests also check that the task survives garbage collection.
  • Return the next message after a streaming subscription reconnects #1231's test_subscribe_topic_with_handler_currently_reconnects_after_handler_error documented the old handler-error behaviour. It is replaced by a test that expects a retry response.

Issue reference

Please reference the issue this PR will close: Fixes #1233

Checklist

  • Code compiles correctly
  • Created/updated tests
  • Extended the documentation

RELEASE NOTE: FIX subscribe_with_handler keeps retrying while the sidecar is unavailable, pauses between reconnects, does not reconnect on handler errors, and the async version no longer stops silently.

🤖 Generated with Claude Code

CasperGN and others added 5 commits September 24, 2026 13:24
After a transient stream error (UNAVAILABLE, UNKNOWN, INTERNAL),
next_message() reconnected and then returned None instead of reading
from the new stream. It now reads from the new stream, and raises after
MAX_RECONNECT_ATTEMPTS reconnects in one call. The async client does the
same and now also retries INTERNAL.

The request iterator of a failed stream kept waiting on the shared send
queue, so after a reconnect it took the next ack and sent it to the dead
stream. Each stream now has its own send queue, and closing a stream
ends its request iterator.

close() now stays closed: a reconnect already in progress no longer
reopens the stream. The close function from subscribe_with_handler waits
up to SUBSCRIPTION_CLOSE_TIMEOUT_SECONDS for the handler thread, so
test_subscribe_topic_with_handler no longer leaves that thread running
and its coverage no longer changes between runs.

Fixes dapr#1230

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Signed-off-by: Casper Nielsen <casper@diagrid.io>
A stream the server ends cleanly (StopIteration in the sync client,
grpc.aio.EOF in the async one) now goes through the same bounded
reconnect as UNAVAILABLE, UNKNOWN and INTERNAL, instead of failing with
"Error while fetching message". The async docstring no longer promises
None at end of stream.

Reconnects inside one next_message() call now wait a short jittered
exponential backoff between attempts. close() ends that wait at once.
After MAX_RECONNECT_ATTEMPTS the clients raise the new
StreamReconnectError, a subclass of Exception.

The async start() now raises StreamInactiveError when close() runs
during the initial read, as the sync one does. A cancellation while the
subscription is still open is re-raised.

Each message remembers the send queue of the stream it came from, and
respond() drops a response to a message from a replaced stream, with a
debug log. The runtime only logs an error for such an ack, but it never
has to see it now.

The sync close function logs a warning when the handler thread is still
running after SUBSCRIPTION_CLOSE_TIMEOUT_SECONDS. The async subscription
logs through a module logger instead of print(). The unreachable
queue.Empty handler is gone.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Signed-off-by: Casper Nielsen <casper@diagrid.io>
Resolve conflicts with dapr#1234 in the sync subscription. Keep this branch's
per-stream send queue with a None sentinel and activation under the lock,
and take dapr#1234's activation before the stream is opened and its reset to
inactive when the initial read fails. Keep dapr#1234's tests.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Signed-off-by: Casper Nielsen <casper@diagrid.io>
The sync handler loop gave up after one failed reconnect: it slept 5
seconds, then read from the still-inactive stream, got
StreamInactiveError and exited. It now retries the reconnect every
SUBSCRIPTION_RECONNECT_BACKOFF_SECONDS until the close function is
called. The close function wakes a loop that is waiting in the backoff,
so it no longer takes up to 5 seconds to return.

An exception from the handler was handled as a stream error, so the
stream was reconnected and the message was delivered again. It is now
logged with its traceback, the message gets a retry response, and the
loop reads the next message.

The async handler loop now handles errors the same way as the sync one.
Before, any error other than StreamInactiveError ended the task, and
the error only showed up later as "Task exception was never retrieved".
The task is now kept referenced, so it can't be garbage-collected while
it runs, and the close function waits up to
SUBSCRIPTION_CLOSE_TIMEOUT_SECONDS for it to finish, except when called
from inside the handler. test_subscribe_topic_with_handler is restored.

Fixes dapr#1233

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Signed-off-by: Casper Nielsen <casper@diagrid.io>
When next_message() gives up or raises a non-retryable error, the handler
loop now waits SUBSCRIPTION_RECONNECT_BACKOFF_SECONDS before reconnecting,
in both clients. Closing ends the wait at once.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Signed-off-by: Casper Nielsen <casper@diagrid.io>
@codecov

codecov Bot commented Sep 24, 2026

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 98.52941% with 3 lines in your changes missing coverage. Please review.
✅ Project coverage is 84.38%. Comparing base (fb229bc) to head (942c58d).

Files with missing lines Patch % Lines
dapr/clients/grpc/subscription.py 96.87% 2 Missing ⚠️
dapr/aio/clients/grpc/subscription.py 98.59% 1 Missing ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main    #1235      +/-   ##
==========================================
+ Coverage   83.89%   84.38%   +0.48%     
==========================================
  Files         123      123              
  Lines       10265    10400     +135     
==========================================
+ Hits         8612     8776     +164     
+ Misses       1653     1624      -29     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[BUG] subscribe_with_handler stops silently: async task can die or be collected, sync retry path exits, handler errors reconnect the stream

1 participant