diff --git a/CHANGELOG.md b/CHANGELOG.md index ba6004a..80355f9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/README.md b/README.md index e13ef38..022e561 100644 --- a/README.md +++ b/README.md @@ -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=, 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 diff --git a/docs/API.md b/docs/API.md index 66d3943..bc7c197 100644 --- a/docs/API.md +++ b/docs/API.md @@ -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 @@ -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 @@ -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 @@ -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 diff --git a/docs/SPEC.md b/docs/SPEC.md index dedf108..1529abc 100644 --- a/docs/SPEC.md +++ b/docs/SPEC.md @@ -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 | diff --git a/src/streamcast/_client.py b/src/streamcast/_client.py index eb8baec..68499d2 100644 --- a/src/streamcast/_client.py +++ b/src/streamcast/_client.py @@ -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 @@ -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". @@ -225,6 +235,7 @@ class Subscription: """A live subscription. Async-iterable, read-only, and offset-aware.""" __slots__ = ( + "_ahead", "_catcher", "_connection", "_cursor", @@ -232,6 +243,7 @@ class Subscription: "_info", "_offset", "_prelude", + "_received", "_stream", "_unsaved", ) @@ -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 ( @@ -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 @@ -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 @@ -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 @@ -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 @@ -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 @@ -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 @@ -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: diff --git a/tests/test_batches.py b/tests/test_batches.py new file mode 100644 index 0000000..50e1b9f --- /dev/null +++ b/tests/test_batches.py @@ -0,0 +1,226 @@ +"""`recv_many` and `batches`: what has arrived, in one batch, never waiting for more. + +Most of these run against a fake connection, because "already arrived" is +exactly what a real socket makes a matter of timing. The fake behaves as +`websockets` does where it matters: a buffered frame comes back without +suspending, an empty queue waits, and a cancelled `recv` loses nothing. The +tests against a real server assert only what timing cannot change: every row, +once, in order, within the limit. +""" + +from __future__ import annotations + +import asyncio +import collections + +import pytest +from websockets.exceptions import ConnectionClosed +from websockets.frames import Close as CloseFrame + +import streamcast +from streamcast._client import Subscription +from streamcast._errors import Close +from streamcast._protocol import encode_projected, greeting, parse_greeting, refusal + + +class FakeConnection: + """The two things `Subscription` asks of a connection: `recv` and `close`.""" + + def __init__(self) -> None: + self._frames: collections.deque[bytes] = collections.deque() + self._closed: ConnectionClosed | None = None + self._arrived = asyncio.Event() + self.reads = 0 + + def arrive(self, *offsets: int) -> None: + for offset in offsets: + self._frames.append(encode_projected(offset, 0, {"i": offset})) + + self._arrived.set() + + def end(self, code: int = 1000, reason: str = "") -> None: + self._closed = ConnectionClosed(CloseFrame(code, reason), None) + self._arrived.set() + + async def recv(self, decode: bool | None = None) -> bytes: # noqa: ARG002 + while not self._frames: + if self._closed is not None: + raise self._closed + + self._arrived.clear() + await self._arrived.wait() # cancelling here takes nothing + + self.reads += 1 + return self._frames.popleft() + + async def close(self, code: int = 1000, reason: str = "") -> None: # noqa: ARG002 + self.end(code) + + +def subscription( + connection: FakeConnection, cursor: streamcast.Cursor | None = None +) -> Subscription: + info = parse_greeting(greeting(stream="t", end_offset=1, replay=None, durable=True)) + return Subscription(connection, info, "t", cursor) # ty: ignore[invalid-argument-type] + + +def offsets(batch: list) -> list[int]: + return [offset for offset, _ts, _row in batch] + + +class TestWhatABatchHolds: + async def test_everything_already_arrived_comes_as_one_batch(self): + connection = FakeConnection() + connection.arrive(1, 2, 3, 4, 5) + sub = subscription(connection) + + assert offsets(await sub.recv_many(limit=10)) == [1, 2, 3, 4, 5] + + async def test_a_batch_stops_at_the_limit(self): + connection = FakeConnection() + connection.arrive(1, 2, 3, 4, 5) + sub = subscription(connection) + + got = [offsets(await sub.recv_many(limit=2)) for _ in range(3)] + assert got == [[1, 2], [3, 4], [5]] + + async def test_it_waits_for_the_first_row_and_never_for_more(self): + connection = FakeConnection() + sub = subscription(connection) + + waiting = asyncio.ensure_future(sub.recv_many(limit=10)) + await asyncio.sleep(0.05) + assert not waiting.done() # nothing has arrived: it waits + + connection.arrive(1) + # One row arrived and nothing behind it: a batch of one, now. + assert offsets(await asyncio.wait_for(waiting, 1)) == [1] + + async def test_the_read_left_waiting_is_the_next_batchs_first_row(self): + """Kept, not cancelled: nothing is lost between batches.""" + connection = FakeConnection() + connection.arrive(1) + sub = subscription(connection) + + assert offsets(await sub.recv_many()) == [1] + connection.arrive(2, 3) + assert offsets(await sub.recv_many()) == [2, 3] + connection.arrive(4) + assert (await sub.recv())[0] == 4 # `recv` takes it too + + def test_a_batch_holds_at_least_one_row(self): + sub = subscription(FakeConnection()) + with pytest.raises(ValueError, match="limit=0"): + asyncio.run(sub.recv_many(limit=0)) + + +class TestWhatTheCallerHasSeen: + """A row read ahead is not a row delivered: `offset`, the cursor and a + refusal's resume point count only what the caller was handed.""" + + async def test_a_batch_is_saved_when_the_next_is_asked_for(self, tmp_path): + cursor = streamcast.Cursor(tmp_path / "cursor", every=0) + connection = FakeConnection() + connection.arrive(1, 2, 3) + sub = subscription(connection, cursor) + + assert offsets(await sub.recv_many()) == [1, 2, 3] + # Handed over, not yet known to be handled: nothing saved. + assert cursor.load() is None + + connection.arrive(4) + await sub.recv_many() + assert cursor.load() == 3 # the whole first batch, and no further + + async def test_a_row_read_ahead_is_not_received(self, tmp_path): + cursor = streamcast.Cursor(tmp_path / "cursor", every=0) + connection = FakeConnection() + connection.arrive(1) + sub = subscription(connection, cursor) + + assert offsets(await sub.recv_many()) == [1] + connection.arrive(2) + await asyncio.sleep(0.01) # the waiting read takes row 2 off the socket + assert connection.reads == 2 + + # But the caller never got it: it is not the offset, and a clean + # exit's commit does not save it. + assert sub.offset == 1 + sub.commit() + assert cursor.load() == 1 + await sub.close() # and closing with it still held is quiet + + async def test_a_failure_mid_batch_hands_over_the_rows_before_it(self): + connection = FakeConnection() + connection.arrive(1, 2) + connection.end(Close.TOO_SLOW, refusal("too_slow", backlog=8)) + sub = subscription(connection) + + assert offsets(await sub.recv_many()) == [1, 2] + with pytest.raises(streamcast.TooSlow) as raised: + await sub.recv_many() + + # Built when raised, so it names the last row the caller was given: + # resuming above it misses nothing and repeats nothing. + assert raised.value.offset == 2 + + async def test_offsets_must_still_increase_within_a_batch(self): + connection = FakeConnection() + connection.arrive(3, 2) + sub = subscription(connection) + + assert offsets(await sub.recv_many()) == [3] + with pytest.raises(streamcast.ProtocolError, match="must increase"): + await sub.recv_many() + + +class TestTheLoop: + async def test_batches_end_on_an_ordinary_close(self): + connection = FakeConnection() + connection.arrive(1, 2, 3) + connection.end(1000) + sub = subscription(connection) + + assert [offsets(b) async for b in sub.batches(limit=2)] == [[1, 2], [3]] + + async def test_a_handler_that_raises_leaves_its_batch_unsaved(self, tmp_path): + cursor = streamcast.Cursor(tmp_path / "cursor", every=0) + connection = FakeConnection() + connection.arrive(1, 2) + sub = subscription(connection, cursor) + + handled = [] + with pytest.raises(RuntimeError): + async for batch in sub.batches(): + handled.append(offsets(batch)) + connection.arrive(3, 4) + if offsets(batch) == [3, 4]: + msg = "the handler failed on this batch" + raise RuntimeError(msg) + + assert handled == [[1, 2], [3, 4]] + assert cursor.load() == 2 # the batch it failed on is read again + + +class TestAgainstAServer: + async def test_every_row_once_in_order_within_the_limit(self, tmp_path, serve): + schema = { + "type": "object", + "properties": {"i": {"type": "integer"}}, + "required": ["i"], + } + stream = streamcast.Stream.new("t", root=tmp_path, schema=schema) + await stream.send_many([{"i": i} for i in range(300)]) + try: + async with serve(stream, maintain=False) as uri: + async with streamcast.connect(uri, offset=streamcast.EARLIEST) as sub: + got: list[list[int]] = [] + async for batch in sub.batches(limit=64): + got.append(offsets(batch)) + if got[-1][-1] == 300: + break + + assert [o for batch in got for o in batch] == list(range(1, 301)) + assert all(1 <= len(batch) <= 64 for batch in got) + finally: + await stream.aclose() diff --git a/tests/test_catchup.py b/tests/test_catchup.py index 86d12a2..83f72d4 100644 --- a/tests/test_catchup.py +++ b/tests/test_catchup.py @@ -162,6 +162,43 @@ async def test_a_long_catch_up_is_not_dropped_for_falling_behind( assert [offset for offset, _ts, _row in got] == list(range(1, TOTAL + 1)) + async def test_batches_cross_the_join_whole_and_in_order( + self, serve, published_log, s3 + ): + """`batches` over a catch-up: published rows, then the socket's. + + The catch-up rows come from a generator a cancelled read would end. + The read after the last published row waits while the generator + connects to the server, so a batch that reaches the join ends there + with that read waiting — the case where it must be kept, never + cancelled. Every row once, in order, across the join. + """ + stream = streamcast.Stream("trades", log=published_log, max_replay=600) + await fill(stream, published_log) + # Not published: these can only come from the server, after the join. + live = 10 + await stream.send_many( + [{"i": i, "pad": PAD} for i in range(TOTAL, TOTAL + live)] + ) + assert published_log.published_through() == TOTAL # inclusive + + async with serve(stream, maintain=False) as uri: + async with streamcast.connect( + uri, offset=1, catch_up=True, s3_options=s3 + ) as sub: + got: list[list[int | None]] = [] + async for batch in sub.batches(limit=1500): + got.append([offset for offset, _ts, _row in batch]) + if got[-1][-1] == TOTAL + live: + break + + flat = [offset for batch in got for offset in batch] + assert flat == list(range(1, TOTAL + live + 1)) + # The published rows are in memory once read, so the first batch is + # full; the second stops at the join, its next read waiting on the + # connection. Deterministic: no handshake completes in one loop turn. + assert [len(batch) for batch in got[:2]] == [1500, TOTAL - 1500] + @pytest.mark.slow async def test_nothing_published_during_the_catch_up_is_lost( self, serve, published_log, s3 diff --git a/tests/test_refusals.py b/tests/test_refusals.py index 94bac6d..0a0c715 100644 --- a/tests/test_refusals.py +++ b/tests/test_refusals.py @@ -393,9 +393,9 @@ async def test_an_out_of_order_frame_stops_the_stream(self, serve, log): assert (await sub.recv())[0] == 1 assert (await sub.recv())[0] == 2 - # Rewind what the subscription believes it has seen, which is + # Rewind what the subscription believes it has read, which is # indistinguishable from the next frame arriving stale. - sub._offset = 99 # noqa: SLF001 + sub._received = 99 # noqa: SLF001 with pytest.raises(streamcast.ProtocolError, match="must increase"): await sub.recv() @@ -409,7 +409,7 @@ async def test_a_gap_in_the_offsets_is_not_out_of_order(self, serve, log): async with streamcast.connect(uri, offset=streamcast.EARLIEST) as sub: assert (await sub.recv())[0] == 1 - sub._offset = -1000 # noqa: SLF001 — a large jump forward + sub._received = -1000 # noqa: SLF001 — a large jump forward assert (await sub.recv())[0] == 2 async def test_a_live_only_stream_has_no_offsets_to_check(self, serve):