diff --git a/CHANGELOG.md b/CHANGELOG.md index 80355f9..5141b58 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,16 @@ there was nothing to have changed from. Everything above it is ordinary. ### Added +- **litelink's read-cache settings on every reader of the published tables** + (#77): `memory_cache`, `disk_cache`, `cache_key` and + `disk_cache_volume_limit` on `Stream.snapshot`, `scan`, `sql`, `live`, and on + `connect` for a catch-up. + - **litelink's defaults:** memory on, disk off. + - **The key is the caller's:** relative to litelink's cache root, absolute, + or its `default`. Nothing is keyed automatically. + - **Separate databases:** readers with different settings read through + different DuckDB databases. + - **`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 @@ -41,6 +51,13 @@ there was nothing to have changed from. Everything above it is ordinary. ### Changed +- **litelink 0.9** (`>=0.9.0,<0.10`), for `litelink.current_metadata`. +- **A published table's current metadata is resolved outside DuckDB**, with + `litelink.current_metadata`. Through a disk-cached connection, + `version-hint.text` pinned a reader to the first snapshot it saw + (litelink#141, fixed in litelink 0.9 by that function). Every file the hint + names is written once and caches safely. + - **The one-shot OTel demo runs without a maintainer**, as the migration demo does. Its maintainer processes were still opening the logs when the run ended, and cost it several seconds: `demo.main` took 3–11 s with them and diff --git a/README.md b/README.md index 022e561..0e498b5 100644 --- a/README.md +++ b/README.md @@ -126,6 +126,8 @@ 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 +# Every reader of the published tables also takes litelink's cache settings: +# memory_cache=True, disk_cache=False, cache_key=None, disk_cache_volume_limit=0.8 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 bc7c197..e61bc59 100644 --- a/docs/API.md +++ b/docs/API.md @@ -653,7 +653,8 @@ arithmetic. ```python streamcast.connect(uri, *, offset=, cursor=None, cursor_uri=None, s3_options=None, upload_every=30.0, catch_up=False, - catch_up_retries=3, metadata=None, + catch_up_retries=3, metadata=None, memory_cache=True, + disk_cache=False, cache_key=None, disk_cache_volume_limit=0.8, **websockets_kwargs) -> Subscription ``` @@ -1273,7 +1274,9 @@ kdb tickerplant has: the feed handler parses, the plant stores typed rows. ```python await streamcast.Stream.snapshot(metadata_uri, *, as_of_offset=None, as_of_ts=None, broker=None, s3_options=None, - max_tail=1_000_000) -> Snapshot + max_tail=1_000_000, memory_cache=True, + disk_cache=False, cache_key=None, + disk_cache_volume_limit=0.8) -> Snapshot await streamcast.Stream.scan(metadata_uri, *, , columns=None, where=None, filters=(), start_offset=None, end_offset=None) -> pa.Table await streamcast.Stream.sql(metadata_uri, query, *, , filters=(), @@ -1389,6 +1392,28 @@ do not hold yet, which it keeps in memory. That stays small while publishing kee `max_tail` rows (1,000,000 by default, counted as read) the snapshot is refused, saying publishing is behind, rather than running the reader out of memory. +**Caching what is read is litelink's, and yours to choose.** The four keywords are +`litelink.duckdb_connection`'s, with its defaults, and the same on `snapshot`, `scan`, +`sql`, `live` and `connect` (for a catch-up): + +| keyword | default | what it does | +|---|---|---| +| `memory_cache` | on | DuckDB's external file cache: what was read stays in memory for the process | +| `disk_cache` | off | `s3://` reads kept on disk by `cache_httpfs`, across restarts | +| `cache_key` | `None` | the disk cache's directory: relative to litelink's cache root (`cache_key=""` is `~/.cache/litelink/`), absolute as given, or `None` for its `default` | +| `disk_cache_volume_limit` | 0.8 | how full the disk cache's volume may get, everything on it counted | + +Nothing chooses a key for you, because only you know what deserves a cache of its own: a +stream id, a team, a job. Readers asking for different settings read through different +DuckDB databases. A disk cache earns its keep where reads cross a network to object +storage; against a store on the same machine, it measured no faster than reading it +again. + +**A cached reader still sees every publish.** Each read resolves the table's current +metadata with `litelink.current_metadata`, outside DuckDB's caches: through a disk-cached +connection, `version-hint.text` would pin a reader to an old snapshot (litelink#141). +Everything the hint names is written once, and caches safely. + **Where the metadata file is.** With an `s3://` published location, the copy beside the tables, which reads from anywhere. Without one, every log publishes under its own directory and the URI is the local `file://` file, which reads only on the server's @@ -1398,7 +1423,9 @@ machine. On another machine the read says so and suggests an `s3://` location. ```python await streamcast.Stream.live(broker, *, s3_options=None, rebase_every=10.0, where=None, - start_offset=None, max_tail=1_000_000) -> Live + start_offset=None, max_tail=1_000_000, memory_cache=True, + disk_cache=False, cache_key=None, + disk_cache_volume_limit=0.8) -> Live async with await streamcast.Stream.live("ws://localhost:8765/trades") as live: await live.wait_for(offset) # until that row is visible diff --git a/pyproject.toml b/pyproject.toml index 87335f1..b34f84f 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -43,12 +43,13 @@ dependencies = [ # 0.5.0 for `column_statistics`, which the manifest is built from, and # for finite-only floats, which is what makes a float column prunable; # 0.5.1 for seals sized by what the Arrow table actually weighs (#84); - # 0.8.1 for Python 3.14. + # 0.8.1 for Python 3.14; 0.9.0 for `current_metadata`, a published + # table's current metadata read outside DuckDB's caches (litelink#141). # # Capped below the next minor: under 0.x a minor is litelink's breaking # release, and a published streamcast must not resolve one it was never # tested against. Each release moves the cap deliberately. - "litelink>=0.8.1,<0.9", + "litelink>=0.9.0,<0.10", # Serialisation is on the hot path in BOTH directions — every publish # encodes a row, every replayed row re-encodes one — which is what earns # a compiled dependency here. Measured against stdlib `json` on a diff --git a/src/streamcast/_catchup.py b/src/streamcast/_catchup.py index fb6c04c..b02cd2b 100644 --- a/src/streamcast/_catchup.py +++ b/src/streamcast/_catchup.py @@ -44,6 +44,7 @@ from streamcast import _snapshot from streamcast._errors import NotReplayable, StreamcastError +from streamcast._published import DEFAULT_CACHE, ReadCache if TYPE_CHECKING: from collections.abc import AsyncGenerator, Awaitable, Callable @@ -135,6 +136,7 @@ class Catcher: """ __slots__ = ( + "_cache", "_first", "_handshake", "_name", @@ -156,8 +158,10 @@ def __init__( start: int, retries: int, handshake: Callable[[int], Awaitable[tuple[Any, Any]]], + cache: ReadCache = DEFAULT_CACHE, ) -> None: self._uri = uri + self._cache = cache self._stream_id = stream_id self._name = name self._s3 = s3_options @@ -173,7 +177,10 @@ def __init__( async def _open(self) -> _snapshot.Snapshot: try: return await _snapshot.snapshot( - self._uri, s3_options=self._s3, stream_id=self._stream_id + self._uri, + s3_options=self._s3, + stream_id=self._stream_id, + cache=self._cache, ) except _snapshot.SnapshotUnavailable as exc: raise CatchUpUnavailable(str(exc)) from exc diff --git a/src/streamcast/_client.py b/src/streamcast/_client.py index 68499d2..44241ce 100644 --- a/src/streamcast/_client.py +++ b/src/streamcast/_client.py @@ -63,6 +63,7 @@ TooSlow, ) from streamcast._protocol import decode, parse_greeting, parse_refusal +from streamcast._published import ReadCache from streamcast._remote import UPLOAD_EVERY, RemoteCursor if TYPE_CHECKING: @@ -718,6 +719,7 @@ class connect: # noqa: N801 — `websockets.connect` is lowercase and this mirr """ __slots__ = ( + "_cache", "_catch_up", "_catch_up_retries", "_cursor", @@ -746,10 +748,22 @@ def __init__( catch_up: bool = False, catch_up_retries: int = CATCH_UP_RETRIES, metadata: str | None = None, + memory_cache: bool = True, + disk_cache: bool = False, + cache_key: str | PathLike[str] | None = None, + disk_cache_volume_limit: float = 0.8, compression: str | None = None, **kwargs: Any, ) -> None: self._uri = uri + # How a catch-up's reads of the published tables are cached: + # litelink's settings, passed through. See `Stream.snapshot`. + self._cache = ReadCache.of( + memory_cache=memory_cache, + disk_cache=disk_cache, + cache_key=cache_key, + disk_cache_volume_limit=disk_cache_volume_limit, + ) # `Any`, so the checker resolves `**self._kwargs` against # `websockets.connect`'s very precise signature rather than against a # value type inferred from `compression`. @@ -851,6 +865,7 @@ async def _recover(self, refused: NotReplayable) -> Subscription: start, self._catch_up_retries, self._handshake, + self._cache, ) # BEFORE handing anything back, so an unreadable table raises here # rather than from whatever line first calls `recv`. Entering the diff --git a/src/streamcast/_live.py b/src/streamcast/_live.py index d756ed2..3f45210 100644 --- a/src/streamcast/_live.py +++ b/src/streamcast/_live.py @@ -51,6 +51,7 @@ from streamcast import _log, _snapshot from streamcast._errors import NotReplayable, TooSlow from streamcast._limits import MAX_TAIL +from streamcast._published import DEFAULT_CACHE, ReadCache if TYPE_CHECKING: from collections.abc import Callable, Sequence @@ -88,8 +89,12 @@ def __init__( where: dict[str, object] | None = None, start: int | None = None, max_tail: int = MAX_TAIL, + cache: ReadCache = DEFAULT_CACHE, ) -> None: self._broker = broker + # How every read of the published tables is cached — the base, each + # rebase, and a catch-up — as `Stream.live` was asked. + self._cache = cache # What the view will hold from the broker, unpublished, before it # stops: a publisher that has stalled must not run this process out # of memory. See `_limits.MAX_TAIL`. @@ -143,6 +148,7 @@ async def _connect(self) -> Any: catch_up=True, s3_options=self._s3, where=self._where, + **self._cache.keywords(), # ty: ignore[invalid-argument-type] ) try: subscription = await connecting @@ -257,6 +263,7 @@ async def rebase(self) -> None: _require(self._broker, self._uri), s3_options=self._s3, stream_id=self._stream_id, + cache=self._cache, ) async with self._lock: # Converted against the old base first, so a pending row is judged @@ -490,6 +497,7 @@ async def live( where: dict[str, object] | None = None, start_offset: int | None = None, max_tail: int = MAX_TAIL, + cache: ReadCache = DEFAULT_CACHE, ) -> Live: """See `Stream.live`.""" from streamcast import _client # noqa: PLC0415 — the client imports `_snapshot` @@ -515,7 +523,7 @@ async def live( raise ValueError(msg) base = await _snapshot.snapshot( - uri, s3_options=s3_options, stream_id=greeting.stream_id + uri, s3_options=s3_options, stream_id=greeting.stream_id, cache=cache ) view = Live( broker, @@ -526,6 +534,7 @@ async def live( where=where, start=start, max_tail=max_tail, + cache=cache, ) try: await view._start() # noqa: SLF001 diff --git a/src/streamcast/_published.py b/src/streamcast/_published.py index d19934e..d71581b 100644 --- a/src/streamcast/_published.py +++ b/src/streamcast/_published.py @@ -15,7 +15,18 @@ **The DuckDB connection is litelink's** (`litelink.duckdb_connection`), with the `iceberg` and `httpfs` extensions litelink provisions rather than installs -at first read, and the S3 secret it creates. +at first read, the S3 secret it creates, and the read caches the caller asks +for (`ReadCache`). + +**The hint is resolved around DuckDB, never through it.** `version-hint.text` +is the one object a reader touches that changes: everything it names — a +`metadata.json`, its manifests, the data files — is written once under a +name of its own, so every cache is safe for them. Read through a connection +with the disk cache on, the hint goes through `cache_httpfs`, whose +file-handle cache keeps its first handle for up to an hour: a reader would +stay on an old snapshot, or since litelink 0.9 fail an ETag check +(litelink#141). So it is resolved with `litelink.current_metadata`, which +reads it outside DuckDB, and every query names the metadata file it returns. **One database per process, a connection per reader.** Loading `iceberg` into a fresh DuckDB database costs 400-580 ms (measured, `just bench-snapshot`), and @@ -26,17 +37,20 @@ from __future__ import annotations +import os import threading -from typing import TYPE_CHECKING +from dataclasses import dataclass +from typing import TYPE_CHECKING, Final from urllib.parse import unquote, urlsplit import pyarrow as pa -from litelink import S3Options, duckdb_connection +from litelink import S3Options, current_metadata, duckdb_connection from streamcast._log import COLUMN if TYPE_CHECKING: from collections.abc import Sequence + from os import PathLike import duckdb @@ -49,18 +63,66 @@ def path(uri: str) -> str: return uri +@dataclass(frozen=True, slots=True) +class ReadCache: + """How a reader caches what it reads: `litelink.duckdb_connection`'s four + settings, with its defaults. Hashable, because it is part of the key a + database is shared under.""" + + memory_cache: bool = True + """DuckDB's external file cache, for the database's life.""" + disk_cache: bool = False + """`cache_httpfs` on disk, surviving restarts. `s3://` tables only.""" + cache_key: str | None = None + """The disk cache's directory: relative to litelink's cache root, absolute + as given, or None for its `default` directory. The caller's choice — only + it knows what deserves a cache of its own, a stream id for one.""" + disk_cache_volume_limit: float = 0.8 + """How full the disk cache's VOLUME may get, everything on it counted.""" + + @classmethod + def of( + cls, + *, + memory_cache: bool, + disk_cache: bool, + cache_key: str | PathLike[str] | None, + disk_cache_volume_limit: float, + ) -> ReadCache: + """From the keywords the public calls take, which are litelink's.""" + return cls( + memory_cache, + disk_cache, + None if cache_key is None else os.fspath(cache_key), + disk_cache_volume_limit, + ) + + def keywords(self) -> dict[str, object]: + """Back to those keywords, for a call that takes them.""" + return { + "memory_cache": self.memory_cache, + "disk_cache": self.disk_cache, + "cache_key": self.cache_key, + "disk_cache_volume_limit": self.disk_cache_volume_limit, + } + + +DEFAULT_CACHE: Final = ReadCache() +"""litelink's defaults: memory on, disk off.""" + + def _quoted(text: str) -> str: return "'" + text.replace("'", "''") + "'" -# The databases readers connect to, by what they can read with. Never evicted: -# a process reads with a handful of credential sets at most. -_DATABASES: dict[S3Options | None, duckdb.DuckDBPyConnection] = {} +# The databases readers connect to, by what they can read with and how they +# cache. Never evicted: a process reads with a handful of each at most. +_DATABASES: dict[tuple[S3Options | None, ReadCache], duckdb.DuckDBPyConnection] = {} _LOCK = threading.Lock() def connection( - s3_options: S3Options | None, *, remote: bool + s3_options: S3Options | None, *, remote: bool, cache: ReadCache = DEFAULT_CACHE ) -> duckdb.DuckDBPyConnection: """A DuckDB connection that can read published tables. The caller closes it. @@ -77,15 +139,30 @@ def connection( `REFRESH auto`), so a database that outlives an STS session keeps reading. Keys rotated in the environment resolve to different options, and so to a database of their own. + + **And one per `cache`**, because a cache belongs to the database too: two + readers asking for different caching, or different directories, cannot + share one. A local database caches only in memory — the disk cache is + for `s3://` reads — so its disk settings do not split it. """ - key = (s3_options or S3Options()).resolved() if remote else None + credentials = (s3_options or S3Options()).resolved() if remote else None + if credentials is None: + cache = ReadCache(memory_cache=cache.memory_cache) + + key = (credentials, cache) with _LOCK: database = _DATABASES.get(key) if database is None: database = ( - duckdb_connection(s3_options=key) - if key is not None - else duckdb_connection() + duckdb_connection( + s3_options=credentials, + memory_cache=cache.memory_cache, + disk_cache=cache.disk_cache, + cache_key=cache.cache_key, + disk_cache_volume_limit=cache.disk_cache_volume_limit, + ) + if credentials is not None + else duckdb_connection(memory_cache=cache.memory_cache) ) _DATABASES[key] = database @@ -161,15 +238,8 @@ def open( else shared ) try: - table = path(uri) - hint = connected.execute( - f"SELECT content FROM read_text({_quoted(f'{table}/metadata/version-hint.text')})" - ).fetchone() - if hint is None: # pragma: no cover — read_text yields a row or raises - msg = f"{uri} has no version-hint.text" - raise FileNotFoundError(msg) - - metadata = f"{table}/metadata/{hint[0].strip()}.metadata.json" + # Not through `connected`: see the module docstring. + metadata = path(current_metadata(uri, s3_options=s3_options)) scan = f"iceberg_scan({_quoted(metadata)})" schema = connected.execute(f"SELECT * FROM {scan} LIMIT 0").arrow().schema low, high, count = connected.execute( diff --git a/src/streamcast/_snapshot.py b/src/streamcast/_snapshot.py index dc2f94a..a7c25f7 100644 --- a/src/streamcast/_snapshot.py +++ b/src/streamcast/_snapshot.py @@ -45,6 +45,7 @@ from streamcast._errors import NotReplayable, StreamcastError from streamcast._limits import MAX_TAIL from streamcast._metadata import Metadata +from streamcast._published import DEFAULT_CACHE, ReadCache if TYPE_CHECKING: from collections.abc import AsyncGenerator, Sequence @@ -440,6 +441,7 @@ async def snapshot( s3_options: S3Options | None = None, stream_id: str | None = None, max_tail: int = MAX_TAIL, + cache: ReadCache = DEFAULT_CACHE, ) -> Snapshot: """See `Stream.snapshot`.""" if as_of_offset is not None and as_of_ts is not None: @@ -457,11 +459,11 @@ async def snapshot( found = await asyncio.to_thread(metadata, metadata_uri, s3_options, stream_id) return ( await asyncio.to_thread( - _assemble, metadata_uri, found, as_of_offset, as_of_ts, s3_options + _assemble, metadata_uri, found, as_of_offset, as_of_ts, s3_options, cache ) if broker is None or as_of_offset is None else await _with_tail( - metadata_uri, found, as_of_offset, broker, s3_options, max_tail + metadata_uri, found, as_of_offset, broker, s3_options, max_tail, cache ) ) @@ -518,10 +520,11 @@ def _assemble( as_of_offset: int | None, as_of_ts: int | None, s3_options: S3Options | None, + cache: ReadCache = DEFAULT_CACHE, ) -> Snapshot: """Everything a snapshot needs but a broker: the pieces, the live pin, the end.""" remote = any((entry.published or "").startswith("s3://") for entry in found.logs) - connection = _published.connection(s3_options, remote=remote) + connection = _published.connection(s3_options, remote=remote, cache=cache) try: pieces = _pieces(found, as_of_ts) live = next((p for p in pieces if p.entry is found.live_log), None) @@ -603,12 +606,13 @@ async def _with_tail( broker: str, s3_options: S3Options | None, max_tail: int, + cache: ReadCache = DEFAULT_CACHE, ) -> Snapshot: """The published snapshot, then the broker's rows above it up to the point.""" from streamcast import _client # noqa: PLC0415 — the client imports this module published = await asyncio.to_thread( - _assemble, metadata_uri, found, None, None, s3_options + _assemble, metadata_uri, found, None, None, s3_options, cache ) start = published.end_offset try: diff --git a/src/streamcast/_stream.py b/src/streamcast/_stream.py index 09e9eec..a38bcb0 100644 --- a/src/streamcast/_stream.py +++ b/src/streamcast/_stream.py @@ -69,6 +69,7 @@ publish_ack, publish_error, ) +from streamcast._published import ReadCache from streamcast._stats import Stats from streamcast._subscriber import Subscriber from streamcast._writer import Job, Writer @@ -714,6 +715,10 @@ async def snapshot( broker: str | None = None, s3_options: S3Options | None = None, max_tail: int = MAX_TAIL, + memory_cache: bool = True, + disk_cache: bool = False, + cache_key: str | PathLike[str] | None = None, + disk_cache_volume_limit: float = 0.8, ) -> _snapshot.Snapshot: """A stream's history as of one point, read from its published tables. @@ -727,6 +732,18 @@ async def snapshot( does not write: `scan`, `sql` and `rows`, then `close` (or `async with`). See `_snapshot` for what each point means, when the broker is consulted, and what is refused rather than answered short. + + **Caching is litelink's, and the caller's choice.** `memory_cache` + (on) keeps what was read in DuckDB's external file cache for the + process; `disk_cache` (off) keeps `s3://` reads on disk with + `cache_httpfs`, across restarts, under `cache_key` — a directory + relative to litelink's cache root (a stream id, say), absolute as + given, or None for its shared `default`. `disk_cache_volume_limit` is + how full that disk may get, everything on it counted. Only the + caller knows what deserves a cache of its own, so nothing is keyed + for it. Readers with different settings read through different + databases; see `litelink.duckdb_connection`. The same keywords are on + `scan`, `sql`, `live` and `connect` (for a catch-up). """ return await _snapshot.snapshot( metadata_uri, @@ -735,6 +752,12 @@ async def snapshot( broker=broker, s3_options=s3_options, max_tail=max_tail, + cache=ReadCache.of( + memory_cache=memory_cache, + disk_cache=disk_cache, + cache_key=cache_key, + disk_cache_volume_limit=disk_cache_volume_limit, + ), ) @staticmethod @@ -746,6 +769,10 @@ async def scan( broker: str | None = None, s3_options: S3Options | None = None, max_tail: int = MAX_TAIL, + memory_cache: bool = True, + disk_cache: bool = False, + cache_key: str | PathLike[str] | None = None, + disk_cache_volume_limit: float = 0.8, columns: Sequence[str] | None = None, where: str | None = None, filters: Sequence[_manifest.Term] = (), @@ -760,6 +787,10 @@ async def scan( broker=broker, s3_options=s3_options, max_tail=max_tail, + memory_cache=memory_cache, + disk_cache=disk_cache, + cache_key=cache_key, + disk_cache_volume_limit=disk_cache_volume_limit, ) as snap: return await snap.scan( columns=columns, @@ -779,6 +810,10 @@ async def sql( broker: str | None = None, s3_options: S3Options | None = None, max_tail: int = MAX_TAIL, + memory_cache: bool = True, + disk_cache: bool = False, + cache_key: str | PathLike[str] | None = None, + disk_cache_volume_limit: float = 0.8, filters: Sequence[_manifest.Term] = (), start_offset: int | None = None, end_offset: int | None = None, @@ -795,6 +830,10 @@ async def sql( broker=broker, s3_options=s3_options, max_tail=max_tail, + memory_cache=memory_cache, + disk_cache=disk_cache, + cache_key=cache_key, + disk_cache_volume_limit=disk_cache_volume_limit, ) as snap: return await snap.sql( query, @@ -812,6 +851,10 @@ async def live( where: dict[str, object] | None = None, start_offset: int | None = None, max_tail: int = MAX_TAIL, + memory_cache: bool = True, + disk_cache: bool = False, + cache_key: str | PathLike[str] | None = None, + disk_cache_volume_limit: float = 0.8, ) -> _live.Live: """A stream's history kept current in memory: `scan` and `sql` as of now. @@ -839,6 +882,12 @@ async def live( where=where, start_offset=start_offset, max_tail=max_tail, + cache=ReadCache.of( + memory_cache=memory_cache, + disk_cache=disk_cache, + cache_key=cache_key, + disk_cache_volume_limit=disk_cache_volume_limit, + ), ) @property diff --git a/tests/test_catchup.py b/tests/test_catchup.py index 83f72d4..84e8fb5 100644 --- a/tests/test_catchup.py +++ b/tests/test_catchup.py @@ -199,6 +199,41 @@ async def test_batches_cross_the_join_whole_and_in_order( # connection. Deterministic: no handshake completes in one loop turn. assert [len(batch) for batch in got[:2]] == [1500, TOTAL - 1500] + async def test_a_catch_up_reads_with_the_callers_cache_settings( + self, serve, published_log, s3, tmp_path, monkeypatch + ): + """`connect`'s cache keywords reach the catch-up's reads (#77).""" + from streamcast import _published # noqa: PLC0415 + + asked: list[_published.ReadCache] = [] + real = _published.connection + + def recording(s3_options, *, remote, cache=_published.DEFAULT_CACHE): # noqa: ANN001, ANN202 + asked.append(cache) + return real(s3_options, remote=remote, cache=cache) + + monkeypatch.setattr(_published, "connection", recording) + stream = streamcast.Stream("trades", log=published_log, max_replay=600) + await fill(stream, published_log) + + async with serve(stream, maintain=False) as uri: + async with streamcast.connect( + uri, + offset=1, + catch_up=True, + s3_options=s3, + disk_cache=True, + cache_key=tmp_path / "cache", + ) as sub: + assert (await sub.recv())[0] == 1 + + assert asked + assert all( + cache + == _published.ReadCache(disk_cache=True, cache_key=str(tmp_path / "cache")) + for cache in asked + ) + @pytest.mark.slow async def test_nothing_published_during_the_catch_up_is_lost( self, serve, published_log, s3 diff --git a/tests/test_published.py b/tests/test_published.py index 706d81b..13e7fd1 100644 --- a/tests/test_published.py +++ b/tests/test_published.py @@ -155,3 +155,123 @@ def test_an_s3_table_reads_with_the_readers_own_credentials(tmp_path, s3, bucket assert table.scan().read_all().num_rows == 4 finally: table.close() + + +def publish(log: litelink.WriteHandle, start: int, count: int) -> None: + log.extend([{"i": i} for i in range(start, start + count)]) + while log.seal(flush=True) is not None: + pass + + log.publish(flush=True) + + +class TestTheReadCache: + """litelink's read caches, as the caller asks for them (#77).""" + + def test_different_caching_never_shares_a_database(self): + cached = _published.connection(None, remote=False) + uncached = _published.connection( + None, remote=False, cache=_published.ReadCache(memory_cache=False) + ) + # Disk settings mean nothing to a local read, so they do not split it. + disk = _published.connection( + None, + remote=False, + cache=_published.ReadCache(disk_cache=True, cache_key="anything"), + ) + setting = "SELECT current_setting('enable_external_file_cache')" + try: + assert cached.execute(setting).fetchone() == (True,) + assert uncached.execute(setting).fetchone() == (False,) + cached.execute("CREATE OR REPLACE TABLE cache_probe AS SELECT 1 AS x") + assert disk.execute("SELECT x FROM cache_probe").fetchall() == [(1,)] + with pytest.raises(Exception, match="cache_probe"): + uncached.execute("SELECT x FROM cache_probe") + finally: + cached.execute("DROP TABLE cache_probe") + for connected in (cached, uncached, disk): + connected.close() + + def test_a_disk_cached_reader_sees_every_publish(self, tmp_path, s3, bucket): + """The hint is the one object that changes, and `cache_httpfs` would + serve an old one for ever (litelink#141): a reader pinned to the + first snapshot it saw. So it is read around the cache.""" + cache = tmp_path / "cache" + connected = _published.connection( + s3, + remote=True, + cache=_published.ReadCache(disk_cache=True, cache_key=str(cache)), + ) + log = litelink.new( + tmp_path / "data", "trades", schema=SCHEMA, published=bucket, s3_options=s3 + ) + try: + with log: + publish(log, 0, 100) + first = _published.Table.open(bucket, "trades", s3, shared=connected) + assert first.record_count == 100 + + publish(log, 100, 50) + again = _published.Table.open(bucket, "trades", s3, shared=connected) + assert again.record_count == 150 + assert again.extent == (1, 151) + + # And the disk cache is on: what was read is under the key. + assert any(path.is_file() for path in cache.rglob("*")) + finally: + connected.close() + + async def test_the_settings_reach_every_read_of_a_stream( + self, tmp_path, serve, monkeypatch + ): + """`Stream.snapshot`, `live` — its base and every rebase — pass the + caller's settings to the connection every published read goes through.""" + import streamcast # noqa: PLC0415 + + asked: list[_published.ReadCache] = [] + real = _published.connection + + def recording(s3_options, *, remote, cache=_published.DEFAULT_CACHE): # noqa: ANN001, ANN202 + asked.append(cache) + return real(s3_options, remote=remote, cache=cache) + + monkeypatch.setattr(_published, "connection", recording) + 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(3)]) + assert stream.log is not None + while stream.log.seal(flush=True) is not None: + pass + + stream.log.publish(flush=True) + key = tmp_path / "k" + wanted = _published.ReadCache(False, True, str(tmp_path / "k"), 0.5) + try: + async with serve(stream, maintain=False) as uri: + assert stream.metadata_uri is not None + async with await streamcast.Stream.snapshot( + stream.metadata_uri, + memory_cache=False, + disk_cache=True, + cache_key=key, + disk_cache_volume_limit=0.5, + ): + pass + + async with await streamcast.Stream.live( + uri, + memory_cache=False, + disk_cache=True, + cache_key=key, + disk_cache_volume_limit=0.5, + ) as live: + await live.rebase() + + # A snapshot, the live base, and the rebase. + assert asked == [wanted] * 3 + finally: + await stream.aclose() diff --git a/uv.lock b/uv.lock index ef70288..9c243b3 100644 --- a/uv.lock +++ b/uv.lock @@ -700,7 +700,7 @@ wheels = [ [[package]] name = "litelink" -version = "0.8.1" +version = "0.9.0" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "duckdb" }, @@ -708,12 +708,12 @@ dependencies = [ { name = "pyiceberg", extra = ["pyarrow", "sql-sqlite"] }, { name = "xxhash" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/68/27/4be7f829057fd2c06fc28307a2ae8453af19ca53d33caadd4eb3e26e4db4/litelink-0.8.1.tar.gz", hash = "sha256:89c61682c3ab0f00336dd8e4f6b104a149cb37117ad392b595901f168d10b06e", size = 333546, upload-time = "2026-10-03T19:48:28.631Z" } +sdist = { url = "https://files.pythonhosted.org/packages/b7/ed/78869573728d3629c81df24dd3336ab78cf39e5359e2f873a5986e42709c/litelink-0.9.0.tar.gz", hash = "sha256:1cc772ef187126244ab4de12c99dd553b0b398fd686d874e1689e4f05b1fbf2b", size = 334817, upload-time = "2026-10-04T08:26:53.513Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/0d/f1/236e3e040d5012dfa7b87fa89592ce3a078e6a566bc6f65abf65373cdae7/litelink-0.8.1-py3-none-macosx_11_0_arm64.whl", hash = "sha256:829be2c6df8c7cfdd57834c713e9916eca246f0d1616888052ffa2dbfbd10695", size = 44218172, upload-time = "2026-10-03T19:48:13.2Z" }, - { url = "https://files.pythonhosted.org/packages/14/dd/aa04e4a5ef450c47d9a5dee2bf445887b2cb81bded1d93e72208b6ac9441/litelink-0.8.1-py3-none-macosx_11_0_x86_64.whl", hash = "sha256:ea16ea2dee1f2ce8807e8b5acabcc89362fded9fd538541b84f1c62e90004f43", size = 46939767, upload-time = "2026-10-03T19:48:17.01Z" }, - { url = "https://files.pythonhosted.org/packages/1a/07/0a162925dcd847d976e747fc989892e163d20bc8d069dec72be377bf35e9/litelink-0.8.1-py3-none-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:14c9df92d2e93a4df3b3702e6306e1968af10bfaffdad192153bf6d8ba3898ef", size = 56088773, upload-time = "2026-10-03T19:48:21.03Z" }, - { url = "https://files.pythonhosted.org/packages/e7/87/4860e2d493f7f08fc805f86153deb57171c1249c4a89cadc2b76b387d5de/litelink-0.8.1-py3-none-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:14f0570bc8bea50fa87336e6deeef4b64b6ec5d73601ccdae5eab53f44c116a9", size = 59986772, upload-time = "2026-10-03T19:48:25.721Z" }, + { url = "https://files.pythonhosted.org/packages/76/d0/ecf290860c88540ef5c83205ac4792a2a4940a7cbf629a623c670c54e7ff/litelink-0.9.0-py3-none-macosx_11_0_arm64.whl", hash = "sha256:907558f6f251da04027b568761eee6e5dd681ab0233bdf0f00edb3b59e65e146", size = 44219326, upload-time = "2026-10-04T08:26:31.605Z" }, + { url = "https://files.pythonhosted.org/packages/fa/f5/e6e38817be74ef44e43c9cd55d553a4aceedbd118e71dc056047d5a01ded/litelink-0.9.0-py3-none-macosx_11_0_x86_64.whl", hash = "sha256:d3a5f5dd26b8ead8d6dd7eea63cdd562dfe77239f1683f4f74c5e43102a44184", size = 46940922, upload-time = "2026-10-04T08:26:37.06Z" }, + { url = "https://files.pythonhosted.org/packages/60/a6/dc7778a650d1c9e4e1509607a5edde406972f25e66764302284860a75ab4/litelink-0.9.0-py3-none-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:42050220a3aeec55454a76a35a4791e3fbaba5ed00db84b3b17c6f383e3eb8af", size = 56089929, upload-time = "2026-10-04T08:26:43.398Z" }, + { url = "https://files.pythonhosted.org/packages/3e/20/2b813be4090c27816fde13e9eab9aa73bb11ac69129016f47778a8a4ad60/litelink-0.9.0-py3-none-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:47e6ca3a605280bb9668bc83ab966c0b0c135de6b7da72d8744bdd6102af9ff9", size = 59987930, upload-time = "2026-10-04T08:26:50.854Z" }, ] [[package]] @@ -2003,7 +2003,7 @@ dev = [ [package.metadata] requires-dist = [ - { name = "litelink", specifier = ">=0.8.1,<0.9" }, + { name = "litelink", specifier = ">=0.9.0,<0.10" }, { name = "msgspec", specifier = ">=0.18" }, { name = "starlette", marker = "extra == 'asgi'", specifier = ">=0.40" }, { name = "websockets", specifier = ">=14" },