[BUG] subscribe_with_handler stops silently: async task can die or be collected, sync retry path exits, handler errors reconnect the stream
Maintainer thường phản hồi trong vòng 1 ngày
Chưa có ai nhận issue này.
Đánh giá
- Độ khó
- 4/5
- Thời gian dự kiến
- 3-5 ngày
- Mức phù hợp với người mới
- 65/100
Hướng nghiên cứu
Start with dapr/aio/clients/grpc/client.py#L580-L599 and dapr/clients/grpc/client.py#L628-L654, then read subscription.py#L70-L74 and the commented test in tests/clients/test_dapr_grpc_client_async.py#L424-L477. Restore async coverage and add sync and async tests for failed-then-successful reconnects and a handler that raises once. Done means subscriptions keep retrying, handler errors do not reconnect the stream, async tasks are retained and awaited on close, and caller closure remains final.
Do mô hình lập chỉ mục viết ra từ nội dung của issue.
Mô tả
Expected Behavior
A subscription made with subscribe_with_handler (sync or async) keeps delivering messages to the handler until the caller closes it. If the stream drops, it reconnects, and it keeps retrying while the sidecar is unavailable. An exception raised by the handler does not tear down the stream. The returned close function stops the handler before it returns.
Actual Behavior
Links are to main at 03eebe1.
Async client: dapr/aio/clients/grpc/client.py#L580-L599
- No reference is kept to the task. It calls
asyncio.create_task(stream_messages(subscription))and discards the result. Python's docs warn that the event loop keeps only a weak reference to such a task, so it can be garbage-collected while it is still running. The subscription then stops with no error. - Any error ends the subscription. The loop catches only
StreamInactiveError. If the handler raises, or the stream raisesStreamCancelledErroror a non-retryable gRPC error, the task dies. Nothing reconnects, and the error only appears later as an "exception was never retrieved" warning. The sync client reconnects in that case. - The close function doesn't wait for the task.
close_subscription()closes the subscription but never awaits the task, so the handler can still be running afterawait close_fn()returns. - No test covers it.
test_subscribe_topic_with_handleris commented out:tests/clients/test_dapr_grpc_client_async.py#L424-L477.
Sync client: dapr/clients/grpc/client.py#L628-L654
-
The "reconnect failed, back off and retry" branch never retries.
- On a stream error the loop calls
sub.reconnect_stream(). reconnect_stream()marks the stream inactive first. Then it waits for the sidecar and callsstart().- If that wait or
start()raises, the loop sleeps 5 seconds and runscontinue. - The next
for message in subcallsnext_message(). The stream is still inactive, so that raisesStreamInactiveError, and the loop treats it as a close andbreaks.
So a sidecar outage longer than the health wait ends the subscription for good, with no error.
- On a stream error the loop calls
-
Handler errors reconnect a healthy stream. The same
except Exceptionalso catches exceptions raised byhandler_fn. One failing handler call tears down and reconnects the stream, and the message is never acked, so it is delivered again.
Steps to Reproduce the Problem
- 5: stop the sidecar while a
subscribe_with_handlersubscription is running. Keep it down for longer thanDaprHealth.wait_for_sidecarwaits, then start it again. The handler thread has exited, and no new messages reach the handler. - 6: have the handler raise for one message. The log shows a stream reconnect, and the message is delivered again.
- 2: do the same with the async client. The task ends, and no further messages are handled.
Suggested fix
Build on #1231, which makes close() final: a subscription that has been closed can no longer be reopened by a reconnect.
- Sync (5): after a failed reconnect, retry the reconnect with the backoff. Leave the loop only when the subscription was closed by the caller.
- Both (6): treat a handler exception separately from a stream error. Log it and respond with retry, instead of reconnecting the stream.
- Async (1–3):
- Keep a reference to the task.
- Handle
StreamCancelledErrorand stream errors the way the sync client does. - Have the close function await the task, with a timeout.
- Tests: restore the async handler test. Add tests for both clients: a reconnect that fails and then succeeds, and a handler that raises once.
Release Note
RELEASE NOTE: FIX subscribe_with_handler keeps retrying while the sidecar is unavailable, does not reconnect on handler errors, and the async version no longer stops silently.
- Ngôn ngữ chính
- Python
- Star
- 272
- Fork
- 152
- Merge trung bình
- 2 ngày 22 giờ
- Pull request đã merge (30 ngày)
- 7
Chuẩn bị môi trường
Bắt đầu từ đâu
- Đọc hết issue, rồi đọc hướng dẫn đóng góp của dự án.
- Bình luận trên issue rằng bạn sẽ nhận — tránh hai người làm cùng một việc.
- Fork repository và làm thay đổi trên một nhánh.
- Mở pull request có tham chiếu số hiệu của issue.
Issue khác của dapr/python-sdk
-
dapr-ext-workflow good first issue kind/enhancement P2
Độ khó 2/5 1-2 ngày Mức phù hợp với người mới 72/100
dapr/python-sdk#853 · 4 bình luận ·
Maintainer thường phản hồi trong vòng 1 ngày
-
kind/bug
Độ khó 4/5 3-5 ngày Mức phù hợp với người mới 68/100
dapr/python-sdk#1232 ·
Maintainer thường phản hồi trong vòng 1 ngày
-
kind/bug
Độ khó 4/5 3-5 ngày Mức phù hợp với người mới 45/100
dapr/python-sdk#1230 ·
Maintainer thường phản hồi trong vòng 1 ngày
-
dapr-ext-workflow kind/enhancement
Độ khó 3/5 1-2 ngày Mức phù hợp với người mới 78/100
dapr/python-sdk#1213 ·
Maintainer thường phản hồi trong vòng 1 ngày
-
Độ khó 3/5 1-2 ngày Mức phù hợp với người mới 76/100
dapr/python-sdk#1200 ·
Maintainer thường phản hồi trong vòng 1 ngày
Tất cả issue của dapr/python-sdk
Issue tương tự
-
correction metadata
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 68/100
acl-org/acl-anthology#10104 · 1 bình luận ·
Maintainer thường phản hồi trong vòng 1 ngày
-
bug status/needs-triage
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 86/100
prowler-cloud/prowler#12885 · 1 bình luận ·
Maintainer thường phản hồi trong vòng 1 ngày
-
Bug in GaussianTailProbabilityCalibrator: running_statistics=False still uses a windowed varianceĐang mởbug good first issue
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 88/100
selimfirat/pysad#107 ·
Maintainer thường phản hồi trong vòng 1 ngày
-
bug ci-failure high priority
Độ khó 1/5 Dưới một giờ Mức phù hợp với người mới 88/100
vllm-project/vllm-omni#8194 · 1 bình luận ·
Maintainer thường phản hồi trong vòng 1 ngày
-
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 88/100
Maintainer thường phản hồi trong vòng 1 ngày