feat(client): recv_many and batches, a subscriber's group commit - #92
Merged
Merged
Conversation
recv_many(limit=500) returns at least one row, then every row already on the connection, up to limit, never waiting for more; batches() loops it until an ordinary close. A consumer that keeps up gets small batches at recv's latency, one that falls behind gets what queued, and a catch-up's batches come full. Whether a row is ready is asked by starting a read and giving it one loop turn. An unready read is kept as the next call's first row, never cancelled: a cancelled read of the catch-up generator would end it. recv's reading moves into _read, which delivers nothing, so a row read ahead never moves offset, the cursor or a refusal's resume point; the order check runs against the new _received. The cursor counts a batch as handled when the next is asked for, and a connection ending mid-batch returns the rows before it, then raises naming the last row delivered. Closes #86. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FsSDkeb5rVAxA1FSmKKfQi
Merged
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.
Closes #86.
What
Subscription.recv_many(limit=500): at least one row, then every row already on the connection, up tolimit. It never waits for more, which is group commit's rule on the subscribing side:recv's latency;Subscription.batches(limit=500):recv_manyuntil an ordinary close, the way iterating a subscription stops on one.limitrather than the issue'smax, becausemaxshadows the builtin (ruff A002).How "already arrived" is asked
Nothing public reports what's buffered, so
recv_manystarts a read and gives it one loop turn, then takes the row if the read finished. An unfinished read is kept as the next call's first row, never cancelled:asynchronous generator is already running.websockets' own bounded queue. SPEC §4's memory table has a row for it.The invariants it keeps
recv's reading moves into_read(), which reads, decodes and checks order but delivers nothing.recv,recv_manyandbatchesdeliver.offset, the clean-exit commit and a refusal's resume point count only rows handed to the caller (_offset). The order check runs against the new_received.commit()'s docstring and API.md's Resuming section now say that batching throughrecv_many/batchesmakes the automatic save correct. Only consumers that build their own batches out ofrecvshould commit by hand.Two existing tests in
test_refusals.pysimulated a stale frame by setting_offset. They now set_received, which the order check reads. Without that change, the forward-gap test would have passed without testing anything.Tests
tests/test_batches.py(12 tests, 0.4 s): mostly against a fake connection that behaves aswebsocketsdoes: a buffered frame comes back without suspending, and a cancelled read loses nothing. That makes "already arrived" exact instead of timing-dependent. It covers:batchesending on an ordinary close, and a raising handler's batch staying unsaved;test_catchup.py: across the join. 2,000 published rows plus 10 unpublished ones, which only the server can deliver.limit=1500makes the second batch stop at the join, with its next read waiting on the catch-up's connection. The test checks that every row arrives once, in order, and that the batch sizes are[1500, 500], which is deterministic.batchesraising on an ordinary close.🤖 Generated with Claude Code
https://claude.ai/code/session_01FsSDkeb5rVAxA1FSmKKfQi