Skip to content

feat(client): recv_many and batches, a subscriber's group commit - #92

Merged
nhobin219 merged 1 commit into
mainfrom
subscriber-batches
Oct 4, 2026
Merged

nhobin219 merged 1 commit into
mainfrom
subscriber-batches

Conversation

@nhobin219

Copy link
Copy Markdown
Owner

Closes #86.

async for batch in sub.batches(limit=500):
    await database.insert_many(row for _offset, _ts, row in batch)

What

  • Subscription.recv_many(limit=500): at least one row, then every row already on the connection, up to limit. It never waits for more, which is group commit's rule on the subscribing side:
    • a consumer that keeps up gets small batches at recv's latency;
    • one that falls behind gets everything that queued, and catches up;
    • catch-up batches come full.
  • Subscription.batches(limit=500): recv_many until an ordinary close, the way iterating a subscription stops on one.
  • limit rather than the issue's max, because max shadows the builtin (ruff A002).

How "already arrived" is asked

Nothing public reports what's buffered, so recv_many starts 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:

  • It costs nothing. The read is the wait the next call would have done anyway.
  • It's the only safe choice in a catch-up. A cancelled read of the catch-up generator ends that generator. The new catch-up test shows it: with cancelling instead of keeping, it fails with asynchronous generator is already running.
  • It needs no new buffer. At most one row is read ahead, outside 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_many and batches deliver.

  • A row read ahead is not received. 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.
  • The cursor never gets ahead of handled work (CLAUDE.md Validate rows on the live-only path, so a stream accepts the same rows with or without a log #9). A batch is saved when the next one is asked for. A handler that raises reads its whole batch again.
  • A connection ending mid-batch returns the rows before the end, then raises the refusal on the next call. The refusal is built at delivery, so its resume offset is the last row delivered.
  • commit()'s docstring and API.md's Resuming section now say that batching through recv_many/batches makes the automatic save correct. Only consumers that build their own batches out of recv should commit by hand.

Two existing tests in test_refusals.py simulated 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 as websockets does: a buffered frame comes back without suspending, and a cancelled read loses nothing. That makes "already arrived" exact instead of timing-dependent. It covers:
    • batch contents and the limit;
    • waiting for the first row only;
    • the kept read becoming the next batch's first row;
    • cursor timing, and a read-ahead row not being received;
    • a refusal mid-batch naming the last delivered row;
    • the order check within a batch;
    • batches ending on an ordinary close, and a raising handler's batch staying unsaved;
    • against a real server: 300 rows, once each, in order, within the limit.
  • test_catchup.py: across the join. 2,000 published rows plus 10 unpublished ones, which only the server can deliver. limit=1500 makes 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.
  • Falsified: six breaks, each caught:
    • cancelling the waiting read;
    • saving the cursor as a batch is handed over;
    • letting a read move the caller's offset;
    • building the refusal where the read failed;
    • checking order against the delivered offset;
    • batches raising on an ordinary close.

🤖 Generated with Claude Code

https://claude.ai/code/session_01FsSDkeb5rVAxA1FSmKKfQi

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
@nhobin219
nhobin219 merged commit 360bed5 into main Oct 4, 2026
6 checks passed
@nhobin219 nhobin219 mentioned this pull request Oct 4, 2026
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.

Subscriber batching: recv_many / batches, the consumer's group commit

1 participant