From 4b080c089212c3db0d7f1ebb0a23f6cb4863f67d Mon Sep 17 00:00:00 2001 From: Tam Nguyen Duc <1218621+tamnd@users.noreply.github.com> Date: Tue, 18 Aug 2026 11:00:49 +0700 Subject: [PATCH] Ctrl-C stops a statement, and so does another thread A press during a long statement did nothing until the statement ended, because Python delivers a signal by setting a flag and raising it at the next bytecode on the main thread, and a thread inside the engine with the GIL released runs no bytecode. Ten seconds of that reads as a hang, and a notebook that looks hung gets its kernel killed. So a statement called from the main thread now runs on a thread the connection keeps, and the calling thread waits for it two milliseconds at a time, asking Python for signals in between. A press raises KeyboardInterrupt from execute at five milliseconds measured against the fifty the milestone asks for. Everywhere else the statement runs inline as it always did, because Python delivers a signal to the main thread and to no other, and a watch there would be a cost buying nothing. The thread is kept rather than made per statement, and made at all only at the first statement that needs it. Starting one costs thirty microseconds against a small statement that costs ten, and handing work to a parked thread costs six, so the thread looks for work for two hundred microseconds before it parks and the waiting side looks for the answer for as long before it sleeps. A statement that finishes on its own is not measurably slower for being watched, which is a test rather than a claim. interrupt() is the same stop from another thread, and raises Interrupted where the statement is. It takes no lock on the connection, because the engine handle it sets lives beside the lock rather than under it and an ask that queued behind the statement could only arrive after it. closed and rows_read are answerable while a statement runs for the same reason, which is what a progress bar is drawn from. --- README.md | 16 +- python/zudb/_zudb.pyi | 7 + src/conn.rs | 145 ++++++++++----- src/interrupt.rs | 391 ++++++++++++++++++++++++++++++++++++++++ src/lib.rs | 1 + tests/conftest.py | 23 +++ tests/test_interrupt.py | 204 +++++++++++++++++++++ 7 files changed, 737 insertions(+), 50 deletions(-) create mode 100644 src/interrupt.rs create mode 100644 tests/test_interrupt.py diff --git a/README.md b/README.md index 4778ad6..654edf5 100644 --- a/README.md +++ b/README.md @@ -85,6 +85,20 @@ result.record_batches() # a reader, for a result larger than memory `Result` implements `__arrow_c_stream__`, so anything that reads the protocol reads a result directly and none of the four methods above is needed: `pyarrow.table(result)` and `polars.DataFrame(result)` both work. Batches are 65,536 rows. A column holds one type, which the values decide, and integers beside floats are the one mixture that widens rather than being refused. Nodes, rels and paths go across as structs. The copy runs with the GIL released, and on this machine 300,000 rows across three columns take 44 ms as Arrow against 67 ms as Python objects, and a single integer column takes 13.8 ms against 44.5 ms. +## Stopping a statement + +A statement that is running can be stopped two ways, and neither of them closes the connection: the session, its plans and its warm readers are all there afterwards, which is the whole difference between stopping a statement and starting again. + +```python +conn.execute(long_one) # Ctrl-C raises KeyboardInterrupt here +conn.interrupt() # from another thread, raises zudb.Interrupted there +conn.rows_read # how far the statement running now has got +``` + +`Ctrl-C` is the one a person presses, and it raises `KeyboardInterrupt` on the thread that called `execute`, measured at 5 ms from the press on this machine against a budget of 50. Python only delivers a signal to the main thread between two bytecodes, so a statement called from the main thread runs on a thread this client keeps for it and the main thread waits and asks for signals while it does. That thread is kept rather than made per statement, because making one costs 30 microseconds against a small statement that costs 10, and a statement called from any other thread runs inline where a signal was never going to arrive anyway. + +`interrupt()` is the one a program calls, from a thread that is not the one inside `execute`, and it raises `zudb.Interrupted` there. It is one of the three calls that may be made on a connection while a statement is running, with `rows_read` and `closed`, and none of the three waits for it: a progress bar drawn from `rows_read` is a poll of an atomic, not a queue behind the executor. + ## Types The wheel carries `py.typed` and a stub for the compiled module, so mypy, pyright and an editor's completion all work with nothing else installed. `zudb.Value` is the union a row holds and a parameter takes, for code that passes rows around and wants to say so. @@ -93,7 +107,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, 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, and the GIL released around every statement, every load and every copy out. `register` and the interrupt are 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, 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. ## Wheels diff --git a/python/zudb/_zudb.pyi b/python/zudb/_zudb.pyi index 166f08e..f88953c 100644 --- a/python/zudb/_zudb.pyi +++ b/python/zudb/_zudb.pyi @@ -63,6 +63,13 @@ class Connection: def closed(self) -> bool: """Whether this connection is still open.""" + @property + def rows_read(self) -> int: + """How many rows the statement running on this connection has read out of storage.""" + + def interrupt(self) -> None: + """Asks the statement running on this connection to stop.""" + def execute(self, statement: str, params: Mapping[str, Value] | None = None) -> Result: """Runs one statement and gives back its rows.""" diff --git a/src/conn.rs b/src/conn.rs index 92668bc..4e77a02 100644 --- a/src/conn.rs +++ b/src/conn.rs @@ -9,16 +9,18 @@ use std::ffi::CStr; use std::path::PathBuf; -use std::sync::Mutex; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::{Arc, Mutex, OnceLock}; use pyo3::prelude::*; use pyo3::types::{PyCapsule, PyDict, PyList, PyTuple}; use zudb::query::{QueryResult, Value}; -use zudb::{Config, Database}; +use zudb::{Config, Database, Interrupt}; use crate::appender::Appender; use crate::columns; use crate::error::{closed, to_py_err}; +use crate::interrupt; use crate::value::{Names, from_py, to_py}; /// What a capsule holding an Arrow stream is called. The name is part @@ -44,7 +46,25 @@ pub struct Connection { /// thread in the process for the length of somebody else's /// statement, and would deadlock against the thread inside that /// statement, which needs the GIL back to return. - pub(crate) inner: Mutex>, + /// + /// Counted rather than plain, because the thread a statement runs + /// on takes a share of it: a job handed to that thread outlives the + /// call that made it in the type system even though it never does + /// in fact. + pub(crate) inner: Arc>>, + /// The thread statements run on where a `Ctrl-C` can arrive, + /// started at the first one that needs it and kept for the rest. + runner: OnceLock, + /// The word the running statement reads, held out here rather than + /// reached through the lock. A stop that had to wait for the + /// connection to be free could only ever arrive after the statement + /// it was meant to stop. + stop: Interrupt, + /// Whether the connection is still open, kept beside the lock for + /// the same reason: asking a connection whether it is closed, or + /// asking it to stop, should not queue behind a ten second + /// statement. + alive: AtomicBool, #[pyo3(get)] path: PathBuf, #[pyo3(get)] @@ -67,27 +87,30 @@ impl Connection { params: Option<&Bound<'_, PyDict>>, ) -> PyResult { let params = bind(params)?; - let borrowed: Vec<(&str, Value)> = params - .iter() - .map(|(name, value)| (name.as_str(), value.clone())) - .collect(); + // 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, and the signal - // handler gets to run too, which is what lets a `Ctrl-C` - // arrive at all. Waiting for the lock with the GIL held would - // be worse than slow: the thread inside the statement has to - // take the GIL back to return, and it could not. - let (result, names) = py - .detach(|| -> std::result::Result<_, Trouble> { - let mut held = self.inner.lock().map_err(|_| Trouble::Closed)?; - let conn = held.as_mut().ok_or(Trouble::Closed)?; + // 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()); - let result = conn.query_with(statement, &borrowed)?; - Ok((result, names)) + conn.query_with(&statement, &borrowed) + .map(|result| (result, names)) }) - .map_err(|trouble| trouble.raise(py))?; + .map_err(|stopped| stopped.raise(py))? + .map_err(|err| to_py_err(py, err))?; Ok(Result { result, names, @@ -119,6 +142,40 @@ impl Connection { Appender::open(py, slf, table) } + /// Asks the statement running on this connection to stop. + /// + /// The one call meant to be made from another thread while the + /// connection is busy, which is why it waits for nothing: the + /// statement it stops is the one holding everything a call that + /// waited would be waiting for. The statement raises + /// `zudb.Interrupted` and the connection is exactly as it was, so + /// the next statement on it starts warm. That is the difference + /// between stopping a statement and closing a connection. + /// + /// With nothing running this does nothing. It does not arm a stop + /// for the next statement, because a statement nobody has run yet + /// is not one anybody has waited too long for. + fn interrupt(&self, py: Python<'_>) -> PyResult<()> { + if !self.alive.load(Ordering::Acquire) { + return Err(closed(py, "this connection")); + } + self.stop.stop(); + Ok(()) + } + + /// How many rows the statement running on this connection has read + /// out of storage, for showing a person that something is + /// happening. + /// + /// Rows read rather than rows answered, because the statement + /// somebody is waiting on is exactly the one that reads a hundred + /// million rows to answer one. It starts at zero at each statement + /// and holds its last value once one ends. + #[getter] + fn rows_read(&self) -> u64 { + self.stop.rows() + } + /// Closes the connection and frees what it held. /// /// Doing it twice is not an error, because a `with` block that @@ -131,13 +188,21 @@ impl Connection { if let Ok(mut held) = self.inner.lock() { drop(held.take()); } + // Written after the drop rather than before it, so that a + // connection reports itself open until it really is not, + // and a poisoned lock reports itself closed because + // nothing can be run on one. + self.alive.store(false, Ordering::Release); }); } /// Whether this connection is still open. + /// + /// Answered from a word beside the lock rather than through it, so + /// that asking does not queue behind whatever is running. #[getter] - fn closed(&self, py: Python<'_>) -> bool { - py.detach(|| self.inner.lock().map(|held| held.is_none()).unwrap_or(true)) + fn closed(&self) -> bool { + !self.alive.load(Ordering::Acquire) } fn __enter__(slf: PyRef<'_, Self>) -> PyRef<'_, Self> { @@ -153,8 +218,8 @@ impl Connection { false } - fn __repr__(&self, py: Python<'_>) -> String { - let state = if self.closed(py) { ", closed" } else { "" }; + fn __repr__(&self) -> String { + let state = if self.closed() { ", closed" } else { "" }; format!("", self.path.display()) } } @@ -189,39 +254,21 @@ impl Connection { } .and_then(|db| db.connect()) }); + let opened = opened.map_err(|err| to_py_err(py, err))?; Ok(Connection { - inner: Mutex::new(Some(opened.map_err(|err| to_py_err(py, err))?)), + // Taken here, once, because every later reader of it wants + // it while the connection is busy and taking it then would + // mean waiting for the statement it is there to stop. + stop: opened.interrupt(), + inner: Arc::new(Mutex::new(Some(opened))), + runner: OnceLock::new(), + alive: AtomicBool::new(true), path, read_only, }) } } -/// What can go wrong inside a statement, with the GIL down and no way -/// to build a Python exception yet. -enum Trouble { - Closed, - Engine(zudb::ZuError), -} - -impl From for Trouble { - fn from(err: zudb::ZuError) -> Trouble { - Trouble::Engine(err) - } -} - -impl Trouble { - fn raise(self, py: Python<'_>) -> PyErr { - match self { - // A connection whose lock a panic left poisoned is a - // connection nothing can be run on again, which is the - // same fact as a closed one and reads better as one. - Trouble::Closed => closed(py, "this connection"), - Trouble::Engine(err) => to_py_err(py, err), - } - } -} - /// The rows a statement gave back. /// /// Held as the engine produced them and turned into Python objects diff --git a/src/interrupt.rs b/src/interrupt.rs new file mode 100644 index 0000000..48f8900 --- /dev/null +++ b/src/interrupt.rs @@ -0,0 +1,391 @@ +//! Stopping a statement while it runs: a `Ctrl-C` from the person +//! waiting for it, and `interrupt()` from another thread. +//! +//! Python answers a signal by setting a flag in the handler and raising +//! the exception the next time the main thread runs bytecode, and a +//! thread inside the engine with the GIL released is not running +//! bytecode. That is why a press during a ten second statement is felt +//! ten seconds later in every binding that does nothing about it, and +//! it is why a statement that could be interrupted runs on a thread of +//! its own while the thread that asked for it waits, asking Python +//! every few milliseconds whether a press arrived. +//! +//! The thread is kept rather than made, and it is only made at all +//! where it buys something. A thread costs some thirty microseconds to +//! start and six to hand a job to once it is parked, and a small +//! statement costs less than either, so the connection keeps one and +//! wakes it. Python delivers a signal to the main thread and to no +//! other, so a statement anywhere else runs on the thread that asked +//! for it, exactly as it did before, and is stopped by `interrupt()` or +//! not at all. +//! +//! Stopping is the engine's [`Interrupt`], which the connection holds a +//! clone of from the moment it is opened. Holding it outside the lock +//! is the whole trick: an ask that had to wait for the connection to be +//! free could only ever arrive after the statement it was meant to +//! stop. + +use std::cell::Cell; +use std::panic::AssertUnwindSafe; + +use std::sync::{Arc, Condvar, Mutex, OnceLock}; +use std::time::{Duration, Instant}; + +use pyo3::prelude::*; +use zudb::Interrupt; + +use crate::error::closed; + +/// How long the waiting thread sleeps between two asks. +/// +/// The budget is 50 ms from the press to the exception, and most of a +/// press is this sleep: the engine answers a stop at the boundary of a +/// chunk, which is a tenth of a millisecond of the budget, and the rest +/// of it is the wait between the press landing and the next ask. Two +/// puts the whole path at five milliseconds measured, forty five of it +/// spare for a machine under load, and a statement running for an hour +/// wakes this thread under two million times, which is a wakeup against +/// two milliseconds of work. +const TICK: Duration = Duration::from_millis(2); + +/// How long the waiting thread looks for the answer before it sleeps +/// for it. +/// +/// A sleep and the wake that ends it cost microseconds, which is more +/// than a small statement costs to run, so a thread that went straight +/// to sleep would add that to every quick statement in the process. +/// Two hundred microseconds is longer than a small statement and +/// nothing beside a statement anybody would reach for the keyboard +/// over, and it is spent on the thread with nothing else to do rather +/// than on the one running the query. +const SPIN: Duration = Duration::from_micros(200); + +/// Why a statement gave back no answer at all. +pub enum Stopped { + /// The connection was closed, or a panic left its lock poisoned, + /// which for anything that would run on it is the same fact. + Closed, + /// A signal arrived while the statement ran. This is the exception + /// Python had waiting, which is a `KeyboardInterrupt` unless the + /// program installed a handler that raises something else. + Signal(PyErr), +} + +impl Stopped { + pub fn raise(self, py: Python<'_>) -> PyErr { + match self { + Stopped::Closed => closed(py, "this connection"), + Stopped::Signal(err) => err, + } + } +} + +/// A statement on its way to the thread that runs it. +/// +/// Type-erased, so that one thread serves statements answering +/// different things: what each job gives back it gives back through the +/// slot it was built around, and by then nobody needs to know what it +/// was. +type Job = Box; + +/// The thread a connection keeps for the statements it wants to be able +/// to interrupt. +/// +/// One thread per connection rather than one per statement, and started +/// at the first statement rather than at the connection, so a program +/// that never runs one on the main thread never pays for it. +pub struct Runner { + post: Arc, + /// Joined when the connection goes, so that nothing of this + /// connection's is still running once nothing points at it. It can + /// only ever be a statement that is already over: a statement being + /// waited for is being waited for by somebody holding the + /// connection. + thread: Option>, +} + +/// Where the statements waiting for that thread sit. +#[derive(Default)] +struct Post { + queue: Mutex, + knock: Condvar, +} + +#[derive(Default)] +struct Queue { + /// A queue rather than a slot, for the one caller who has two at + /// once: a signal handler that runs a statement of its own while + /// the statement it interrupted is still being waited for. + jobs: std::collections::VecDeque, + /// Set when the connection is dropped, which is what ends the + /// thread. + shut: bool, +} + +impl Post { + /// The next statement, or `None` once the connection is gone. + /// + /// Looked for before it is slept for. A statement handed to a + /// thread that is still awake costs a lock, and one handed to a + /// thread that has parked costs a wake, which on this machine is + /// the difference between one microsecond and eighty. A program + /// running statements in a loop hands them over inside the look + /// every time, and one that has gone quiet pays the wake once. + fn next(&self) -> Option { + let looking = Instant::now(); + loop { + if let Ok(mut queue) = self.queue.try_lock() { + if let Some(job) = queue.jobs.pop_front() { + return Some(job); + } + if queue.shut { + return None; + } + } + if looking.elapsed() >= SPIN { + break; + } + std::hint::spin_loop(); + } + let mut queue = self.queue.lock().unwrap_or_else(|held| held.into_inner()); + loop { + if let Some(job) = queue.jobs.pop_front() { + return Some(job); + } + if queue.shut { + return None; + } + queue = self + .knock + .wait(queue) + .unwrap_or_else(|held| held.into_inner()); + } + } + + /// Hands a statement over, or says that the thread is gone. + fn send(&self, job: Job) -> bool { + let mut queue = self.queue.lock().unwrap_or_else(|held| held.into_inner()); + if queue.shut { + return false; + } + queue.jobs.push_back(job); + self.knock.notify_one(); + true + } +} + +impl Runner { + fn start() -> Runner { + let post = Arc::new(Post::default()); + let taking = Arc::clone(&post); + let thread = std::thread::Builder::new() + .name("zudb statement".into()) + .spawn(move || { + while let Some(job) = taking.next() { + job(); + } + }) + // A process that cannot start a thread cannot run a + // statement on one, and there is nothing to say about that + // which the operating system's own message does not. + .expect("a thread to run statements on"); + Runner { + post, + thread: Some(thread), + } + } +} + +impl Drop for Runner { + fn drop(&mut self) { + self.post + .queue + .lock() + .unwrap_or_else(|held| held.into_inner()) + .shut = true; + self.post.knock.notify_all(); + if let Some(thread) = self.thread.take() { + // A panic inside a statement was caught where it happened + // and raised again on the thread that asked for it, so + // there is nothing here for one to have left behind. + let _ = thread.join(); + } + } +} + +/// What a statement leaves for the thread waiting on it. +struct Answer { + /// `None` until the statement is over. A panic inside the engine + /// arrives here as the payload it unwound with, to be raised again + /// on the thread that asked for the statement, which is where a + /// Python caller can see it. + slot: Mutex>>, + over: Condvar, +} + +impl Answer { + fn new() -> Answer { + Answer { + slot: Mutex::new(None), + over: Condvar::new(), + } + } + + fn done(&self, out: std::thread::Result) { + *self.slot.lock().unwrap_or_else(|held| held.into_inner()) = Some(out); + self.over.notify_all(); + } + + /// The answer if it is there, without waiting for it. + fn peek(&self) -> Option> { + self.slot.try_lock().ok().and_then(|mut slot| slot.take()) + } + + /// The answer, waiting up to `patience` for it, or forever when + /// there is none. + fn wait(&self, patience: Option) -> Option> { + let slot = self.slot.lock().unwrap_or_else(|held| held.into_inner()); + let waiting = |slot: &mut Option>| slot.is_none(); + let mut slot = match patience { + Some(patience) => { + self.over + .wait_timeout_while(slot, patience, waiting) + .unwrap_or_else(|held| held.into_inner()) + .0 + } + None => self + .over + .wait_while(slot, waiting) + .unwrap_or_else(|held| held.into_inner()), + }; + slot.take() + } +} + +/// Runs `work` on the connection, with a `Ctrl-C` able to stop it. +/// +/// `stop` is the connection's own handle, taken when it was opened +/// rather than fetched from inside the lock, because the point of it is +/// to be reachable while the lock is held by the statement it stops. +/// +/// A press that arrives while the statement is still running raises, +/// and whatever the statement was about to answer is dropped: the +/// person pressed the key, and the answer to that is the exception +/// rather than the rows they interrupted. A press with nothing running +/// is Python's to deliver as it always was. +pub fn watched( + py: Python<'_>, + runner: &OnceLock, + inner: &Arc>>, + stop: &Interrupt, + work: impl FnOnce(&mut zudb::Connection) -> T + Send + 'static, +) -> Result { + let held = Arc::clone(inner); + let run = move || -> Result { + let mut held = held.lock().map_err(|_| Stopped::Closed)?; + let conn = held.as_mut().ok_or(Stopped::Closed)?; + // Cleared here, holding the lock, rather than before the wait. + // An ask that arrived while nothing was running must not end + // the statement about to start, and an ask meant for the + // statement another thread is running must not be cleared by + // this one queuing up behind it. + conn.interrupt().clear(); + Ok(work(conn)) + }; + if !on_main_thread(py) { + // The GIL goes down for the whole statement, waiting for the + // connection's own lock included, which is what lets another + // thread run while this one is inside the executor. Waiting + // for the lock with the GIL held would be worse than slow: the + // thread inside the statement has to take the GIL back to + // return, and it could not. + return py.detach(run); + } + let answer: Arc>> = Arc::new(Answer::new()); + let done = Arc::clone(&answer); + let job: Job = Box::new(move || { + // The panic is caught rather than left to end the thread this + // connection keeps: an engine that panicked on one statement + // is still a connection somebody is going to call `close` on, + // and a thread that died holding the answer would leave the + // caller waiting for it forever. + done.done(std::panic::catch_unwind(AssertUnwindSafe(run))); + }); + if !runner.get_or_init(Runner::start).post.send(job) { + // The thread is gone, which happens when the connection it + // belongs to is dropped, which is a connection nothing can run + // on anyway. + return Err(Stopped::Closed); + } + let out = py.detach(|| { + // Most statements are over in less time than a sleep and a wake + // cost, so the first wait is a look. + let looking = Instant::now(); + loop { + if let Some(out) = answer.peek() { + return Some(out); + } + if looking.elapsed() >= SPIN { + return None; + } + std::hint::spin_loop(); + } + }); + let out = match out { + Some(out) => out, + None => loop { + if let Some(out) = py.detach(|| answer.wait(Some(TICK))) { + break out; + } + // Where the press is felt. Python runs the handler here, on + // the thread it is meant to run on, and hands back what it + // raised. + if let Err(signal) = py.check_signals() { + stop.stop(); + // Waited for rather than abandoned, so that the next + // statement finds the connection free rather than + // locked by one nobody is listening to. + py.detach(|| answer.wait(None)); + return Err(Stopped::Signal(signal)); + } + }, + }; + // Raised again here, on the thread that asked for the statement, + // where PyO3 turns it into an exception a caller can see rather + // than into a process that stops. + out.unwrap_or_else(|panicked| std::panic::resume_unwind(panicked)) +} + +thread_local! { + /// Whether this thread is the one a signal can reach, asked once + /// per thread and kept, because the answer cannot change while the + /// thread lives. + static IS_MAIN: Cell> = const { Cell::new(None) }; +} + +/// Whether a `Ctrl-C` can arrive on this thread at all. +/// +/// `threading` rather than the interpreter's own thread state, because +/// the limited ABI these wheels are built against does not expose the +/// second one, and asking Python is a couple of attribute lookups on a +/// module that is imported long before anything here runs. +fn on_main_thread(py: Python<'_>) -> bool { + IS_MAIN.with(|known| { + if let Some(known) = known.get() { + return known; + } + let asked = idents(py).is_ok_and(|(this, main)| this == main); + known.set(Some(asked)); + asked + }) +} + +fn idents(py: Python<'_>) -> PyResult<(u64, u64)> { + let threading = py.import("threading")?; + let this: u64 = threading.call_method0("get_ident")?.extract()?; + let main: u64 = threading + .call_method0("main_thread")? + .getattr("ident")? + .extract()?; + Ok((this, main)) +} diff --git a/src/lib.rs b/src/lib.rs index 862dd03..d0f54aa 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -17,6 +17,7 @@ mod buffer; mod columns; mod conn; mod error; +mod interrupt; mod load; mod value; diff --git a/tests/conftest.py b/tests/conftest.py index ebcc1e5..5d57a77 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -76,6 +76,29 @@ def loaded(tmp_path: Path) -> zudb.Connection: conn.close() +@pytest.fixture +def throng(tmp_path: Path) -> zudb.Connection: + """Enough people that the statement over the pairs of them runs for + about ten seconds, which is the statement the milestone asks to be + interruptible. + + Loaded rather than inserted because twelve thousand rows one + statement at a time is twelve thousand commits, and this is a + fixture rather than the thing under test. + """ + path = tmp_path / "throng.zu1" + zudb.load( + path, + nodes="person", + rels="knows", + columns={"uid": list(range(12_000))}, + edges=[(0, 1)], + ) + conn = zudb.connect(path, read_only=True) + yield conn + conn.close() + + @pytest.fixture def crowd(tmp_path: Path) -> zudb.Connection: """Enough people that a statement over the pairs of them takes long diff --git a/tests/test_interrupt.py b/tests/test_interrupt.py new file mode 100644 index 0000000..5c4cce8 --- /dev/null +++ b/tests/test_interrupt.py @@ -0,0 +1,204 @@ +"""Stopping a statement that is already running. + +Two ways in, one word underneath. A person presses `Ctrl-C` and the +statement raises `KeyboardInterrupt` on the thread that asked for it; a +program calls `interrupt()` from another thread and the statement raises +`zudb.Interrupted`. Either way the connection is exactly as it was +afterwards, which is the difference between stopping a statement and +closing a connection. + +The budget is 50 ms from the press to the exception, and it is asserted +rather than described: a long statement in a notebook that ignores the +keyboard is a kernel somebody kills. +""" + +from __future__ import annotations + +import os +import signal +import threading +import time + +import pytest +import zudb + +# Every pair of twelve thousand people, filtered, which is about ten +# seconds of work on the machine this was written on. Nothing here waits +# for it to finish: the point of each test is that it does not. +WORK = "MATCH (a:person), (b:person) WHERE a.uid < b.uid RETURN count(a) AS n" + +# What the milestone asks for, from the press to the exception. +BUDGET = 0.05 + +# Long enough that the statement is certainly inside the executor rather +# than still being parsed when the ask arrives. +SETTLED = 0.2 + + +def after(delay: float, do) -> threading.Thread: + """Runs `do` on a thread of its own, `delay` from now, and answers + the thread so a test can join it.""" + + def wait() -> None: + time.sleep(delay) + do() + + thread = threading.Thread(target=wait) + thread.start() + return thread + + +def test_a_press_raises_keyboard_interrupt_within_the_budget(throng: zudb.Connection) -> None: + pressed: list[float] = [] + + def press() -> None: + pressed.append(time.perf_counter()) + # The signal a terminal sends, sent the way a terminal sends it: + # to the process, for the main thread to answer. + os.kill(os.getpid(), signal.SIGINT) + + presser = after(SETTLED, press) + with pytest.raises(KeyboardInterrupt): + throng.execute(WORK) + felt = time.perf_counter() + presser.join(timeout=30) + assert felt - pressed[0] < BUDGET, f"the press took {felt - pressed[0]:.3f}s to arrive" + + +def test_the_connection_is_the_same_afterwards(throng: zudb.Connection) -> None: + presser = after(SETTLED, lambda: os.kill(os.getpid(), signal.SIGINT)) + with pytest.raises(KeyboardInterrupt): + throng.execute(WORK) + presser.join(timeout=30) + # Not reopened, not reconnected: the same connection, which still + # knows what it knew. + assert throng.closed is False + assert throng.execute("MATCH (p:person) RETURN count(p) AS n").fetchone() == (12_000,) + + +def test_interrupt_stops_a_statement_from_another_thread(throng: zudb.Connection) -> None: + asked: list[float] = [] + + def ask() -> None: + asked.append(time.perf_counter()) + throng.interrupt() + + asker = after(SETTLED, ask) + with pytest.raises(zudb.Interrupted): + throng.execute(WORK) + felt = time.perf_counter() + asker.join(timeout=30) + assert felt - asked[0] < BUDGET, f"the ask took {felt - asked[0]:.3f}s to arrive" + + +def test_a_statement_on_another_thread_is_stopped_too(throng: zudb.Connection) -> None: + """A worker thread's statement runs on the worker thread, since a + signal cannot reach one, and `interrupt()` still ends it.""" + raised: list[BaseException] = [] + running = threading.Event() + + def run() -> None: + running.set() + try: + throng.execute(WORK) + except BaseException as why: # noqa: BLE001 + raised.append(why) + + worker = threading.Thread(target=run) + worker.start() + running.wait(timeout=30) + time.sleep(SETTLED) + throng.interrupt() + worker.join(timeout=30) + assert not worker.is_alive() + assert isinstance(raised[0], zudb.Interrupted) + + +def test_an_ask_with_nothing_running_does_not_end_the_next_statement( + social: zudb.Connection, +) -> None: + social.interrupt() + assert social.execute("MATCH (p:person) RETURN count(p) AS n").fetchone() == (3,) + + +def test_a_press_with_nothing_running_is_pythons_to_deliver(social: zudb.Connection) -> None: + with pytest.raises(KeyboardInterrupt): + os.kill(os.getpid(), signal.SIGINT) + # Python raises the press at the next thing this thread does, + # which is this call, and the statement it stopped is none. + for _ in range(1000): + pass + assert social.execute("MATCH (p:person) RETURN count(p) AS n").fetchone() == (3,) + + +def test_interrupt_on_a_closed_connection_says_so(social: zudb.Connection) -> None: + social.close() + with pytest.raises(zudb.ProgrammingError, match="closed"): + social.interrupt() + + +def test_a_connection_answers_while_it_is_busy(throng: zudb.Connection) -> None: + """Asking a connection how it is going does not queue behind the + statement it is going through.""" + running = threading.Event() + + def run() -> None: + running.set() + with pytest.raises(zudb.Interrupted): + throng.execute(WORK) + + worker = threading.Thread(target=run) + worker.start() + running.wait(timeout=30) + time.sleep(SETTLED) + asked = time.perf_counter() + assert throng.closed is False + assert throng.rows_read > 0, "a statement that has been running has read something" + assert time.perf_counter() - asked < BUDGET, "asking waited for the statement" + throng.interrupt() + worker.join(timeout=30) + assert not worker.is_alive() + + +def test_rows_read_holds_what_the_last_statement_cost(social: zudb.Connection) -> None: + social.execute("MATCH (p:person) RETURN p.name AS n") + assert social.rows_read >= 3 + + +def test_a_statement_that_finishes_is_not_slowed_by_being_watched(social: zudb.Connection) -> None: + """The thread a statement runs on is kept rather than made, so a + small statement on the main thread costs about what it costs on any + other.""" + query = "MATCH (p:person) RETURN p.name AS n" + social.execute(query) + + def timed() -> float: + best = float("inf") + for _ in range(5): + start = time.perf_counter() + for _ in range(200): + social.execute(query) + best = min(best, (time.perf_counter() - start) / 200) + return best + + here = timed() + elsewhere: list[float] = [] + worker = threading.Thread(target=lambda: elsewhere.append(timed())) + worker.start() + worker.join(timeout=120) + # Generous, because this is a gate against paying for a thread per + # statement rather than a benchmark: starting one costs some thirty + # microseconds against a statement that costs ten. + assert here < elsewhere[0] * 3, f"{here * 1e6:.1f}us here against {elsewhere[0] * 1e6:.1f}us" + + +def test_a_press_during_a_statement_that_finishes_first_is_still_raised( + social: zudb.Connection, +) -> None: + """A press is never swallowed. The statement was over before it + arrived, so Python raises it at the next thing this thread does.""" + with pytest.raises(KeyboardInterrupt): + os.kill(os.getpid(), signal.SIGINT) + social.execute("MATCH (p:person) RETURN p.name AS n") + for _ in range(1000): + pass