From 68bf34db33264539def577f31010e9384c024b9c Mon Sep 17 00:00:00 2001 From: Tam Nguyen Duc <1218621+tamnd@users.noreply.github.com> Date: Wed, 19 Aug 2026 17:17:15 +0700 Subject: [PATCH 1/2] Hand the rows over as the engine makes them `execute` runs a statement to the end and gives back every row, which is the wrong shape for a result bigger than the memory meant for it and for a reader that wants the first rows before the last are made. `conn.stream` is the other shape: rows a batch at a time, while the statement is still running. The statement goes on a thread of its own rather than the connection's runner, because a stream lasts as long as its reader and the runner is for statements that end. The rows are copied out of the borrowed batch on that thread with the GIL down and turned into Python objects on the reader's, and the two are joined by a queue two batches deep, so the engine fills the next batch while the loop reads this one. Over a million people in two columns on this machine the first row arrives in 0.7 ms against 30 ms for `execute`, reading all of them takes 214 ms against 256, a batch at a time takes 176, and the Python side peaks at nothing worth measuring against 153 MB. A stream holds the connection until it ends, since a connection runs one statement at a time, so a statement run on it while one is open is refused rather than queued behind a loop that may never finish. Closing the connection hangs the stream up first, which is what lets the close return rather than wait for a scan nobody is reading. A reader that stopped a stream early can still ask what it stopped, because the summary is kept beside the queue rather than in it. `zudb.aio` gets the same surface as `AsyncStream`, where waiting for a batch is handed off the loop like every other wait there. Twenty seven tests: twenty two on the stream and five on the async spelling, covering the rows, the columns before a row is read, a first row while the summary is still `None`, the batch sizes, the refusal, the summary of a run that ended and one that was stopped, the buffered case that `ORDER BY` is, the failures, and the memory. --- README.md | 33 +- python/zudb/__init__.py | 6 + python/zudb/_zudb.pyi | 68 ++++ python/zudb/aio.py | 170 ++++++++- src/conn.rs | 59 +++ src/interrupt.rs | 2 +- src/lib.rs | 4 + src/stream.rs | 815 ++++++++++++++++++++++++++++++++++++++++ src/value.rs | 2 +- tests/test_aio.py | 80 ++++ tests/test_stream.py | 330 ++++++++++++++++ 11 files changed, 1565 insertions(+), 4 deletions(-) create mode 100644 src/stream.rs create mode 100644 tests/test_stream.py diff --git a/README.md b/README.md index 8d13e30..36eaf80 100644 --- a/README.md +++ b/README.md @@ -126,6 +126,37 @@ result.record_batches() # a reader, for a result larger than memory `Result` implements `__arrow_c_stream__`, so anything that reads the protocol reads a result directly and none of the four methods above is needed: `pyarrow.table(result)` and `polars.DataFrame(result)` both work. Batches are 65,536 rows. A column holds one type, which the values decide, and integers beside floats are the one mixture that widens rather than being refused. Nodes, rels and paths go across as structs. The copy runs with the GIL released, and on this machine 300,000 rows across three columns take 44 ms as Arrow against 67 ms as Python objects, and a single integer column takes 13.8 ms against 44.5 ms. +## Reading a result as it arrives + +`execute` runs the statement to the end and hands back every row. `stream` hands back the rows as the engine makes them, which is what you want when the result is bigger than the memory you meant to spend on it, or when the first rows are worth having before the last are made. + +```python +with conn.stream("MATCH (p:person) RETURN p.uid AS uid, p.name AS name") as rows: + for uid, name in rows: + write(uid, name) +``` + +```python +with conn.stream("MATCH (p:person) RETURN p.uid AS uid", batch_rows=10_000) as rows: + for batch in rows.batches(): + write_many(batch) + rows.summary.rows # how many were read +``` + +The statement runs on a thread of its own with the GIL down and hands each batch over a queue two batches deep, so the engine is filling the next batch while your loop is reading this one and neither waits for the other for long. On this machine, over a million people in two columns, the first row arrives 0.7 ms after the call against 30 ms for `execute`, and reading all of them takes 214 ms streamed against 256 ms in one piece, which is 4.7 million rows a second against 3.9 million. Reading them a batch at a time takes 176 ms, or 5.7 million a second, because a batch is one call where a row is one call. Streaming being the faster of the two is not a trick: the rows are made and consumed while the cache still has them, and nothing has to hold a million tuples at once. Held is the whole difference: the Python side of `fetchall` peaks at 153 MB for that result and the same read streamed peaks at nothing worth measuring, since every tuple is freed as the loop moves past it. + +A stream holds the connection until it ends, because a connection runs one statement at a time. Reading it to the end frees the connection, so does `close`, and so does the end of a `with` block however it was left. A statement run on the connection while a stream is open is refused with a message that says so rather than queued behind a loop that may never finish, which is the deadlock the refusal exists to prevent. If a program needs to read a stream and run statements at the same time, that is two connections. + +`summary` is `None` while the statement runs and afterwards says what it did: the columns, how many rows were handed over, whether the reader stopped it early, and whether the rows arrived as they were made. That last one is worth reading. A statement that has to see every row before it can give one, which is `ORDER BY`, `DISTINCT` and the aggregates, is run whole by the engine and handed over in batches afterwards, and `streamed` is `False` for it. The loop over it is the same loop and what differs is what it cost. + +`zudb.aio` has the same thing, where the waits go off the loop like every other wait there: + +```python +async with conn.stream("MATCH (p:person) RETURN p.name AS name") as rows: + async for (name,) in rows: + await write(name) +``` + ## Preparing a statement A statement a program runs many times with different values can be compiled once and kept. @@ -286,7 +317,7 @@ Half of this package is compiled, which is the one thing an inspection cannot se ## What works today -The list above is what this client is for. What it does so far is the core of it: `connect`, `execute` and `sql` with named parameters, results that iterate and fetch, values as Python objects both ways including dates, times, datetimes and durations, `Node`, `Rel` and `Path` as classes, `load` for building a graph with edges in it, an appender for growing one, transactions as a context manager that commits at the end of a block and rolls back when it raises, every condition as an exception class carrying its code, its position and its documentation link, results as Arrow columns and as pandas and polars frames, `register` for putting a frame under a name a statement can match on and reading it where it lies, stubs inside the wheel with a gate that keeps them true, the GIL released around every statement, every load and every copy out, `Ctrl-C` and `interrupt()` stopping a statement without touching the connection under it, `zudb.aio` for the same calls awaited on an event loop, results, nodes, rels and paths that draw themselves in a notebook with `%gql` and `%%gql` to run statements in one, `zudb.dbapi` for code written against PEP 249, and `prepare`, `explain` and `profile` for a statement compiled once, the plan it would run and the plan it did. Each one landed with the tests that say it works. +The list above is what this client is for. What it does so far is the core of it: `connect`, `execute` and `sql` with named parameters, results that iterate and fetch, values as Python objects both ways including dates, times, datetimes and durations, `Node`, `Rel` and `Path` as classes, `load` for building a graph with edges in it, an appender for growing one, transactions as a context manager that commits at the end of a block and rolls back when it raises, every condition as an exception class carrying its code, its position and its documentation link, results as Arrow columns and as pandas and polars frames, `register` for putting a frame under a name a statement can match on and reading it where it lies, stubs inside the wheel with a gate that keeps them true, the GIL released around every statement, every load and every copy out, `Ctrl-C` and `interrupt()` stopping a statement without touching the connection under it, `zudb.aio` for the same calls awaited on an event loop, results, nodes, rels and paths that draw themselves in a notebook with `%gql` and `%%gql` to run statements in one, `zudb.dbapi` for code written against PEP 249, `prepare`, `explain` and `profile` for a statement compiled once, the plan it would run and the plan it did, and `stream` for a result read as the engine makes it rather than after it has made all of it. Each one landed with the tests that say it works. ## Wheels diff --git a/python/zudb/__init__.py b/python/zudb/__init__.py index d04594e..ed29345 100644 --- a/python/zudb/__init__.py +++ b/python/zudb/__init__.py @@ -34,6 +34,9 @@ Rel, Result, ScalarPlan, + Stream, + StreamBatches, + StreamSummary, Transaction, __abi_version__, connect, @@ -61,6 +64,9 @@ "Appender", "Prepared", "Result", + "Stream", + "StreamBatches", + "StreamSummary", "Plan", "PlanNode", "ScalarPlan", diff --git a/python/zudb/_zudb.pyi b/python/zudb/_zudb.pyi index f43d95b..3070332 100644 --- a/python/zudb/_zudb.pyi +++ b/python/zudb/_zudb.pyi @@ -85,6 +85,15 @@ class Connection: def profile(self, statement: str, params: Mapping[str, Value] | None = None) -> Profile: """Runs the statement with the counters on and answers what its operators really did.""" + def stream( + self, + statement: str, + params: Mapping[str, Value] | None = None, + *, + batch_rows: int | None = None, + ) -> Stream: + """Runs one statement and hands back its rows as the engine makes them.""" + def transaction(self, *, read_only: bool = False) -> Transaction: """Starts a transaction and hands it back for a `with` block.""" @@ -199,6 +208,65 @@ class Prepared: def __exit__(self, *_exception: object) -> bool: ... def __repr__(self) -> str: ... +class Stream: + """Rows arriving one batch at a time, while the statement is still running.""" + + @property + def columns(self) -> list[str]: + """The column names, in the order the statement projects them.""" + + @property + def summary(self) -> StreamSummary | None: + """What the statement did, once it has done it, and `None` while it is still running.""" + + def batches(self) -> StreamBatches: + """The rows in the batches they arrived in, as lists of tuples.""" + + def close(self) -> None: + """Stops the statement and gives the connection back.""" + + @property + def closed(self) -> bool: + """Whether the statement is over, by running out of rows or by being closed.""" + + def __iter__(self) -> Iterator[tuple[Value, ...]]: ... + def __next__(self) -> tuple[Value, ...]: ... + def __enter__(self) -> Stream: ... + def __exit__(self, *_exception: object) -> bool: ... + def __repr__(self) -> str: ... + +class StreamBatches: + """The same rows, in the batches they arrived in.""" + + def __iter__(self) -> Iterator[list[tuple[Value, ...]]]: ... + def __next__(self) -> list[tuple[Value, ...]]: ... + def __repr__(self) -> str: ... + +class StreamSummary: + """What a streamed statement did, known once it has ended.""" + + @property + def columns(self) -> list[str]: + """The column names, in the order the statement projected them.""" + + @property + def rows(self) -> int: + """How many rows were handed over.""" + + @property + def stopped(self) -> bool: + """Whether the reader stopped it before it ran out of rows.""" + + @property + def streamed(self) -> bool: + """Whether the rows arrived as they were made.""" + + @property + def notices(self) -> list[dict[str, str]]: + """The warnings the statement raised, in the shape a result reports them.""" + + def __repr__(self) -> str: ... + class PlanNode: """One operator of a plan.""" diff --git a/python/zudb/aio.py b/python/zudb/aio.py index 23d0dd8..7b4d680 100644 --- a/python/zudb/aio.py +++ b/python/zudb/aio.py @@ -49,7 +49,18 @@ from typing import Any, Generic, TypeVar from . import _zudb -from ._zudb import Appender, Connection, Plan, Prepared, Profile, Result, Transaction +from ._zudb import ( + Appender, + Connection, + Plan, + Prepared, + Profile, + Result, + Stream, + StreamBatches, + StreamSummary, + Transaction, +) from .types import Value __all__ = [ @@ -58,10 +69,18 @@ "AsyncTransaction", "AsyncAppender", "AsyncPrepared", + "AsyncStream", + "AsyncStreamBatches", ] T = TypeVar("T") +#: Handed to `next` so that the end of a stream arrives as a value on the +#: connection's thread. A `StopIteration` raised there would cross a +#: future on its way back, and a `StopIteration` crossing a future is the +#: one exception asyncio cannot let through. +_NOTHING = object() + def connect( path: str | os.PathLike[str], @@ -252,6 +271,45 @@ async def profile(self, statement: str, params: Mapping[str, Value] | None = Non """ return await self._call(functools.partial(self._conn.profile, statement, params)) + def stream( + self, + statement: str, + params: Mapping[str, Value] | None = None, + *, + batch_rows: int | None = None, + ) -> _Opening[AsyncStream]: + """Runs one statement and hands back its rows as the engine + makes them. + + async with conn.stream("MATCH (p:person) RETURN p.name AS name") as rows: + async for (name,) in rows: + await write(name) + + This is the call for a result too big to want in memory and for + one whose first rows are worth having before the last are made. + The statement runs on a thread of its own, not this connection's, + and every wait for a batch is handed off the loop like every + other wait here, so a task reading a stream leaves the loop free + between batches. + + A stream holds the connection until it ends, so a statement run + on the same connection while one is open is refused rather than + queued behind a loop that may never finish. Read it to the end, + close it, or open it with `async with`. + """ + return _Opening(functools.partial(self._stream, statement, params, batch_rows)) + + async def _stream( + self, + statement: str, + params: Mapping[str, Value] | None, + batch_rows: int | None, + ) -> AsyncStream: + opened = await self._call( + functools.partial(self._conn.stream, statement, params, batch_rows=batch_rows) + ) + return AsyncStream(self, opened) + def transaction(self, *, read_only: bool = False) -> _Opening[AsyncTransaction]: """Starts a transaction and hands it back for an `async with` block. @@ -597,3 +655,113 @@ async def __aexit__(self, *_exception: Any) -> bool: def __repr__(self) -> str: closed = ", closed" if self.closed else "" return f"" + + +class AsyncStream: + """Rows arriving one batch at a time, awaited a row at a time. + + Take one with `AsyncConnection.stream`. The stream underneath is + `zudb.Stream` and the rules are its rules: it holds the connection + until it ends, closing it twice does nothing, and the summary is + there once the statement is over whether it ran out of rows or was + stopped. + + Waiting for the next batch is handed to the connection's thread, so + the loop is free while the engine is working. The rows of a batch + already in hand are turned into tuples on the loop's thread, which is + the same work `Result` does when it is read, and there is nothing to + wait for in it. + """ + + __slots__ = ("_conn", "_stream") + + def __init__(self, conn: AsyncConnection, stream: Stream) -> None: + self._conn = conn + self._stream = stream + + async def columns(self) -> list[str]: + """The column names, in the order the statement projects them. + + A method where the sync client has a property, because answering + reads the first batch and reading a batch can wait. + """ + return await self._conn._call(lambda: self._stream.columns) + + @property + def summary(self) -> StreamSummary | None: + """What the statement did, once it has done it, and `None` while + it is still running. + + A property and not a coroutine: the answer is beside the queue + rather than behind the engine, so asking reaches nothing that + could wait. + """ + return self._stream.summary + + @property + def closed(self) -> bool: + """Whether the statement is over, by running out of rows or by + being closed. + """ + return self._stream.closed + + def batches(self) -> AsyncStreamBatches: + """The rows in the batches they arrived in, as lists of tuples. + + For a writer with a size of its own: one await per batch rather + than one per row, and one call into whatever is being written to. + """ + return AsyncStreamBatches(self._conn, self._stream.batches()) + + async def close(self) -> None: + """Stops the statement and gives the connection back. + + Awaited because it waits for the statement to stop, so the + connection is free by the time the call returns. Doing it twice + is not an error. + """ + await self._conn._call(self._stream.close) + + def __aiter__(self) -> AsyncStream: + return self + + async def __anext__(self) -> tuple[Value, ...]: + row = await self._conn._call(functools.partial(next, self._stream, _NOTHING)) + if row is _NOTHING: + raise StopAsyncIteration + return row # type: ignore[return-value] + + async def __aenter__(self) -> AsyncStream: + return self + + async def __aexit__(self, *_exception: Any) -> bool: + """Closes on the way out, whether the block ended well or badly, + so the connection comes back either way. + """ + await self.close() + return False + + def __repr__(self) -> str: + return f"" + + +class AsyncStreamBatches: + """The same rows, in the batches they arrived in.""" + + __slots__ = ("_conn", "_batches") + + def __init__(self, conn: AsyncConnection, batches: StreamBatches) -> None: + self._conn = conn + self._batches = batches + + def __aiter__(self) -> AsyncStreamBatches: + return self + + async def __anext__(self) -> list[tuple[Value, ...]]: + batch = await self._conn._call(functools.partial(next, self._batches, _NOTHING)) + if batch is _NOTHING: + raise StopAsyncIteration + return batch # type: ignore[return-value] + + def __repr__(self) -> str: + return f"" diff --git a/src/conn.rs b/src/conn.rs index 54e1600..9709e1b 100644 --- a/src/conn.rs +++ b/src/conn.rs @@ -25,6 +25,7 @@ use crate::interrupt; use crate::plan; use crate::prepared::Prepared; use crate::register; +use crate::stream; use crate::txn::Transaction; use crate::value::{Names, from_py, to_py}; @@ -81,6 +82,12 @@ pub struct Connection { /// asking it to stop, should not queue behind a ten second /// statement. alive: AtomicBool, + /// The stream holding the connection, if one is. A stream ends when + /// its reader says so, and a statement that queued behind a + /// half-read one would be waiting for a loop that is waiting for + /// it, so every other statement asks this and is told no rather + /// than left to deadlock. + pub(crate) feeding: Arc, #[pyo3(get)] path: PathBuf, #[pyo3(get)] @@ -117,6 +124,46 @@ impl Connection { self.execute(py, statement, params) } + /// Runs one statement and hands back its rows a batch at a time, as + /// the executor makes them. + /// + /// ```python + /// with conn.stream("MATCH (p:person) RETURN p.name AS name") as rows: + /// for (name,) in rows: + /// print(name) + /// ``` + /// + /// What is in memory is a batch and not the answer, which is what + /// this is for: a statement over ten million rows read by a program + /// holding a thousand. A reader that stops early stops the scan, + /// which is the other half of it, and `batch_rows` says how many + /// rows a batch may hold when the rows are going somewhere with a + /// size of its own. + /// + /// The stream holds the connection until it ends, because a + /// connection runs one statement at a time. Read it to the end, + /// close it, or open it in a `with` block, which closes it however + /// the block is left. + #[pyo3(signature = (statement, params = None, *, batch_rows = None))] + fn stream( + &self, + py: Python<'_>, + statement: &str, + params: Option<&Bound<'_, PyDict>>, + batch_rows: Option, + ) -> PyResult { + stream::open( + py, + &self.inner, + &self.stop, + &self.feeding, + self.alive.load(Ordering::Acquire), + statement.to_string(), + bind(params)?, + batch_rows, + ) + } + /// Compiles a statement now and hands back something that runs it /// later, as often as you like, with different values bound each /// time. @@ -340,6 +387,11 @@ impl Connection { /// Doing it twice is not an error, because a `with` block that /// closed early would otherwise fail on the way out. fn close(&self, py: Python<'_>) { + // A stream is told first. Closing waits for the connection's + // lock and a stream holds it for as long as its reader takes, + // so a connection closed under a half-read one would wait for a + // loop nobody is going to run again. + self.feeding.hang_up(); // Released here too, because closing waits for the statement // another thread is running and frees caches worth megabytes // once it has. @@ -422,6 +474,7 @@ impl Connection { inner: Arc::new(Mutex::new(Some(opened))), runner: OnceLock::new(), alive: AtomicBool::new(true), + feeding: Arc::new(stream::Feeding::new()), path, read_only, }) @@ -476,6 +529,9 @@ impl Connection { source: Source, params: Vec<(String, Value)>, ) -> PyResult { + if self.feeding.busy() { + return Err(programming(py, stream::STREAMING)); + } // The GIL goes down for the whole statement, waiting for the // connection's own lock included. That is the point of a // compiled engine in a Python process: another thread runs @@ -538,6 +594,9 @@ impl Connection { if !self.alive.load(Ordering::Acquire) { return Err(closed(py, "this connection")); } + if self.feeding.busy() { + return Err(programming(py, stream::STREAMING)); + } py.detach(|| { let mut held = self.inner.lock().map_err(|_| ())?; let conn = held.as_mut().ok_or(())?; diff --git a/src/interrupt.rs b/src/interrupt.rs index 48f8900..327cac4 100644 --- a/src/interrupt.rs +++ b/src/interrupt.rs @@ -369,7 +369,7 @@ thread_local! { /// the limited ABI these wheels are built against does not expose the /// second one, and asking Python is a couple of attribute lookups on a /// module that is imported long before anything here runs. -fn on_main_thread(py: Python<'_>) -> bool { +pub(crate) fn on_main_thread(py: Python<'_>) -> bool { IS_MAIN.with(|known| { if let Some(known) = known.get() { return known; diff --git a/src/lib.rs b/src/lib.rs index 6c5bd75..7004c0a 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -24,6 +24,7 @@ mod load; mod plan; mod prepared; mod register; +mod stream; mod txn; mod value; @@ -62,6 +63,9 @@ fn _zudb(module: &Bound<'_, PyModule>) -> PyResult<()> { module.add_class::()?; module.add_class::()?; module.add_class::()?; + module.add_class::()?; + module.add_class::()?; + module.add_class::()?; module.add_class::()?; module.add_class::()?; module.add_class::()?; diff --git a/src/stream.rs b/src/stream.rs new file mode 100644 index 0000000..f02183f --- /dev/null +++ b/src/stream.rs @@ -0,0 +1,815 @@ +//! A statement read a batch at a time, instead of all at once. +//! +//! A result that does not fit in memory is the reason this exists, and +//! a reader that will not read all of it is the reason it stops +//! properly. The engine's shape for both is a sink: it hands over a +//! batch of rows, the sink says whether it wants more, and a sink that +//! says no ends the scan at the boundary an interrupt is answered at. +//! That is a push and Python wants a pull, so the two are joined by a +//! queue of two batches and a thread of this statement's own. +//! +//! A thread of its own rather than the one the connection keeps for +//! interruptible statements, because a stream lasts as long as its +//! reader takes and that thread is the connection's: a statement parked +//! on it waiting for the body of a `for` loop to come round again would +//! be every later statement on the connection waiting too. The queue is +//! bounded because that is what backpressure is. A reader that stops +//! reading stops the scan two batches later rather than pulling a +//! database into memory behind it. +//! +//! The rows are copied out of the batch on the statement's thread, +//! where the GIL is not held, and turned into Python objects on the +//! reader's thread, where it is. So the copy is a batch and never the +//! answer, and the interpreter is free for the whole time the executor +//! is working. + +use std::collections::VecDeque; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::{Arc, Condvar, Mutex}; +use std::time::Duration; + +use pyo3::prelude::*; +use pyo3::types::{PyDict, PyList, PyTuple}; +use zudb::query::Value; +use zudb::{Batch, Flow, Interrupt, Streamed}; + +use crate::error::{closed, programming, to_py_err}; +use crate::interrupt::on_main_thread; +use crate::value::{Names, to_py}; + +/// How many batches may sit between the statement and the reader. +/// +/// Two, so that one is being turned into Python objects while the next +/// is being made, and no more, because everything past that is memory +/// spent to hide a reader slower than the scan. It is the whole of the +/// backpressure and it is deliberately small. +const QUEUE: usize = 2; + +/// How long a reader on the main thread waits for a batch before it +/// asks Python whether a `Ctrl-C` arrived. +/// +/// The same two milliseconds every other statement waits in, for the +/// same reason: the budget from the press to the exception is fifty +/// milliseconds and most of a press is this wait. +const TICK: Duration = Duration::from_millis(2); + +/// How long `close` waits for the statement to notice that nobody is +/// reading before it stops it outright. +/// +/// A stream that is handing rows over notices at its next batch, which +/// is immediately, and ends with a summary that says it was stopped. A +/// statement reading a million rows that answer nothing hands over no +/// batch to notice at, so after this it is interrupted instead, which +/// is the difference between a close that waits for a scan and one that +/// does not. +const GRACE: Duration = Duration::from_millis(50); + +/// What every batch of one statement shares. +/// +/// The column names and the table names belong to the statement rather +/// than to the batch, so they are made once and each batch carries a +/// share of them instead of a copy. +struct Head { + columns: Vec, + names: Names, +} + +/// What the thread running the statement sends back. +enum Chunk { + /// Rows, copied out of the batch they were borrowed in. + Rows { + head: Arc, + rows: Vec>, + }, + /// The statement ended, and this is what it ended as. Always the + /// last thing sent. + Done(Box), +} + +/// How a statement finished, from the reader's side. +enum Ending { + /// It ran to the end, or to where the reader stopped it. + Ran(Streamed), + /// It failed, and this is what a caller is told. Held as the + /// exception rather than as the engine's error so that it can be + /// raised more than once, since an iterator that is asked again + /// after it failed should say the same thing again. + Failed(PyErr), + /// The connection was closed, or a panic left its lock poisoned, + /// which for anything that would run on it is the same fact. + Closed, +} + +/// The queue between the statement and the reader, with both ends able +/// to walk away. +/// +/// A bounded channel is what this is. The standard library's own would +/// do but for the wait: a reader on the main thread has to give up +/// every couple of milliseconds to ask Python whether a signal arrived, +/// and a receive with a deadline is not on the stable side of the +/// channel. So the queue is written out, with one condition variable +/// making room and arrival exact rather than polled. +struct Pipe { + held: Mutex, + /// One variable for both directions. There is one writer and, in + /// any program worth calling correct, one reader, so waking both is + /// waking at most one of each. + moved: Condvar, +} + +struct Held { + queue: VecDeque<(Arc, Vec>)>, + /// How the statement finished, taken by the reader that asks after + /// the last batch. + ending: Option, + /// Nobody is reading any more, so the statement stops. + shut: bool, + /// The statement is over and has said so, which is how a reader + /// tells the end from a pause. + over: bool, +} + +impl Pipe { + fn new() -> Pipe { + Pipe { + held: Mutex::new(Held { + queue: VecDeque::with_capacity(QUEUE), + ending: None, + shut: false, + over: false, + }), + moved: Condvar::new(), + } + } + + /// Hands a batch over, waiting for room. `false` means nobody is + /// reading any more, which is what ends the scan. + fn send(&self, head: Arc, rows: Vec>) -> bool { + let Ok(mut held) = self.held.lock() else { + return false; + }; + while held.queue.len() >= QUEUE && !held.shut { + let Ok(next) = self.moved.wait(held) else { + return false; + }; + held = next; + } + if held.shut { + return false; + } + held.queue.push_back((head, rows)); + self.moved.notify_all(); + true + } + + /// Says the statement is over. Kept beside the queue rather than + /// put in it, so that a reader which hung up and threw the queued + /// batches away still finds out how the statement ended. + fn finish(&self, ending: Ending) { + let Ok(mut held) = self.held.lock() else { + return; + }; + held.ending = Some(ending); + held.over = true; + self.moved.notify_all(); + } + + /// The next chunk, waiting up to `patience` for one, or forever + /// when there is none. `None` is a wait that ran out. + fn recv(&self, patience: Option) -> Option { + let held = self.held.lock().ok()?; + let waiting = |held: &mut Held| held.queue.is_empty() && !held.over; + let mut held = match patience { + Some(patience) => { + self.moved + .wait_timeout_while(held, patience, waiting) + .ok()? + .0 + } + None => self.moved.wait_while(held, waiting).ok()?, + }; + if let Some((head, rows)) = held.queue.pop_front() { + self.moved.notify_all(); + return Some(Chunk::Rows { head, rows }); + } + if held.over { + // Taken rather than copied, and a second reader asking gets + // the same answer a closed connection gets, because there + // is no second reader in a program worth calling correct. + return Some(Chunk::Done(Box::new( + held.ending.take().unwrap_or(Ending::Closed), + ))); + } + None + } + + /// Says that nobody is reading, and throws away what was waiting to + /// be read. The statement finds out at its next batch, which is the + /// same boundary an interrupt is answered at. + fn hang_up(&self) { + if let Ok(mut held) = self.held.lock() { + held.shut = true; + held.queue.clear(); + self.moved.notify_all(); + } + } + + /// How the statement ended, for the reader that stopped it rather + /// than read to it. `None` while it is still running. + fn ending(&self) -> Option { + self.held.lock().ok()?.ending.take() + } + + /// Waits for the statement to say it is over, up to `patience`. + fn wait_over(&self, patience: Option) -> bool { + let Ok(held) = self.held.lock() else { + return true; + }; + let waiting = |held: &mut Held| !held.over; + match patience { + Some(patience) => self + .moved + .wait_timeout_while(held, patience, waiting) + .map(|(held, _)| held.over) + .unwrap_or(true), + None => self + .moved + .wait_while(held, waiting) + .map(|held| held.over) + .unwrap_or(true), + } + } +} + +/// What a connection keeps of the stream it is feeding. +/// +/// A connection runs one statement at a time and a stream holds the +/// connection for as long as its reader takes, so this is the thing +/// every other statement asks before it queues: a statement that waited +/// for a half-read stream would wait for a loop that is waiting for it, +/// and a program deadlocked on itself is worse than a program that is +/// told no. +pub(crate) struct Feeding { + now: Mutex>>, + /// Read without the lock, because refusing is the common answer and + /// asking should cost a load. + busy: AtomicBool, +} + +impl Feeding { + pub(crate) fn new() -> Feeding { + Feeding { + now: Mutex::new(None), + busy: AtomicBool::new(false), + } + } + + /// Whether a stream is holding the connection now. + pub(crate) fn busy(&self) -> bool { + self.busy.load(Ordering::Acquire) + } + + /// Tells the stream that is running, if one is, that nobody is + /// reading. For a connection on its way out: the statement ends at + /// its next batch and gives the lock back, which is what closing is + /// waiting for. + pub(crate) fn hang_up(&self) { + if let Ok(now) = self.now.lock() + && let Some(pipe) = now.as_ref() + { + pipe.hang_up(); + } + } + + fn took(&self, pipe: &Arc) { + if let Ok(mut now) = self.now.lock() { + *now = Some(Arc::clone(pipe)); + } + self.busy.store(true, Ordering::Release); + } + + fn gave_back(&self) { + if let Ok(mut now) = self.now.lock() { + *now = None; + } + self.busy.store(false, Ordering::Release); + } +} + +/// What the reader has in hand between two calls. +#[derive(Default)] +struct State { + /// Rows taken out of a batch and not yet handed to Python. + pending: VecDeque>, + /// What the last batch was made of, kept for the rows still in + /// `pending` and for `columns` after the statement is over. + head: Option>, + /// `Some` once the statement has said it is over. + ending: Option, +} + +/// The pieces of one stream that both sides hold. +struct Live { + pipe: Arc, + /// The connection's own, so that a `Ctrl-C` or a close reaches a + /// statement that is producing nothing to notice at. + stop: Interrupt, + feeding: Arc, + state: Mutex, + /// Set by the reader, so a second `close` costs nothing. + hung_up: AtomicBool, +} + +/// A statement's rows, read as the executor makes them. +/// +/// Iterate it for rows and `batches()` for the lists they arrive in. +/// Either way what is in memory is a batch and not the answer, which is +/// the whole point: a statement over ten million rows is read by a +/// program holding a thousand. +/// +/// It holds the connection until it ends, because a connection runs one +/// statement at a time. Read it to the end, close it, or open it in a +/// `with` block, which closes it however the block is left. +#[pyclass(module = "zudb")] +pub struct Stream { + live: Arc, + statement: String, +} + +/// What a streamed statement did, known once it has ended. +/// +/// The rows are gone by then, which is the point of streaming, so this +/// is what is worth keeping about a result nobody held: what it +/// projected, how much of it was read, whether the reader stopped it +/// early, and what the engine wanted to say along the way. +#[pyclass(module = "zudb", frozen)] +pub struct StreamSummary { + /// The column names, in the order the statement projected them. + #[pyo3(get)] + columns: Vec, + /// How many rows were handed over, which is fewer than the + /// statement would have returned when the reader stopped early. + #[pyo3(get)] + rows: u64, + /// Whether the reader stopped it before it ran out of rows. + #[pyo3(get)] + stopped: bool, + /// Whether the rows arrived as they were made, rather than the + /// statement running whole and being handed over in batches + /// afterwards. A statement that has to see every row before it can + /// give one, which is `ORDER BY`, `DISTINCT` and the aggregates, is + /// the second kind. The loop over it reads the same either way and + /// what differs is what it cost. + #[pyo3(get)] + streamed: bool, + notices: Vec<(String, String, String)>, +} + +#[pymethods] +impl StreamSummary { + /// The warnings the statement raised, in the shape a result reports + /// them. + #[getter] + fn notices<'py>(&self, py: Python<'py>) -> PyResult> { + let out = PyList::empty(py); + for (code, condition, detail) in &self.notices { + let one = PyDict::new(py); + one.set_item("code", code)?; + one.set_item("condition", condition)?; + one.set_item("detail", detail)?; + out.append(one)?; + } + Ok(out) + } + + fn __repr__(&self) -> String { + format!( + "", + self.rows, + if self.stopped { ", stopped" } else { "" }, + if self.streamed { "" } else { ", buffered" } + ) + } +} + +/// The same rows, in the batches they arrived in. +/// +/// One object rather than a method that yields, so that closing the +/// stream is the one call it always was and a batch reader is the same +/// stream seen a list at a time. +#[pyclass(module = "zudb")] +pub struct StreamBatches { + stream: Py, +} + +#[pymethods] +impl StreamBatches { + fn __iter__(slf: PyRef<'_, Self>) -> PyRef<'_, Self> { + slf + } + + fn __next__<'py>(&self, py: Python<'py>) -> PyResult> { + let stream = self.stream.borrow(py); + let live = Arc::clone(&stream.live); + drop(stream); + loop { + { + let mut state = live.state.lock().map_err(|_| panicked())?; + if !state.pending.is_empty() { + let head = state.head.clone().ok_or_else(panicked)?; + let rows: Vec> = state.pending.drain(..).collect(); + drop(state); + let out = PyList::empty(py); + for row in &rows { + out.append(tuple(py, row, &head.names)?)?; + } + return Ok(out); + } + } + if !pull(&live, py)? { + return Err(pyo3::exceptions::PyStopIteration::new_err(())); + } + } + } + + fn __repr__(&self, py: Python<'_>) -> String { + format!( + "", + self.stream.borrow(py).__repr__() + ) + } +} + +#[pymethods] +impl Stream { + /// The column names, in the order the statement projects them. + /// + /// Answering this reads the first batch and keeps it, because the + /// names are the statement's and the statement does not say them + /// until it has made a row. Nothing is lost by that: the batch is + /// handed to the reader that asks next. + #[getter] + fn columns(&self, py: Python<'_>) -> PyResult> { + loop { + { + let state = self.live.state.lock().map_err(|_| panicked())?; + if let Some(head) = state.head.as_ref() { + return Ok(head.columns.clone()); + } + if let Some(Ending::Ran(streamed)) = state.ending.as_ref() { + return Ok(streamed.columns.clone()); + } + } + if !pull(&self.live, py)? { + return Ok(Vec::new()); + } + } + } + + /// What the statement did, once it has done it, and `None` while it + /// is still running. + /// + /// A statement that failed has raised its failure at the row it + /// happened on, so a stream that ended badly reports no summary + /// rather than a summary of half a run. + #[getter] + fn summary(&self, py: Python<'_>) -> PyResult> { + let state = self.live.state.lock().map_err(|_| panicked())?; + let _ = py; + Ok(match state.ending.as_ref() { + Some(Ending::Ran(streamed)) => Some(summarised(streamed)), + _ => None, + }) + } + + /// The rows in the batches they arrived in, as lists of tuples. + /// + /// For a reader that writes what it reads somewhere with a size of + /// its own. A batch holds at most what `batch_rows` asked for, and + /// the last one holds whatever was left. + fn batches(slf: Py) -> StreamBatches { + StreamBatches { stream: slf } + } + + /// Stops the statement and gives the connection back. + /// + /// Doing it twice is not an error, and doing it to a stream that + /// has already ended does nothing, so a `with` block that read to + /// the end leaves the same way one that read a page does. + fn close(&self, py: Python<'_>) { + hang_up(&self.live, py); + } + + /// Whether the statement is over, either because it ran out of rows + /// or because this stream was closed. + #[getter] + fn closed(&self, py: Python<'_>) -> PyResult { + let ended = self + .live + .state + .lock() + .map_err(|_| panicked())? + .ending + .is_some(); + let _ = py; + Ok(ended || self.live.hung_up.load(Ordering::Acquire)) + } + + fn __iter__(slf: PyRef<'_, Self>) -> PyRef<'_, Self> { + slf + } + + fn __next__<'py>(&self, py: Python<'py>) -> PyResult> { + loop { + { + let mut state = self.live.state.lock().map_err(|_| panicked())?; + if let Some(row) = state.pending.pop_front() { + let head = state.head.clone().ok_or_else(panicked)?; + drop(state); + return tuple(py, &row, &head.names); + } + } + if !pull(&self.live, py)? { + return Err(pyo3::exceptions::PyStopIteration::new_err(())); + } + } + } + + fn __enter__(slf: PyRef<'_, Self>) -> PyRef<'_, Self> { + slf + } + + #[pyo3(signature = (*_exception))] + fn __exit__(&self, py: Python<'_>, _exception: &Bound<'_, PyTuple>) -> bool { + self.close(py); + false + } + + fn __repr__(&self) -> String { + let over = self + .live + .state + .lock() + .map(|state| state.ending.is_some()) + .unwrap_or(true); + format!( + "", + self.statement, + if over { ", ended" } else { "" } + ) + } +} + +impl Drop for Stream { + /// A stream nobody holds any more is a reader that walked away, and + /// the statement behind it should not go on reading for one. The + /// wait belongs to `close`, though, not here: a collector running + /// this cannot afford to wait for a scan, so the word is said and + /// the connection comes free a batch later. + fn drop(&mut self) { + if !self.live.hung_up.swap(true, Ordering::AcqRel) { + self.live.pipe.hang_up(); + self.live.stop.stop(); + } + } +} + +/// Starts a statement on a thread of its own and hands back the stream +/// that reads it. +#[allow(clippy::too_many_arguments)] +pub(crate) fn open( + py: Python<'_>, + inner: &Arc>>, + stop: &Interrupt, + feeding: &Arc, + alive: bool, + statement: String, + params: Vec<(String, Value)>, + batch_rows: Option, +) -> PyResult { + if !alive { + return Err(closed(py, "this connection")); + } + if feeding.busy() { + return Err(programming(py, STREAMING)); + } + if batch_rows == Some(0) { + return Err(pyo3::exceptions::PyValueError::new_err( + "batch_rows is how many rows a batch may hold, so it starts at one", + )); + } + let live = Arc::new(Live { + pipe: Arc::new(Pipe::new()), + stop: stop.clone(), + feeding: Arc::clone(feeding), + state: Mutex::new(State::default()), + hung_up: AtomicBool::new(false), + }); + feeding.took(&live.pipe); + let held = Arc::clone(inner); + let running = Arc::clone(&live); + let text = statement.clone(); + let started = std::thread::Builder::new() + .name("zudb stream".into()) + .spawn(move || run(&running, &held, &text, ¶ms, batch_rows)); + if started.is_err() { + feeding.gave_back(); + return Err(pyo3::exceptions::PyRuntimeError::new_err( + "a thread to stream a statement on could not be started", + )); + } + Ok(Stream { live, statement }) +} + +/// What a connection a stream is still reading says to the next +/// statement somebody runs on it. +pub(crate) const STREAMING: &str = "a stream on this connection has not finished, and a connection runs one statement at a \ + time: read the stream to the end, close it, or open a second connection"; + +/// The statement, on its own thread, from the moment it takes the +/// connection to the moment it gives it back. +fn run( + live: &Arc, + inner: &Arc>>, + statement: &str, + params: &[(String, Value)], + batch_rows: Option, +) { + let pipe = Arc::clone(&live.pipe); + let ending = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + let Ok(mut held) = inner.lock() else { + return Ending::Closed; + }; + let Some(conn) = held.as_mut() else { + return Ending::Closed; + }; + // Cleared holding the lock, for the reason every other + // statement clears it there: a stop that arrived while nothing + // was running must not end the statement about to start. + conn.interrupt().clear(); + let names = Names::of(conn.session_mut().catalog()); + let borrowed: Vec<(&str, Value)> = params + .iter() + .map(|(name, value)| (name.as_str(), value.clone())) + .collect(); + let mut head: Option> = None; + let mut sink = |batch: Batch<'_>| -> zudb::Result { + let head = head.get_or_insert_with(|| { + Arc::new(Head { + columns: batch.columns().to_vec(), + names: names.clone(), + }) + }); + // Copied out here, on this thread, with the GIL nowhere in + // sight. The batch is borrowed for the length of this call + // and the reader is going to hold it for the length of a + // loop body, so there is nothing to keep instead. + Ok(if pipe.send(Arc::clone(head), batch.rows().to_vec()) { + Flow::More + } else { + Flow::Stop + }) + }; + let out = match batch_rows { + Some(rows) => conn.query_stream_batched(statement, &borrowed, rows, &mut sink), + None => conn.query_stream(statement, &borrowed, &mut sink), + }; + match out { + Ok(streamed) => Ending::Ran(streamed), + // Built here rather than on the reader's thread, which + // means taking the GIL for as long as one exception takes + // to make. It is the last thing this thread does and it is + // what lets the same failure be raised twice. + Err(err) => Python::attach(|py| Ending::Failed(to_py_err(py, err))), + } + })); + // The lock is gone by here, so the connection is free before + // anybody is told the statement is over. + live.feeding.gave_back(); + let ending = ending.unwrap_or_else(|_| { + // Raised again on the reader's thread rather than left to end + // this one silently, because a panic inside the engine is + // something the program that ran the statement should see. + Ending::Failed(pyo3::exceptions::PyRuntimeError::new_err( + "the statement behind this stream ended in a panic", + )) + }); + live.pipe.finish(ending); +} + +/// Waits for the next chunk and puts it where the reader will find it. +/// +/// `false` means the statement is over and there is nothing more +/// coming, which is the only answer a caller has to act on. +fn pull(live: &Arc, py: Python<'_>) -> PyResult { + { + let state = live.state.lock().map_err(|_| panicked())?; + if let Some(ending) = state.ending.as_ref() { + return match ending { + Ending::Failed(err) => Err(err.clone_ref(py)), + Ending::Closed => Err(closed(py, "this connection")), + Ending::Ran(_) => Ok(false), + }; + } + } + // Off the main thread a signal was never going to arrive, so the + // wait is one wait. On it the wait is a couple of milliseconds at a + // time with a question to Python between them, which is where a + // press is felt. + let chunk = if on_main_thread(py) { + loop { + if let Some(chunk) = py.detach(|| live.pipe.recv(Some(TICK))) { + break chunk; + } + if let Err(signal) = py.check_signals() { + // Stopped rather than left running, so the connection + // is free by the time the exception reaches the caller + // rather than busy with rows nobody wants. + hang_up(live, py); + return Err(signal); + } + } + } else { + py.detach(|| live.pipe.recv(None)) + .ok_or_else(|| closed(py, "this connection"))? + }; + let mut state = live.state.lock().map_err(|_| panicked())?; + match chunk { + Chunk::Rows { head, rows } => { + state.head = Some(head); + state.pending.extend(rows); + Ok(true) + } + Chunk::Done(ending) => { + let out = match ending.as_ref() { + Ending::Failed(err) => Err(err.clone_ref(py)), + Ending::Closed => Err(closed(py, "this connection")), + Ending::Ran(_) => Ok(false), + }; + state.ending = Some(*ending); + out + } + } +} + +/// Says that nobody is reading and waits for the connection to come +/// back, stopping the statement outright if it does not. +fn hang_up(live: &Arc, py: Python<'_>) { + if live.hung_up.swap(true, Ordering::AcqRel) { + return; + } + live.pipe.hang_up(); + // Waited for rather than abandoned, so the connection is free by + // the time this call returns and the next statement on it is not + // told that a stream nobody holds is still reading. + py.detach(|| { + if !live.pipe.wait_over(Some(GRACE)) { + live.stop.stop(); + live.pipe.wait_over(None); + } + }); + // Kept, so that a reader which stopped a stream can still ask what + // it stopped: the summary of a run cut short is exactly the thing + // that says how much of it was read. + if let Some(ending) = live.pipe.ending() + && let Ok(mut state) = live.state.lock() + && state.ending.is_none() + { + state.ending = Some(ending); + } +} + +fn tuple<'py>(py: Python<'py>, row: &[Value], names: &Names) -> PyResult> { + PyTuple::new( + py, + row.iter() + .map(|value| to_py(py, value, names)) + .collect::>>()?, + ) +} + +fn summarised(streamed: &Streamed) -> StreamSummary { + StreamSummary { + columns: streamed.columns.clone(), + rows: streamed.rows, + stopped: streamed.stopped, + streamed: streamed.streamed, + notices: streamed + .notices + .iter() + .map(|notice| { + ( + notice.status.code().to_string(), + notice.status.standard_text().to_string(), + notice.detail.clone(), + ) + }) + .collect(), + } +} + +fn panicked() -> PyErr { + pyo3::exceptions::PyRuntimeError::new_err( + "this stream was left in an unknown state by a thread that panicked", + ) +} diff --git a/src/value.rs b/src/value.rs index 9d5f6f8..fd49a86 100644 --- a/src/value.rs +++ b/src/value.rs @@ -30,7 +30,7 @@ const MICROS_PER_DAY: i64 = 86_400 * 1_000_000; /// statement runs and carried with the rows. A catalog holds tens of /// tables, so this is a copy of a few short strings and not a /// structure worth sharing. -#[derive(Default)] +#[derive(Default, Clone)] pub struct Names { nodes: HashMap, rels: HashMap, diff --git a/tests/test_aio.py b/tests/test_aio.py index 5b59693..0eded6a 100644 --- a/tests/test_aio.py +++ b/tests/test_aio.py @@ -32,6 +32,12 @@ # scheduled. WORK = "MATCH (a:person), (b:person) WHERE a.uid < b.uid RETURN count(a) AS n" +# The same pairs, returned rather than counted, which is millions of +# rows and so a statement that is still running however long a test +# takes to look at it. Every streaming test that says something about a +# stream that has not ended is written on this one. +PAIRS = "MATCH (a:person), (b:person) WHERE a.uid < b.uid RETURN a.uid AS a, b.uid AS b" + # Three thousand people is about a second of that work on the machine # this was written on, and six thousand about three seconds. The first # is for the tests that wait for the statement to end and the second @@ -414,3 +420,77 @@ async def test_sql_is_execute_under_another_name(tmp_path: Path) -> None: await conn.sql("INSERT (p:person {uid: 10, name: 'ada'})") rows = await conn.sql("MATCH (p:person) RETURN p.name AS name") assert rows.fetchall() == [("ada",)] + + +@run +async def test_a_stream_gives_back_its_rows_awaited_one_at_a_time(tmp_path: Path) -> None: + conn = await crowded(tmp_path / "streamed.zu1", WATCHED) + async with conn: + async with conn.stream("MATCH (p:person) RETURN p.uid AS uid") as rows: + assert await rows.columns() == ["uid"] + seen = [uid async for (uid,) in rows] + assert seen == list(range(WATCHED)) + assert rows.summary is not None + assert rows.summary.rows == WATCHED + + +@run +@pytest.mark.timing +async def test_the_loop_runs_while_a_stream_waits_for_a_batch(tmp_path: Path) -> None: + """The reason streaming is here at all rather than in the sync + client only: a task reading rows off a scan leaves the loop free + between the batches, so everything else on it keeps its turn. + """ + conn = await crowded(tmp_path / "free.zu1", LONG) + ticks = 0 + + async def tick() -> None: + nonlocal ticks + while True: + ticks += 1 + await asyncio.sleep(0.001) + + async with conn: + counting = asyncio.create_task(tick()) + async with conn.stream(PAIRS) as rows: + await rows.__anext__() + await asyncio.sleep(0.05) + counting.cancel() + + assert ticks > 10 + + +@run +async def test_a_stream_holds_the_connection_until_it_ends(tmp_path: Path) -> None: + conn = await crowded(tmp_path / "held.zu1", WATCHED) + async with conn: + rows = await conn.stream(PAIRS) + await rows.__anext__() + + with pytest.raises(zudb.ProgrammingError, match="runs one statement at a time"): + await conn.execute("MATCH (p:person) RETURN count(*) AS n") + + await rows.close() + counted = await conn.execute("MATCH (p:person) RETURN count(*) AS n") + assert counted.fetchall() == [(WATCHED,)] + assert rows.closed is True + + +@run +async def test_a_stream_reads_in_batches_when_it_is_asked_to(tmp_path: Path) -> None: + conn = await crowded(tmp_path / "batched.zu1", WATCHED) + async with conn: + async with conn.stream("MATCH (p:person) RETURN p.uid AS uid", batch_rows=250) as rows: + sizes = [len(batch) async for batch in rows.batches()] + + assert sum(sizes) == WATCHED + assert max(sizes) <= 250 + + +@run +async def test_a_stream_that_fails_raises_at_the_row_it_failed_on(tmp_path: Path) -> None: + async with zudb.aio.connect(tmp_path / "broken.zu1") as conn: + await conn.execute("INSERT (p:person {uid: 10, name: 'ada'})") + rows = await conn.stream("MATCH (") + with pytest.raises(zudb.SyntaxError): + await rows.__anext__() diff --git a/tests/test_stream.py b/tests/test_stream.py new file mode 100644 index 0000000..c826424 --- /dev/null +++ b/tests/test_stream.py @@ -0,0 +1,330 @@ +"""Rows read as the engine makes them. + +A result is rows already in memory and a stream is rows on their way, so +what these assert is the difference between the two: that the first row +arrives before the last is made, that no more than a batch or two is +held at once, and that the connection a stream is reading is a +connection nothing else can run on until it ends. The rest is the +behaviour a result already has, checked once through the streaming +spelling to show it survived the crossing. +""" + +from __future__ import annotations + +import gc +import itertools +import tracemalloc +from pathlib import Path + +import pytest +import zudb + +# Enough people that a scan over them takes more than one batch, which +# is what makes a batch worth talking about at all. +MANY = 5_000 +# Enough people that the pairs of them are more rows than anything here +# reads. A statement that fits in the queue is one the engine can finish +# while the reader is between two rows, so every test about a stream +# still running is written on the pairs rather than on the people. +PAIRED = 2_000 +PAIRS = PAIRED * (PAIRED - 1) // 2 +PAIRS_OF_PEOPLE = "MATCH (a:person), (b:person) WHERE a.uid < b.uid RETURN a.uid AS a, b.uid AS b" + + +@pytest.fixture +def many(tmp_path: Path) -> zudb.Connection: + """Five thousand people, which is a couple of batches and a bit. + + Loaded rather than inserted because a row at a time is a commit at a + time, and this is scaffolding rather than the thing under test. + """ + path = tmp_path / "many.zu1" + zudb.load( + path, + nodes="person", + rels="knows", + columns={"uid": list(range(MANY)), "name": [f"p{uid}" for uid in range(MANY)]}, + edges=[(0, 1)], + ) + conn = zudb.connect(path, read_only=True) + yield conn + conn.close() + + +@pytest.fixture +def paired(tmp_path: Path) -> zudb.Connection: + """Enough people that the pairs of them are millions of rows. + + Two thousand people is two million pairs, which no test here reads + to the end: it is the statement to open a stream on when what is + being asserted is about a statement that is still running. + """ + path = tmp_path / "paired.zu1" + zudb.load( + path, + nodes="person", + rels="knows", + columns={"uid": list(range(PAIRED))}, + edges=[(0, 1)], + ) + conn = zudb.connect(path, read_only=True) + yield conn + conn.close() + + +def test_a_stream_gives_back_every_row_the_statement_made(many: zudb.Connection) -> None: + rows = list(many.stream("MATCH (p:person) RETURN p.uid AS uid")) + + assert len(rows) == MANY + assert rows[0] == (0,) + assert rows[-1] == (MANY - 1,) + + +def test_a_stream_says_its_columns_before_a_row_is_read(many: zudb.Connection) -> None: + with many.stream("MATCH (p:person) RETURN p.uid AS uid, p.name AS name") as stream: + assert stream.columns == ["uid", "name"] + assert next(iter(stream)) == (0, "p0") + + +def test_rows_arrive_before_the_statement_is_over(paired: zudb.Connection) -> None: + """The claim the whole module is for. + + Two million pairs take long enough that a statement which had to + finish before it gave a row would not have given one yet, so a first + row in hand while the summary is still `None` is the proof that the + rows are coming as they are made. + """ + with paired.stream(PAIRS_OF_PEOPLE) as stream: + first = next(iter(stream)) + + assert first == (0, 1) + assert stream.summary is None + assert stream.closed is False + + +def test_a_stream_holds_a_batch_or_two_and_not_the_result(many: zudb.Connection) -> None: + """What the memory of a stream looks like against the memory of a + result. + + Only the Python side is measured, which is the side a caller sees: + the tuples. A result makes every one of them and holds them, a + stream makes them a batch at a time and lets each go. + """ + statement = "MATCH (p:person) RETURN p.uid AS uid, p.name AS name" + gc.collect() + tracemalloc.start() + held = many.execute(statement).fetchall() + whole = tracemalloc.get_traced_memory()[1] + tracemalloc.stop() + assert len(held) == MANY + del held + gc.collect() + + tracemalloc.start() + counted = 0 + for _ in many.stream(statement): + counted += 1 + streaming = tracemalloc.get_traced_memory()[1] + tracemalloc.stop() + + assert counted == MANY + # A tenth would hold on this machine and half is the assertion that + # survives a machine where an allocator rounds differently. + assert streaming < whole / 2 + + +def test_the_batches_are_the_size_that_was_asked_for(many: zudb.Connection) -> None: + with many.stream("MATCH (p:person) RETURN p.uid AS uid", batch_rows=500) as stream: + sizes = [len(batch) for batch in stream.batches()] + + assert sum(sizes) == MANY + assert max(sizes) <= 500 + assert sizes[0] == 500 + + +def test_a_batch_is_a_list_of_the_rows_the_loop_would_have_given(social: zudb.Connection) -> None: + with social.stream("MATCH (p:person) RETURN p.uid AS uid") as stream: + batches = list(stream.batches()) + + assert batches == [[(10,), (20,), (30,)]] + + +def test_a_stream_holds_the_connection_until_it_ends(paired: zudb.Connection) -> None: + """A statement run on a connection a stream is reading is told no + rather than queued, because queueing it behind a loop that may never + finish is a program deadlocked on itself. + """ + stream = paired.stream(PAIRS_OF_PEOPLE) + next(iter(stream)) + + with pytest.raises(zudb.ProgrammingError, match="runs one statement at a time"): + paired.execute("MATCH (p:person) RETURN count(*) AS n") + + stream.close() + assert paired.execute("MATCH (p:person) RETURN count(*) AS n").fetchall() == [(PAIRED,)] + + +def test_a_stream_read_to_the_end_gives_the_connection_back(social: zudb.Connection) -> None: + assert list(social.stream("MATCH (p:person) RETURN p.uid AS uid")) == [(10,), (20,), (30,)] + + assert social.execute("MATCH (p:person) RETURN count(*) AS n").fetchall() == [(3,)] + + +def test_a_stream_that_ran_out_says_what_it_did(many: zudb.Connection) -> None: + stream = many.stream("MATCH (p:person) RETURN p.uid AS uid") + assert stream.summary is None + for _ in stream: + pass + summary = stream.summary + + assert summary is not None + assert summary.columns == ["uid"] + assert summary.rows == MANY + assert summary.stopped is False + assert summary.streamed is True + assert summary.notices == [] + + +def test_a_stream_stopped_early_says_how_much_of_it_was_read(many: zudb.Connection) -> None: + """The summary of a run cut short is the thing that says how much of + it happened, so a stream that was closed has one rather than none. + """ + stream = many.stream("MATCH (p:person) RETURN p.uid AS uid", batch_rows=100) + next(iter(stream)) + stream.close() + summary = stream.summary + + assert summary is not None + assert summary.stopped is True + assert 0 < summary.rows < MANY + + +def test_a_statement_that_sorts_reads_the_same_and_says_it_was_buffered( + many: zudb.Connection, +) -> None: + """`ORDER BY` cannot give a row before it has seen every row, so the + engine runs it whole and hands it over in batches afterwards. The + loop is the loop either way and what differs is what it cost, which + is why the summary says which of the two happened. + """ + with many.stream("MATCH (p:person) RETURN p.uid AS uid ORDER BY p.uid DESC") as stream: + rows = list(itertools.islice(stream, 3)) + stream.close() + summary = stream.summary + + assert rows == [(MANY - 1,), (MANY - 2,), (MANY - 3,)] + assert summary is not None + assert summary.streamed is False + + +def test_closing_twice_is_not_an_error(social: zudb.Connection) -> None: + stream = social.stream("MATCH (p:person) RETURN p.uid AS uid") + stream.close() + stream.close() + + assert stream.closed is True + + +def test_a_stream_nobody_read_still_frees_the_connection(many: zudb.Connection) -> None: + with many.stream("MATCH (p:person) RETURN p.uid AS uid"): + pass + + assert many.execute("MATCH (p:person) RETURN count(*) AS n").fetchall() == [(MANY,)] + + +def test_reading_a_closed_stream_gives_nothing_rather_than_rows(many: zudb.Connection) -> None: + stream = many.stream("MATCH (p:person) RETURN p.uid AS uid") + stream.close() + + assert list(stream) == [] + + +def test_a_stream_of_a_statement_that_does_not_compile_fails(social: zudb.Connection) -> None: + """The compile happens on the statement's own thread, so the failure + arrives at the first row asked for rather than at the call that + opened the stream, and it is the failure the same statement would + have raised. + """ + stream = social.stream("MATCH (") + + with pytest.raises(zudb.SyntaxError): + next(iter(stream)) + + +def test_the_same_failure_is_raised_at_every_row_asked_for(social: zudb.Connection) -> None: + """A failure is kept rather than spent, because a loop that asked + twice would otherwise be told the second time that the rows had run + out, which is the one thing that did not happen. + """ + stream = social.stream("MATCH (") + with pytest.raises(zudb.SyntaxError): + next(iter(stream)) + + with pytest.raises(zudb.SyntaxError): + next(iter(stream)) + + +def test_a_stream_binds_its_parameters(social: zudb.Connection) -> None: + statement = "MATCH (p:person) WHERE p.name = $name RETURN p.uid AS uid" + + assert list(social.stream(statement, {"name": "grace"})) == [(20,)] + + +def test_a_batch_of_no_rows_is_refused_rather_than_looped_on(social: zudb.Connection) -> None: + with pytest.raises(ValueError, match="starts at one"): + social.stream("MATCH (p:person) RETURN p.uid AS uid", batch_rows=0) + + +def test_streaming_on_a_closed_connection_is_refused(tmp_path: Path) -> None: + conn = zudb.connect(tmp_path / "shut.zu1") + conn.close() + + with pytest.raises(zudb.ProgrammingError, match="closed"): + conn.stream("MATCH (p:person) RETURN p.uid AS uid") + + +def test_closing_the_connection_under_a_stream_stops_the_stream(paired: zudb.Connection) -> None: + """Closing waits for the statement, so a connection cannot be closed + out from under a stream that is still reading. What it does instead + is hang the stream up first, which is why the close returns at all + rather than waiting for a scan nobody is reading, and why what the + stream reports afterwards is a run that was stopped. + """ + stream = paired.stream(PAIRS_OF_PEOPLE) + next(iter(stream)) + paired.close() + + assert list(stream) != [] + assert stream.closed is True + summary = stream.summary + assert summary is not None + assert summary.stopped is True + assert summary.rows < PAIRS + + +def test_a_stream_carries_the_values_a_result_would(loaded: zudb.Connection) -> None: + """A node and an edge come back as the objects a result gives, named + for their tables rather than numbered, which is the conversion this + module had to carry over from the connection it borrowed the + catalog from. + """ + statement = "MATCH (a:person)-[r:knows]->(b:person) RETURN a, r, b" + rows = list(loaded.stream(statement)) + + assert rows == loaded.execute(statement).fetchall() + node, rel, other = rows[0] + assert isinstance(node, zudb.Node) + assert isinstance(rel, zudb.Rel) + assert isinstance(other, zudb.Node) + assert node.table == "person" + assert rel.table == "knows" + + +def test_a_stream_says_what_it_is(social: zudb.Connection) -> None: + stream = social.stream("MATCH (p:person) RETURN p.uid AS uid") + + assert repr(stream) == '' + assert repr(stream.batches()).startswith("") + assert repr(stream.summary) == "" From f02c760c9cad9ef8a2d932cc5dbc0680db951a83 Mon Sep 17 00:00:00 2001 From: Tam Nguyen Duc <1218621+tamnd@users.noreply.github.com> Date: Wed, 19 Aug 2026 17:20:12 +0700 Subject: [PATCH 2/2] Say what batch_rows really does It is a ceiling rather than a size: the engine fills whole vectors and a batch holds as many of them as fit under the number, so asking for ten thousand gives batches of 9,216 and asking for five hundred gives batches of exactly five hundred. A caller who asked for a round figure and counted something else would think the client had lost rows, so the README, the method's own documentation and a test all say it. --- README.md | 2 ++ src/conn.rs | 5 ++++- tests/test_stream.py | 15 +++++++++++++++ 3 files changed, 21 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index 36eaf80..10e8198 100644 --- a/README.md +++ b/README.md @@ -143,6 +143,8 @@ with conn.stream("MATCH (p:person) RETURN p.uid AS uid", batch_rows=10_000) as r rows.summary.rows # how many were read ``` +`batch_rows` is a ceiling rather than a size. The engine fills whole vectors and a batch holds as many of them as fit under the number, so asking for 10,000 gives batches of 9,216 and asking for 500 gives batches of exactly 500, since a vector divides evenly by that. + The statement runs on a thread of its own with the GIL down and hands each batch over a queue two batches deep, so the engine is filling the next batch while your loop is reading this one and neither waits for the other for long. On this machine, over a million people in two columns, the first row arrives 0.7 ms after the call against 30 ms for `execute`, and reading all of them takes 214 ms streamed against 256 ms in one piece, which is 4.7 million rows a second against 3.9 million. Reading them a batch at a time takes 176 ms, or 5.7 million a second, because a batch is one call where a row is one call. Streaming being the faster of the two is not a trick: the rows are made and consumed while the cache still has them, and nothing has to hold a million tuples at once. Held is the whole difference: the Python side of `fetchall` peaks at 153 MB for that result and the same read streamed peaks at nothing worth measuring, since every tuple is freed as the loop moves past it. A stream holds the connection until it ends, because a connection runs one statement at a time. Reading it to the end frees the connection, so does `close`, and so does the end of a `with` block however it was left. A statement run on the connection while a stream is open is refused with a message that says so rather than queued behind a loop that may never finish, which is the deadlock the refusal exists to prevent. If a program needs to read a stream and run statements at the same time, that is two connections. diff --git a/src/conn.rs b/src/conn.rs index 9709e1b..a163717 100644 --- a/src/conn.rs +++ b/src/conn.rs @@ -138,7 +138,10 @@ impl Connection { /// holding a thousand. A reader that stops early stops the scan, /// which is the other half of it, and `batch_rows` says how many /// rows a batch may hold when the rows are going somewhere with a - /// size of its own. + /// size of its own. It is a ceiling rather than a size: the engine + /// fills whole vectors and a batch holds as many of them as fit + /// under the number, so a round figure comes back a little short + /// unless a vector divides by it. /// /// The stream holds the connection until it ends, because a /// connection runs one statement at a time. Read it to the end, diff --git a/tests/test_stream.py b/tests/test_stream.py index c826424..65f810a 100644 --- a/tests/test_stream.py +++ b/tests/test_stream.py @@ -142,6 +142,21 @@ def test_the_batches_are_the_size_that_was_asked_for(many: zudb.Connection) -> N assert sizes[0] == 500 +def test_the_size_asked_for_is_a_ceiling_rather_than_a_size(many: zudb.Connection) -> None: + """The engine fills whole vectors and a batch holds as many of them + as fit under the number, so a round figure comes back a little short + unless a vector divides by it. Written down because a caller who + asked for three thousand and counted 2,048 would otherwise think + something was wrong. + """ + with many.stream("MATCH (p:person) RETURN p.uid AS uid", batch_rows=3_000) as stream: + sizes = [len(batch) for batch in stream.batches()] + + assert sum(sizes) == MANY + assert max(sizes) < 3_000 + assert len(set(sizes[:-1])) == 1 + + def test_a_batch_is_a_list_of_the_rows_the_loop_would_have_given(social: zudb.Connection) -> None: with social.stream("MATCH (p:person) RETURN p.uid AS uid") as stream: batches = list(stream.batches())