diff --git a/Cargo.lock b/Cargo.lock index 18f8988..6921392 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1647,7 +1647,7 @@ dependencies = [ [[package]] name = "zu" version = "0.0.1" -source = "git+https://github.com/tamnd/zu?rev=8009a961f463d6b576509e0f752b60f44071cd33#8009a961f463d6b576509e0f752b60f44071cd33" +source = "git+https://github.com/tamnd/zu?rev=6f38950b9f1c4b6cbca37f404ee660370a945c12#6f38950b9f1c4b6cbca37f404ee660370a945c12" dependencies = [ "zu-common", "zu-encoding", @@ -1663,7 +1663,7 @@ dependencies = [ [[package]] name = "zu-common" version = "0.0.1" -source = "git+https://github.com/tamnd/zu?rev=8009a961f463d6b576509e0f752b60f44071cd33#8009a961f463d6b576509e0f752b60f44071cd33" +source = "git+https://github.com/tamnd/zu?rev=6f38950b9f1c4b6cbca37f404ee660370a945c12#6f38950b9f1c4b6cbca37f404ee660370a945c12" dependencies = [ "thiserror", ] @@ -1671,7 +1671,7 @@ dependencies = [ [[package]] name = "zu-encoding" version = "0.0.1" -source = "git+https://github.com/tamnd/zu?rev=8009a961f463d6b576509e0f752b60f44071cd33#8009a961f463d6b576509e0f752b60f44071cd33" +source = "git+https://github.com/tamnd/zu?rev=6f38950b9f1c4b6cbca37f404ee660370a945c12#6f38950b9f1c4b6cbca37f404ee660370a945c12" dependencies = [ "ruzstd", "zu-common", @@ -1680,7 +1680,7 @@ dependencies = [ [[package]] name = "zu-exec" version = "0.0.1" -source = "git+https://github.com/tamnd/zu?rev=8009a961f463d6b576509e0f752b60f44071cd33#8009a961f463d6b576509e0f752b60f44071cd33" +source = "git+https://github.com/tamnd/zu?rev=6f38950b9f1c4b6cbca37f404ee660370a945c12#6f38950b9f1c4b6cbca37f404ee660370a945c12" dependencies = [ "zu-common", "zu-query", @@ -1690,7 +1690,7 @@ dependencies = [ [[package]] name = "zu-query" version = "0.0.1" -source = "git+https://github.com/tamnd/zu?rev=8009a961f463d6b576509e0f752b60f44071cd33#8009a961f463d6b576509e0f752b60f44071cd33" +source = "git+https://github.com/tamnd/zu?rev=6f38950b9f1c4b6cbca37f404ee660370a945c12#6f38950b9f1c4b6cbca37f404ee660370a945c12" dependencies = [ "crossbeam-deque", "zu-common", @@ -1701,7 +1701,7 @@ dependencies = [ [[package]] name = "zu-s3" version = "0.0.1" -source = "git+https://github.com/tamnd/zu?rev=8009a961f463d6b576509e0f752b60f44071cd33#8009a961f463d6b576509e0f752b60f44071cd33" +source = "git+https://github.com/tamnd/zu?rev=6f38950b9f1c4b6cbca37f404ee660370a945c12#6f38950b9f1c4b6cbca37f404ee660370a945c12" dependencies = [ "crc32c", "object_store", @@ -1712,7 +1712,7 @@ dependencies = [ [[package]] name = "zu-sqlite" version = "0.0.1" -source = "git+https://github.com/tamnd/zu?rev=8009a961f463d6b576509e0f752b60f44071cd33#8009a961f463d6b576509e0f752b60f44071cd33" +source = "git+https://github.com/tamnd/zu?rev=6f38950b9f1c4b6cbca37f404ee660370a945c12#6f38950b9f1c4b6cbca37f404ee660370a945c12" dependencies = [ "rusqlite", "zu-common", @@ -1722,7 +1722,7 @@ dependencies = [ [[package]] name = "zu-storage" version = "0.0.1" -source = "git+https://github.com/tamnd/zu?rev=8009a961f463d6b576509e0f752b60f44071cd33#8009a961f463d6b576509e0f752b60f44071cd33" +source = "git+https://github.com/tamnd/zu?rev=6f38950b9f1c4b6cbca37f404ee660370a945c12#6f38950b9f1c4b6cbca37f404ee660370a945c12" dependencies = [ "zu-common", "zu-encoding", @@ -1731,7 +1731,7 @@ dependencies = [ [[package]] name = "zu-vector" version = "0.0.1" -source = "git+https://github.com/tamnd/zu?rev=8009a961f463d6b576509e0f752b60f44071cd33#8009a961f463d6b576509e0f752b60f44071cd33" +source = "git+https://github.com/tamnd/zu?rev=6f38950b9f1c4b6cbca37f404ee660370a945c12#6f38950b9f1c4b6cbca37f404ee660370a945c12" dependencies = [ "zu-common", ] @@ -1739,7 +1739,7 @@ dependencies = [ [[package]] name = "zu-zu1" version = "0.0.1" -source = "git+https://github.com/tamnd/zu?rev=8009a961f463d6b576509e0f752b60f44071cd33#8009a961f463d6b576509e0f752b60f44071cd33" +source = "git+https://github.com/tamnd/zu?rev=6f38950b9f1c4b6cbca37f404ee660370a945c12#6f38950b9f1c4b6cbca37f404ee660370a945c12" dependencies = [ "crc32c", "loom", diff --git a/Cargo.toml b/Cargo.toml index bc41f04..f344e5d 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -18,8 +18,8 @@ crate-type = ["cdylib"] # with (ADR 0002), so a revision is the honest way to say which one. # A local checkout is used instead with a `paths` override in # `.cargo/config.toml`, which is untracked on purpose. -zudb = { package = "zu", git = "https://github.com/tamnd/zu", rev = "8009a961f463d6b576509e0f752b60f44071cd33" } -zu-common = { git = "https://github.com/tamnd/zu", rev = "8009a961f463d6b576509e0f752b60f44071cd33" } +zudb = { package = "zu", git = "https://github.com/tamnd/zu", rev = "6f38950b9f1c4b6cbca37f404ee660370a945c12" } +zu-common = { git = "https://github.com/tamnd/zu", rev = "6f38950b9f1c4b6cbca37f404ee660370a945c12" } # `extension-module` is asked for by maturin, in pyproject.toml, and # not here. Only the build backend knows how an extension is linked on # the platform it is building for, and a crate that turns the feature diff --git a/README.md b/README.md index fdeaa8a..2974824 100644 --- a/README.md +++ b/README.md @@ -102,13 +102,15 @@ conn.register("people", frame) conn.execute("MATCH (p:people) WHERE p.age > 40 RETURN p.name AS name") ``` -Anything that speaks Arrow goes in, which is a pandas or polars DataFrame, a pyarrow table or a reader over one, and a dictionary of lists is there for a caller with none of them installed. The frame arrives over the same C Data Interface a result leaves by, so a column costs a memcpy and no Python object per cell: a million rows of one integer column are read in 5 ms, and a million rows of an integer, a float and a string take 138 ms, which is the string column being the only one that allocates. +Anything that speaks Arrow goes in, which is a pandas or polars DataFrame, a pyarrow table or a reader over one, and a dictionary of lists is there for a caller with none of them installed. Nothing is copied. The frame arrives over the same C Data Interface a result leaves by, and what the engine is told is where each column is, how wide its values are and what they mean; a statement that names it builds vectors pointing straight at the caller's buffers. So registering costs what describing the columns costs and not what the rows cost: 2.6 microseconds for ten rows and 2.3 for ten million. -What this is not is a scan of the frame where it sits. DuckDB's `register` is zero copy because its executor can call back out to the Python object holding the data, and this engine has no such callback yet, so the frame is copied in and becomes a table like any other. That makes a registered frame a snapshot: changing the DataFrame afterwards changes nothing until it is registered again. The call is the one it would be either way, and the day the engine grows a scan it becomes the cheap thing under the same name. +The one column that is walked is a string column, and it is walked once. Every offset is checked at registration so that reading the frame afterwards cannot fail, which is 362 microseconds for a million strings. Two other things copy and both are said rather than hidden: a stream that arrives as several batches is concatenated into one, because a column of a table is one run of bytes and two batches are two of them, and a dictionary of Python lists is read into buffers of this client's own, because a list holds objects rather than numbers and there is nothing in it to point at. -The copy is not the slow kind, but the write is a write. Those million rows take 9.3 seconds in all, of which 0.14 is reading the frame and the rest is the engine, which is the same 0.1 to 0.2 million rows a second `load` and the appender manage on this machine. Registering is the right way to get a frame in and it is not a way to make the v0 write path faster. +Because it is not a copy, a registered frame is a view and not a snapshot. Write into the array behind it and the next statement answers what is there now, which is the thing to know about the call and the reason it is worth having. Reading one is as fast as reading a table of the database and faster where the database has to decode: over a million rows on this machine, summing an integer column takes 662 microseconds against a stored table's 833, and finding one row by a string takes 1.8 ms against 3.1. -`unregister(name)` takes the rows back out and `conn.registered` says what is registered here. Registering the same name again replaces the rows under it, which is what rerunning a cell means by it, and registering over a table this connection did not register is refused, since a frame knows nothing about the rows that were already there. Two things show the engine through and are said rather than worked around: a name that has been used keeps the columns it was used with, because a table's columns are declared by its first row and no statement alters them, and `unregister` empties the table rather than removing it, because no statement drops one. A null anywhere is refused by column and row, since a property that is null is one no row of this engine holds, and registering inside a transaction is refused for the reason an appender is. +The frame belongs to the connection it was registered on and goes when that connection does. Nothing is written to the file, so another program opening the same database has never heard of it, and nothing writes to it either: a statement that inserts into or deletes from a registered name is refused with the reason, because that memory is the caller's DataFrame. `unregister(name)` takes the name away and hands the bytes back, which is not always that instant, since a statement still reading the frame holds it until it ends. `conn.registered` says what is registered here. + +Registering the same name again replaces what it stands for, columns and all, which is what rerunning a cell means by it. Registering over a table the database already holds is refused, since a statement naming it would mean the stored one. A frame with no rows is a table to match on and answers nothing, because a frame knows its columns without being told by a row. A null anywhere is refused by column and row, since a property that is null is one no row of this engine holds, and registering inside a transaction is refused because a frame is registered on the session, which is the thing the transaction is running on. ## Reading a result as columns @@ -152,7 +154,7 @@ The stub is checked against the module it describes in CI: griffe reads the stub ## 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, stubs inside the wheel with a gate that keeps them true, the GIL released around every statement, every load and every copy out, and `Ctrl-C` and `interrupt()` stopping a statement without touching the connection under it. `zudb.aio` is next, and each one lands 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, and `Ctrl-C` and `interrupt()` stopping a statement without touching the connection under it. `zudb.aio` is next, and each one lands with the tests that say it works. ## Wheels diff --git a/python/zudb/_zudb.pyi b/python/zudb/_zudb.pyi index 8afe57d..2a8fff5 100644 --- a/python/zudb/_zudb.pyi +++ b/python/zudb/_zudb.pyi @@ -90,7 +90,7 @@ class Connection: """Puts a DataFrame under a name a statement can match on.""" def unregister(self, name: str) -> None: - """Takes the rows of a registered frame back out.""" + """Takes a registered frame's name away.""" @property def registered(self) -> list[str]: diff --git a/src/buffer.rs b/src/buffer.rs index fcde8c6..af0479d 100644 --- a/src/buffer.rs +++ b/src/buffer.rs @@ -21,7 +21,6 @@ use pyo3::types::{PyBool, PyBytes, PyDate, PyDateTime, PyDelta, PyTime}; use zu_common::temporal::days_from_civil; use zu_common::{DurationKind, FloatBits, IntBits, LogicalType, Temporal}; use zudb::Field; -use zudb::query::Value; use crate::value::{Duration, clock_nanos}; @@ -286,27 +285,6 @@ impl Column { } } - /// One value of this column as a statement's parameter takes it. - /// - /// Owned rather than borrowed, unlike `field`, because a parameter - /// is bound for the length of a statement and the statement is - /// where the copy was always going to happen. `None` for a column - /// of byte strings, which no statement carries a literal or a - /// parameter for. - pub fn value(&self, row: usize) -> Option { - Some(match self { - Column::Int(v) => Value::Int(v[row] as i64), - Column::Float(v) => Value::Float(v[row]), - Column::Bool(v) => Value::Bool(v[row]), - Column::Str(v) => Value::Str(v[row].clone()), - Column::Bytes(_) => return None, - Column::Date(v) => Value::Temporal(Temporal::Date(v[row])), - Column::LocalTime(v) => Value::Temporal(Temporal::LocalTime(v[row])), - Column::LocalDatetime(v) => Value::Temporal(Temporal::LocalDatetime(v[row])), - Column::Duration(kind, v) => Value::Temporal(Temporal::Duration(*kind, v[row])), - }) - } - /// One value of this column as the engine's appender takes it. /// /// A field borrows rather than owning, which is the point of it on diff --git a/src/conn.rs b/src/conn.rs index 7db13b9..9428676 100644 --- a/src/conn.rs +++ b/src/conn.rs @@ -7,7 +7,6 @@ //! to run statements at once want two connections. It is there so that //! a program which shares one by accident waits rather than corrupts. -use std::collections::BTreeMap; use std::ffi::CStr; use std::path::PathBuf; use std::sync::atomic::{AtomicBool, Ordering}; @@ -68,18 +67,6 @@ pub struct Connection { /// asking it to stop, should not queue behind a ten second /// statement. alive: AtomicBool, - /// The names this connection has registered frames under, against - /// whether one is registered under them now. A name that was - /// unregistered stays here as `false`, because the empty table it - /// left behind is still a table this connection made and is still - /// the one a later `register` under the same name may fill: forget - /// it entirely and the name would be refused as somebody else's. - /// - /// Kept per connection rather than per database, because a - /// registered frame is a caller's data under a caller's name and - /// another program opening the same file has no idea it is anything - /// but a table. - frames: Mutex>, #[pyo3(get)] path: PathBuf, #[pyo3(get)] @@ -196,48 +183,38 @@ impl Connection { /// /// Anything that speaks Arrow goes in: a pandas or polars /// DataFrame, a pyarrow Table or RecordBatchReader, or a dictionary - /// of lists. The frame is copied into the database rather than - /// scanned where it sits, because the engine has no way yet to call - /// back out to a Python object mid statement, so what a program - /// registers is a snapshot: changing the DataFrame afterwards - /// changes nothing here until it is registered again. + /// of lists. Nothing is copied. The engine is told where the + /// columns are and reads them where they lie, so registering costs + /// the same whether the frame has ten rows or ten million, and a + /// frame is a view rather than a snapshot: write into the array + /// behind it and the next statement answers what is there now. A + /// dictionary of Python lists is the one exception, since a list + /// holds objects rather than numbers and there is nothing in it to + /// point at. /// - /// Registering the same name again replaces the rows under it. - /// Registering over a table this connection did not register is - /// refused, since a frame knows nothing about the rows that were - /// already there. Gives back the number of rows written. - fn register( - slf: Py, - py: Python<'_>, - name: &str, - data: &Bound<'_, PyAny>, - ) -> PyResult { - register::register(py, slf, name, data) + /// The frame belongs to this connection and goes when it does. + /// Nothing is written to the database and no other program sees it. + /// Registering the same name again replaces what it stands for, + /// columns and all. Registering over a table the database already + /// holds is refused, since a statement naming it would mean the + /// stored one. Gives back the number of rows the frame has. + fn register(&self, py: Python<'_>, name: &str, data: &Bound<'_, PyAny>) -> PyResult { + register::register(py, self, name, data) } - /// Takes the rows of a registered frame back out. + /// Takes a registered frame's name away. /// - /// The name stops being a frame of this connection's and the rows - /// go. The empty table stays in the catalog, because no statement - /// of this engine drops one yet, so the name cannot afterwards be - /// used for something that is not a table. + /// The name stops standing for anything and the bytes go back to + /// the caller, which is not necessarily this instant: a statement + /// still reading the frame holds it until it ends. fn unregister(&self, py: Python<'_>, name: &str) -> PyResult<()> { register::unregister(py, self, name) } /// The names this connection has registered frames under, sorted. #[getter] - fn registered(&self) -> Vec { - self.frames - .lock() - .map(|frames| { - frames - .iter() - .filter(|&(_, ®istered)| registered) - .map(|(name, _)| name.clone()) - .collect() - }) - .unwrap_or_default() + fn registered(&self, py: Python<'_>) -> PyResult> { + register::registered(py, self) } /// Asks the statement running on this connection to stop. @@ -361,7 +338,6 @@ impl Connection { inner: Arc::new(Mutex::new(Some(opened))), runner: OnceLock::new(), alive: AtomicBool::new(true), - frames: Mutex::new(BTreeMap::new()), path, read_only, }) @@ -448,35 +424,6 @@ impl Connection { }) .map_err(|()| closed(py, "this connection")) } - - /// Whether a frame is registered here under this name right now. - pub(crate) fn registered_here(&self, name: &str) -> bool { - self.frames - .lock() - .is_ok_and(|frames| frames.get(name) == Some(&true)) - } - - /// Whether this connection made the table of this name, whether or - /// not a frame is registered under it now. - pub(crate) fn made_here(&self, name: &str) -> bool { - self.frames - .lock() - .is_ok_and(|frames| frames.contains_key(name)) - } - - pub(crate) fn remember(&self, name: &str) { - if let Ok(mut frames) = self.frames.lock() { - frames.insert(name.to_string(), true); - } - } - - pub(crate) fn forget(&self, name: &str) { - if let Ok(mut frames) = self.frames.lock() - && let Some(registered) = frames.get_mut(name) - { - *registered = false; - } - } } /// The rows a statement gave back. diff --git a/src/frame.rs b/src/frame.rs index a872042..33e3121 100644 --- a/src/frame.rs +++ b/src/frame.rs @@ -1,14 +1,28 @@ -//! A frame of columns, read out of whatever the caller is holding. +//! A frame of columns, described where the caller keeps them. //! //! The way in for data that is already in columns. pandas, polars and //! pyarrow all hand out an Arrow stream through the PyCapsule protocol, -//! and reading one costs a memcpy per column and no Python objects at -//! all, which is the same trade the way out makes and for the same -//! reason. +//! and what comes back from one is buffers: eight-byte words back to +//! back, one bit a row for a boolean, characters end to end with +//! offsets cutting them up. That is how this engine lays a column out +//! too, so what this module produces is not a copy of any of it but a +//! description: where each column is, how wide its values are, and what +//! they mean. //! -//! What comes back is the same [`Column`] buffers a load and an -//! appender already use, so nothing downstream of here knows whether -//! the rows arrived as Arrow or as Python lists. +//! Two things do copy, and both are said rather than hidden. A stream +//! of several batches is concatenated into one, because a column of a +//! table is one run of bytes and two batches are two of them; that is a +//! memcpy per column and it happens once. A dictionary of Python lists +//! is read into buffers of this module's own, because a list holds +//! objects and a column holds numbers, so there is nothing there to +//! point at. +//! +//! What keeps the bytes alive is [`Held`], which the frame holds and +//! which the engine drops when the last table naming those bytes goes. +//! Dropping it releases the Arrow arrays back to whoever exported them, +//! and that release may reach into an interpreter, so it takes the GIL +//! first: the engine drops a frame on whichever thread finished with +//! it, and that is not always a thread of Python's. //! //! There is no null anywhere in it. A property that is null is one no //! row of this engine can hold, so a column with a gap in it can only @@ -17,25 +31,20 @@ //! knows only that something somewhere was empty. use std::ffi::CStr; +use std::ptr::NonNull; +use std::sync::Arc; -use arrow::array::{ - Array, ArrayRef, BooleanArray, LargeStringArray, PrimitiveArray, StringArray, StringViewArray, -}; -use arrow::datatypes::{ - DataType, Date32Type, DurationMicrosecondType, DurationMillisecondType, DurationNanosecondType, - DurationSecondType, Float32Type, Float64Type, Int8Type, Int16Type, Int32Type, Int64Type, - IntervalUnit, IntervalYearMonthType, Time32MillisecondType, Time32SecondType, - Time64MicrosecondType, Time64NanosecondType, TimeUnit, TimestampMicrosecondType, - TimestampMillisecondType, TimestampNanosecondType, TimestampSecondType, UInt8Type, UInt16Type, - UInt32Type, UInt64Type, -}; +use arrow::array::{Array, ArrayRef, LargeStringArray, StringArray}; +use arrow::compute::{concat, concat_batches}; +use arrow::datatypes::{DataType, IntervalUnit, TimeUnit}; use arrow::ffi_stream::{ArrowArrayStreamReader, FFI_ArrowArrayStream}; use arrow::record_batch::{RecordBatch, RecordBatchReader}; use pyo3::prelude::*; use pyo3::types::{PyCapsule, PyDict}; -use zu_common::DurationKind; +use zu_common::{DurationKind, FloatBits, IntBits, LogicalType}; +use zudb::{Column, Layout}; -use crate::buffer::{Column, type_name}; +use crate::buffer::{self, type_name}; use crate::columns::Snag; use crate::load; @@ -43,13 +52,92 @@ use crate::load; /// checks before it reads the pointer. const STREAM: &CStr = c"arrow_array_stream"; -/// Columns, in the order the frame holds them, and how long they are. -pub struct Frame { - pub columns: Vec<(String, Column)>, - pub rows: usize, +/// What a registered frame's bytes are, and the thing whose life is +/// their life. +/// +/// One of the two fields is filled. The batch is the arrays a producer +/// handed over, held so the buffers underneath them stay where they +/// are; the owned vectors are what this client built out of Python +/// objects, for a caller with no frame library installed. +pub struct Held { + /// An `Option` only so that [`Drop`] can take it, which is where + /// the GIL has to be held. + batch: Option, + owned: Vec, +} + +/// Releasing an imported Arrow array calls the callback the producer +/// gave with it, and a producer that is pyarrow drops Python objects in +/// that callback. The engine drops a frame when the last table naming +/// it goes, which may be inside a statement on a thread holding no GIL, +/// so this takes one. Attaching on a thread that already has it costs a +/// check. +impl Drop for Held { + fn drop(&mut self) { + if self.batch.is_some() { + Python::attach(|_| drop(self.batch.take())); + } + } +} + +/// One column this client built, because a list of Python objects is +/// not a column and something has to hold the bytes. +/// +/// The variants are the layouts the engine reads rather than the types +/// a value has: a date and a count of days are one variant, and what +/// tells them apart is the logical type recorded beside the pointer. +enum Bytes { + /// Signed 64-bit values, whatever they count. + Counts(Vec), + /// Signed 32-bit values, which is what a date is. + Days(Vec), + Floats(Vec), + /// One bit a row, low bit of the first byte first. + Bits(Vec), + /// Arrow's `Utf8`: characters end to end and `rows + 1` offsets. + Text { + data: Vec, + offsets: Vec, + }, } -/// Reads whatever the caller handed over into columns. +/// A frame as the engine is about to be told about it. +/// +/// The columns hold raw pointers into what `held` keeps alive, which is +/// why the two travel together and why neither is any use without the +/// other. +pub struct Described { + pub columns: Vec, + pub rows: u64, + held: Arc, +} + +// A pointer is not `Send`, because Rust cannot know what it addresses. +// These address buffers the `Arc` in the same struct keeps alive, and +// that `Arc` is `Send` and `Sync`, so a description travels wherever +// the thing it describes does. Nothing writes through them. +unsafe impl Send for Described {} + +impl Described { + /// Registers this as a table named `name` on `engine`. + /// + /// Every pointer in it addresses a buffer the `Arc` handed over + /// with them keeps alive, which is what building one of these + /// promises and what nothing between there and here undoes, so the + /// `unsafe` that [`zudb::Frame::new`] asks for is discharged where + /// the buffers were described rather than here. + pub fn register(self, engine: &mut zudb::Connection, name: &str) -> zudb::Result<()> { + let Described { + columns, + rows, + held, + } = self; + let frame = unsafe { zudb::Frame::new(name, rows, columns, held) }?; + engine.register(frame) + } +} + +/// Reads whatever the caller handed over into a description of it. /// /// Two shapes, and the first of them is the one that matters: anything /// that speaks Arrow, which is every frame library worth naming, and a @@ -57,14 +145,12 @@ pub struct Frame { /// Anything else is refused here rather than iterated hopefully, /// because the message a caller wants is the list of what would have /// worked. -pub fn read(py: Python<'_>, data: &Bound<'_, PyAny>) -> PyResult { +pub fn read(py: Python<'_>, data: &Bound<'_, PyAny>) -> PyResult { if data.hasattr("__arrow_c_stream__")? { return from_arrow(py, data); } if let Ok(dict) = data.cast::() { - let columns = load::build(Some(dict))?; - let rows = columns.first().map(|(_, column)| column.len()).unwrap_or(0); - return Ok(Frame { columns, rows }); + return from_lists(dict); } Err(pyo3::exceptions::PyTypeError::new_err(format!( "a frame is anything with `__arrow_c_stream__`, which a pandas, polars or pyarrow table \ @@ -78,11 +164,10 @@ pub fn read(py: Python<'_>, data: &Bound<'_, PyAny>) -> PyResult { /// /// The GIL is held to pull each batch, because the producer on the far /// side of the stream may be Python and calling into an interpreter -/// nobody is holding is how a process ends. It goes down again for the -/// conversion, which is the part that touches every value and is pure -/// Rust: a ten million row frame is a memcpy per column rather than ten -/// million reference counts. -fn from_arrow(py: Python<'_>, data: &Bound<'_, PyAny>) -> PyResult { +/// nobody is holding is how a process ends. It goes down for everything +/// after that, which is pure Rust and, on the path that matters, walks +/// the pointers rather than the rows. +fn from_arrow(py: Python<'_>, data: &Bound<'_, PyAny>) -> PyResult { // No requested schema. Asking for one would mean casting on the // producer's side, and a producer that cannot cast is entitled to // refuse: what is wanted here is whatever it already has. @@ -112,280 +197,436 @@ fn from_arrow(py: Python<'_>, data: &Bound<'_, PyAny>) -> PyResult { }; let reader = ArrowArrayStreamReader::try_new(stream).map_err(|err| Snag::Arrow(err).raise(py))?; - let schema = reader.schema(); - let mut columns: Vec<(String, Column)> = Vec::with_capacity(schema.fields().len()); - for field in schema.fields() { - columns.push(( - field.name().clone(), - start(field.name(), field.data_type())?, - )); - } - let mut rows = 0usize; + let mut batches = Vec::new(); for batch in reader { - let batch = batch.map_err(|err| Snag::Arrow(err).raise(py))?; - let base = rows; - py.detach(|| take(&mut columns, &batch, base)) - .map_err(|snag| snag.raise(py))?; - rows += batch.num_rows(); + batches.push(batch.map_err(|err| Snag::Arrow(err).raise(py))?); + } + + let batch = py + .detach(|| -> Result { + // One batch is the frame where it lies. Several are one + // memcpy per column, because a column of a table is one run + // of bytes and a table of two batches is two of them. None + // is an empty frame, which the caller refuses under the + // name it was asked to register. + let batch = match batches.len() { + 1 => batches.pop().expect("the one batch"), + _ => concat_batches(&schema, &batches)?, + }; + settled(batch) + }) + .map_err(|snag| snag.raise(py))?; + + let rows = batch.num_rows() as u64; + let mut columns = Vec::with_capacity(batch.num_columns()); + for (at, array) in batch.columns().iter().enumerate() { + let name = batch.schema_ref().field(at).name(); + let (ty, layout) = described(name, array)?; + columns.push(Column { + name: name.clone(), + ty, + layout, + }); + } + Ok(Described { + columns, + rows, + held: Arc::new(Held { + batch: Some(batch), + owned: Vec::new(), + }), + }) +} + +/// A batch with everything about it that a pointer cannot express taken +/// out of it. +/// +/// A producer may hand over a slice of a longer array, and a slice does +/// not start where its buffers do: a bitmap that begins partway into a +/// byte and offsets that begin partway into their data are both things +/// a bare pointer does not say. A skewed column is copied down to +/// itself, which costs the rows it actually holds and no more, and +/// every other column is left where it is. +/// +/// Runs with the GIL down, so what is wrong comes back rather than +/// being raised. +fn settled(batch: RecordBatch) -> Result { + let mut columns = batch.columns().to_vec(); + let mut skewed = false; + for column in &mut columns { + if offset(column) { + // An empty slice of it goes in front, because `concat` of + // one array hands that array straight back: the right + // answer for a concatenation and the wrong one here, where + // the copy is the whole point. + let nothing = column.slice(0, 0); + *column = concat(&[nothing.as_ref(), column.as_ref()])?; + skewed = true; + } + } + let batch = match skewed { + true => RecordBatch::try_new(batch.schema(), columns)?, + false => batch, + }; + for (at, column) in batch.columns().iter().enumerate() { + whole(batch.schema_ref().field(at).name(), column)?; + } + Ok(batch) +} + +/// Whether this column starts somewhere other than where its buffers +/// do. +/// +/// Two ways it can, because a slice reaches this by two roads. An array +/// may carry the row offset itself, which is what `offset` is; or the +/// producer may have handed the offset over already applied to the +/// buffers, which is what an Arrow import does to a string column, and +/// then the array counts from zero and its first offset does not. +fn offset(array: &ArrayRef) -> bool { + if array.offset() != 0 { + return true; } - Ok(Frame { columns, rows }) + let starts_at = |from: Option| from.is_some_and(|from| from != 0); + match array.data_type() { + DataType::Utf8 => starts_at( + array + .as_any() + .downcast_ref::() + .and_then(|text| text.value_offsets().first().map(|&from| from as i64)), + ), + DataType::LargeUtf8 => starts_at( + array + .as_any() + .downcast_ref::() + .and_then(|text| text.value_offsets().first().copied()), + ), + _ => false, + } +} + +/// That a column has a value in every row of it. +/// +/// Names the first row rather than the count, because a caller with a +/// gap in a column wants to go and look at it. +fn whole(name: &str, array: &ArrayRef) -> Result<(), Snag> { + if array.null_count() == 0 { + return Ok(()); + } + let row = (0..array.len()) + .find(|&row| array.is_null(row)) + .unwrap_or(0); + Err(Snag::Value(format!( + "column '{name}' has no value at row {row}, and every column of a row holds one" + ))) +} + +/// Reads a dictionary of Python lists into buffers of this module's +/// own. +/// +/// The one path that copies, and it copies because there is nothing to +/// point at: a Python list holds objects and a column holds numbers. +/// What decides a column's type is its first value, which is +/// [`buffer::Column`]'s rule and the loader's, so a dictionary and a +/// load read the same way and refuse the same things. +fn from_lists(dict: &Bound<'_, PyDict>) -> PyResult { + let built = load::build(Some(dict))?; + let rows = built.first().map(|(_, column)| column.len()).unwrap_or(0) as u64; + let mut named = Vec::with_capacity(built.len()); + let mut owned = Vec::with_capacity(built.len()); + for (name, column) in built { + let (ty, bytes) = packed(&name, column)?; + named.push((name, ty)); + owned.push(bytes); + } + let held = Arc::new(Held { batch: None, owned }); + let columns = named + .into_iter() + .zip(&held.owned) + .map(|((name, ty), bytes)| Column { + name, + ty, + layout: lent(bytes), + }) + .collect(); + Ok(Described { + columns, + rows, + held, + }) +} + +/// One column of Python values, as the bytes a frame reads and what +/// they mean. +fn packed(name: &str, column: buffer::Column) -> PyResult<(LogicalType, Bytes)> { + let counts = LogicalType::Int { + signed: true, + bits: IntBits::B64, + precision: None, + }; + let characters = LogicalType::Str { + min: None, + max: None, + fixed: false, + }; + // The temporal buffers count in `i64` and the integer one holds the + // bits of one in a `u64`, and both of them are the eight-byte lane + // the engine reads through, so this cast is the identity on the + // bytes and the logical type beside them is what says how to read + // them. + let same = |v: Vec| Bytes::Counts(v.into_iter().map(|n| n as u64).collect()); + Ok(match column { + buffer::Column::Int(v) => (counts, Bytes::Counts(v)), + buffer::Column::Float(v) => ( + LogicalType::Float { + bits: FloatBits::B64, + precision: None, + }, + Bytes::Floats(v), + ), + buffer::Column::Bool(v) => { + // At least one byte, so the pointer is an allocation and + // not the dangling address an empty vector lends out. A + // frame of no rows never reaches a read, but it does reach + // the check that the pointer is not null. + let mut bits = vec![0u8; v.len().div_ceil(8).max(1)]; + for (row, &yes) in v.iter().enumerate() { + if yes { + bits[row / 8] |= 1 << (row % 8); + } + } + (LogicalType::Bool, Bytes::Bits(bits)) + } + buffer::Column::Str(v) => { + let mut data = Vec::with_capacity(v.iter().map(String::len).sum::().max(1)); + let mut offsets = Vec::with_capacity(v.len() + 1); + offsets.push(0i32); + for word in &v { + data.extend_from_slice(word.as_bytes()); + let end = i32::try_from(data.len()).map_err(|_| { + pyo3::exceptions::PyValueError::new_err(format!( + "column '{name}' holds more than two gigabytes of characters, which is \ + further than the offsets of a frame reach" + )) + })?; + offsets.push(end); + } + (characters, Bytes::Text { data, offsets }) + } + // Refused by the loader that built the column, so this arm is + // here to be exhaustive rather than to be reached. + buffer::Column::Bytes(_) => { + return Err(pyo3::exceptions::PyTypeError::new_err(format!( + "column '{name}' holds byte strings, and no statement can read a column of bytes \ + back yet" + ))); + } + buffer::Column::Date(v) => (LogicalType::Date, Bytes::Days(v)), + buffer::Column::LocalTime(v) => (LogicalType::LocalTime, same(v)), + buffer::Column::LocalDatetime(v) => (LogicalType::LocalDatetime, same(v)), + buffer::Column::Duration(kind, v) => (LogicalType::Duration(kind), same(v)), + }) } -/// The buffer a column of this Arrow type fills, or the reason there -/// is none. +/// Where one of this module's own buffers is, as a layout. +fn lent(bytes: &Bytes) -> Layout { + match bytes { + Bytes::Counts(v) => Layout::Int { + ptr: at(v.as_ptr()), + bits: IntBits::B64, + signed: true, + scale: 1, + }, + Bytes::Days(v) => Layout::Int { + ptr: at(v.as_ptr()), + bits: IntBits::B32, + signed: true, + scale: 1, + }, + Bytes::Floats(v) => Layout::Float { + ptr: at(v.as_ptr()), + bits: FloatBits::B64, + }, + Bytes::Bits(v) => Layout::Bool { + ptr: at(v.as_ptr()), + }, + Bytes::Text { data, offsets } => Layout::Str { + offsets: at(offsets.as_ptr()), + wide: false, + data: at(data.as_ptr()), + data_len: data.len(), + }, + } +} + +/// A pointer to the start of something this process is holding. +/// +/// Never null: a vector's pointer is an allocation or a dangling +/// aligned address, and an Arrow buffer's is an allocation. Neither of +/// them is address zero. +fn at(ptr: *const T) -> NonNull { + NonNull::new(ptr as *mut u8).expect("a buffer of this process is never at address zero") +} + +/// Where one Arrow column is and what it means. +/// +/// The buffers come off the array data rather than off a downcast per +/// type, because what is wanted is the same thing every time: the run +/// of bytes Arrow put the values in. The row offset was taken out of +/// the array before this, so buffer zero starts at row zero. /// -/// The type decides it and not the first value, which is the one thing -/// a frame gives that a list of Python objects does not: an empty frame -/// still knows what its columns are, and a column of floats that -/// happens to start with a whole number is still a column of floats. -fn start(name: &str, ty: &DataType) -> PyResult { +/// The scale is what one value is multiplied by to reach the unit its +/// meaning counts in, which is where Arrow's microseconds meet this +/// engine's nanoseconds. Nothing is converted here: the multiplication +/// happens per scanned chunk, on the rows a statement actually reads. +fn described(name: &str, array: &ArrayRef) -> PyResult<(LogicalType, Layout)> { + let ty = array.data_type(); let refused = |instead: &str| { - Err(pyo3::exceptions::PyTypeError::new_err(format!( - "column '{name}' is {ty}, and {instead}" - ))) + pyo3::exceptions::PyTypeError::new_err(format!("column '{name}' is {ty}, and {instead}")) + }; + let data = array.to_data(); + let bufs = data.buffers(); + let word = |bits: IntBits, signed: bool, scale: i64, means: LogicalType| { + ( + means, + Layout::Int { + ptr: at(bufs[0].as_ptr()), + bits, + signed, + scale, + }, + ) + }; + let plain = |signed: bool, bits: IntBits| LogicalType::Int { + signed, + bits, + precision: None, + }; + let float = |bits: FloatBits| { + ( + LogicalType::Float { + bits, + precision: None, + }, + Layout::Float { + ptr: at(bufs[0].as_ptr()), + bits, + }, + ) + }; + let characters = || LogicalType::Str { + min: None, + max: None, + fixed: false, + }; + let text = |wide: bool| { + ( + characters(), + Layout::Str { + offsets: at(bufs[0].as_ptr()), + wide, + data: at(bufs[1].as_ptr()), + data_len: bufs[1].len(), + }, + ) }; Ok(match ty { - DataType::Boolean => Column::Bool(Vec::new()), - DataType::Int8 - | DataType::Int16 - | DataType::Int32 - | DataType::Int64 - | DataType::UInt8 - | DataType::UInt16 - | DataType::UInt32 - | DataType::UInt64 => Column::Int(Vec::new()), - DataType::Float32 | DataType::Float64 => Column::Float(Vec::new()), + DataType::Boolean => ( + LogicalType::Bool, + Layout::Bool { + ptr: at(bufs[0].as_ptr()), + }, + ), + DataType::Int8 => word(IntBits::B8, true, 1, plain(true, IntBits::B8)), + DataType::Int16 => word(IntBits::B16, true, 1, plain(true, IntBits::B16)), + DataType::Int32 => word(IntBits::B32, true, 1, plain(true, IntBits::B32)), + DataType::Int64 => word(IntBits::B64, true, 1, plain(true, IntBits::B64)), + DataType::UInt8 => word(IntBits::B8, false, 1, plain(false, IntBits::B8)), + DataType::UInt16 => word(IntBits::B16, false, 1, plain(false, IntBits::B16)), + DataType::UInt32 => word(IntBits::B32, false, 1, plain(false, IntBits::B32)), + DataType::UInt64 => word(IntBits::B64, false, 1, plain(false, IntBits::B64)), + DataType::Float32 => float(FloatBits::B32), + DataType::Float64 => float(FloatBits::B64), // The three string layouts of Arrow, which are three ways of // holding the same characters: offsets into one buffer, wider // offsets into one buffer, and views into several. polars hands - // out the third by default and pandas the first, and a column of - // this table holds strings either way. - DataType::Utf8 | DataType::LargeUtf8 | DataType::Utf8View => Column::Str(Vec::new()), - DataType::Date32 => Column::Date(Vec::new()), - DataType::Time32(_) | DataType::Time64(_) => Column::LocalTime(Vec::new()), - DataType::Timestamp(_, None) => Column::LocalDatetime(Vec::new()), + // out the third by default and pandas the first, and a column + // of this table holds strings either way. A short view is + // already this engine's own, byte for byte. + DataType::Utf8 => text(false), + DataType::LargeUtf8 => text(true), + DataType::Utf8View => ( + characters(), + Layout::View { + views: at(bufs[0].as_ptr()), + data: bufs[1..] + .iter() + .map(|buf| (at(buf.as_ptr()), buf.len())) + .collect(), + }, + ), + DataType::Date32 => word(IntBits::B32, true, 1, LogicalType::Date), + DataType::Time32(TimeUnit::Second) => { + word(IntBits::B32, true, 1_000_000_000, LogicalType::LocalTime) + } + DataType::Time32(TimeUnit::Millisecond) => { + word(IntBits::B32, true, 1_000_000, LogicalType::LocalTime) + } + DataType::Time64(TimeUnit::Microsecond) => { + word(IntBits::B64, true, 1_000, LogicalType::LocalTime) + } + DataType::Time64(TimeUnit::Nanosecond) => { + word(IntBits::B64, true, 1, LogicalType::LocalTime) + } + DataType::Timestamp(unit, None) => { + word(IntBits::B64, true, nanos(unit), LogicalType::LocalDatetime) + } DataType::Timestamp(_, Some(zone)) => { - return refused(&format!( + return Err(refused(&format!( "a column of this table has nowhere to keep '{zone}', so drop the zone once the \ values are in the zone you want them in, or write it as a string" - )); - } - DataType::Duration(_) => Column::Duration(DurationKind::DayTime, Vec::new()), - DataType::Interval(IntervalUnit::YearMonth) => { - Column::Duration(DurationKind::YearMonth, Vec::new()) + ))); } + DataType::Duration(unit) => word( + IntBits::B64, + true, + nanos(unit), + LogicalType::Duration(DurationKind::DayTime), + ), + DataType::Interval(IntervalUnit::YearMonth) => word( + IntBits::B32, + true, + 1, + LogicalType::Duration(DurationKind::YearMonth), + ), DataType::Dictionary(_, _) => { - return refused( + return Err(refused( "a dictionary is a layout rather than a type here, so cast it to what it holds \ first", - ); + )); } DataType::Binary | DataType::LargeBinary | DataType::BinaryView => { - return refused( - "no statement can read a column of bytes back yet, so writing one would be \ - writing data the caller cannot get at", - ); + return Err(refused( + "no statement can read a column of bytes back yet, so registering one would be \ + naming data the caller cannot get at", + )); } _ => { - return refused( + return Err(refused( "a column holds booleans, integers, floats, strings, dates, times, datetimes or \ durations", - ); + )); } }) } -/// Every column of one batch, appended to the buffers it belongs in. -/// -/// Runs with the GIL released and so cannot raise: what goes wrong -/// comes back as a [`Snag`] and is raised by the caller. -fn take(columns: &mut [(String, Column)], batch: &RecordBatch, base: usize) -> Result<(), Snag> { - for (at, (name, column)) in columns.iter_mut().enumerate() { - let array = batch.column(at); - if array.null_count() > 0 { - let row = (0..array.len()) - .find(|&row| array.is_null(row)) - .unwrap_or(0); - return Err(Snag::Value(format!( - "column '{name}' has no value at row {}, and every column of a row holds one", - base + row - ))); - } - one(name, column, array, base)?; - } - Ok(()) -} - -/// Every value of a primitive array, converted and pushed. -/// -/// A macro rather than a function because the conversion differs per -/// arm in both directions: a `Time32` in seconds and an `Int32` are -/// the same bits and not the same value, and an integer column keeps -/// its bits where a temporal one keeps a count. Twenty arms of the same -/// three lines is what this saves. -macro_rules! pushed { - ($array:expr, $ty:ty, $out:expr, $changed:expr, $convert:expr) => {{ - let values = $array - .as_any() - .downcast_ref::>() - .ok_or_else($changed)?; - $out.extend(values.values().iter().copied().map($convert)); - }}; -} - -/// One column of one batch. -/// -/// The pair is matched rather than the array alone, so a stream that -/// changed a column's type between batches is caught here instead of -/// filling one buffer with values of two kinds. It cannot happen -/// through a producer that keeps to its own schema, which is why the -/// arm says so rather than explaining itself. -fn one(name: &str, column: &mut Column, array: &ArrayRef, base: usize) -> Result<(), Snag> { - let changed = || { - Snag::Type(format!( - "column '{name}' arrived as {} in a later batch than the one that declared it, which \ - is a producer that broke its own schema", - array.data_type() - )) - }; - match (array.data_type(), column) { - (DataType::Boolean, Column::Bool(out)) => { - let values = array - .as_any() - .downcast_ref::() - .ok_or_else(changed)?; - out.extend(values.values().iter()); - } - (DataType::Int8, Column::Int(out)) => { - pushed!(array, Int8Type, out, changed, |n| n as i64 as u64) - } - (DataType::Int16, Column::Int(out)) => { - pushed!(array, Int16Type, out, changed, |n| n as i64 as u64) - } - (DataType::Int32, Column::Int(out)) => { - pushed!(array, Int32Type, out, changed, |n| n as i64 as u64) - } - (DataType::Int64, Column::Int(out)) => { - pushed!(array, Int64Type, out, changed, |n| n as u64) - } - (DataType::UInt8, Column::Int(out)) => { - pushed!(array, UInt8Type, out, changed, u64::from) - } - (DataType::UInt16, Column::Int(out)) => { - pushed!(array, UInt16Type, out, changed, u64::from) - } - (DataType::UInt32, Column::Int(out)) => { - pushed!(array, UInt32Type, out, changed, u64::from) - } - (DataType::UInt64, Column::Int(out)) => { - // The one integer width that does not fit. A value above - // what an INT64 holds is refused where it sits rather than - // written as a negative number, which is the kind of thing - // a caller finds out about a year later. - let values = array - .as_any() - .downcast_ref::>() - .ok_or_else(changed)?; - for (row, &value) in values.values().iter().enumerate() { - if value > i64::MAX as u64 { - return Err(Snag::Value(format!( - "column '{name}' holds {value} at row {}, which is larger than the largest \ - integer a column of a table holds", - base + row - ))); - } - out.push(value); - } - } - (DataType::Float32, Column::Float(out)) => { - pushed!(array, Float32Type, out, changed, f64::from) - } - (DataType::Float64, Column::Float(out)) => { - pushed!(array, Float64Type, out, changed, |f| f) - } - (DataType::Utf8, Column::Str(out)) => { - let values = array - .as_any() - .downcast_ref::() - .ok_or_else(changed)?; - out.extend(values.iter().map(|s| s.unwrap_or_default().to_string())); - } - (DataType::LargeUtf8, Column::Str(out)) => { - let values = array - .as_any() - .downcast_ref::() - .ok_or_else(changed)?; - out.extend(values.iter().map(|s| s.unwrap_or_default().to_string())); - } - (DataType::Utf8View, Column::Str(out)) => { - let values = array - .as_any() - .downcast_ref::() - .ok_or_else(changed)?; - out.extend(values.iter().map(|s| s.unwrap_or_default().to_string())); - } - (DataType::Date32, Column::Date(out)) => { - pushed!(array, Date32Type, out, changed, |days| days) - } - (DataType::Time32(TimeUnit::Second), Column::LocalTime(out)) => { - pushed!(array, Time32SecondType, out, changed, |n| i64::from(n) - * 1_000_000_000) - } - (DataType::Time32(TimeUnit::Millisecond), Column::LocalTime(out)) => { - pushed!(array, Time32MillisecondType, out, changed, |n| i64::from(n) - * 1_000_000) - } - (DataType::Time64(TimeUnit::Microsecond), Column::LocalTime(out)) => { - pushed!(array, Time64MicrosecondType, out, changed, |n| n * 1_000) - } - (DataType::Time64(TimeUnit::Nanosecond), Column::LocalTime(out)) => { - pushed!(array, Time64NanosecondType, out, changed, |n| n) - } - (DataType::Timestamp(TimeUnit::Second, None), Column::LocalDatetime(out)) => { - pushed!(array, TimestampSecondType, out, changed, |n| n - * 1_000_000_000) - } - (DataType::Timestamp(TimeUnit::Millisecond, None), Column::LocalDatetime(out)) => { - pushed!(array, TimestampMillisecondType, out, changed, |n| n - * 1_000_000) - } - (DataType::Timestamp(TimeUnit::Microsecond, None), Column::LocalDatetime(out)) => { - pushed!(array, TimestampMicrosecondType, out, changed, |n| n * 1_000) - } - (DataType::Timestamp(TimeUnit::Nanosecond, None), Column::LocalDatetime(out)) => { - pushed!(array, TimestampNanosecondType, out, changed, |n| n) - } - (DataType::Duration(TimeUnit::Second), Column::Duration(DurationKind::DayTime, out)) => { - pushed!(array, DurationSecondType, out, changed, |n| n - * 1_000_000_000) - } - ( - DataType::Duration(TimeUnit::Millisecond), - Column::Duration(DurationKind::DayTime, out), - ) => { - pushed!(array, DurationMillisecondType, out, changed, |n| n - * 1_000_000) - } - ( - DataType::Duration(TimeUnit::Microsecond), - Column::Duration(DurationKind::DayTime, out), - ) => { - pushed!(array, DurationMicrosecondType, out, changed, |n| n * 1_000) - } - ( - DataType::Duration(TimeUnit::Nanosecond), - Column::Duration(DurationKind::DayTime, out), - ) => { - pushed!(array, DurationNanosecondType, out, changed, |n| n) - } - ( - DataType::Interval(IntervalUnit::YearMonth), - Column::Duration(DurationKind::YearMonth, out), - ) => { - pushed!(array, IntervalYearMonthType, out, changed, i64::from) - } - _ => return Err(changed()), +/// How many nanoseconds one count of this unit is, which is the scale a +/// temporal column is read through. +fn nanos(unit: &TimeUnit) -> i64 { + match unit { + TimeUnit::Second => 1_000_000_000, + TimeUnit::Millisecond => 1_000_000, + TimeUnit::Microsecond => 1_000, + TimeUnit::Nanosecond => 1, } - Ok(()) } diff --git a/src/register.rs b/src/register.rs index 02ae4c0..a7c4a28 100644 --- a/src/register.rs +++ b/src/register.rs @@ -5,157 +5,88 @@ //! conn.execute("MATCH (p:people) WHERE p.age > 40 RETURN p.name AS name") //! ``` //! -//! What this is not is a scan of the frame where it sits. DuckDB's -//! `register` is zero copy because its executor can call back out to -//! the Python object holding the data; this engine has no such callback -//! yet, so a registered frame is copied in and becomes a table like any -//! other. The call is the one it would be either way, and the day the -//! engine grows a scan it becomes the cheap thing under the same name. +//! This is the replacement scan, and it copies nothing. What the engine +//! is told is where the caller's columns are, how wide their values are +//! and what they mean; a statement that matches the name builds vectors +//! that point straight at those buffers, so a frame of ten million rows +//! is registered in the time it takes to describe its columns and read +//! at the speed of the memory it already sits in. //! -//! The copy is not the slow kind. A frame arrives as Arrow buffers and -//! goes in through the same appender a bulk load uses, which is one -//! commit for the batch and no Python object per cell, so what it costs -//! is a memcpy per column and the write. +//! Because it is not a copy, a registered frame is a view and not a +//! snapshot: write into the array behind it and the next statement +//! answers what is there now. The one exception is a dictionary of +//! Python lists, which has to be read into buffers of this client's own +//! because a list holds objects rather than numbers. //! -//! Registering owns the name. A name that is already a table of the -//! database is refused rather than overwritten, since a frame knows -//! nothing about the rows that were there, and a name this connection -//! registered before is emptied and rewritten, which is what rerunning -//! a cell means by it. -//! -//! What a name cannot do is change shape. A table's columns are the -//! ones its first row declared and no statement of this engine alters -//! or drops one, so a frame with different columns under a used name is -//! refused rather than half written, and `unregister` empties the table -//! rather than removing it. Both of those are the engine showing -//! through, and both are better said than worked around quietly. +//! A frame belongs to the connection it was registered on and goes when +//! that connection does. It is not written to the database, no other +//! program opening the same file sees it, and nothing writes to it: a +//! statement that tries to insert into or delete from one is refused +//! with the reason. `unregister` takes the name away entirely, and the +//! bytes go back to the caller when the last statement reading them has +//! finished with them. use pyo3::prelude::*; -use zu_common::{DurationKind, Temporal}; -use zudb::query::Value; +use zudb::ZuError; use zudb::zu1::catalog::Catalog; -use zudb::{Field, ZuError}; use crate::conn::Connection; use crate::error::{programming, to_py_err}; -use crate::frame::{self, Frame}; - -/// What can go wrong with the GIL down, where an exception cannot be -/// built yet. -enum Snag { - Engine(ZuError), - /// A row the engine's appender would not take, named by where it - /// sits in the frame, since that is the row a caller can go and - /// look at. - Row(u64, ZuError), - /// A column the table declares and the frame has not got, which is - /// a table that was written to while it was being registered. - Missing(String), -} - -impl Snag { - fn raise(self, py: Python<'_>, name: &str) -> PyErr { - match self { - Snag::Engine(err) => to_py_err(py, err), - Snag::Row(row, err) => { - let raised = to_py_err(py, err); - programming( - py, - &format!("row {row} of this frame: {}", raised.value(py)), - ) - } - Snag::Missing(column) => programming( - py, - &format!( - "'{name}' declares a column '{column}' that this frame has not got, which is \ - a table that was written to while it was being registered" - ), - ), - } - } -} +use crate::frame; -/// Copies a frame in as a node table called `name`. +/// Registers a frame as a table called `name`. pub fn register( py: Python<'_>, - conn: Py, + conn: &Connection, name: &str, data: &Bound<'_, PyAny>, ) -> PyResult { identifier(py, name, "a registered frame")?; - let frame = frame::read(py, data)?; - if frame.columns.is_empty() { + let described = frame::read(py, data)?; + if described.columns.is_empty() { return Err(programming( py, "this frame has no columns, and a table whose rows hold nothing is not a table", )); } - if frame.rows == 0 { - return Err(programming( - py, - "this frame has no rows, and the columns of a table are declared by the first row \ - written into it", - )); + for column in &described.columns { + identifier(py, &column.name, "a column of a registered frame")?; } - for (column, _) in &frame.columns { - identifier(py, column, "a column of a registered frame")?; - } - - let held = conn.bind(py).borrow(); - let held: &Connection = &held; - // Refused inside a transaction for the reason an appender is: the - // rows after the first go in through the appender, whose batches are - // commits of their own, so half the frame would survive a rollback - // and half of it would not. - if held.in_transaction(py)? { - return Err(programming( - py, - "registering a frame writes its own commits, which a rollback does not take back, so \ - it cannot be done inside a transaction", - )); - } - let mine = held.made_here(name); - if !mine && has_table(py, held, name)? { + // Refused here rather than at the statement that would have hit it. + // The engine keeps frames in an id space of their own and would + // take this name happily; what it could not do is bind it, because + // a label in a statement is one thing and the stored table would + // win. Better said at the call that made the clash. + if has_table(py, conn, name)? { return Err(programming( py, &format!( "'{name}' is already a table of this database, and registering over one would \ - write across rows this frame knows nothing about" + hide rows this frame knows nothing about" ), )); } - // A name registered again is the same name holding something else, - // which is what a caller who reruns a cell means by it. The rows go - // first, so the frame that replaces them is the whole table rather - // than the half of it that is new. - // - // The columns cannot change, though, and that is worth saying at the - // call rather than at the row the engine refuses: a table's columns - // are the ones the first row declared, no statement of this engine - // alters or drops one, and the empty table left behind by an - // `unregister` still has them. - if mine { - same_columns(py, held, name, &frame)?; - held.run(py, &empty_out(name))?; - } - // The first row through a statement, because a table nothing - // declares is declared by the row that is written into it and no - // statement of this engine declares one on its own. Every row after - // it goes through the appender, which is the fast way and needs the - // table to be there. - held.run(py, &first(py, name, &frame)?)?; - if frame.rows > 1 { - rest(py, held, name, &frame)?; - } - held.remember(name); - Ok(frame.rows) + let rows = described.rows as usize; + // The description carries raw pointers, so it travels into the + // detached call inside the type that says who keeps them alive. The + // walk `Frame::new` does over the unsigned, scaled and string + // columns happens in there, with the GIL down. + conn.engine(py, move |engine| described.register(engine, name))? + .map_err(|err| to_py_err(py, err))?; + Ok(rows) } -/// Empties the table a registered frame was copied into and forgets the -/// name. +/// Takes a registered frame's name away and gives the bytes back. +/// +/// The bytes go when the last statement reading them lets go, which is +/// usually now and is never before: a frame a running statement is +/// still scanning is held until it ends. pub fn unregister(py: Python<'_>, conn: &Connection, name: &str) -> PyResult<()> { - if !conn.registered_here(name) { + let dropped = conn + .engine(py, |engine| engine.unregister(name))? + .map_err(|err| to_py_err(py, err))?; + if !dropped { return Err(programming( py, &format!( @@ -164,167 +95,13 @@ pub fn unregister(py: Python<'_>, conn: &Connection, name: &str) -> PyResult<()> ), )); } - conn.run(py, &empty_out(name))?; - conn.forget(name); Ok(()) } -/// The statement that takes the rows of a registered frame back out. -/// -/// `DETACH`, because a registered table's rows may have been joined to -/// since, and a delete that refused would leave the name registered -/// over rows nobody can replace. -fn empty_out(name: &str) -> String { - format!("MATCH (frame:{name}) DETACH DELETE frame") -} - -/// The `INSERT` that writes the first row and, in writing it, declares -/// the table. -/// -/// The values go in as literals rather than as parameters, because a -/// column is declared by what is written into it and a parameter is -/// worked out rather than written: the engine refuses `{uid: $uid}` on -/// a table that does not exist yet, and it is right to, since the -/// statement alone does not say what the column would hold. -fn first(py: Python<'_>, name: &str, frame: &Frame) -> PyResult { - let mut fields = Vec::with_capacity(frame.columns.len()); - for (column, values) in &frame.columns { - let value = values - .value(0) - .and_then(|value| literal(&value)) - .ok_or_else(|| { - programming( - py, - &format!( - "column '{column}' holds values no statement can write down, and the \ - columns of a table are declared by the first row written into it" - ), - ) - })?; - fields.push(format!("{column}: {value}")); - } - Ok(format!("INSERT (frame:{name} {{{}}})", fields.join(", "))) -} - -/// One value as a statement writes it, or `None` when no statement -/// writes one at all. -/// -/// The two that come back `None` are byte strings, which have no -/// literal, and year-month durations, whose only spelling is the one a -/// day-time duration already has. -fn literal(value: &Value) -> Option { - Some(match value { - Value::Int(n) => n.to_string(), - // Debug rather than Display, because Display writes `1` for a - // float that holds one and the column would come out an - // integer. Non-finite is refused for want of a literal: no - // spelling of infinity parses, and a column cannot be declared - // by a value that cannot be written. - Value::Float(f) if f.is_finite() => format!("{f:?}"), - Value::Bool(b) => b.to_string(), - Value::Str(s) => quoted(s), - Value::Temporal(t) => { - let kind = match t { - Temporal::Date(_) => "DATE", - Temporal::LocalTime(_) => "LOCAL TIME", - Temporal::LocalDatetime(_) => "LOCAL DATETIME", - Temporal::Duration(DurationKind::DayTime, _) => "DURATION", - _ => return None, - }; - format!("{kind} {}", quoted(&t.to_string())) - } - _ => return None, - }) -} - -/// A string as a statement carries it, with the five characters that -/// end a literal or start an escape written as escapes. -fn quoted(text: &str) -> String { - let mut out = String::with_capacity(text.len() + 2); - out.push('\''); - for ch in text.chars() { - match ch { - '\\' => out.push_str("\\\\"), - '\'' => out.push_str("\\'"), - '\n' => out.push_str("\\n"), - '\r' => out.push_str("\\r"), - '\t' => out.push_str("\\t"), - _ => out.push(ch), - } - } - out.push('\''); - out -} - -/// Every row after the first, through the engine's appender. -/// -/// The columns are taken in the order the table declares them rather -/// than the order the frame holds them, because the appender writes a -/// row as a list and the two orders are only the same by luck. -fn rest(py: Python<'_>, conn: &Connection, name: &str, frame: &Frame) -> PyResult<()> { - conn.engine(py, |engine| { - let order = column_order(engine, name)?; - // The frame's columns in the table's order, worked out once - // rather than looked up per value, since the whole point of the - // appender is that a value costs nothing but a copy. - let mut taken: Vec = Vec::with_capacity(order.len()); - for column in &order { - let at = frame - .columns - .iter() - .position(|(held, _)| held == column) - .ok_or_else(|| Snag::Missing(column.clone()))?; - taken.push(at); - } - let mut appender = engine.appender(name).map_err(Snag::Engine)?; - let mut row: Vec> = Vec::with_capacity(taken.len()); - for at in 1..frame.rows { - row.clear(); - row.extend( - taken - .iter() - .map(|&column| frame.columns[column].1.field(at)), - ); - appender - .append_row(&row[..]) - .map_err(|err| Snag::Row(at as u64, err))?; - } - appender.close().map(|_| ()).map_err(Snag::Engine) - })? - .map_err(|snag| snag.raise(py, name)) -} - -/// That a frame going in under a name this connection already used has -/// the columns that name was declared with. -/// -/// Names rather than types, because the engine refuses a value that does -/// not fit a column at the row that carries it and says so well, where a -/// column that is simply not there reads as a mistake about the data -/// rather than about the name. -fn same_columns(py: Python<'_>, conn: &Connection, name: &str, frame: &Frame) -> PyResult<()> { - let mut declared = conn - .engine(py, |engine| column_order(engine, name))? - .map_err(|snag| snag.raise(py, name))?; - let mut holds: Vec = frame - .columns - .iter() - .map(|(column, _)| column.clone()) - .collect(); - declared.sort(); - holds.sort(); - if declared != holds { - return Err(programming( - py, - &format!( - "'{name}' was registered with the columns {}, and this frame has {}, which no \ - statement of this engine can turn one into the other: register it under another \ - name", - declared.join(", "), - holds.join(", ") - ), - )); - } - Ok(()) +/// The names frames are registered under on this connection, sorted. +pub fn registered(py: Python<'_>, conn: &Connection) -> PyResult> { + conn.engine(py, |engine| Ok::<_, ZuError>(engine.registered()))? + .map_err(|err| to_py_err(py, err)) } /// Whether the database already holds a table of this name. @@ -334,45 +111,13 @@ fn same_columns(py: Python<'_>, conn: &Connection, name: &str, frame: &Frame) -> /// the same mistake. fn has_table(py: Python<'_>, conn: &Connection, name: &str) -> PyResult { conn.engine(py, |engine| { - let catalog = catalog(engine)?; - Ok(catalog.node_by_name(name).is_some() || catalog.rel_by_name(name).is_some()) + let file = engine.session_mut().file_mut()?; + let catalog = Catalog::load(file)?; + Ok::<_, ZuError>( + catalog.node_by_name(name).is_some() || catalog.rel_by_name(name).is_some(), + ) })? - .map_err(|snag: Snag| snag.raise(py, name)) -} - -/// The columns of a node table, in the order it declares them, which is -/// the order the appender writes a row in. -fn column_order(engine: &mut zudb::Connection, table: &str) -> Result, Snag> { - let id = catalog(engine)? - .node_by_name(table) - .map(|node| node.id) - .ok_or_else(|| { - Snag::Engine(ZuError::InvalidArgument(format!( - "no node table '{table}', which is a table that went away while it was being \ - registered" - ))) - })?; - let file = engine.session_mut().file_mut().map_err(Snag::Engine)?; - let directory = zudb::zu1::props::load_props(file, id) - .map_err(Snag::Engine)? - .ok_or_else(|| { - Snag::Engine(ZuError::InvalidArgument(format!( - "'{table}' stores no properties, so it has no columns to write into" - ))) - })?; - Ok(directory - .columns - .iter() - .map(|column| column.name.clone()) - .collect()) -} - -/// The catalog as the file has it, which is not always the one the -/// session is holding: an appender's commit folds rows in without the -/// session hearing about it. -fn catalog(engine: &mut zudb::Connection) -> Result { - let file = engine.session_mut().file_mut().map_err(Snag::Engine)?; - Catalog::load(file).map_err(Snag::Engine) + .map_err(|err| to_py_err(py, err)) } /// Refuses a name a statement could not carry. diff --git a/tests/test_register.py b/tests/test_register.py index a13a350..06ecc2c 100644 --- a/tests/test_register.py +++ b/tests/test_register.py @@ -1,11 +1,11 @@ """Frames, registered under a name a statement can match on. The point of the call is that a DataFrame a program already has becomes -something a statement can read, so most of these register one and then -match it. The rest are the refusals, which are where the engine shows -through: a table's columns are declared by its first row and nothing -alters or drops one, so a name that has been used holds the shape it was -used with. +something a statement can read, and that reading it costs nothing: the +engine is told where the columns are and reads them where they lie. So +most of these register one and then match it, and the ones that matter +most prove the two halves of that claim, which are that the bytes are +never copied and that they are handed back when the name goes. pandas and polars are asked for by the tests that need them and skipped when they are not installed, because the wheel depends on neither. @@ -14,7 +14,10 @@ from __future__ import annotations import datetime +import gc +import struct import time +from pathlib import Path import pytest import zudb @@ -66,8 +69,8 @@ def test_a_pyarrow_table_goes_in(empty: zudb.Connection) -> None: def test_a_stream_of_several_batches_is_one_frame(empty: zudb.Connection) -> None: - """Read a batch at a time, so a frame bigger than memory is not two - frames.""" + """A column of a table is one run of bytes, so this is the one shape + that costs a memcpy per column on the way in.""" schema = pa.schema([("uid", pa.int64()), ("name", pa.string())]) batches = [ pa.record_batch([pa.array([1, 2]), pa.array(["ada", "grace"])], schema=schema), @@ -78,6 +81,14 @@ def test_a_stream_of_several_batches_is_one_frame(empty: zudb.Connection) -> Non assert names(empty, "people") == ["ada", "grace", "lynn"] +def test_a_sliced_column_is_read_from_the_row_it_starts_at(empty: zudb.Connection) -> None: + """A slice is an array with a row offset, which is the one thing a + bare pointer cannot say, so it is copied down to itself first.""" + frame = pa.table({"uid": pa.array([1, 2, 3, 4]), "name": pa.array(["a", "b", "c", "d"])}) + assert empty.register("people", frame.slice(1, 2)) == 2 + assert names(empty, "people") == ["b", "c"] + + def test_every_kind_of_column_a_row_can_hold_arrives_as_itself(empty: zudb.Connection) -> None: frame = pa.table( { @@ -113,14 +124,68 @@ def test_every_kind_of_column_a_row_can_hold_arrives_as_itself(empty: zudb.Conne ) -def test_a_frame_is_a_snapshot_and_not_a_view(empty: zudb.Connection) -> None: - """The one thing about this being a copy that a caller has to know.""" - held = {"uid": [1], "name": ["ada"]} - empty.register("people", held) - held["name"][0] = "grace" +def held_column(values: list[int]) -> tuple[bytearray, object]: + """An Arrow column over memory this test keeps and can write into. + + `from_buffers` is the way to hand pyarrow bytes that are already + laid out, so what comes back points at the bytearray rather than at + a copy of it. There is no validity buffer, which is what `None` + first says, because a column of a row of this engine holds a value + everywhere. + """ + held = bytearray(struct.pack(f"<{len(values)}q", *values)) + return held, pa.Array.from_buffers(pa.int64(), len(values), [None, pa.py_buffer(held)]) + + +def test_a_frame_is_read_where_it_lies_and_not_copied(empty: zudb.Connection) -> None: + """The whole point of the call, proved the only way it can be. + + The column is written into between two statements and the second one + answers the new number, which no copy taken at registration could + do. + """ + held, column = held_column([10, 20, 30]) + empty.register("numbers", pa.table({"n": column})) + assert empty.execute("MATCH (x:numbers) RETURN sum(x.n) AS total").fetchone() == (60,) + struct.pack_into(" None: + """A bytearray refuses to resize while anything is holding a buffer + of it, so whether it will is exactly the question of whether the + engine has let go.""" + held, column = held_column([1, 2, 3]) + empty.register("numbers", pa.table({"n": column})) + del column + with pytest.raises(BufferError): + held.append(0) + empty.unregister("numbers") + gc.collect() + held.append(0) + + +def test_a_dictionary_of_lists_is_copied_because_a_list_is_not_a_column( + empty: zudb.Connection, +) -> None: + """The one way in that does copy, and the reason it has to.""" + lists = {"uid": [1], "name": ["ada"]} + empty.register("people", lists) + lists["name"][0] = "grace" assert names(empty, "people") == ["ada"] +def test_a_frame_belongs_to_the_connection_that_registered_it( + empty: zudb.Connection, tmp_path: Path +) -> None: + """Nothing is written to the database, so another program opening the + same file has never heard of it.""" + empty.register("people", {"uid": [1], "name": ["ada"]}) + with zudb.connect(tmp_path / "empty.zu1") as other: + assert other.registered == [] + assert names(other, "people") == [] + + def test_registered_says_what_is_registered_here(empty: zudb.Connection) -> None: assert empty.registered == [] empty.register("second", {"a": [1]}) @@ -128,14 +193,23 @@ def test_registered_says_what_is_registered_here(empty: zudb.Connection) -> None assert empty.registered == ["first", "second"] -def test_registering_a_name_again_replaces_the_rows(empty: zudb.Connection) -> None: +def test_registering_a_name_again_replaces_what_it_stands_for(empty: zudb.Connection) -> None: """Which is what rerunning a cell means by it.""" empty.register("people", {"uid": [1, 2], "name": ["ada", "grace"]}) assert empty.register("people", {"uid": [3], "name": ["lynn"]}) == 1 assert names(empty, "people") == ["lynn"] -def test_unregister_takes_the_rows_back_out(empty: zudb.Connection) -> None: +def test_a_name_registered_again_may_hold_a_different_shape(empty: zudb.Connection) -> None: + """A frame is not a table, so nothing about the first registration + survives the second.""" + empty.register("people", {"uid": [1], "name": ["ada"]}) + empty.register("people", {"name": ["grace"], "age": [45]}) + row = empty.execute("MATCH (p:people) RETURN p.name AS name, p.age AS age").fetchone() + assert row == ("grace", 45) + + +def test_unregister_takes_the_name_away(empty: zudb.Connection) -> None: empty.register("people", {"uid": [1, 2], "name": ["ada", "grace"]}) empty.unregister("people") assert empty.registered == [] @@ -143,7 +217,6 @@ def test_unregister_takes_the_rows_back_out(empty: zudb.Connection) -> None: def test_a_name_that_was_unregistered_can_be_registered_again(empty: zudb.Connection) -> None: - """The empty table it left behind is still this connection's.""" empty.register("people", {"uid": [1], "name": ["ada"]}) empty.unregister("people") assert empty.register("people", {"uid": [2], "name": ["grace"]}) == 1 @@ -163,20 +236,19 @@ def test_unregistering_a_table_nobody_registered_is_refused(social: zudb.Connect def test_registering_over_a_table_of_the_database_is_refused(social: zudb.Connection) -> None: - """A frame knows nothing about the rows that were already there.""" + """A statement naming it would mean the stored one.""" with pytest.raises(zudb.ProgrammingError, match="already a table of this database"): social.register("person", {"uid": [1]}) -def test_a_name_that_has_been_used_keeps_the_columns_it_was_used_with( - empty: zudb.Connection, -) -> None: - """No statement of this engine alters a table, so this is refused at - the call rather than halfway through the rows.""" +def test_nothing_writes_to_a_registered_frame(empty: zudb.Connection) -> None: + """It is the caller's memory, read where it lies, and a statement + that wrote into it would be writing into the DataFrame.""" empty.register("people", {"uid": [1], "name": ["ada"]}) - with pytest.raises(zudb.ProgrammingError, match="register it under another name"): - empty.register("people", {"uid": [1], "name": ["ada"], "age": [36]}) - assert names(empty, "people") == ["ada"] + with pytest.raises(zudb.TransactionError, match="never written"): + empty.execute("INSERT (p:people {uid: 2, name: 'grace'})") + with pytest.raises(zudb.TransactionError, match="never written"): + empty.execute("MATCH (p:people) DETACH DELETE p") def test_a_null_anywhere_is_refused_by_column_and_row(empty: zudb.Connection) -> None: @@ -186,11 +258,11 @@ def test_a_null_anywhere_is_refused_by_column_and_row(empty: zudb.Connection) -> empty.register("people", frame) -def test_a_frame_with_no_rows_is_refused(empty: zudb.Connection) -> None: - """The columns of a table are declared by the first row written into - it, and there is none.""" - with pytest.raises(zudb.ProgrammingError, match="no rows"): - empty.register("people", pa.table({"uid": pa.array([], pa.int64())})) +def test_a_frame_with_no_rows_registers_and_matches_nothing(empty: zudb.Connection) -> None: + """A frame knows what its columns are without being told by a row, so + a filter that came back empty is still a table to match on.""" + assert empty.register("people", pa.table({"uid": pa.array([], pa.int64())})) == 0 + assert empty.execute("MATCH (p:people) RETURN count(*) AS n").fetchone() == (0,) def test_a_frame_with_no_columns_is_refused(empty: zudb.Connection) -> None: @@ -212,14 +284,16 @@ def test_a_zoned_timestamp_is_refused_with_what_to_do_about_it(empty: zudb.Conne def test_a_column_of_bytes_is_refused(empty: zudb.Connection) -> None: - """Writing one would be writing data no statement reads back.""" + """Naming one would be naming data no statement reads back.""" with pytest.raises(TypeError, match="column of bytes"): empty.register("blobs", pa.table({"raw": pa.array([b"x"])})) def test_an_integer_too_large_for_a_column_is_refused_by_row(empty: zudb.Connection) -> None: + """Checked once, at registration, so that reading a frame cannot + fail: the engine's lane is signed and this value is not in it.""" frame = pa.table({"big": pa.array([1, 2**63], pa.uint64())}) - with pytest.raises(ValueError, match="at row 1"): + with pytest.raises(zudb.ProgrammingError, match="at row 1"): empty.register("numbers", frame) @@ -231,10 +305,11 @@ def test_something_that_is_not_a_frame_at_all_is_refused_with_the_list( def test_registering_inside_a_transaction_is_refused(empty: zudb.Connection) -> None: - """Its rows go in through the appender, whose batches are commits of - their own that no rollback reaches.""" + """A frame is registered on the session, which is the thing a + transaction is running on, and a rollback has nothing to say about + memory the caller owns.""" with empty.transaction(): - with pytest.raises(zudb.ProgrammingError, match="own commits"): + with pytest.raises(zudb.TransactionError, match="not inside a transaction"): empty.register("people", {"a": [1, 2]}) @@ -244,23 +319,20 @@ def test_a_closed_connection_registers_nothing(empty: zudb.Connection) -> None: empty.register("people", {"a": [1]}) -def test_reading_the_frame_costs_a_memcpy_and_not_a_python_object_per_cell( - empty: zudb.Connection, -) -> None: - """The read is the part this client owns, so it is the part measured. +def test_registering_costs_the_same_whatever_the_frame_holds(empty: zudb.Connection) -> None: + """Nothing is copied, so nothing about the call is per row. - A frame whose column name no statement could carry is read in full - and then refused, which times the way in without the write behind - it. The budget is loose by a factor of fifty against the 1 ms this - takes on a laptop, because it is here to catch a way in that started - making a Python object per cell rather than to hold a number. + Five million rows against ten, and the budget is a millisecond + against the 30 microseconds either of them takes on a laptop. It is + here to catch a way in that started walking the rows rather than to + hold a number, which is why it is loose by a factor of thirty. """ - rows = 200_000 - frame = pa.table({"two words": pa.array(range(rows))}) - best = float("inf") - for _ in range(3): - started = time.perf_counter() - with pytest.raises(zudb.ProgrammingError): + frames = {rows: pa.table({"n": pa.array(range(rows))}) for rows in (10, 5_000_000)} + best = {} + for rows, frame in frames.items(): + best[rows] = float("inf") + for _ in range(5): + started = time.perf_counter() empty.register("numbers", frame) - best = min(best, time.perf_counter() - started) - assert best < 50e-3, f"reading {rows} rows took {best * 1e3:.0f} ms" + best[rows] = min(best[rows], time.perf_counter() - started) + assert best[5_000_000] < 1e-3, f"registering 5m rows took {best[5_000_000] * 1e3:.1f} ms"