fix(ps): mark streaming sub active before opening grpc stream - #1234
Merged
Merged
Conversation
…tream Signed-off-by: Samantha Coyle <sam@diagrid.io>
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## main #1234 +/- ##
==========================================
- Coverage 86.63% 83.91% -2.72%
==========================================
Files 84 123 +39
Lines 4473 10260 +5787
==========================================
+ Hits 3875 8610 +4735
- Misses 598 1650 +1052 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
CasperGN
approved these changes
Sep 24, 2026
nelson-parente
added a commit
to nelson-parente/python-sdk
that referenced
this pull request
Sep 24, 2026
Resolves the conflict in tests/integration/test_pubsub.py. The next_message None loop from dapr#1234 moves into the shared _next_message helper, so the bulk metadata test gets it too. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Nelson Parente <nelson_parente@live.com.pt>
3 tasks
CasperGN
added a commit
to CasperGN/python-sdk
that referenced
this pull request
Sep 24, 2026
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>
3 tasks
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Description
Fixes a race in the sync streaming
Subscriptionthat madetest_streaming_subscribe_receives_published_messagefail intermittently in CI. The same race can drop real user subscriptions on connect.Subscription.start()created the bidi stream and only then marked it active. gRPC starts reading the request iterator on its own thread beforeSubscribeTopicEventsAlpha1()returns. When that thread won the race, the iterator sent the initial request, saw the stream as inactive, and ended. That half-closed the stream, so the sidecar unsubscribed right away and the client gotUNKNOWN: EOF. The reconnect hit the same race,next_message()returnedNone, and the test failed with'NoneType' object has no attribute 'data'.The async
Subscriptionis not affected, since it marks the stream active before its firstawait.Changes
dapr/clients/grpc/subscription.py: mark the stream active before creating it, and mark it inactive again if setup fails.dapr/clients/grpc/subscription.py: give each stream its own send queue, and wait on it with a timeout. Before this, an iterator left over afterclose()or a reconnect blocked forever on a shared queue and could take acks meant for the new stream.tests/clients/test_subscription.py: add unit tests with a stub that reads the request iterator beforestart()returns, which reproduces the race every time. Also check that the iterator exits afterclose().tests/integration/test_pubsub.py: keep reading whennext_message()returnsNone, since the API documents that as the normal result of a reconnect.Testing
masterand passes with this fix.ruffandmypypass.Issue reference
We strive to have all PR being opened based on an issue, where the problem or feature have been discussed prior to implementation.
Please reference the issue this PR will close: #[issue number]
Checklist
Please make sure you've completed the relevant tasks for this PR, out of the following list: