Repository navigation
Add opt-in receipt-confirmed sends - #207
Merged
Merged
Conversation
Client.send(receipt_timeout=...) now waits for the broker's STOMP RECEIPT instead of returning once the frame is written to the socket. Failures are reported as SendError with reason rejected, timeout or connection_lost; a confirmed send is attempted once and never replayed, so an ambiguous outcome stays the caller's decision. Pending receipts move from ActiveSubscriptions into a connection-scoped, epoch-aware PendingReceipts registry shared by subscriptions and sends. The frame reader's handler-concurrency backpressure consults the same registry, so a confirmed send still completes while message handlers are saturated. faststream-stomp exposes receipt_timeout on the broker, on declared publishers and per call, with the same precedence as add_content_length. Batch publishing goes through a transaction and does not claim confirmation.
The error only ever arises from receipt confirmation, so the name says which failure it reports.
anyio, faker, hypothesis and ruff move to their current releases. ruff 0.16.8 selects two preview rules through `select = ["ALL"]` that rewrite suppression comments to `ruff: ignore[...]` and rule codes in this configuration to rule names. Both are stylistic, so ignore RUF105 and RUF201 and keep `noqa` comments and rule codes.
lesnik512
force-pushed
the
feature/send-receipt-confirmation
branch
from
September 17, 2026 14:19
968765a to
26b1022
Compare
StompProducer.publish passes receipt_timeout to Client.send on every publish, so an older stompman raises TypeError for every message, not only for opted-in ones. The workspace source hides this in CI, where the local stompman is always used.
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.
Follows #204:
Client.send()can now wait for the broker to accept the message, the waysubscribe()already can.Why
A successful
send()means only that the frame was written to the socket. Callers that acknowledge upstream work on the strength of a publish (a Kafka consumer acknowledging a command, say) are acknowledging something the broker may never have received.stompmanThe client generates the
receiptheader, waits for the matchingRECEIPT, and raisesSendReceiptErrorwithreasonrejected,timeoutorconnection_lost.receipt_timeout=None(the default) is byte-for-byte the old path; a timeout must be finite and positive, and passing your ownreceiptheader alongside it is aValueError.Only
rejectedis definitive. A frame is fully buffered before the socket is drained, so ontimeoutorconnection_lostthe broker may well have taken the message. The confirmed write is therefore attempted exactly once:connect_retry_attemptsstill applies while acquiring a connection, butwrite_retry_attemptsdoes not, because neitherConnectionnorWebSocketConnectioncan distinguish a failure before the bytes left the process from one after they reached the broker. Nothing is replayed silently; retrying is the caller's decision, and the README says so and points at deduplication by a stable application event ID.Transaction.send()and batch publication are deliberately untouched: a receipt for a transactionalSENDsays nothing about whether theCOMMITsucceeded.Design: one receipt registry, not two
Pending receipts move out of
ActiveSubscriptionsintoPendingReceipts(receipts.py): connection-scoped, epoch-keyed, dispatchingRECEIPT/ERROR/ connection-loss to whichever waiter registered the ID. Subscriptions keep their own failure semantics as one waiter;send.pyis a much simpler second one. The subscription state machine itself is unchanged.Sharing the registry is not only tidiness.
Client._reserve_handler_slotlets the frame reader run ahead ofmax_concurrent_handlersonly while receipts are pending. Had sends tracked their receipts separately, a confirmedsend()issued while handlers were saturated would have hung until its own timeout, because the reader was parked on the semaphore and never read theRECEIPT. There is a regression test for exactly that.An
ERRORwithoutreceipt-idfails every pending confirmation on that connection, matching the policy #204 established for subscriptions.on_error_framestill sees every broker error.faststream-stompreceipt_timeouton the broker, on declared publishers, and per call, with the same precedence asadd_content_length(per-call, then publisher, then broker), defaulting toNoneeverywhere.publish_batchdoes not take it and ignores any configured default rather than claiming a confirmation it cannot make.TestStompBrokernever waits for a real receipt. AsyncAPI output is unaffected.Tests
Deterministic tests at the real client/frame-dispatch seam (the queue-backed connection harness from #204, hoisted into
conftest.pyso both suites share it): receipt matching, unrelated receipts, concurrent sends, correlated and uncorrelatedERROR, timeout, connection loss and cancellation both during and after the write, a receipt that wins the race against the deadline, stale-epoch receipts, invalid timeouts, and the header conflict.receipts.pyandsend.pyare at full line coverage. An Artemis/ActiveMQ Classic integration test proves a confirmed send returns and the message arrives.Full suite with both brokers: 476 passed, 1 skipped.
mypyandruffclean.Note on the last commit
Bump dev dependencies and disable ruff's suppression-style rulesis independent maintenance and can be dropped or split out. Beside the version bumps it turns off two preview rules thatselect = ["ALL"]pulls in: RUF105 rewrites# noqa:comments to# ruff: ignore[...], and RUF201 rewrites rule codes in this configuration to rule names. Both are stylistic, and withfix = trueandunsafe-fixes = truethey rewrite the repo on everyjust lint, so they are now ignored and the existingnoqacomments and rule codes stay as they are.just lintis idempotent and reports no errors.faststream-stomp's floor moves tostompman>=3.16.0, on the assumption that this ships asstompman-3.16.0; adjust that line before tagging if you pick a different number, and release stompman before faststream-stomp. The floor is load-bearing:StompProducer.publishpassesreceipt_timeouttoClient.sendon every publish, so an older stompman would raiseTypeErrorfor every message. CI cannot catch a wrong floor, because the workspace source always resolves stompman locally.