mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-01 02:04:24 +00:00
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 <noreply@anthropic.com>
This commit is contained in:
parent
7eb5ca502c
commit
3ffe7e00ce
8 changed files with 876 additions and 7 deletions
7
Cargo.lock
generated
7
Cargo.lock
generated
|
|
@ -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]]
|
||||
|
|
|
|||
|
|
@ -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"] }
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
588
lib/components/fabro-petri/src/run_store.rs
Normal file
588
lib/components/fabro-petri/src/run_store.rs
Normal file
|
|
@ -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 <n>` 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<LogId> {
|
||||
match text {
|
||||
"coordinator" => Some(LogId::Coordinator),
|
||||
"resources" => Some(LogId::Resources),
|
||||
other => other
|
||||
.strip_prefix(EXECUTION_LOG_PREFIX)?
|
||||
.parse::<u64>()
|
||||
.ok()
|
||||
.map(|id| LogId::Execution(ExecutionId::new(id))),
|
||||
}
|
||||
}
|
||||
|
||||
/// Petri's run store over Fabro's SQLite database.
|
||||
pub struct SqliteRunStore {
|
||||
shared: Arc<Shared>,
|
||||
}
|
||||
|
||||
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<HashMap<RunKey, Weak<SqliteRunLogs>>>,
|
||||
/// The releases dropped handles spawned, awaited before the next open.
|
||||
releases: Mutex<Vec<JoinHandle<()>>>,
|
||||
}
|
||||
|
||||
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<Option<OwnerId>, StoreError> {
|
||||
self.shared.drain_releases().await;
|
||||
let holder: Option<Option<String>> =
|
||||
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<dyn RunLogs> {
|
||||
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<Box<dyn Error + Send + Sync>>,
|
||||
) -> 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<Option<String>> =
|
||||
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<Option<String>> =
|
||||
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<Arc<dyn RunLogs>, 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<Shared>,
|
||||
key: RunKey,
|
||||
owner: Option<OwnerId>,
|
||||
}
|
||||
|
||||
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<Box<dyn Error + Send + Sync>>,
|
||||
) -> StoreError {
|
||||
self.shared.backend(&self.key, action, cause)
|
||||
}
|
||||
|
||||
/// A stored `record_json` back into the record it was.
|
||||
fn decode(&self, json: &str) -> Result<Record, StoreError> {
|
||||
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, StoreError> {
|
||||
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<String> = 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<Vec<Record>, StoreError> {
|
||||
let rows: Vec<String> = 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<Digest, StoreError> {
|
||||
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<Option<Vec<u8>>, 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<T>(mutex: &Mutex<T>) -> 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);
|
||||
}
|
||||
}
|
||||
217
lib/components/fabro-petri/tests/sqlite_store.rs
Normal file
217
lib/components/fabro-petri/tests/sqlite_store.rs
Normal file
|
|
@ -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<dyn RunStore> {
|
||||
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
|
||||
);
|
||||
}
|
||||
|
|
@ -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 <n>`),
|
||||
-- `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))
|
||||
);
|
||||
|
|
@ -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 =
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue