SQLite engine open: SqliteOpts, connection architecture, spawn_blocking seam posture (task sqlite-engine-open-opts)
- SqliteOpts (poll_interval: Option<Duration>, None = 1 ms shipping default ADR-023 §4; max_readers, DEFAULT_MAX_READERS = 8; opts exemption from non_exhaustive per ADR-017 §3) - open(path, opts) -> Box<dyn Store>: writer slot, reader pool, SharedUpdateWatcher with fallible spawn mapped W-2-style into Error::Database; plain-path posture (ADR-023 §3) pinned by test - Substrate delta D-30: open_conn_bootstrapped — every connection (writer and pooled readers) carries the full bootstrap surface (pragmas, notify, alkstore functions, schema) per engine-sqlite.md; lineage opened readers pragmas-only. Registered in PROVENANCE.md - spawn_blocking seam helper + error mappings (string/rusqlite -> Database, source chains preserved); helper cfg(test)-gated until sqlite-engine-seam-tx wires the trait impls - Store trait stubbed with Database errors (no panics); close()/Drop join the watcher, clear subscribers (death-guard), close pool+writer - 15 new tests (open boot, watcher fanout/death-close, cadence accounting, plain-path literal filename, :memory:, pool bounding, drop teardown, seam smoke incl. panic mapping); gates green
This commit is contained in:
1 parent
5474141f2b
commit
7d400906f5
11 files changed
+937
-13
No files matched your search
@@ -11,6 +11,7 @@ parking_lot = "0.12"
|
||||
rusqlite = { version = "0.40", features = ["bundled", "functions"] }
|
||||
serde_json = "1"
|
||||
thiserror = "2"
|
||||
tokio = { version = "1", features = ["rt", "sync"] }
|
||||
|
||||
# file-id only supports unix and windows. Other targets (WASI, Redox,
|
||||
# illumos, etc.) get the `(0, 0)` fallback in substrate::watcher's
|
||||
@@ -19,4 +20,5 @@ thiserror = "2"
|
||||
file-id = "0.2"
|
||||
|
||||
[dev-dependencies]
|
||||
alkstore-contract-suite = { path = "../alkstore-contract-suite" }
|
||||
alkstore-contract-suite = { path = "../alkstore-contract-suite" }
|
||||
tokio = { version = "1", features = ["rt-multi-thread", "macros", "time"] }
|
||||
@@ -1 +1,26 @@
|
||||
//! alkstore-sqlite — the SQLite engine: the core contract implemented
|
||||
//! on rusqlite over the forked honker-core substrate carried in-tree
|
||||
//! (ADR-001, ADR-011, ADR-013).
|
||||
//!
|
||||
//! Single-host by nature: file-backed, one machine; NFS
|
||||
//! two-writers-unsupported (the lineage's honesty posture, inherited —
|
||||
//! ADR-016). Durability knobs, pool sizing, and the watcher cadence
|
||||
//! are engine-configuration concerns on [`SqliteOpts`] (ADR-008 §6) —
|
||||
//! never contract surface.
|
||||
//!
|
||||
//! Long transactions park the writer: the writer slot is a lease, and
|
||||
//! holding a caller-held [`alkstore::TxHandle`] across `await` points
|
||||
//! holds it (ADR-007's honest model, surfaced so consumers budget
|
||||
//! transactions).
|
||||
//!
|
||||
//! Every engine call round-trips a blocking thread (the
|
||||
//! `spawn_blocking` seam — rusqlite connections are not
|
||||
//! `Send`-across-await; the family REQ-TTY-01 posture).
|
||||
|
||||
mod opts;
|
||||
mod seam;
|
||||
mod store;
|
||||
mod substrate;
|
||||
|
||||
pub use opts::{DEFAULT_MAX_READERS, SqliteOpts};
|
||||
pub use store::{SqliteStore, open};
|
||||
@@ -0,0 +1,68 @@
|
||||
//! The SQLite engine's option struct and the watcher-config plumbing
|
||||
//! (ADR-008 §6 — constructors and their option structs live in the
|
||||
//! engine crates; engine-configuration concerns never touch the
|
||||
//! contract surface).
|
||||
//!
|
||||
//! `#[non_exhaustive]` deliberately does not apply (ADR-017 §3): opts
|
||||
//! structs are consumer-*constructed*, and the attribute on them would
|
||||
//! push every construction through a builder. New fields may be added
|
||||
//! with `Default` fallbacks, per the core opts modules' posture.
|
||||
|
||||
use std::time::Duration;
|
||||
|
||||
use crate::substrate::WatcherConfig;
|
||||
|
||||
/// The default reader-pool size (the shared watcher's poll cadence
|
||||
/// rides [`SqliteOpts::poll_interval`]; this is the pool's own knob).
|
||||
pub const DEFAULT_MAX_READERS: usize = 8;
|
||||
|
||||
/// Construction options for the SQLite engine's
|
||||
/// [`open`](crate::open) constructor.
|
||||
///
|
||||
/// Durability knobs, pool sizing, and the watcher cadence are
|
||||
/// engine-configuration concerns — they never appear on the contract
|
||||
/// surface (ADR-008 §6); deployment.md carries the documented tuning
|
||||
/// facts.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct SqliteOpts {
|
||||
/// The watcher's `PRAGMA data_version` poll cadence. `None` = the
|
||||
/// substrate's 1 ms shipping default (ADR-023 §4 — the default the
|
||||
/// published wake-latency numbers were measured at); `Some(d)` sets
|
||||
/// an explicit cadence (raise it for latency-tolerant deployments —
|
||||
/// the documented idle cost at 1 ms is ~1000 poll reads/sec per
|
||||
/// instance). Zero is rejected at `open` with
|
||||
/// [`Error::Database`](alkstore::Error::Database) (the substrate's
|
||||
/// `WatcherConfig::with_poll_interval` rejects zero; the mapping is
|
||||
/// engine-layer work per ADR-012 §4's W-2 posture).
|
||||
pub poll_interval: Option<Duration>,
|
||||
/// Upper bound on pooled reader connections (claims/lookups ride
|
||||
/// these). Fresh readers open on demand up to the bound; `0` clamps
|
||||
/// to 1.
|
||||
pub max_readers: usize,
|
||||
}
|
||||
|
||||
impl Default for SqliteOpts {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
poll_interval: None,
|
||||
max_readers: DEFAULT_MAX_READERS,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl SqliteOpts {
|
||||
/// Resolve this opts into the substrate's
|
||||
/// [`WatcherConfig`](crate::substrate::WatcherConfig), mapping the
|
||||
/// substrate's string error into the contract taxonomy's
|
||||
/// [`Error::Database`](alkstore::Error) arm (the W-2 posture —
|
||||
/// engine-layer mapping, ADR-012 §4).
|
||||
pub(crate) fn watcher_config(&self) -> Result<WatcherConfig, alkstore::Error> {
|
||||
let config = WatcherConfig::default();
|
||||
match self.poll_interval {
|
||||
None => Ok(config),
|
||||
Some(d) => config
|
||||
.with_poll_interval(d)
|
||||
.map_err(alkstore::Error::database),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,45 @@
|
||||
//! The `spawn_blocking` seam (ADR-003, ADR-007) and the error
|
||||
//! mapping — the two engine-wide postures established here that every
|
||||
//! later mechanism task reuses.
|
||||
//!
|
||||
//! - Every engine call round-trips a blocking thread: rusqlite
|
||||
//! connections are not `Send`-across-await, and the sync substrate
|
||||
//! bridges at the trait seam (the family REQ-TTY-01 posture).
|
||||
//! - Substrate errors map into the contract taxonomy's opaque
|
||||
//! [`Error::Database`](alkstore::Error) fallback with the source
|
||||
//! chain preserved (ADR-008 §5). No panics cross the seam: a panic
|
||||
//! inside a bridged closure is the `JoinError` arm's `Database`
|
||||
//! mapping, family-standard no-panics discipline maintained.
|
||||
|
||||
use alkstore::Error;
|
||||
|
||||
/// The blocking bridge every trait method's substrate round trip goes
|
||||
/// through. Private: the seam is engine-internal — consumers meet it
|
||||
/// only through the async trait surface. `cfg(test)`-gated until the
|
||||
/// seam task (`sqlite-engine-seam-tx`) wires the trait impls; the
|
||||
/// helper's shape and error path are pinned by this wave's smoke test.
|
||||
#[cfg(test)]
|
||||
pub(crate) async fn blocking<T, F>(f: F) -> Result<T, Error>
|
||||
where
|
||||
T: Send + 'static,
|
||||
F: FnOnce() -> Result<T, Error> + Send + 'static,
|
||||
{
|
||||
match tokio::task::spawn_blocking(f).await {
|
||||
Ok(result) => result,
|
||||
Err(e) => Err(database_error(format!("blocking task failed: {e}"))),
|
||||
}
|
||||
}
|
||||
|
||||
/// Map a substrate string failure into the contract taxonomy. The
|
||||
/// substrate is contract-blind — its fallible API surfaces strings
|
||||
/// (the W-2 spawn posture, ADR-012 §4); the engine layer owns the
|
||||
/// mapping, ADR-012 §2.
|
||||
pub(crate) fn database_error(message: impl Into<String>) -> Error {
|
||||
Error::database(std::io::Error::other(message.into()))
|
||||
}
|
||||
|
||||
/// Map a `rusqlite::Error` into the contract taxonomy, source chain
|
||||
/// preserved (ADR-008 §5 — no engine-specific variants minted).
|
||||
pub(crate) fn sqlite_error(e: rusqlite::Error) -> Error {
|
||||
Error::database(e)
|
||||
}
|
||||
@@ -0,0 +1,220 @@
|
||||
//! The SQLite engine's store: the `open` constructor, the connection
|
||||
//! architecture (engine-sqlite.md's "Connection architecture"), and
|
||||
//! the [`Store`] trait wiring.
|
||||
//!
|
||||
//! Connection architecture: one writer connection behind the
|
||||
//! substrate's `Writer` slot (WAL single-writer modeled honestly), a
|
||||
//! pool of reader connections (fresh opens on demand up to
|
||||
//! [`SqliteOpts::max_readers`]), and one `SharedUpdateWatcher` poll
|
||||
//! thread over the db path. Every connection — writer and pooled
|
||||
//! readers alike — runs the substrate's full bootstrap at open:
|
||||
//! pragmas (WAL, `synchronous=NORMAL`, busy timeout), the notify
|
||||
//! machinery, the substrate's SQL function set, and the schema
|
||||
//! bootstrap.
|
||||
//!
|
||||
//! Watcher spawn is fallible in the substrate (ADR-012 §4's W-2); the
|
||||
//! failure maps into [`Error::Database`](alkstore::Error::Database)
|
||||
//! engine-side and `open` fails if the watcher cannot start. On
|
||||
//! watcher death the substrate's `WatcherDeathGuard` closes every
|
||||
//! subscriber (ADR-006) — the death-guard plumbing exists from open
|
||||
//! onward; the consumer-visible receiver bridging arrives with
|
||||
//! `listen`.
|
||||
//!
|
||||
//! Drop/close: the watcher's poll thread is joined (releasing its
|
||||
//! connection and the db file handle) and the writer slot + reader
|
||||
//! pool are closed; [`Drop`](SqliteStore::drop) delegates to
|
||||
//! [`SqliteStore::close`].
|
||||
//!
|
||||
//! Every engine call round-trips the `spawn_blocking` seam
|
||||
//! ([`crate::seam`]); rusqlite connections are not
|
||||
//! `Send`-across-await. Single-host by nature: file-backed, one
|
||||
//! machine; NFS two-writers-unsupported (the lineage's honesty
|
||||
//! posture, inherited — ADR-016).
|
||||
|
||||
use alkstore::Store;
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
|
||||
use crate::opts::SqliteOpts;
|
||||
use crate::seam::{database_error, sqlite_error};
|
||||
use crate::substrate::{Readers, SharedUpdateWatcher, Writer};
|
||||
|
||||
/// The SQLite engine's store handle. Consumers hold
|
||||
/// [`alkstore::Store`] — the trait is the contract surface; the type
|
||||
/// exists so the engine can own its machinery and so tests can reach
|
||||
/// the close path directly.
|
||||
pub struct SqliteStore {
|
||||
writer: Arc<Writer>,
|
||||
readers: Arc<Readers>,
|
||||
watcher: Arc<SharedUpdateWatcher>,
|
||||
db_path: PathBuf,
|
||||
}
|
||||
|
||||
impl std::fmt::Debug for SqliteStore {
|
||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
f.debug_struct("SqliteStore")
|
||||
.field("db_path", &self.db_path)
|
||||
.finish_non_exhaustive()
|
||||
}
|
||||
}
|
||||
|
||||
impl SqliteStore {
|
||||
/// Explicit close: stops the watcher, disconnects every subscriber
|
||||
/// (the death-guard close — receivers see them close), closes the
|
||||
/// reader pool, and closes the writer slot. Idempotent. [`Drop`]
|
||||
/// delegates here.
|
||||
pub fn close(&self) {
|
||||
let _ = self.watcher.close();
|
||||
self.readers.close();
|
||||
self.writer.close();
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for SqliteStore {
|
||||
fn drop(&mut self) {
|
||||
self.close();
|
||||
}
|
||||
}
|
||||
|
||||
/// Open a SQLite-backed store at `path`.
|
||||
///
|
||||
/// The path is a **plain filesystem path** (ADR-023 §3): no URI
|
||||
/// interpretation — a `?name=value` suffix is part of the literal
|
||||
/// filename (`:memory:` works without URI mode). The full connection
|
||||
/// architecture boots: the writer slot, the reader pool, and the
|
||||
/// update watcher (whose spawn is fallible — open fails with
|
||||
/// [`Error::Database`](alkstore::Error::Database) if it cannot start).
|
||||
pub fn open(path: &str, opts: SqliteOpts) -> alkstore::Result<Box<dyn Store>> {
|
||||
Ok(Box::new(open_store(path, opts)?))
|
||||
}
|
||||
|
||||
fn open_store(path: &str, opts: SqliteOpts) -> alkstore::Result<SqliteStore> {
|
||||
let config = opts.watcher_config()?;
|
||||
|
||||
let writer_conn = crate::substrate::open_conn_bootstrapped(path)
|
||||
.map_err(|e| match e {
|
||||
crate::substrate::SchemaError::Sqlite(inner) => sqlite_error(inner),
|
||||
})
|
||||
.map_err(|e| database_error(format!("failed to open/bootstrap writer connection: {e}")))?;
|
||||
|
||||
let watcher = SharedUpdateWatcher::new_with_config(PathBuf::from(path), config)
|
||||
.map_err(|e| database_error(format!("failed to spawn update watcher: {e}")))?;
|
||||
|
||||
Ok(SqliteStore {
|
||||
writer: Arc::new(Writer::new(writer_conn)),
|
||||
readers: Arc::new(Readers::new(path.to_string(), opts.max_readers)),
|
||||
watcher: Arc::new(watcher),
|
||||
db_path: PathBuf::from(path),
|
||||
})
|
||||
}
|
||||
|
||||
impl Store for SqliteStore {
|
||||
fn begin_tx(
|
||||
&self,
|
||||
) -> alkstore::BoxedFuture<'_, alkstore::Result<Box<dyn alkstore::TxHandle + Send>>> {
|
||||
Box::pin(async { Err(database_error("begin_tx wiring lands with the seam task")) })
|
||||
}
|
||||
|
||||
fn notify<'a>(
|
||||
&'a self,
|
||||
_channel: &str,
|
||||
_payload: serde_json::Value,
|
||||
) -> alkstore::BoxedFuture<'a, alkstore::Result<()>> {
|
||||
Box::pin(async {
|
||||
Err(database_error(
|
||||
"notify wiring lands with the mechanism tasks",
|
||||
))
|
||||
})
|
||||
}
|
||||
|
||||
fn listen<'a>(
|
||||
&'a self,
|
||||
_channel: &str,
|
||||
) -> alkstore::BoxedFuture<'a, alkstore::Result<Box<dyn alkstore::WakeReceiver>>> {
|
||||
Box::pin(async {
|
||||
Err(database_error(
|
||||
"listen wiring lands with the mechanism tasks",
|
||||
))
|
||||
})
|
||||
}
|
||||
|
||||
fn stream<'a>(
|
||||
&'a self,
|
||||
_name: &str,
|
||||
) -> alkstore::BoxedFuture<'a, alkstore::Result<Box<dyn alkstore::StreamHandle>>> {
|
||||
Box::pin(async {
|
||||
Err(database_error(
|
||||
"stream wiring lands with the mechanism tasks",
|
||||
))
|
||||
})
|
||||
}
|
||||
|
||||
fn queue<'a>(
|
||||
&'a self,
|
||||
_name: &str,
|
||||
_opts: alkstore::QueueOpts,
|
||||
) -> alkstore::BoxedFuture<'a, alkstore::Result<Box<dyn alkstore::Queue>>> {
|
||||
Box::pin(async {
|
||||
Err(database_error(
|
||||
"queue wiring lands with the mechanism tasks",
|
||||
))
|
||||
})
|
||||
}
|
||||
|
||||
fn outbox<'a>(
|
||||
&'a self,
|
||||
_name: &str,
|
||||
) -> alkstore::BoxedFuture<'a, alkstore::Result<Box<dyn alkstore::Outbox>>> {
|
||||
Box::pin(async {
|
||||
Err(database_error(
|
||||
"outbox wiring lands with the mechanism tasks",
|
||||
))
|
||||
})
|
||||
}
|
||||
|
||||
fn try_lock<'a>(
|
||||
&'a self,
|
||||
_name: &str,
|
||||
_owner: &str,
|
||||
_ttl: i64,
|
||||
) -> alkstore::BoxedFuture<'a, alkstore::Result<Option<Box<dyn alkstore::Lock>>>> {
|
||||
Box::pin(async { Err(database_error("lock wiring lands with the mechanism tasks")) })
|
||||
}
|
||||
|
||||
fn schedule<'a>(
|
||||
&'a self,
|
||||
_name: &str,
|
||||
_spec: &str,
|
||||
_queue: &str,
|
||||
_payload: serde_json::Value,
|
||||
_opts: alkstore::ScheduleOpts,
|
||||
) -> alkstore::BoxedFuture<'a, alkstore::Result<alkstore::Schedule>> {
|
||||
Box::pin(async {
|
||||
Err(database_error(
|
||||
"schedule wiring lands with the mechanism tasks",
|
||||
))
|
||||
})
|
||||
}
|
||||
|
||||
fn unschedule<'a>(&'a self, _name: &str) -> alkstore::BoxedFuture<'a, alkstore::Result<bool>> {
|
||||
Box::pin(async {
|
||||
Err(database_error(
|
||||
"unschedule wiring lands with the mechanism tasks",
|
||||
))
|
||||
})
|
||||
}
|
||||
|
||||
fn run_schedules<'a>(
|
||||
&'a self,
|
||||
_stop: alkstore::StopToken,
|
||||
) -> alkstore::BoxedFuture<'a, alkstore::Result<()>> {
|
||||
Box::pin(async {
|
||||
Err(database_error(
|
||||
"scheduler wiring lands with the mechanism tasks",
|
||||
))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod open_tests;
|
||||
@@ -0,0 +1,415 @@
|
||||
//! The `open` constructor's acceptance tests (`sqlite-engine-open-opts`):
|
||||
//! the connection architecture boots, opts flow, the plain-path
|
||||
//! posture holds, and close/drop teardown releases the file.
|
||||
|
||||
use std::path::PathBuf;
|
||||
use std::time::Duration;
|
||||
|
||||
use crate::opts::{DEFAULT_MAX_READERS, SqliteOpts};
|
||||
use crate::seam::blocking;
|
||||
use crate::store::{open, open_store};
|
||||
|
||||
fn temp_dir(tag: &str) -> PathBuf {
|
||||
let d = std::env::temp_dir().join(format!(
|
||||
"alkstore-open-{tag}-{}-{:?}",
|
||||
std::process::id(),
|
||||
std::thread::current().id()
|
||||
));
|
||||
let _ = std::fs::remove_dir_all(&d);
|
||||
std::fs::create_dir_all(&d).unwrap();
|
||||
d
|
||||
}
|
||||
|
||||
fn cleanup(dir: &PathBuf) {
|
||||
let _ = std::fs::remove_dir_all(dir);
|
||||
}
|
||||
|
||||
/// Acceptance: `open(path, opts)` boots the full connection
|
||||
/// architecture — the writer slot takes writes, the reader pool
|
||||
/// serves reads, and the substrate's bootstrap ran on every
|
||||
/// connection (the `__alkstore_*` table family exists; the notify and
|
||||
/// ops SQL functions respond on a fresh pooled reader).
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn open_boots_the_full_connection_architecture() {
|
||||
let dir = temp_dir("boot");
|
||||
let path = dir.join("store.db");
|
||||
let store = open_store(path.to_str().unwrap(), SqliteOpts::default()).unwrap();
|
||||
|
||||
let writer_conn = store.writer.acquire().unwrap();
|
||||
let mode: String = writer_conn
|
||||
.pragma_query_value(None, "journal_mode", |r| r.get(0))
|
||||
.unwrap();
|
||||
assert!(
|
||||
mode.eq_ignore_ascii_case("wal"),
|
||||
"writer not in WAL: {mode}"
|
||||
);
|
||||
store.writer.release(writer_conn);
|
||||
|
||||
let reader_conn = store.readers.acquire().unwrap();
|
||||
let reader_mode: String = reader_conn
|
||||
.pragma_query_value(None, "journal_mode", |r| r.get(0))
|
||||
.unwrap();
|
||||
assert!(
|
||||
reader_mode.eq_ignore_ascii_case("wal"),
|
||||
"pooled reader not in WAL: {reader_mode}"
|
||||
);
|
||||
let notify_fn: i64 = reader_conn
|
||||
.query_row("SELECT notify('__alkstore_probe', 'x')", [], |r| r.get(0))
|
||||
.unwrap();
|
||||
assert!(notify_fn >= 1, "pooled reader missing the function set");
|
||||
let publish_offset: i64 = reader_conn
|
||||
.query_row(
|
||||
"SELECT alkstore_stream_publish('__alkstore_probe', NULL, '{}')",
|
||||
[],
|
||||
|r| r.get(0),
|
||||
)
|
||||
.unwrap();
|
||||
assert!(publish_offset >= 1, "pooled reader missing the ops surface");
|
||||
store.readers.release(reader_conn);
|
||||
|
||||
assert!(store.db_path.exists(), "db file at the plain path");
|
||||
drop(store);
|
||||
cleanup(&dir);
|
||||
}
|
||||
|
||||
/// Acceptance: the watcher is up after open — a commit from the
|
||||
/// writer slot fans out to a raw substrate subscriber (the guard the
|
||||
/// death plumbing rides; the consumer-visible receiver bridge lands
|
||||
/// with `listen`).
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn open_leaves_the_watcher_running_and_fanning_out() {
|
||||
let dir = temp_dir("watcher");
|
||||
let path = dir.join("store.db");
|
||||
let store = open_store(path.to_str().unwrap(), SqliteOpts::default()).unwrap();
|
||||
|
||||
let (_id, rx) = store.watcher.subscribe();
|
||||
|
||||
let writer_conn = store.writer.acquire().unwrap();
|
||||
writer_conn
|
||||
.execute("CREATE TABLE t(x INTEGER PRIMARY KEY)", [])
|
||||
.unwrap();
|
||||
writer_conn.execute("INSERT INTO t VALUES (1)", []).unwrap();
|
||||
store.writer.release(writer_conn);
|
||||
|
||||
let deadline = std::time::Instant::now() + Duration::from_secs(2);
|
||||
let got = loop {
|
||||
match rx.try_recv() {
|
||||
Ok(()) => break true,
|
||||
Err(std::sync::mpsc::TryRecvError::Empty) => {
|
||||
assert!(
|
||||
std::time::Instant::now() < deadline,
|
||||
"watcher not fanning out after a commit"
|
||||
);
|
||||
std::thread::sleep(Duration::from_millis(5));
|
||||
}
|
||||
Err(std::sync::mpsc::TryRecvError::Disconnected) => break false,
|
||||
}
|
||||
};
|
||||
assert!(
|
||||
got,
|
||||
"subscriber channel disconnected instead of fanning out"
|
||||
);
|
||||
|
||||
store.close();
|
||||
cleanup(&dir);
|
||||
}
|
||||
|
||||
/// Acceptance: watcher death handling wired — closing the shared
|
||||
/// watcher disconnects every subscriber receiver (the death-guard
|
||||
/// plumbing `listen` bridges into `recv() -> None`). The mechanism
|
||||
/// lands with `listen`; the plumbing exists here.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn closing_the_store_closes_watcher_subscribers() {
|
||||
let dir = temp_dir("death");
|
||||
let path = dir.join("store.db");
|
||||
let store = open_store(path.to_str().unwrap(), SqliteOpts::default()).unwrap();
|
||||
|
||||
let (_id, rx) = store.watcher.subscribe();
|
||||
store.close();
|
||||
assert!(
|
||||
rx.recv().is_err(),
|
||||
"subscriber receiver must close with the watcher (WatcherDeathGuard)"
|
||||
);
|
||||
cleanup(&dir);
|
||||
}
|
||||
|
||||
/// Acceptance: `poll_interval` flows to the watcher's config (the
|
||||
/// config-surface accounting test — asserting the configured cadence,
|
||||
/// not wall-clock timing; ADR-023 §4's backlog row; `None` = the 1 ms
|
||||
/// default, zero rejected at open).
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn poll_interval_flows_to_the_watcher_config() {
|
||||
use crate::substrate::DEFAULT_WATCHER_POLL_INTERVAL;
|
||||
|
||||
assert_eq!(
|
||||
DEFAULT_WATCHER_POLL_INTERVAL,
|
||||
Duration::from_millis(1),
|
||||
"the substrate's shipping default is 1 ms (ADR-023 §4)"
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
SqliteOpts::default().poll_interval,
|
||||
None,
|
||||
"None = the substrate's 1 ms default"
|
||||
);
|
||||
|
||||
let opts = SqliteOpts {
|
||||
poll_interval: Some(Duration::from_millis(25)),
|
||||
..Default::default()
|
||||
};
|
||||
let config = opts.watcher_config().unwrap();
|
||||
assert_eq!(
|
||||
config.poll_interval,
|
||||
Duration::from_millis(25),
|
||||
"the knob must flow to the watcher's configured cadence"
|
||||
);
|
||||
|
||||
let zero_err = SqliteOpts {
|
||||
poll_interval: Some(Duration::ZERO),
|
||||
..Default::default()
|
||||
}
|
||||
.watcher_config()
|
||||
.unwrap_err();
|
||||
assert!(
|
||||
matches!(zero_err, alkstore::Error::Database(_)),
|
||||
"the W-2 mapping lands in Error::Database, got: {zero_err:?}"
|
||||
);
|
||||
let source = match &zero_err {
|
||||
alkstore::Error::Database(source) => source.to_string(),
|
||||
_ => unreachable!("matched above"),
|
||||
};
|
||||
assert!(
|
||||
source.contains("poll interval must be positive"),
|
||||
"the substrate's zero-interval rejection must ride the source chain, got: {source}"
|
||||
);
|
||||
}
|
||||
|
||||
/// The 1 ms default resolves at open without an explicit interval,
|
||||
/// and a generous explicit interval opens just the same.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn open_succeeds_with_default_and_explicit_poll_intervals() {
|
||||
let dir = temp_dir("cadence");
|
||||
for opts in [
|
||||
SqliteOpts::default(),
|
||||
SqliteOpts {
|
||||
poll_interval: Some(Duration::from_millis(50)),
|
||||
..Default::default()
|
||||
},
|
||||
] {
|
||||
let path = dir.join(format!(
|
||||
"cadence-{}.db",
|
||||
opts.poll_interval.map(|d| d.as_millis()).unwrap_or(0)
|
||||
));
|
||||
let store = open_store(path.to_str().unwrap(), opts).unwrap();
|
||||
assert!(store.db_path.exists());
|
||||
store.close();
|
||||
}
|
||||
cleanup(&dir);
|
||||
}
|
||||
|
||||
/// Acceptance: the plain-path posture holds (ADR-023 §3 pin — a
|
||||
/// `?name=value`-suffixed path creates a file with that literal
|
||||
/// name; no URI flag anywhere in the open path).
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn uri_shaped_open_path_is_a_literal_filename() {
|
||||
let dir = temp_dir("plain-path");
|
||||
let path = dir.join("plain.db?mode=memory&cache=shared");
|
||||
let store = open_store(path.to_str().unwrap(), SqliteOpts::default()).unwrap();
|
||||
let writer_conn = store.writer.acquire().unwrap();
|
||||
writer_conn.execute("CREATE TABLE t(x)", []).unwrap();
|
||||
store.writer.release(writer_conn);
|
||||
store.close();
|
||||
assert!(
|
||||
path.exists(),
|
||||
"the file must exist under its literal `?name=value` name (no URI interpretation)"
|
||||
);
|
||||
cleanup(&dir);
|
||||
}
|
||||
|
||||
/// `:memory:` works without URI mode (ADR-023 §3): each connection is
|
||||
/// its own private memory db, so the writer carries the bootstrap and
|
||||
/// a fresh pooled reader opens the same way.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn memory_paths_open_without_uri_mode() {
|
||||
let store = open_store(":memory:", SqliteOpts::default()).unwrap();
|
||||
let reader_conn = store.readers.acquire().unwrap();
|
||||
let tables: i64 = reader_conn
|
||||
.query_row(
|
||||
"SELECT count(*) FROM sqlite_master
|
||||
WHERE type = 'table' AND name LIKE '__alkstore_%'",
|
||||
[],
|
||||
|r| r.get(0),
|
||||
)
|
||||
.unwrap();
|
||||
store.readers.release(reader_conn);
|
||||
assert!(tables >= 7, ":memory: open must carry the bootstrap");
|
||||
store.close();
|
||||
}
|
||||
|
||||
/// Acceptance: reader-pool sizing — `max_readers` bounds the pool
|
||||
/// (excess acquirers block until release) and `0` clamps to 1
|
||||
/// (the substrate's max discipline).
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn reader_pool_size_is_bounded_by_opts() {
|
||||
let dir = temp_dir("pool");
|
||||
let path = dir.join("store.db");
|
||||
|
||||
let store = open_store(
|
||||
path.to_str().unwrap(),
|
||||
SqliteOpts {
|
||||
max_readers: 2,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let c1 = store.readers.acquire().unwrap();
|
||||
let c2 = store.readers.acquire().unwrap();
|
||||
|
||||
let readers = store.readers.clone();
|
||||
let blocked = std::thread::spawn(move || readers.acquire());
|
||||
std::thread::sleep(Duration::from_millis(100));
|
||||
assert!(
|
||||
!blocked.is_finished(),
|
||||
"acquire above max_readers must block (pool bound enforced)"
|
||||
);
|
||||
|
||||
store.readers.release(c2);
|
||||
let c3 = blocked.join().unwrap().unwrap();
|
||||
store.readers.release(c3);
|
||||
store.readers.release(c1);
|
||||
store.close();
|
||||
|
||||
let store = open_store(
|
||||
path.to_str().unwrap(),
|
||||
SqliteOpts {
|
||||
max_readers: 0,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.unwrap();
|
||||
let only = store.readers.acquire().unwrap();
|
||||
store.readers.release(only);
|
||||
store.close();
|
||||
cleanup(&dir);
|
||||
}
|
||||
|
||||
/// The default reader count is the documented constant.
|
||||
#[test]
|
||||
fn default_max_readers_is_the_documented_constant() {
|
||||
assert_eq!(DEFAULT_MAX_READERS, 8, "the documented default pool size");
|
||||
assert_eq!(
|
||||
SqliteOpts::default().max_readers,
|
||||
DEFAULT_MAX_READERS,
|
||||
"`Default` resolves the reader-pool knob to the documented default"
|
||||
);
|
||||
}
|
||||
|
||||
/// A pool beyond its bound opens fresh readers on demand up to the
|
||||
/// bound; reads on those readers see a writer's committed data
|
||||
/// (the WAL cross-connection read path the reader pool exists for).
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn pooled_readers_see_writer_commits() {
|
||||
let dir = temp_dir("read-visibility");
|
||||
let path = dir.join("store.db");
|
||||
let store = open_store(path.to_str().unwrap(), SqliteOpts::default()).unwrap();
|
||||
|
||||
let writer_conn = store.writer.acquire().unwrap();
|
||||
writer_conn
|
||||
.execute("CREATE TABLE t(x INTEGER)", [])
|
||||
.unwrap();
|
||||
writer_conn
|
||||
.execute_batch("INSERT INTO t VALUES (42);")
|
||||
.unwrap();
|
||||
store.writer.release(writer_conn);
|
||||
|
||||
let reader_conn = store.readers.acquire().unwrap();
|
||||
let x: i64 = reader_conn
|
||||
.query_row("SELECT x FROM t", [], |r| r.get(0))
|
||||
.unwrap();
|
||||
store.readers.release(reader_conn);
|
||||
assert_eq!(x, 42, "pooled reader reads the writer's committed row");
|
||||
|
||||
store.close();
|
||||
cleanup(&dir);
|
||||
}
|
||||
|
||||
/// Explicit close: the watcher joins (its thread exits), the file
|
||||
/// handle is released, and the store is gone-but-reopenable.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn explicit_close_releases_the_db_file_and_permits_reopen() {
|
||||
let dir = temp_dir("close");
|
||||
let path = dir.join("store.db");
|
||||
let store = open_store(path.to_str().unwrap(), SqliteOpts::default()).unwrap();
|
||||
store.close();
|
||||
|
||||
let reopened = open(path.to_str().unwrap(), SqliteOpts::default()).unwrap();
|
||||
drop(reopened);
|
||||
cleanup(&dir);
|
||||
}
|
||||
|
||||
/// `Drop` stops the watcher and closes connections — the same teardown
|
||||
/// as explicit close, reached through the drop path.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn drop_runs_the_full_teardown() {
|
||||
let dir = temp_dir("drop");
|
||||
let path = dir.join("store.db");
|
||||
{
|
||||
let store = open_store(path.to_str().unwrap(), SqliteOpts::default()).unwrap();
|
||||
let (_id, rx) = store.watcher.subscribe();
|
||||
drop(store);
|
||||
assert!(
|
||||
rx.recv().is_err(),
|
||||
"drop must run the same teardown as close (subscriber closed)"
|
||||
);
|
||||
}
|
||||
let reopened = open(path.to_str().unwrap(), SqliteOpts::default()).unwrap();
|
||||
drop(reopened);
|
||||
cleanup(&dir);
|
||||
}
|
||||
|
||||
/// Acceptance: the `spawn_blocking` seam — every engine call's
|
||||
/// substrate round trip rides this helper; smoked here through the
|
||||
/// real bridge with a rusqlite-shaped closure (the pattern the
|
||||
/// mechanism tasks reuse one-for-one).
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn the_blocking_bridge_round_trips_and_maps_errors() {
|
||||
let ok = blocking(|| Ok::<_, alkstore::Error>(7i64)).await;
|
||||
assert_eq!(ok.unwrap(), 7);
|
||||
|
||||
let err = blocking(|| -> Result<(), alkstore::Error> {
|
||||
Err(alkstore::Error::database(std::io::Error::other(
|
||||
"smoke failure",
|
||||
)))
|
||||
})
|
||||
.await;
|
||||
assert!(matches!(err, Err(alkstore::Error::Database(_))));
|
||||
|
||||
let panicking: Result<(), alkstore::Error> =
|
||||
blocking(|| -> Result<(), alkstore::Error> { panic!("bridge smoke panic") }).await;
|
||||
assert!(
|
||||
matches!(&panicking, Err(alkstore::Error::Database(_))),
|
||||
"a panic inside a bridged closure maps to Database (no panics cross the seam), got {panicking:?}"
|
||||
);
|
||||
}
|
||||
|
||||
/// `open` on an unopenable path fails cleanly — the substrate's
|
||||
/// bootstrap errors map into `Error::Database` at the constructor, not
|
||||
/// panics.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn open_fails_cleanly_on_an_unopenable_path() {
|
||||
let result = open(
|
||||
"/this/parent/does/not/exist/alkstore-open-fail.db",
|
||||
SqliteOpts::default(),
|
||||
);
|
||||
match result {
|
||||
Err(alkstore::Error::Database(source)) => {
|
||||
assert!(
|
||||
!source.to_string().is_empty(),
|
||||
"the source chain carries the substrate's detail"
|
||||
);
|
||||
}
|
||||
Ok(_) => panic!("expected Err(Database) from an unopenable path, got Ok"),
|
||||
Err(e) => panic!("expected Err(Database) from an unopenable path, got {e}"),
|
||||
}
|
||||
}
|
||||
@@ -118,6 +118,7 @@ fixes. Cherry-picks: empty at scaffold, append-only forever.
|
||||
| D-27 | re-derivation | wave-2 review fix — `scheduler_tick`'s fire enqueue resolved the schedule row's relative `expires_s` into the absolute row `expires_at` (`now + s` at the fire transaction, one clock read) instead of passing the relative value through; regression test pins a fired job's `expires_at = fire-instant + s` and its claimability. The pre-fix pass-through would have produced rows expired in 1970 under any `expires_s` schedule | `queue_ops.rs` (`scheduler_tick`, fire-expiry resolution + test) | ADR-020 §2 (relative `expires` → resolved absolute row value, resolved at enqueue against the single clock) — D-12's landed state carried the lineage's `expires_s` field shape without carrying the lineage's enqueue-side resolution; wave-2 review (task review-wave-2) re-derivation spot-check found it | ours |
|
||||
| D-28 | hygiene | wave-2 review fix — the cross-mechanism pressure test's lock-churn branch exercised a per-producer lock name (literal `"p{producer}"` had no interpolation, so all producers contended on one name) with its result swallowed via `.unwrap_or(0)`; now interpolates, releases periodically, propagates errors | `ops.rs` (pressure test, lock branch) | test-quality only (no lineage counterpart, no library code touched); `fork-provenance-and-floor`'s ported-surface pressure intent | ours |
|
||||
| D-29 | port delta | `SQLITE_OPEN_URI` dropped from `open_conn`'s flags (the watcher's own opens never carried it — the inherited posture was internally inconsistent); the path argument is a plain filesystem path, `?name=value` suffixes are literal filenames, no connection-semantics mutation outside the engine constructor's opts; pinned by test (`open_conn_treats_uri_shaped_path_as_plain_filename`) | `schema.rs` (`open_conn`) | ADR-023 §3 (plain-path open — the URI parameter surface is SQLite's, not this engine's; no consumer-inventory row names URI features) | ours |
|
||||
| D-30 | port delta | `open_conn_bootstrapped` — new function stacking the engine's full per-connection open posture (pragmas, notify, `attach_alkstore_functions`, `bootstrap_schema`), and `Readers::acquire`'s open path switched from bare `open_conn` to it (fresh connections only — the pool never re-bootstraps a reused connection). The lineage's pools opened readers pragmas+notify only, the function set + schema living only on the writer; this engine's reader-slot operations (claim/ack, stream reads, lock ops, scheduler checks, leadership probes) run this substrate's SQL functions on pooled connections, so every pooled connection carries the full surface | `schema.rs` (`open_conn_bootstrapped`, `Readers::acquire` open path) | engine-sqlite.md's connection architecture ("each connection runs the substrate's bootstrap at open — the writer and each pooled reader") — engine-layer obligation of `sqlite-engine-open-opts`; landed with the engine's `open` wiring | ours |
|
||||
|
||||
### Cherry-picks
|
||||
|
||||
|
||||
@@ -53,5 +53,8 @@ mod queue_ops;
|
||||
mod schema;
|
||||
mod watcher;
|
||||
|
||||
pub(crate) use schema::{attach_notify, bootstrap_schema, open_conn};
|
||||
pub(crate) use watcher::SharedUpdateWatcher;
|
||||
pub(crate) use schema::{
|
||||
Error as SchemaError, Readers, Writer, apply_default_pragmas, attach_notify, bootstrap_schema,
|
||||
open_conn, open_conn_bootstrapped,
|
||||
};
|
||||
pub(crate) use watcher::{DEFAULT_WATCHER_POLL_INTERVAL, SharedUpdateWatcher, WatcherConfig};
|
||||
@@ -324,6 +324,27 @@ pub(crate) fn open_conn(path: &str, install_notify: bool) -> Result<Connection,
|
||||
Ok(conn)
|
||||
}
|
||||
|
||||
/// Open a connection carrying the engine's full per-connection
|
||||
/// bootstrap: pragmas, notify machinery, this substrate's function
|
||||
/// set, and schema bootstrap. The engine's `open` constructor runs
|
||||
/// this on the writer and each pooled reader (engine-sqlite.md's
|
||||
/// connection architecture: "each connection runs the substrate's
|
||||
/// bootstrap at open").
|
||||
///
|
||||
/// D-30: the lineage's pools open readers through bare `open_conn`
|
||||
/// (pragmas + notify only), the function set + schema living only on
|
||||
/// the writer. This substrate extends the pool open path with
|
||||
/// `attach_alkstore_functions` + `bootstrap_schema` — the engine's
|
||||
/// reader-slot operations (claim/ack, stream reads, lock ops,
|
||||
/// scheduler checks) run this substrate's SQL functions on pooled
|
||||
/// connections, so every pooled connection must carry them.
|
||||
pub(crate) fn open_conn_bootstrapped(path: &str) -> Result<Connection, Error> {
|
||||
let conn = open_conn(path, true)?;
|
||||
super::ops::attach_alkstore_functions(&conn)?;
|
||||
bootstrap_schema(&conn)?;
|
||||
Ok(conn)
|
||||
}
|
||||
|
||||
pub(crate) struct Writer {
|
||||
slot: Mutex<Option<Connection>>,
|
||||
available: Condvar,
|
||||
@@ -411,7 +432,7 @@ impl Readers {
|
||||
*out += 1;
|
||||
drop(out);
|
||||
drop(pool);
|
||||
match open_conn(&self.path, false) {
|
||||
match open_conn_bootstrapped(&self.path) {
|
||||
Ok(conn) => {
|
||||
if self.closed.load(Ordering::Acquire) {
|
||||
*self.outstanding.lock() -= 1;
|
||||
@@ -577,6 +598,55 @@ mod writer_reader_tests {
|
||||
let _ = std::fs::remove_dir_all(&dir);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn open_conn_bootstrapped_carries_the_full_surface_on_every_connection() {
|
||||
let dir = std::env::temp_dir().join(format!(
|
||||
"alkstore-bootstrapped-{}-{:?}",
|
||||
std::process::id(),
|
||||
std::thread::current().id()
|
||||
));
|
||||
std::fs::create_dir_all(&dir).unwrap();
|
||||
let path = dir.join("bootstrapped.db");
|
||||
|
||||
let first = open_conn_bootstrapped(path.to_str().unwrap()).unwrap();
|
||||
drop(first);
|
||||
|
||||
for i in 0..3 {
|
||||
let conn = open_conn_bootstrapped(path.to_str().unwrap())
|
||||
.unwrap_or_else(|e| panic!("bootstrapped open {i} failed: {e}"));
|
||||
let mode: String = conn
|
||||
.pragma_query_value(None, "journal_mode", |r| r.get(0))
|
||||
.unwrap();
|
||||
assert!(mode.eq_ignore_ascii_case("wal"), "open {i}: not WAL");
|
||||
let notify_fn: i64 = conn
|
||||
.query_row("SELECT notify('__alkstore_probe', 'x')", [], |r| r.get(0))
|
||||
.unwrap_or(0);
|
||||
assert!(notify_fn >= 1, "open {i}: notify function missing");
|
||||
let tables: i64 = conn
|
||||
.query_row(
|
||||
"SELECT count(*) FROM sqlite_master
|
||||
WHERE type = 'table' AND name LIKE '__alkstore_%'",
|
||||
[],
|
||||
|r| r.get(0),
|
||||
)
|
||||
.unwrap();
|
||||
assert!(
|
||||
tables >= 7,
|
||||
"open {i}: bootstrap schema missing (saw {tables} tables)"
|
||||
);
|
||||
let publish_ready: i64 = conn
|
||||
.query_row(
|
||||
"SELECT alkstore_stream_publish('__alkstore_probe', NULL, '{}')",
|
||||
[],
|
||||
|r| r.get(0),
|
||||
)
|
||||
.unwrap_or(0);
|
||||
assert!(publish_ready >= 1, "open {i}: function set missing");
|
||||
drop(conn);
|
||||
}
|
||||
let _ = std::fs::remove_dir_all(&dir);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn concurrent_opens_never_return_database_is_locked() {
|
||||
const ROUNDS: usize = 8;
|
||||
|
||||
Reference in new issue
Block a user