From 3ffe7e00cedffc907811246de6b347ece2d01158 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 17 Sep 2026 19:30:30 -0400 Subject: [PATCH] Implement Petri's run store over Fabro's SQLite database `SqliteRunStore` implements Petri's `RunStore` and `RunLogs` on the pool Fabro's other stores share. A run's existence and writer lease live in `petri_runs`; every record of every log lives in `petri_records`, keyed by (run, log, seq) with the record stored as JSON and read back unchanged; blobs share the `blobs` table with `BlobStore`. The lease is taken idempotently per owner, ends when the last handle drops or when an operator releases it, and never by timeout; every write checks it inside its own transaction. An append is one `BEGIN IMMEDIATE` transaction per batch: a repeated record is accepted, a different record at a taken seq or a seq past the head is a conflict that stores nothing. Petri's store conformance suite passes against it, with the operator release, lease exclusivity, a crash between appends, and blob interoperation checked beside it. Co-Authored-By: Claude Fable 5.1 --- Cargo.lock | 7 + lib/components/fabro-petri/Cargo.toml | 9 + lib/components/fabro-petri/README.md | 21 +- lib/components/fabro-petri/src/lib.rs | 8 +- lib/components/fabro-petri/src/run_store.rs | 588 ++++++++++++++++++ .../fabro-petri/tests/sqlite_store.rs | 217 +++++++ .../migrations/2026091701_petri_records.sql | 27 + lib/foundation/fabro-db/src/lib.rs | 6 + 8 files changed, 876 insertions(+), 7 deletions(-) create mode 100644 lib/components/fabro-petri/src/run_store.rs create mode 100644 lib/components/fabro-petri/tests/sqlite_store.rs create mode 100644 lib/foundation/fabro-db/migrations/2026091701_petri_records.sql diff --git a/Cargo.lock b/Cargo.lock index cada5b097..70342d49b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2886,6 +2886,10 @@ dependencies = [ name = "fabro-petri" version = "0.357.0-nightly.0" dependencies = [ + "async-trait", + "fabro-db", + "fabro-store", + "fabro-types", "petri-attractor-steps", "petri-execution", "petri-frontend-attractor", @@ -2893,8 +2897,11 @@ dependencies = [ "petri-runtime", "petri-store", "petri-testkit", + "serde_json", + "sqlx", "tempfile", "tokio", + "tracing", ] [[package]] diff --git a/lib/components/fabro-petri/Cargo.toml b/lib/components/fabro-petri/Cargo.toml index 5259e102e..9398c5481 100644 --- a/lib/components/fabro-petri/Cargo.toml +++ b/lib/components/fabro-petri/Cargo.toml @@ -13,14 +13,23 @@ doctest = false workspace = true [dependencies] +fabro-db = { path = "../../foundation/fabro-db" } +fabro-store = { path = "../fabro-store" } +fabro-types = { path = "../../foundation/fabro-types" } petri_runtime.workspace = true petri_execution.workspace = true petri_store.workspace = true petri_attractor_steps.workspace = true petri_frontend_attractor.workspace = true petri_frontend_fabro.workspace = true +async-trait.workspace = true +serde_json.workspace = true +sqlx.workspace = true +tokio.workspace = true +tracing.workspace = true [dev-dependencies] +fabro-store = { path = "../fabro-store", features = ["test-support"] } petri_testkit.workspace = true tempfile = "3" tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } diff --git a/lib/components/fabro-petri/README.md b/lib/components/fabro-petri/README.md index 8e2a0eb43..727afb430 100644 --- a/lib/components/fabro-petri/README.md +++ b/lib/components/fabro-petri/README.md @@ -12,9 +12,15 @@ to this crate and the lockfile, nothing else. ## What it holds -Every adapter the integration plan describes lands here: the run store over -Fabro's SQLite database, then the platform adapters (hooks, interviews, -secrets, output storage, run tools, the event projection). +Every adapter the integration plan describes lands here. + +- `SqliteRunStore`: Petri's `RunStore` and `RunLogs` over Fabro's SQLite + database, so a run's records live in Fabro's tables (`petri_runs` for the + run and its writer lease, `petri_records` for every record of every log, + and the shared `blobs` table). The module docs state the lease and append + rules. +- The platform adapters the plan adds after it: hooks, interviews, secrets, + output storage, run tools, the event projection. ## How it is tested @@ -23,8 +29,13 @@ Integration tests live under `tests/`: - `runs.rs` runs the `hello` bundle in memory through `Runtime::standard()` with the Fabro frontend and the model-free stub registry, then a command-only workflow on the host sandbox through the real step registry. - The sandbox test skips, and says why, when the `sandbox-driver-host` plugin - executable is not on `PATH`. + Both skip, and say why, when the `sandbox-driver-host` plugin executable + is not on `PATH` (every run takes its scope's environment through it); + the sandbox-plugins CI job requires them. +- `sqlite_store.rs` runs Petri's store conformance suite + (`petri_testkit::run_store::conformance`) against `SqliteRunStore`, plus the + operator release, lease exclusivity, a crash between appends, and blob + interoperation with Fabro's `BlobStore`. Run them with: diff --git a/lib/components/fabro-petri/src/lib.rs b/lib/components/fabro-petri/src/lib.rs index f1c1546cd..dd3f77139 100644 --- a/lib/components/fabro-petri/src/lib.rs +++ b/lib/components/fabro-petri/src/lib.rs @@ -8,10 +8,14 @@ //! //! What lives here, as the integration plan lands it: //! -//! - the run store over Fabro's SQLite database, so Petri's records are the -//! run's source of truth in Fabro's tables; +//! - [`SqliteRunStore`]: Petri's run store over Fabro's SQLite database, so a +//! run's records are its source of truth in Fabro's tables; //! - the platform adapters: hooks, interviews, secrets, output storage, the run //! tools, the event projection. //! //! The Petri packages are pinned by revision in the workspace `Cargo.toml` //! under `petri_*` keys. + +pub mod run_store; + +pub use run_store::SqliteRunStore; diff --git a/lib/components/fabro-petri/src/run_store.rs b/lib/components/fabro-petri/src/run_store.rs new file mode 100644 index 000000000..89d55e21d --- /dev/null +++ b/lib/components/fabro-petri/src/run_store.rs @@ -0,0 +1,588 @@ +//! Petri's run store over Fabro's SQLite database: the Petri store crate's +//! `RunStore` and `RunLogs`, implemented on the pool Fabro's other stores +//! share. +//! +//! # Tables +//! +//! - `petri_runs` is a run's existence and its writer lease: one row per run +//! key, with the owner holding the lease and when it took it. This is the +//! lease row the integration plan describes as "a row in `runs`". Petri opens +//! runs by keys of its own, with no Fabro run row behind them, and `runs` has +//! columns only a Fabro run can fill, so the lease lives in a table of its +//! own. The create handler inserts the Fabro `runs` row separately. +//! - `petri_records` holds every record of every log, keyed by `(run_id, log, +//! seq)`: `recorded_at` lifted out for indexing, and the record itself as +//! JSON, stored and read back unchanged. +//! - `blobs` is Fabro's content-addressed blob table, shared with +//! [`BlobStore`]. Petri's digest is the same SHA-256 hex. +//! +//! # The log column +//! +//! `log` is the `LogId` rendered with its `Display`: `coordinator`, +//! `resources`, or `execution ` for execution `n`. [`log_id_text`] and +//! [`parse_log_id`] are the two directions, and a test pins the strings. +//! +//! # The lease +//! +//! `Create` inserts the run row and takes the lease in one statement, and +//! refuses an existing key with `Exists`. `Write` takes the lease of an +//! existing key when nobody holds it or when the same owner holds it (a retry +//! after a lost reply gets the same lease), and refuses a live different +//! owner with `Leased`. `Read` takes no lease and never blocks a writer. A +//! same-owner reopen in this process shares the live handle, so the lease +//! lasts while any handle of the owner does. +//! +//! The lease ends when the last handle drops, by an operator's +//! [`SqliteRunStore::release_lease`], or when the server observes the +//! worker's exit and calls the same method. Never by timeout. Dropping a +//! handle spawns the release on the current Tokio runtime, because sqlx has +//! no synchronous path; the store awaits every spawned release before its +//! next `open`, so a drop followed by an open observes the release. With no +//! runtime at drop, the row stays leased until an operator releases it, and +//! the drop says so in the log. +//! +//! Every write checks, inside its own transaction, that the handle's owner +//! still holds the lease, and fails with `StaleOwner` otherwise. +//! +//! # Appends +//! +//! One `BEGIN IMMEDIATE` transaction per batch. A record at a seq below the +//! log's head must equal the stored record as a JSON value, and is then +//! accepted without a second insert (a lost-reply retry is safe). A different +//! record at a taken seq, or a seq past the head, is `Conflict`, and the +//! batch stores nothing. A committed transaction is durable past a process +//! crash: Fabro's pool runs SQLite in WAL mode with `synchronous = NORMAL`. + +use std::collections::HashMap; +use std::error::Error; +use std::sync::{Arc, Mutex, MutexGuard, PoisonError, Weak}; +use std::time::{SystemTime, UNIX_EPOCH}; +use std::{fmt, mem, ptr}; + +use fabro_db::DbPool; +use fabro_store::BlobStore; +use fabro_types::BlobHash; +use petri_store::{ + Access, Digest, ExecutionId, LogId, OwnerId, Record, RunKey, RunLogs, RunStore, StoreError, +}; +use serde_json::Value; +use sqlx::{Executor, Sqlite}; +use tokio::runtime::Handle; +use tokio::task::JoinHandle; +use tracing::{debug, warn}; + +const EXECUTION_LOG_PREFIX: &str = "execution "; + +/// The `log` column value of a log id: its `Display`. +pub fn log_id_text(log: &LogId) -> String { + log.to_string() +} + +/// The log id a `log` column value names, or `None` when the text is not +/// one [`log_id_text`] produces. +pub fn parse_log_id(text: &str) -> Option { + match text { + "coordinator" => Some(LogId::Coordinator), + "resources" => Some(LogId::Resources), + other => other + .strip_prefix(EXECUTION_LOG_PREFIX)? + .parse::() + .ok() + .map(|id| LogId::Execution(ExecutionId::new(id))), + } +} + +/// Petri's run store over Fabro's SQLite database. +pub struct SqliteRunStore { + shared: Arc, +} + +impl fmt::Debug for SqliteRunStore { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("SqliteRunStore") + .field("database", &self.shared.database) + .finish_non_exhaustive() + } +} + +/// What the store and every handle it opens share. +struct Shared { + pool: DbPool, + blobs: BlobStore, + /// The database file, for locators. + database: String, + /// The writer handle alive in this process per run, so a same-owner + /// reopen shares it and the lease lasts while any handle does. + live: Mutex>>, + /// The releases dropped handles spawned, awaited before the next open. + releases: Mutex>>, +} + +impl SqliteRunStore { + /// A store over a pool whose migrations have run. + #[must_use] + pub fn new(pool: DbPool) -> Self { + let database = pool.connect_options().get_filename().display().to_string(); + Self { + shared: Arc::new(Shared { + blobs: BlobStore::new(pool.clone()), + pool, + database, + live: Mutex::default(), + releases: Mutex::default(), + }), + } + } + + /// End the writer lease of `key` from outside, as an operator does, or + /// as the server does when it observes the worker that held it exit. + /// The holder's handles turn stale, and the next `Write` open takes the + /// run. `NotFound` when the store does not hold the key. + pub async fn release_lease(&self, key: &RunKey) -> Result<(), StoreError> { + self.shared.drain_releases().await; + let result = sqlx::query( + "UPDATE petri_runs SET owner_id = NULL, acquired_at_ms = NULL WHERE run_id = ?", + ) + .bind(key.as_str()) + .execute(&self.shared.pool) + .await + .map_err(|cause| self.shared.backend(key, "release the run's lease", cause))?; + if result.rows_affected() == 0 { + return Err(self.shared.not_found(key)); + } + lock(&self.shared.live).remove(key); + debug!(run_id = %key, "Petri run lease released from outside"); + Ok(()) + } + + /// The owner holding the writer lease of `key`, if any. `NotFound` when + /// the store does not hold the key. + pub async fn owner(&self, key: &RunKey) -> Result, StoreError> { + self.shared.drain_releases().await; + let holder: Option> = + sqlx::query_scalar("SELECT owner_id FROM petri_runs WHERE run_id = ?") + .bind(key.as_str()) + .fetch_optional(&self.shared.pool) + .await + .map_err(|cause| self.shared.backend(key, "read the run's lease", cause))?; + match holder { + None => Err(self.shared.not_found(key)), + Some(holder) => Ok(holder.map(OwnerId::new)), + } + } + + /// The writer handle for `owner`, once the lease is taken: the live one + /// when this owner already holds a handle here, else a new one. + fn writer(&self, key: &RunKey, owner: OwnerId) -> Arc { + let mut live = lock(&self.shared.live); + if let Some(handle) = live.get(key).and_then(Weak::upgrade) { + if handle.owner.as_ref() == Some(&owner) { + return handle; + } + } + let handle = Arc::new(SqliteRunLogs { + shared: self.shared.clone(), + key: key.clone(), + owner: Some(owner), + }); + live.insert(key.clone(), Arc::downgrade(&handle)); + handle + } +} + +impl Shared { + fn locator(&self, key: &RunKey) -> String { + format!("sqlite database {}, run `{key}`", self.database) + } + + fn backend( + &self, + key: &RunKey, + action: &'static str, + cause: impl Into>, + ) -> StoreError { + StoreError::backend(self.locator(key), action, cause) + } + + fn not_found(&self, key: &RunKey) -> StoreError { + StoreError::NotFound { + key: key.clone(), + locator: self.locator(key), + } + } + + /// Await every release a dropped handle spawned, so what follows sees + /// the lease as the drops left it. + async fn drain_releases(&self) { + let pending = mem::take(&mut *lock(&self.releases)); + for release in pending { + // A release task never panics: it reports its own failure. + let _ = release.await; + } + } + + /// Take the lease of an existing run for `owner`, or share it when the + /// same owner holds it. + async fn take_lease(&self, key: &RunKey, owner: &OwnerId) -> Result<(), StoreError> { + let mut tx = self + .pool + .begin_with("BEGIN IMMEDIATE") + .await + .map_err(|cause| self.backend(key, "take the run's lease", cause))?; + let holder: Option> = + sqlx::query_scalar("SELECT owner_id FROM petri_runs WHERE run_id = ?") + .bind(key.as_str()) + .fetch_optional(&mut *tx) + .await + .map_err(|cause| self.backend(key, "take the run's lease", cause))?; + match holder { + None => return Err(self.not_found(key)), + Some(Some(holder)) if holder != owner.as_str() => { + return Err(StoreError::Leased { + locator: self.locator(key), + owner: OwnerId::new(holder), + }); + } + Some(Some(_)) => { + debug!(run_id = %key, owner = %owner, "Petri run lease shared with its holder"); + } + Some(None) => { + sqlx::query( + "UPDATE petri_runs SET owner_id = ?, acquired_at_ms = ? WHERE run_id = ?", + ) + .bind(owner.as_str()) + .bind(now_ms()) + .bind(key.as_str()) + .execute(&mut *tx) + .await + .map_err(|cause| self.backend(key, "take the run's lease", cause))?; + debug!(run_id = %key, owner = %owner, "Petri run lease taken"); + } + } + tx.commit() + .await + .map_err(|cause| self.backend(key, "take the run's lease", cause)) + } + + /// Whether `owner` still holds the lease of `key`, read through + /// `executor` so a write's check sits in the write's own transaction. + async fn check_owner<'c, E>( + &self, + executor: E, + key: &RunKey, + owner: &OwnerId, + ) -> Result<(), StoreError> + where + E: Executor<'c, Database = Sqlite>, + { + let holder: Option> = + sqlx::query_scalar("SELECT owner_id FROM petri_runs WHERE run_id = ?") + .bind(key.as_str()) + .fetch_optional(executor) + .await + .map_err(|cause| self.backend(key, "check the run's lease", cause))?; + match holder { + Some(Some(holder)) if holder == owner.as_str() => Ok(()), + _ => Err(StoreError::StaleOwner), + } + } + + /// End the lease of `key` when `owner` still holds it: what a dropped + /// handle does. A lease that already moved is left alone. + async fn release_owner(&self, key: &RunKey, owner: &OwnerId) { + let released = sqlx::query( + "UPDATE petri_runs SET owner_id = NULL, acquired_at_ms = NULL \ + WHERE run_id = ? AND owner_id = ?", + ) + .bind(key.as_str()) + .bind(owner.as_str()) + .execute(&self.pool) + .await; + match released { + Ok(result) if result.rows_affected() == 1 => { + debug!(run_id = %key, owner = %owner, "Petri run lease released at drop"); + } + Ok(_) => { + debug!(run_id = %key, owner = %owner, "Petri run lease had already moved at drop"); + } + Err(error) => { + warn!( + run_id = %key, + owner = %owner, + error = %error, + "Petri run lease not released at drop; release it from outside" + ); + } + } + } +} + +#[async_trait::async_trait] +impl RunStore for SqliteRunStore { + async fn open(&self, key: &RunKey, access: Access) -> Result, StoreError> { + let shared = &self.shared; + shared.drain_releases().await; + match access { + Access::Create { owner } => { + let now = now_ms(); + let result = sqlx::query( + "INSERT INTO petri_runs (run_id, created_at_ms, owner_id, acquired_at_ms) \ + VALUES (?, ?, ?, ?) ON CONFLICT(run_id) DO NOTHING", + ) + .bind(key.as_str()) + .bind(now) + .bind(owner.as_str()) + .bind(now) + .execute(&shared.pool) + .await + .map_err(|cause| shared.backend(key, "create the run", cause))?; + if result.rows_affected() == 0 { + return Err(StoreError::Exists { + key: key.clone(), + locator: shared.locator(key), + }); + } + debug!(run_id = %key, owner = %owner, "Petri run created"); + Ok(self.writer(key, owner)) + } + Access::Write { owner } => { + shared.take_lease(key, &owner).await?; + Ok(self.writer(key, owner)) + } + Access::Read => { + let exists: bool = + sqlx::query_scalar("SELECT EXISTS(SELECT 1 FROM petri_runs WHERE run_id = ?)") + .bind(key.as_str()) + .fetch_one(&shared.pool) + .await + .map_err(|cause| shared.backend(key, "open the run", cause))?; + if !exists { + return Err(shared.not_found(key)); + } + Ok(Arc::new(SqliteRunLogs { + shared: shared.clone(), + key: key.clone(), + owner: None, + })) + } + } + } +} + +/// One run in the database, opened. A writer handle carries the owner it +/// was opened with; a reader handle refuses every mutation. +struct SqliteRunLogs { + shared: Arc, + key: RunKey, + owner: Option, +} + +impl SqliteRunLogs { + fn owner(&self) -> Result<&OwnerId, StoreError> { + self.owner.as_ref().ok_or(StoreError::ReadOnly) + } + + fn backend( + &self, + action: &'static str, + cause: impl Into>, + ) -> StoreError { + self.shared.backend(&self.key, action, cause) + } + + /// A stored `record_json` back into the record it was. + fn decode(&self, json: &str) -> Result { + let value: Value = serde_json::from_str(json) + .map_err(|cause| self.backend("decode a stored record", cause))?; + Record::from_value(value).map_err(|cause| self.backend("decode a stored record", cause)) + } + + /// A seq as SQLite stores it. + fn column(&self, value: u64) -> Result { + i64::try_from(value).map_err(|cause| self.backend("encode a record", cause)) + } +} + +impl Drop for SqliteRunLogs { + fn drop(&mut self) { + let Some(owner) = self.owner.clone() else { + return; + }; + { + let mut live = lock(&self.shared.live); + let this: *const Self = self; + if live + .get(&self.key) + .is_some_and(|weak| ptr::eq(weak.as_ptr(), this)) + { + live.remove(&self.key); + } + } + match Handle::try_current() { + Ok(runtime) => { + let shared = self.shared.clone(); + let key = self.key.clone(); + let release = runtime.spawn(async move { + shared.release_owner(&key, &owner).await; + }); + lock(&self.shared.releases).push(release); + } + Err(_) => { + warn!( + run_id = %self.key, + owner = %owner, + "Petri run lease not released at drop: no async runtime; release it from outside" + ); + } + } + } +} + +#[async_trait::async_trait] +impl RunLogs for SqliteRunLogs { + fn locator(&self) -> String { + self.shared.locator(&self.key) + } + + async fn append(&self, log: &LogId, records: &[Record]) -> Result<(), StoreError> { + let owner = self.owner()?; + let mut tx = self + .shared + .pool + .begin_with("BEGIN IMMEDIATE") + .await + .map_err(|cause| self.backend("begin an append", cause))?; + self.shared.check_owner(&mut *tx, &self.key, owner).await?; + let log_text = log_id_text(log); + let head: i64 = sqlx::query_scalar( + "SELECT COALESCE(MAX(seq) + 1, 0) FROM petri_records WHERE run_id = ? AND log = ?", + ) + .bind(self.key.as_str()) + .bind(&log_text) + .fetch_one(&mut *tx) + .await + .map_err(|cause| self.backend("read the log's head", cause))?; + let mut next = + u64::try_from(head).map_err(|cause| self.backend("read the log's head", cause))?; + for record in records { + let conflict = || StoreError::Conflict { + log: *log, + seq: record.seq, + }; + if record.seq < next { + let stored: Option = sqlx::query_scalar( + "SELECT record_json FROM petri_records WHERE run_id = ? AND log = ? AND seq = ?", + ) + .bind(self.key.as_str()) + .bind(&log_text) + .bind(self.column(record.seq)?) + .fetch_optional(&mut *tx) + .await + .map_err(|cause| self.backend("read a stored record", cause))?; + let same = match stored { + Some(json) => self.decode(&json)? == *record, + None => false, + }; + if same { + continue; + } + return Err(conflict()); + } + if record.seq != next { + return Err(conflict()); + } + let json = serde_json::to_string(&record.record) + .map_err(|cause| self.backend("encode a record", cause))?; + sqlx::query( + "INSERT INTO petri_records (run_id, log, seq, recorded_at, record_json) \ + VALUES (?, ?, ?, ?, ?)", + ) + .bind(self.key.as_str()) + .bind(&log_text) + .bind(self.column(record.seq)?) + .bind(self.column(record.recorded_at)?) + .bind(json) + .execute(&mut *tx) + .await + .map_err(|cause| self.backend("append a record", cause))?; + next += 1; + } + tx.commit() + .await + .map_err(|cause| self.backend("commit an append", cause)) + } + + async fn read(&self, log: &LogId) -> Result, StoreError> { + let rows: Vec = sqlx::query_scalar( + "SELECT record_json FROM petri_records WHERE run_id = ? AND log = ? ORDER BY seq", + ) + .bind(self.key.as_str()) + .bind(log_id_text(log)) + .fetch_all(&self.shared.pool) + .await + .map_err(|cause| self.backend("read a log", cause))?; + rows.iter().map(|json| self.decode(json)).collect() + } + + async fn put_blob(&self, bytes: &[u8]) -> Result { + let owner = self.owner()?; + self.shared + .check_owner(&self.shared.pool, &self.key, owner) + .await?; + self.shared + .blobs + .write(bytes) + .await + .map_err(|cause| self.backend("store a blob", cause))?; + Ok(Digest::of(bytes)) + } + + async fn get_blob(&self, digest: Digest) -> Result>, StoreError> { + let hash: BlobHash = digest + .to_hex() + .parse() + .map_err(|cause| self.backend("read a blob", cause))?; + let bytes = self + .shared + .blobs + .read(&hash) + .await + .map_err(|cause| self.backend("read a blob", cause))?; + Ok(bytes.map(|bytes| bytes.to_vec())) + } +} + +fn lock(mutex: &Mutex) -> MutexGuard<'_, T> { + mutex.lock().unwrap_or_else(PoisonError::into_inner) +} + +/// Milliseconds since the Unix epoch, as SQLite stores them. +fn now_ms() -> i64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .ok() + .and_then(|elapsed| i64::try_from(elapsed.as_millis()).ok()) + .unwrap_or(0) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn log_ids_round_trip_through_their_text() { + let logs = [ + (LogId::Coordinator, "coordinator"), + (LogId::Resources, "resources"), + (LogId::Execution(ExecutionId::new(0)), "execution 0"), + (LogId::Execution(ExecutionId::new(42)), "execution 42"), + ]; + for (log, text) in logs { + assert_eq!(log_id_text(&log), text); + assert_eq!(parse_log_id(text), Some(log)); + } + assert_eq!(parse_log_id("execution"), None); + assert_eq!(parse_log_id("execution x"), None); + assert_eq!(parse_log_id("engine 1"), None); + } +} diff --git a/lib/components/fabro-petri/tests/sqlite_store.rs b/lib/components/fabro-petri/tests/sqlite_store.rs new file mode 100644 index 000000000..21512e1b7 --- /dev/null +++ b/lib/components/fabro-petri/tests/sqlite_store.rs @@ -0,0 +1,217 @@ +//! The SQLite run store against Petri's store contract: the conformance +//! suite, the operator release, lease exclusivity, a crash between appends, +//! and blob interoperation with Fabro's own blob store. + +use std::mem; +use std::path::Path; +use std::sync::Arc; + +use fabro_db::Database; +use fabro_petri::SqliteRunStore; +use fabro_store::{BlobStore, test_support}; +use fabro_types::BlobHash; +use petri_store::{Access, Digest, LogId, OwnerId, RunKey, RunStore, StoreError}; +use petri_testkit::run_store::{self, conformance, stale_owner_conformance}; +use tokio::runtime::Handle; +use tokio::task; + +/// A store over a fresh in-memory database with the production blob and +/// Petri record schemas. +fn fresh_in_memory() -> Arc { + Arc::new(SqliteRunStore::new(test_support::in_memory_pool_with(&[ + fabro_db::BLOBS_MIGRATION_SQL, + fabro_db::PETRI_RECORDS_MIGRATION_SQL, + ]))) +} + +/// A migrated database file, as the server opens it. +async fn migrated(path: &Path) -> Database { + let database = Database::connect(path).await.expect("the database opens"); + database.migrate().await.expect("the migrations run"); + database +} + +#[tokio::test] +async fn the_sqlite_store_passes_the_conformance_suite() { + conformance(fresh_in_memory).await; +} + +/// The operator release ends the lease from outside: the old owner's handle +/// turns stale and the next writer takes the run. The release is async, so +/// the suite's synchronous closure blocks on it in place. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn an_operator_release_makes_the_old_owner_stale() { + let dir = tempfile::tempdir().expect("a temp dir"); + let database = migrated(&dir.path().join("fabro.sqlite3")).await; + let store = SqliteRunStore::new(database.clone_pool()); + let release = |key: &RunKey| { + task::block_in_place(|| Handle::current().block_on(store.release_lease(key))) + .expect("the lease releases"); + }; + stale_owner_conformance(&store, release).await; +} + +/// Two owners never hold one run's lease at the same time, in either order, +/// and the store reports who holds it. +#[tokio::test] +async fn two_owners_cannot_both_hold_the_lease() { + let dir = tempfile::tempdir().expect("a temp dir"); + let database = migrated(&dir.path().join("fabro.sqlite3")).await; + let store = SqliteRunStore::new(database.clone_pool()); + let key = RunKey::new("exclusive"); + let first = OwnerId::new("first"); + let second = OwnerId::new("second"); + + let held = store + .open(&key, Access::Create { + owner: first.clone(), + }) + .await + .expect("the first owner creates"); + assert_eq!(store.owner(&key).await.expect("reads"), Some(first.clone())); + let refused = store + .open(&key, Access::Write { + owner: second.clone(), + }) + .await + .err() + .expect("the second owner is refused while the first is live"); + assert!( + matches!(&refused, StoreError::Leased { owner, .. } if *owner == first), + "{refused}" + ); + assert!( + refused.to_string().contains("fabro.sqlite3") + && refused.to_string().contains("`exclusive`"), + "the message names the database and the run: {refused}" + ); + + drop(held); + let taken = store + .open(&key, Access::Write { + owner: second.clone(), + }) + .await + .expect("the second owner takes the run once the first handle drops"); + assert_eq!(store.owner(&key).await.expect("reads"), Some(second)); + let refused = store + .open(&key, Access::Write { owner: first }) + .await + .err() + .expect("the first owner is refused in turn"); + assert!(matches!(refused, StoreError::Leased { .. }), "{refused}"); + drop(taken); + assert_eq!(store.owner(&key).await.expect("reads"), None); +} + +/// A worker that crashes between two appends leaves a readable prefix and a +/// lease that only an operator ends: a second process opens the file, reads +/// the first batch back intact, is refused the lease until it releases it, +/// and then continues the log past the prefix. +#[tokio::test] +async fn a_crash_between_appends_leaves_a_readable_prefix() { + let dir = tempfile::tempdir().expect("a temp dir"); + let path = dir.path().join("fabro.sqlite3"); + let key = RunKey::new("crashed"); + let log = LogId::Execution(petri_store::ExecutionId::new(0)); + let prefix = [ + run_store::record(0, "execution.started"), + run_store::record(1, "step.started"), + ]; + + // The worker's process: its own pool over the file. + let worker = SqliteRunStore::new(migrated(&path).await.clone_pool()); + let handle = worker + .open(&key, Access::Create { + owner: OwnerId::new("worker"), + }) + .await + .expect("the worker creates"); + handle + .append(&log, &prefix) + .await + .expect("the first batch is durable"); + // The crash: the handle never drops, so nothing releases the lease. + mem::forget(handle); + + // The server's process: a second pool over the same file. + let server = SqliteRunStore::new(migrated(&path).await.clone_pool()); + let reader = server + .open(&key, Access::Read) + .await + .expect("a reader never blocks on the lease"); + assert_eq!(reader.read(&log).await.expect("reads"), prefix); + let refused = server + .open(&key, Access::Write { + owner: OwnerId::new("resumer"), + }) + .await + .err() + .expect("the crashed worker's lease does not time out"); + assert!( + matches!(&refused, StoreError::Leased { owner, .. } if owner.as_str() == "worker"), + "{refused}" + ); + + server + .release_lease(&key) + .await + .expect("the server releases the lease it observed the worker lose"); + let resumed = server + .open(&key, Access::Write { + owner: OwnerId::new("resumer"), + }) + .await + .expect("the resumer takes the run"); + assert_eq!(resumed.read(&log).await.expect("reads"), prefix); + let error = resumed + .append(&log, &[run_store::record(0, "different")]) + .await + .expect_err("the prefix cannot be rewritten"); + assert!( + matches!(error, StoreError::Conflict { seq: 0, .. }), + "{error}" + ); + resumed + .append(&log, &[run_store::record(2, "step.finished")]) + .await + .expect("the log continues past the prefix"); + assert_eq!(resumed.read(&log).await.expect("reads").len(), 3); +} + +/// Petri's blobs and Fabro's blob store are one table: a blob either side +/// writes, the other reads by the same SHA-256 hex. +#[tokio::test] +async fn blobs_interoperate_with_the_blob_store() { + let dir = tempfile::tempdir().expect("a temp dir"); + let database = migrated(&dir.path().join("fabro.sqlite3")).await; + let store = SqliteRunStore::new(database.clone_pool()); + let blobs = BlobStore::new(database.clone_pool()); + let logs = store + .open(&RunKey::new("blobs"), Access::Create { + owner: OwnerId::new("owner"), + }) + .await + .expect("creates"); + + let graph = br#"{"nodes":[],"edges":[]}"#; + let digest = logs.put_blob(graph).await.expect("stores"); + assert_eq!(digest.to_hex(), BlobHash::new(graph).to_string()); + let hash: BlobHash = digest.to_hex().parse().expect("the digest is a blob hash"); + assert_eq!( + blobs.read(&hash).await.expect("reads").as_deref(), + Some(graph.as_slice()) + ); + + let output = b"large step output"; + let hash = blobs.write(output).await.expect("stores"); + let digest: Digest = hash.to_string().parse().expect("the blob hash is a digest"); + assert_eq!( + logs.get_blob(digest).await.expect("reads"), + Some(output.to_vec()) + ); + assert_eq!( + logs.get_blob(Digest::of(b"missing")).await.expect("reads"), + None + ); +} diff --git a/lib/foundation/fabro-db/migrations/2026091701_petri_records.sql b/lib/foundation/fabro-db/migrations/2026091701_petri_records.sql new file mode 100644 index 000000000..19ba7eaed --- /dev/null +++ b/lib/foundation/fabro-db/migrations/2026091701_petri_records.sql @@ -0,0 +1,27 @@ +-- Petri's durable run record in Fabro's database. +-- +-- `petri_runs` is a run's existence and its writer lease: one row per run +-- key, with the owner holding the lease and when it took it. Petri opens runs +-- by keys of its own, so this row is separate from the Fabro `runs` summary +-- row the create handler writes. +CREATE TABLE petri_runs ( + run_id TEXT PRIMARY KEY NOT NULL, + created_at_ms INTEGER NOT NULL, + owner_id TEXT NULL, + acquired_at_ms INTEGER NULL +); + +-- Every record of every log of a run, keyed by (run, log, seq). `log` is the +-- Petri log id as text (`coordinator`, `resources`, `execution `), +-- `recorded_at` is lifted out of the record for indexing, and `record_json` +-- is the record itself, stored and read back unchanged. Blobs share the +-- `blobs` table. +CREATE TABLE petri_records ( + run_id TEXT NOT NULL, + log TEXT NOT NULL, + seq INTEGER NOT NULL, + recorded_at INTEGER NOT NULL, + record_json TEXT NOT NULL, + PRIMARY KEY (run_id, log, seq), + CHECK (json_valid(record_json)) +); diff --git a/lib/foundation/fabro-db/src/lib.rs b/lib/foundation/fabro-db/src/lib.rs index e44ec0ace..a6eda1e57 100644 --- a/lib/foundation/fabro-db/src/lib.rs +++ b/lib/foundation/fabro-db/src/lib.rs @@ -41,6 +41,12 @@ pub const RUN_EVENT_SESSION_OWNER_MIGRATION_SQL: &str = pub const RUN_SESSION_RECORDS_MIGRATION_SQL: &str = include_str!("../migrations/2026091101_run_session_records.sql"); +/// The Petri run record migration (`petri_runs`, `petri_records`), exposed +/// so fixtures in other crates can install the production schema without a +/// filesystem path into this crate. +pub const PETRI_RECORDS_MIGRATION_SQL: &str = + include_str!("../migrations/2026091701_petri_records.sql"); + /// The temporary run-history activation migration, exposed so fixtures in /// other crates can install the production compatibility schema. pub const RUN_HISTORY_ACTIVATION_MIGRATION_SQL: &str =