Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 15 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,21 @@ there was nothing to have changed from. Everything above it is ordinary.

### Added

- **`Subscription.recv_many(limit=500)` and `Subscription.batches(limit=500)`**
(#86): a subscriber's version of group commit.
- **What a batch holds:** at least one row, then every row already on the
connection, up to `limit`. It never waits for more. A consumer that keeps
up gets small batches at `recv`'s latency; one that falls behind gets
everything that queued, and catch-up batches come full.
- **The same guarantees as `recv`:** the same rows in the same order, with
the same offset check.
- **The cursor:** it counts a batch as handled when the next is asked for,
so the automatic save is now right for batching consumers.
- **A connection that ends mid-batch:** the rows before the end are
returned, and the refusal is raised next, naming the last row delivered.
- **No new buffer:** one row at most is read ahead, outside `websockets`'
own bounded queue.

- **Metrics in the OTel example** (`examples/otel/metrics.py`):
- **`StreamMetricExporter`**, an OTel metric exporter that publishes each
collection to a stream, one row per data point. It covers sums, gauges
Expand Down
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,8 @@ streamcast.serve(streams, host, port, *, maintain=True, replicate=False,
max_in_flight=64, ...) -> Server # an int, or {stream: int}
streamcast.connect(uri, *, offset=<unset>, cursor=None, cursor_uri=None,
catch_up=False, ...) -> Subscription
await sub.recv() · async for offset, ts, row in sub
async for batch in sub.batches(limit=500) # what has arrived, never waiting for more
streamcast.publish(uri, ...) -> Publication # server needs publish=True
await producer.send(row) · await producer.send_many(rows)
await producer.submit(row) -> Future # pipelined: up to max_in_flight=64 unanswered
Expand Down
28 changes: 24 additions & 4 deletions docs/API.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ offsets, subscribers, replay — and holds no socket. `serve` puts it behind a p
Stream send · send_many · end_offset · subscribers · durable · aclose
│
├─ serve(…) ─► websockets.Server
└─ connect(…) ─► Subscription recv · __aiter__ · offset · info · close
└─ connect(…) ─► Subscription recv · recv_many · batches · __aiter__ · offset · info · close
```

`SCHEMA` and `EARLIEST` are exported because they appear in calls you write: the first
Expand Down Expand Up @@ -687,6 +687,8 @@ some `recv`.
```python
await sub.recv() -> tuple[int | None, int | None, dict[str, object]] # (offset, ts, msg)
async for offset, ts, msg in sub: ...
await sub.recv_many(limit=500) -> list[tuple] # at least one row, then what has arrived
async for batch in sub.batches(limit=500): ... # recv_many until an ordinary close
await sub.close(code=1000, reason="") -> None # drains what is in flight
sub.commit(offset=None) -> None # save the cursor now — see Resuming

Expand All @@ -713,6 +715,23 @@ Iteration **stops** on a normal close (1000/1001) and **raises** on anything els
same contract as iterating a `websockets` connection, with the refusals below filling in
for what a bare code cannot say.

**`recv_many` takes what has already arrived, and never waits for more.** It waits for
the first row as `recv` does, then takes every row already on the connection, up to
`limit` — the rule group commit follows on the publishing side. A consumer that keeps up
gets small batches at `recv`'s latency; one that falls behind gets everything that
queued, which costs it less per row, so it catches up and the batches shrink again.
Catch-up rows are in memory once read, so a catch-up's batches come full.

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

The same rows in the same order as `recv`, checked the same way. The cursor counts a
batch as handled when the next one is asked for, so a handler that raises reads the whole
batch again. If the connection ends partway through, the rows before it are returned and
the next call raises the refusal, whose resume offset is the last row you were given.

`info` is the greeting: `end_offset` (the server's frontier at subscribe), `replay` (the
`[start, end)` about to be replayed, or `None`), `durable` (whether these offsets survive
a server restart), `group_commit` (whether sends may share a commit; false means each
Expand Down Expand Up @@ -834,9 +853,10 @@ which is neither. So four things are arranged around that:
durable and correcting an optimistic value downward is the point. It also cancels the
clean-exit save, so leaving the block does not undo it.

**A batching consumer should not use the automatic save at all.** It advances as you read,
which for a batch is ahead of what you have flushed. Drive `streamcast.Cursor` yourself —
it is exported for exactly that — and save only what you have committed.
**A consumer that batches with `recv_many` or `batches` can use the automatic save**: it
counts a batch as handled when the next is asked for. One that builds its own batches out
of `recv` should not — that save advances per row read, ahead of what it has flushed — and
calls `commit()` once each batch is durable, or drives `streamcast.Cursor` itself.

```python
sub.commit() # force a save now — for a batching or non-idempotent consumer
Expand Down
1 change: 1 addition & 0 deletions docs/SPEC.md
Original file line number Diff line number Diff line change
Expand Up @@ -490,6 +490,7 @@ Neither end can be run out of memory by the other, or by a stall:
| broker | a replay | one batch at a time, into the bounded subscriber queue | — |
| client | a publisher's unanswered sends | `publish(max_in_flight=)`, 64 frames | `submit` waits |
| client | a catch-up | one batch at a time | — |
| client | a `recv_many` batch, and the read it leaves waiting | `limit` rows, plus one | — |
| client | a snapshot's rows from the broker | `max_tail`, 1,000,000 rows | the snapshot is refused |
| client | a live view's unpublished rows | `max_tail`, 1,000,000 rows | the view stops; its next query raises |

Expand Down
169 changes: 140 additions & 29 deletions src/streamcast/_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,7 @@
from streamcast._remote import UPLOAD_EVERY, RemoteCursor

if TYPE_CHECKING:
from collections.abc import Callable, Mapping
from collections.abc import AsyncIterator, Callable, Mapping
from os import PathLike
from types import TracebackType
from typing import Self
Expand All @@ -76,6 +76,16 @@
from streamcast._filter import Where
from streamcast._protocol import Greeting

Row = tuple[int | None, int | None, dict[str, object]]
"""`(offset, ts, row)`, as `recv` returns it."""

BATCH: Final = 500
"""`recv_many`'s and `batches`' default `limit`: rows in one batch at most.

A bound on what one batch holds and on the transaction it hands a consumer,
not a target — a batch is what has arrived, up to this.
"""

_UNSET: Final = object()
"""Tells `offset=None` apart from "offset not given".

Expand Down Expand Up @@ -225,13 +235,15 @@ class Subscription:
"""A live subscription. Async-iterable, read-only, and offset-aware."""

__slots__ = (
"_ahead",
"_catcher",
"_connection",
"_cursor",
"_inbound",
"_info",
"_offset",
"_prelude",
"_received",
"_stream",
"_unsaved",
)
Expand Down Expand Up @@ -260,6 +272,15 @@ def __init__(
# saved value when the caller comes back for another message, which is
# the only evidence this library has that the last one was finished.
self._unsaved: int | None = None
# The last offset READ, which runs ahead of `_offset`, the last one
# handed to the caller, by the rows of a batch being built and by
# `_ahead`. The order check is against this; everything a caller can
# see — `offset`, the cursor, a refusal's resume point — is `_offset`.
self._received: int | None = None
# A read `recv_many` started and did not wait for: the next call's
# first row. Kept rather than cancelled, because a cancelled read of
# the catch-up rows would end their generator (see `recv_many`).
self._ahead: asyncio.Task[Row] | None = None

def __repr__(self) -> str:
return (
Expand Down Expand Up @@ -324,7 +345,7 @@ def connection(self) -> ClientConnection:
details, and anything else this class deliberately does not wrap."""
return self._live()

async def recv(self) -> tuple[int | None, int | None, dict[str, object]]:
async def recv(self) -> Row:
"""The next `(offset, ts, row)`.

`ts` is `streamcast_ts`, when the server took the row, in UTC
Expand All @@ -341,14 +362,100 @@ async def recv(self) -> tuple[int | None, int | None, dict[str, object]]:
if self._cursor is not None:
self._cursor.save(self._unsaved)

row = await self._take()
self._offset = self._unsaved = row[0]
return row

async def recv_many(self, limit: int = BATCH) -> list[Row]:
"""At least one row, then every row already arrived, up to `limit`.

async for batch in stream.batches():
await database.insert_many(batch)

**It never waits for a batch to fill.** It waits for the first row as
`recv` does, then takes only what has already arrived — so a consumer
that keeps up gets small batches at `recv`'s latency, and one that
falls behind gets everything that queued, which costs it less per row
and lets it catch up. Group commit's rule, on the other end of the
stream. Rows from the published tables during a catch-up are already
in memory, so those batches start full.

The same rows in the same order as `recv`, with the same check that
offsets increase. The cursor counts a batch as handled when the next
is asked for, never sooner: a consumer whose handler raises re-reads
the whole batch. If the connection ends partway, the rows read before
it are returned, and the next call raises the refusal.
"""
if limit < 1:
msg = f"limit={limit}: a batch holds at least one row"
raise ValueError(msg)

if self._cursor is not None:
self._cursor.save(self._unsaved)

batch = [await self._take()]
while len(batch) < limit:
# Whether a row is ready, asked the only way the public APIs
# allow: start a read and give it one turn of the loop. A ready
# row is taken; an unready read is KEPT as the next call's first
# row rather than cancelled, so it costs nothing, and a read of
# the catch-up generator — which a cancel would end — is never
# cancelled. One row at most is held outside `websockets`' queue.
ahead = asyncio.ensure_future(self._read())
await asyncio.sleep(0)
if not ahead.done() or ahead.exception() is not None:
# A failed read is raised by the NEXT call, after the rows
# before it are handed over.
self._ahead = ahead
break

batch.append(ahead.result())

self._offset = self._unsaved = batch[-1][0]
return batch

async def batches(self, limit: int = BATCH) -> AsyncIterator[list[Row]]:
"""`recv_many` as a loop: batches until an ordinary close, as iterating
the subscription yields rows until one."""
while True:
try:
batch = await self.recv_many(limit)
except ConnectionClosed as exc:
if exc.rcvd is not None and exc.rcvd.code in _ENDED:
return

raise

yield batch

async def _take(self) -> Row:
"""The next row for the caller: a read already started, or a new one.

A refusal is built HERE, at delivery, rather than where the read
failed: its resume point is the last row the caller was given, and
rows read before the failure may still have been on their way to it.
"""
ahead, self._ahead = self._ahead, None
try:
return await (ahead if ahead is not None else self._read())
except ConnectionClosed as exc:
refusal = _refusal(exc, stream=self._stream, offset=self._offset)
if refusal is None:
raise

raise refusal from None

async def _read(self) -> Row:
"""Read one row, from the catch-up or the socket. Delivers nothing.

Checks the order and advances `_received`; `_offset`, the cursor and
a refusal belong to the caller's side, `recv` and `recv_many`.
"""
if self._prelude is not None:
caught = await anext(self._prelude, None)
if caught is not None:
offset, ts, row = caught
self._offset = offset
self._unsaved = offset

return offset, ts, row
self._received = caught[0]
return caught

# Exhausted, which means `Catcher.stream` returned — and it does
# not return until it has a live connection. Everything below the
Expand All @@ -361,17 +468,9 @@ async def recv(self) -> tuple[int | None, int | None, dict[str, object]]:
self._info = catcher.info
self._inbound = _inbound(catcher.info)

try:
# `decode=False`: a data frame is a TEXT frame of JSON, and msgspec
# parses the bytes directly — no `str` built only to be parsed.
frame = await self._live().recv(decode=False)
except ConnectionClosed as exc:
refusal = _refusal(exc, stream=self._stream, offset=self._offset)
if refusal is None:
raise

raise refusal from None

# `decode=False`: a data frame is a TEXT frame of JSON, and msgspec
# parses the bytes directly — no `str` built only to be parsed.
frame = await self._live().recv(decode=False)
offset, ts, row = decode(frame)
# Binary columns arrive as text in their encoding, and a row that
# came out of the published tables by catch-up carries bytes: decoded here so
Expand All @@ -388,7 +487,7 @@ async def recv(self) -> tuple[int | None, int | None, dict[str, object]]:
# What it actually guards is the two places the offsets come from
# somewhere TCP says nothing about:
#
# * **The catch-up join.** `self._offset` above is set by rows read
# * **The catch-up join.** `self._received` is set by rows read
# out of OBJECT STORAGE, and the first frame off the socket is
# compared against it. Those are two different sources spliced into
# one stream, and the splice is computed by `Catcher.start` and the
Expand All @@ -408,28 +507,30 @@ async def recv(self) -> tuple[int | None, int | None, dict[str, object]]:
# is usually assumed to cover cannot happen.
#
# It does NOT span a reconnect: a new `Subscription` starts with
# `_offset = None`, so nothing is compared across the gap. The log is
# `_received = None`, so nothing is compared across the gap. The log is
# what makes that safe, not this.
#
# `<=` rather than `!= previous + 1`: litelink's offset space has
# legitimate GAPS (a `restore` fences 2**20 of them), so a jump
# forward is ordinary and only a step backwards is wrong.
if offset is not None and self._offset is not None and offset <= self._offset:
if (
offset is not None
and self._received is not None
and offset <= self._received
):
msg = (
f"offsets must increase within a subscription; received "
f"{offset} after {self._offset}"
f"{offset} after {self._received}"
)
raise ProtocolError(msg)

self._offset = offset
self._unsaved = offset

self._received = offset
return offset, ts, row

def __aiter__(self) -> Self:
return self

async def __anext__(self) -> tuple[int | None, int | None, dict[str, object]]:
async def __anext__(self) -> Row:
"""Stops on an ordinary close; raises on anything else.

Same contract as iterating a `websockets` connection — a normal
Expand All @@ -455,9 +556,11 @@ def commit(self, offset: int | None = None) -> None:
it is you stating what is durable, and correcting an optimistic value
downward is the point. The automatic saves never do — see `_cursor`.

A consumer that batches should not rely on the automatic save at all:
it advances as you read, which for a batch is ahead of what you have
flushed. Drive `streamcast.Cursor` yourself there.
A consumer that batches with `recv_many` or `batches` can rely on the
automatic save: it counts a batch as handled when the next is asked
for. One that builds its own batches out of `recv` should not — that
save advances per row read, ahead of what it has flushed — and calls
this once each batch is durable.
"""
if self._cursor is None:
return
Expand All @@ -478,6 +581,14 @@ def commit(self, offset: int | None = None) -> None:

async def close(self, code: int = 1000, reason: str = "") -> None:
"""End the subscription. Idempotent, and awaits the close handshake."""
# First: a read still running is reading the catch-up generator or the
# socket, and both are closed below. Its row was never handed over,
# so dropping it loses nothing the caller saw.
ahead, self._ahead = self._ahead, None
if ahead is not None:
ahead.cancel()
await asyncio.gather(ahead, return_exceptions=True)

prelude, self._prelude = self._prelude, None
catcher, self._catcher = self._catcher, None
if prelude is not None:
Expand Down
Loading
Loading