From 6d6f0fdb9d7da2ca880edf10a5007757dbee1aab Mon Sep 17 00:00:00 2001 From: Tam Nguyen Duc <1218621+tamnd@users.noreply.github.com> Date: Tue, 18 Aug 2026 22:29:34 +0700 Subject: [PATCH] a frame under a name a statement can match on `conn.register("people", frame)` puts a DataFrame into the database under a name, and `MATCH (p:people)` reads it. Anything that speaks Arrow goes in, which is pandas, polars, pyarrow and a reader over one, and a dictionary of lists is there for a caller with none of them installed. `unregister` takes the rows back out and `conn.registered` says what is registered here. 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 read in 5 ms, and a million rows of an integer, a float and a string in 138 ms, which is the string column being the only one that allocates. The write behind it is the engine's own, and those million rows take 9.3 seconds in all. It is a copy and not a scan, which is the one thing about it a caller has to know, and the module says why: 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 a registered frame is a snapshot, and the day the engine grows a scan it becomes the cheap thing under the same name. Two more places the engine shows through, both 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. The first row goes in as a statement with its values written out, since a table nothing declares is declared by the row written into it and a parameter is worked out rather than written; every row after it goes through the appender. Registering inside a transaction is refused for the reason an appender is, since its batches are commits of their own. A null anywhere is refused by column and row, a zoned timestamp with what to do about it, a column of bytes because no statement reads one back, and a name a statement could not carry before anything is written. Twenty-eight tests over the four ways in, every column kind a row can hold, a stream of several batches, the snapshot, the replacements and every refusal, plus a budget on the read. --- README.md | 19 +- python/zudb/_zudb.pyi | 10 ++ src/buffer.rs | 22 +++ src/conn.rs | 211 ++++++++++++++++++---- src/frame.rs | 391 ++++++++++++++++++++++++++++++++++++++++ src/lib.rs | 2 + src/load.rs | 2 +- src/register.rs | 399 +++++++++++++++++++++++++++++++++++++++++ tests/test_register.py | 266 +++++++++++++++++++++++++++ 9 files changed, 1284 insertions(+), 38 deletions(-) create mode 100644 src/frame.rs create mode 100644 src/register.rs create mode 100644 tests/test_register.py diff --git a/README.md b/README.md index d42cbaf..fdeaa8a 100644 --- a/README.md +++ b/README.md @@ -93,6 +93,23 @@ An appender is refused inside a transaction. Its batches are commits of their ow The wrapper costs 5 microseconds for an empty transaction, so what it costs is what the engine charges. On this machine that is more rather than less: 200 `INSERT`s cost 2.2 seconds each committing on its own and 3.3 seconds inside one transaction, and reads cost the same either way. A transaction here is worth taking for the span it holds and not for the time it saves, and the v0 write path is where that number has to change. +## Bringing a DataFrame in + +A frame a program already has becomes something a statement can match on, under a name the program picks. + +```python +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. + +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 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. + +`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. + ## Reading a result as columns A result is rows to iterate and columns to hand to something else. The columns go out over the Arrow C Data Interface, so pyarrow, pandas and polars each read the same buffers and none of them gets a Python object per cell. @@ -135,7 +152,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, 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. `register` 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, 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 03a053a..8afe57d 100644 --- a/python/zudb/_zudb.pyi +++ b/python/zudb/_zudb.pyi @@ -86,6 +86,16 @@ class Connection: def appender(self, table: str) -> Appender: """Opens an appender on `table`, for loading rows into a database that already exists.""" + def register(self, name: str, data: Any) -> int: + """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.""" + + @property + def registered(self) -> list[str]: + """The names this connection has registered frames under, sorted.""" + def close(self) -> None: """Closes the connection and frees what it held.""" diff --git a/src/buffer.rs b/src/buffer.rs index af0479d..fcde8c6 100644 --- a/src/buffer.rs +++ b/src/buffer.rs @@ -21,6 +21,7 @@ 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}; @@ -285,6 +286,27 @@ 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 1a4a514..7db13b9 100644 --- a/src/conn.rs +++ b/src/conn.rs @@ -7,6 +7,7 @@ //! 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}; @@ -21,6 +22,7 @@ use crate::appender::Appender; use crate::columns; use crate::error::{closed, programming, to_py_err}; use crate::interrupt; +use crate::register; use crate::txn::Transaction; use crate::value::{Names, from_py, to_py}; @@ -66,6 +68,18 @@ 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)] @@ -87,36 +101,7 @@ impl Connection { statement: &str, params: Option<&Bound<'_, PyDict>>, ) -> PyResult { - let params = bind(params)?; - // Owned rather than borrowed, both of them, because the thread - // that runs the statement is not this one and the borrow would - // have to say so. It is two allocations against a statement, - // which is nothing beside the parse it is about to have. - let statement = statement.to_string(); - // 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 - // while this one is inside the executor. Where a `Ctrl-C` can - // arrive the statement goes on the connection's own thread and - // this one waits for it, which is the only way a press is felt - // before the statement ends. - let (result, names) = - interrupt::watched(py, &self.runner, &self.inner, &self.stop, move |conn| { - let borrowed: Vec<(&str, Value)> = params - .iter() - .map(|(name, value)| (name.as_str(), value.clone())) - .collect(); - let names = Names::of(conn.session_mut().catalog()); - conn.query_with(&statement, &borrowed) - .map(|result| (result, names)) - }) - .map_err(|stopped| stopped.raise(py))? - .map_err(|err| to_py_err(py, err))?; - Ok(Result { - result, - names, - next: Mutex::new(0), - }) + self.query(py, statement, bind(params)?) } /// The same call, named for the way it reads in a notebook: @@ -164,7 +149,7 @@ impl Connection { /// statement, since the answer belongs to the session and the /// session is what the statement is holding. #[getter] - fn in_transaction(&self, py: Python<'_>) -> PyResult { + pub(crate) fn in_transaction(&self, py: Python<'_>) -> PyResult { if !self.alive.load(Ordering::Acquire) { return Err(closed(py, "this connection")); } @@ -202,6 +187,59 @@ impl Connection { Appender::open(py, slf, table) } + /// Puts a DataFrame under a name a statement can match on. + /// + /// ```python + /// conn.register("people", frame) + /// conn.execute("MATCH (p:people) WHERE p.age > 40 RETURN p.name AS name") + /// ``` + /// + /// 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. + /// + /// 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) + } + + /// Takes the rows of a registered frame back out. + /// + /// 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. + 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() + } + /// Asks the statement running on this connection to stop. /// /// The one call meant to be made from another thread while the @@ -323,20 +361,121 @@ impl Connection { inner: Arc::new(Mutex::new(Some(opened))), runner: OnceLock::new(), alive: AtomicBool::new(true), + frames: Mutex::new(BTreeMap::new()), path, read_only, }) } + /// Runs one statement, with its parameters already engine values. + /// + /// The one path every statement takes, whether the parameters came + /// from a caller's dictionary or were built here. + pub(crate) fn query( + &self, + py: Python<'_>, + statement: &str, + params: Vec<(String, Value)>, + ) -> PyResult { + // Owned rather than borrowed, because the thread that runs the + // statement is not this one and the borrow would have to say + // so. It is one allocation against a statement, which is + // nothing beside the parse it is about to have. + let statement = statement.to_string(); + // 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 + // while this one is inside the executor. Where a `Ctrl-C` can + // arrive the statement goes on the connection's own thread and + // this one waits for it, which is the only way a press is felt + // before the statement ends. + let (result, names) = + interrupt::watched(py, &self.runner, &self.inner, &self.stop, move |conn| { + let borrowed: Vec<(&str, Value)> = params + .iter() + .map(|(name, value)| (name.as_str(), value.clone())) + .collect(); + let names = Names::of(conn.session_mut().catalog()); + conn.query_with(&statement, &borrowed) + .map(|result| (result, names)) + }) + .map_err(|stopped| stopped.raise(py))? + .map_err(|err| to_py_err(py, err))?; + Ok(Result { + result, + names, + next: Mutex::new(0), + }) + } + /// Runs a statement that takes no parameters and gives back /// nothing, which is what the three transaction words are. /// - /// It goes through `execute` rather than around it, so a `COMMIT` - /// waits for the connection's lock, releases the GIL and feels a - /// `Ctrl-C` exactly as every other statement does. The result is - /// dropped, since the words return no columns. + /// It goes through the same path as every other statement, so a + /// `COMMIT` waits for the connection's lock, releases the GIL and + /// feels a `Ctrl-C` exactly as they do. The result is dropped, + /// since the words return no columns. pub(crate) fn run(&self, py: Python<'_>, statement: &str) -> PyResult<()> { - self.execute(py, statement, None).map(|_| ()) + self.query(py, statement, Vec::new()).map(|_| ()) + } + + /// Work done on the engine connection itself, with the GIL down. + /// + /// For the things a statement cannot say: reading the catalog, + /// opening the engine's own appender. The two failures are kept + /// apart because they belong to different people. A connection that + /// is closed or locked is this call's own answer and comes back as + /// the outer error; anything the work itself decides is wrong comes + /// back as the inner one, for the caller to turn into an exception + /// now that there is an interpreter to build one with. + pub(crate) fn engine( + &self, + py: Python<'_>, + work: F, + ) -> PyResult> + where + T: Send, + S: Send, + F: FnOnce(&mut zudb::Connection) -> std::result::Result + Send, + { + if !self.alive.load(Ordering::Acquire) { + return Err(closed(py, "this connection")); + } + py.detach(|| { + let mut held = self.inner.lock().map_err(|_| ())?; + let conn = held.as_mut().ok_or(())?; + Ok(work(conn)) + }) + .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; + } } } diff --git a/src/frame.rs b/src/frame.rs new file mode 100644 index 0000000..a872042 --- /dev/null +++ b/src/frame.rs @@ -0,0 +1,391 @@ +//! A frame of columns, read out of whatever the caller is holding. +//! +//! 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. +//! +//! 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. +//! +//! 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 +//! ever be refused, and refusing it by name and row number is the +//! difference between a caller who knows which cell to fix and one who +//! knows only that something somewhere was empty. + +use std::ffi::CStr; + +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::ffi_stream::{ArrowArrayStreamReader, FFI_ArrowArrayStream}; +use arrow::record_batch::{RecordBatch, RecordBatchReader}; +use pyo3::prelude::*; +use pyo3::types::{PyCapsule, PyDict}; +use zu_common::DurationKind; + +use crate::buffer::{Column, type_name}; +use crate::columns::Snag; +use crate::load; + +/// What a capsule holding an Arrow stream is called, which a consumer +/// 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, +} + +/// Reads whatever the caller handed over into columns. +/// +/// 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 +/// dictionary of lists for a caller with none of them installed. +/// 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 { + 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 }); + } + Err(pyo3::exceptions::PyTypeError::new_err(format!( + "a frame is anything with `__arrow_c_stream__`, which a pandas, polars or pyarrow table \ + has, or a dictionary of column name to values, and this is a '{}'", + type_name(data) + ))) +} + +/// Reads the Arrow stream a frame hands out through the capsule +/// protocol. +/// +/// 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 { + // 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. + let capsule = data.call_method1("__arrow_c_stream__", (py.None(),))?; + let capsule = capsule.cast_into::().map_err(|_| { + pyo3::exceptions::PyTypeError::new_err( + "`__arrow_c_stream__` gave back something that is not a capsule, so this frame does \ + not speak the protocol it says it speaks", + ) + })?; + // Checked by name, which is what tells a stream from every other + // capsule a library might hand out, and the pointer comes back + // from the same call so there is no way to read one without the + // other. + let held = capsule.pointer_checked(Some(STREAM)).map_err(|_| { + pyo3::exceptions::PyTypeError::new_err( + "`__arrow_c_stream__` gave back a capsule of some other kind, which is not a stream to \ + read", + ) + })?; + // Moved out and replaced with an empty one, which is the protocol's + // own rule: the consumer owns the stream from here and the + // capsule's destructor has to find nothing left to release. + let stream = unsafe { + let held = held.as_ptr() as *mut FFI_ArrowArrayStream; + std::ptr::replace(held, FFI_ArrowArrayStream::empty()) + }; + 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; + 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(); + } + Ok(Frame { columns, rows }) +} + +/// The buffer a column of this Arrow type fills, or the reason there +/// is none. +/// +/// 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 { + let refused = |instead: &str| { + Err(pyo3::exceptions::PyTypeError::new_err(format!( + "column '{name}' is {ty}, and {instead}" + ))) + }; + 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()), + // 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()), + DataType::Timestamp(_, Some(zone)) => { + return 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::Dictionary(_, _) => { + return 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 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()), + } + Ok(()) +} diff --git a/src/lib.rs b/src/lib.rs index 23066ae..f3b0864 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -17,8 +17,10 @@ mod buffer; mod columns; mod conn; mod error; +mod frame; mod interrupt; mod load; +mod register; mod txn; mod value; diff --git a/src/load.rs b/src/load.rs index 8badc89..304dd4b 100644 --- a/src/load.rs +++ b/src/load.rs @@ -127,7 +127,7 @@ pub fn load( /// Every column, in the order the dictionary holds them, which is the /// order they were written. -fn build(columns: Option<&Bound<'_, PyDict>>) -> PyResult> { +pub(crate) fn build(columns: Option<&Bound<'_, PyDict>>) -> PyResult> { let Some(columns) = columns else { return Ok(Vec::new()); }; diff --git a/src/register.rs b/src/register.rs new file mode 100644 index 0000000..02ae4c0 --- /dev/null +++ b/src/register.rs @@ -0,0 +1,399 @@ +//! Frames, registered under a name a statement can match on. +//! +//! ```python +//! conn.register("people", frame) +//! 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. +//! +//! 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. +//! +//! 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. + +use pyo3::prelude::*; +use zu_common::{DurationKind, Temporal}; +use zudb::query::Value; +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" + ), + ), + } + } +} + +/// Copies a frame in as a node table called `name`. +pub fn register( + py: Python<'_>, + conn: Py, + name: &str, + data: &Bound<'_, PyAny>, +) -> PyResult { + identifier(py, name, "a registered frame")?; + let frame = frame::read(py, data)?; + if frame.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 &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)? { + 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" + ), + )); + } + // 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) +} + +/// Empties the table a registered frame was copied into and forgets the +/// name. +pub fn unregister(py: Python<'_>, conn: &Connection, name: &str) -> PyResult<()> { + if !conn.registered_here(name) { + return Err(programming( + py, + &format!( + "nothing is registered here as '{name}', and a name this connection did not \ + register is a table of the database rather than a frame of the caller's" + ), + )); + } + 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(()) +} + +/// Whether the database already holds a table of this name. +/// +/// Node tables and rel tables both, because a name is a label in a +/// statement either way and registering over either of them would be +/// 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()) + })? + .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) +} + +/// Refuses a name a statement could not carry. +/// +/// A table name goes into a statement as itself, so a name that is not +/// an identifier is a name that would either fail to parse or parse as +/// something else, and the second of those is the one worth refusing +/// for. +fn identifier(py: Python<'_>, name: &str, what: &str) -> PyResult<()> { + let mut chars = name.chars(); + let starts = chars + .next() + .is_some_and(|first| first.is_ascii_alphabetic() || first == '_'); + if !starts || !chars.all(|c| c.is_ascii_alphanumeric() || c == '_') { + return Err(programming( + py, + &format!( + "'{name}' is not a name a statement can carry, and {what} is named by a letter or \ + an underscore followed by letters, digits or underscores" + ), + )); + } + Ok(()) +} diff --git a/tests/test_register.py b/tests/test_register.py new file mode 100644 index 0000000..a13a350 --- /dev/null +++ b/tests/test_register.py @@ -0,0 +1,266 @@ +"""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. + +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. +""" + +from __future__ import annotations + +import datetime +import time + +import pytest +import zudb + +pa = pytest.importorskip("pyarrow") + + +def names(conn: zudb.Connection, table: str) -> list[str]: + """The `name` column of a registered frame, in the order it went in.""" + return [name for (name,) in conn.execute(f"MATCH (f:{table}) RETURN f.name AS name")] + + +def test_a_dictionary_of_lists_is_a_frame(empty: zudb.Connection) -> None: + """The way in for a caller with no frame library installed.""" + assert empty.register("people", {"uid": [1, 2, 3], "name": ["ada", "grace", "lynn"]}) == 3 + assert names(empty, "people") == ["ada", "grace", "lynn"] + + +def test_a_statement_reads_a_registered_frame_like_any_other_table( + empty: zudb.Connection, +) -> None: + empty.register( + "people", + {"uid": [1, 2, 3], "name": ["ada", "grace", "lynn"], "age": [36, 45, 52]}, + ) + rows = empty.execute("MATCH (p:people) WHERE p.age > 40 RETURN p.name AS name").fetchall() + assert rows == [("grace",), ("lynn",)] + + +def test_a_pandas_frame_goes_in(empty: zudb.Connection) -> None: + pd = pytest.importorskip("pandas") + frame = pd.DataFrame({"uid": [1, 2], "name": ["ada", "grace"]}) + assert empty.register("people", frame) == 2 + assert names(empty, "people") == ["ada", "grace"] + + +def test_a_polars_frame_goes_in(empty: zudb.Connection) -> None: + """polars holds strings as views, which is a third Arrow layout and + the same characters.""" + pl = pytest.importorskip("polars") + frame = pl.DataFrame({"uid": [1, 2], "name": ["ada", "grace"]}) + assert empty.register("people", frame) == 2 + assert names(empty, "people") == ["ada", "grace"] + + +def test_a_pyarrow_table_goes_in(empty: zudb.Connection) -> None: + assert empty.register("people", pa.table({"uid": [1, 2], "name": ["ada", "grace"]})) == 2 + assert names(empty, "people") == ["ada", "grace"] + + +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.""" + schema = pa.schema([("uid", pa.int64()), ("name", pa.string())]) + batches = [ + pa.record_batch([pa.array([1, 2]), pa.array(["ada", "grace"])], schema=schema), + pa.record_batch([pa.array([3]), pa.array(["lynn"])], schema=schema), + ] + reader = pa.RecordBatchReader.from_batches(schema, batches) + assert empty.register("people", reader) == 3 + assert names(empty, "people") == ["ada", "grace", "lynn"] + + +def test_every_kind_of_column_a_row_can_hold_arrives_as_itself(empty: zudb.Connection) -> None: + frame = pa.table( + { + "yes": pa.array([True, False]), + "small": pa.array([1, 2], pa.int8()), + "wide": pa.array([3, 4], pa.uint32()), + "narrow": pa.array([1.5, 2.5], pa.float32()), + "word": pa.array(["a", "b"]), + "day": pa.array([datetime.date(2024, 1, 1), datetime.date(2024, 2, 1)]), + "clock": pa.array([datetime.time(1, 2, 3), datetime.time(4, 5, 6)], pa.time64("us")), + "moment": pa.array( + [datetime.datetime(2024, 1, 1, 1, 2, 3), datetime.datetime(2024, 2, 1)] + ), + "span": pa.array([datetime.timedelta(seconds=90)] * 2, pa.duration("us")), + } + ) + assert empty.register("kinds", frame) == 2 + row = empty.execute( + "MATCH (k:kinds) RETURN k.yes AS yes, k.small AS small, k.wide AS wide, " + "k.narrow AS narrow, k.word AS word, k.day AS day, k.clock AS clock, " + "k.moment AS moment, k.span AS span" + ).fetchone() + assert row == ( + True, + 1, + 3, + 1.5, + "a", + datetime.date(2024, 1, 1), + datetime.time(1, 2, 3), + datetime.datetime(2024, 1, 1, 1, 2, 3), + zudb.Duration(nanoseconds=90_000_000_000), + ) + + +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" + assert names(empty, "people") == ["ada"] + + +def test_registered_says_what_is_registered_here(empty: zudb.Connection) -> None: + assert empty.registered == [] + empty.register("second", {"a": [1]}) + empty.register("first", {"a": [1]}) + assert empty.registered == ["first", "second"] + + +def test_registering_a_name_again_replaces_the_rows(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: + empty.register("people", {"uid": [1, 2], "name": ["ada", "grace"]}) + empty.unregister("people") + assert empty.registered == [] + assert names(empty, "people") == [] + + +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 + assert names(empty, "people") == ["grace"] + + +def test_unregistering_twice_is_refused(empty: zudb.Connection) -> None: + empty.register("people", {"a": [1]}) + empty.unregister("people") + with pytest.raises(zudb.ProgrammingError, match="nothing is registered here"): + empty.unregister("people") + + +def test_unregistering_a_table_nobody_registered_is_refused(social: zudb.Connection) -> None: + with pytest.raises(zudb.ProgrammingError, match="nothing is registered here"): + social.unregister("person") + + +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.""" + 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.""" + 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"] + + +def test_a_null_anywhere_is_refused_by_column_and_row(empty: zudb.Connection) -> None: + """A property that is null is one no row of this engine can hold.""" + frame = pa.table({"uid": pa.array([1, 2, 3]), "name": pa.array(["ada", None, "lynn"])}) + with pytest.raises(ValueError, match="column 'name' has no value at row 1"): + 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_columns_is_refused(empty: zudb.Connection) -> None: + with pytest.raises(zudb.ProgrammingError, match="no columns"): + empty.register("people", {}) + + +def test_a_name_a_statement_could_not_carry_is_refused(empty: zudb.Connection) -> None: + with pytest.raises(zudb.ProgrammingError, match="not a name a statement can carry"): + empty.register("two words", {"a": [1]}) + with pytest.raises(zudb.ProgrammingError, match="a column of a registered frame"): + empty.register("people", {"two words": [1]}) + + +def test_a_zoned_timestamp_is_refused_with_what_to_do_about_it(empty: zudb.Connection) -> None: + frame = pa.table({"when": pa.array([1_700_000_000_000_000], pa.timestamp("us", "UTC"))}) + with pytest.raises(TypeError, match="nowhere to keep"): + empty.register("moments", frame) + + +def test_a_column_of_bytes_is_refused(empty: zudb.Connection) -> None: + """Writing one would be writing 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: + frame = pa.table({"big": pa.array([1, 2**63], pa.uint64())}) + with pytest.raises(ValueError, match="at row 1"): + empty.register("numbers", frame) + + +def test_something_that_is_not_a_frame_at_all_is_refused_with_the_list( + empty: zudb.Connection, +) -> None: + with pytest.raises(TypeError, match="__arrow_c_stream__"): + empty.register("people", [1, 2, 3]) + + +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.""" + with empty.transaction(): + with pytest.raises(zudb.ProgrammingError, match="own commits"): + empty.register("people", {"a": [1, 2]}) + + +def test_a_closed_connection_registers_nothing(empty: zudb.Connection) -> None: + empty.close() + with pytest.raises(zudb.ProgrammingError, match="closed"): + 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. + + 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. + """ + 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): + empty.register("numbers", frame) + best = min(best, time.perf_counter() - started) + assert best < 50e-3, f"reading {rows} rows took {best * 1e3:.0f} ms"