Merge branch 'petri-followup-projection' into petri-integration

This commit is contained in:
Bryan Helmkamp 2026-09-18 21:23:12 -04:00
commit a4ffba0d0f
No known key found for this signature in database
7 changed files with 548 additions and 98 deletions

52
Cargo.lock generated
View file

@ -1873,7 +1873,7 @@ dependencies = [
"libc",
"option-ext",
"redox_users",
"windows-sys 0.61.2",
"windows-sys 0.59.0",
]
[[package]]
@ -1987,7 +1987,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb"
dependencies = [
"libc",
"windows-sys 0.61.2",
"windows-sys 0.59.0",
]
[[package]]
@ -3893,7 +3893,7 @@ dependencies = [
"js-sys",
"log",
"wasm-bindgen",
"windows-core 0.61.2",
"windows-core 0.62.2",
]
[[package]]
@ -4735,7 +4735,7 @@ version = "0.50.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5"
dependencies = [
"windows-sys 0.61.2",
"windows-sys 0.59.0",
]
[[package]]
@ -5324,7 +5324,7 @@ checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220"
[[package]]
name = "petri-attractor-steps"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?rev=a5906f6554b94f608ede3902daab106d5f8062c3#a5906f6554b94f608ede3902daab106d5f8062c3"
source = "git+https://github.com/lithoscomputer/petri.git?rev=c216cf226f2a6d072f9e82040d1fab5e8c7f3091#c216cf226f2a6d072f9e82040d1fab5e8c7f3091"
dependencies = [
"async-trait",
"globset",
@ -5355,7 +5355,7 @@ dependencies = [
[[package]]
name = "petri-driver"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?rev=a5906f6554b94f608ede3902daab106d5f8062c3#a5906f6554b94f608ede3902daab106d5f8062c3"
source = "git+https://github.com/lithoscomputer/petri.git?rev=c216cf226f2a6d072f9e82040d1fab5e8c7f3091#c216cf226f2a6d072f9e82040d1fab5e8c7f3091"
dependencies = [
"async-trait",
"getrandom 0.3.4",
@ -5375,7 +5375,7 @@ dependencies = [
[[package]]
name = "petri-engine"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?rev=a5906f6554b94f608ede3902daab106d5f8062c3#a5906f6554b94f608ede3902daab106d5f8062c3"
source = "git+https://github.com/lithoscomputer/petri.git?rev=c216cf226f2a6d072f9e82040d1fab5e8c7f3091#c216cf226f2a6d072f9e82040d1fab5e8c7f3091"
dependencies = [
"petri-ir",
"serde",
@ -5387,7 +5387,7 @@ dependencies = [
[[package]]
name = "petri-execution"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?rev=a5906f6554b94f608ede3902daab106d5f8062c3#a5906f6554b94f608ede3902daab106d5f8062c3"
source = "git+https://github.com/lithoscomputer/petri.git?rev=c216cf226f2a6d072f9e82040d1fab5e8c7f3091#c216cf226f2a6d072f9e82040d1fab5e8c7f3091"
dependencies = [
"async-trait",
"petri-driver",
@ -5411,7 +5411,7 @@ dependencies = [
[[package]]
name = "petri-executor"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?rev=a5906f6554b94f608ede3902daab106d5f8062c3#a5906f6554b94f608ede3902daab106d5f8062c3"
source = "git+https://github.com/lithoscomputer/petri.git?rev=c216cf226f2a6d072f9e82040d1fab5e8c7f3091#c216cf226f2a6d072f9e82040d1fab5e8c7f3091"
dependencies = [
"async-trait",
"libc",
@ -5426,7 +5426,7 @@ dependencies = [
[[package]]
name = "petri-executor-sandbox"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?rev=a5906f6554b94f608ede3902daab106d5f8062c3#a5906f6554b94f608ede3902daab106d5f8062c3"
source = "git+https://github.com/lithoscomputer/petri.git?rev=c216cf226f2a6d072f9e82040d1fab5e8c7f3091#c216cf226f2a6d072f9e82040d1fab5e8c7f3091"
dependencies = [
"async-trait",
"petri-executor",
@ -5448,7 +5448,7 @@ dependencies = [
[[package]]
name = "petri-frontend"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?rev=a5906f6554b94f608ede3902daab106d5f8062c3#a5906f6554b94f608ede3902daab106d5f8062c3"
source = "git+https://github.com/lithoscomputer/petri.git?rev=c216cf226f2a6d072f9e82040d1fab5e8c7f3091#c216cf226f2a6d072f9e82040d1fab5e8c7f3091"
dependencies = [
"marked-yaml",
"petri-ir",
@ -5462,7 +5462,7 @@ dependencies = [
[[package]]
name = "petri-frontend-attractor"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?rev=a5906f6554b94f608ede3902daab106d5f8062c3#a5906f6554b94f608ede3902daab106d5f8062c3"
source = "git+https://github.com/lithoscomputer/petri.git?rev=c216cf226f2a6d072f9e82040d1fab5e8c7f3091#c216cf226f2a6d072f9e82040d1fab5e8c7f3091"
dependencies = [
"minijinja",
"petri-frontend",
@ -5479,7 +5479,7 @@ dependencies = [
[[package]]
name = "petri-frontend-fabro"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?rev=a5906f6554b94f608ede3902daab106d5f8062c3#a5906f6554b94f608ede3902daab106d5f8062c3"
source = "git+https://github.com/lithoscomputer/petri.git?rev=c216cf226f2a6d072f9e82040d1fab5e8c7f3091#c216cf226f2a6d072f9e82040d1fab5e8c7f3091"
dependencies = [
"petri-frontend",
"petri-frontend-attractor",
@ -5495,7 +5495,7 @@ dependencies = [
[[package]]
name = "petri-frontend-native"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?rev=a5906f6554b94f608ede3902daab106d5f8062c3#a5906f6554b94f608ede3902daab106d5f8062c3"
source = "git+https://github.com/lithoscomputer/petri.git?rev=c216cf226f2a6d072f9e82040d1fab5e8c7f3091#c216cf226f2a6d072f9e82040d1fab5e8c7f3091"
dependencies = [
"petri-frontend",
"petri-ir",
@ -5506,7 +5506,7 @@ dependencies = [
[[package]]
name = "petri-ir"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?rev=a5906f6554b94f608ede3902daab106d5f8062c3#a5906f6554b94f608ede3902daab106d5f8062c3"
source = "git+https://github.com/lithoscomputer/petri.git?rev=c216cf226f2a6d072f9e82040d1fab5e8c7f3091#c216cf226f2a6d072f9e82040d1fab5e8c7f3091"
dependencies = [
"regex",
"serde",
@ -5519,7 +5519,7 @@ dependencies = [
[[package]]
name = "petri-runtime"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?rev=a5906f6554b94f608ede3902daab106d5f8062c3#a5906f6554b94f608ede3902daab106d5f8062c3"
source = "git+https://github.com/lithoscomputer/petri.git?rev=c216cf226f2a6d072f9e82040d1fab5e8c7f3091#c216cf226f2a6d072f9e82040d1fab5e8c7f3091"
dependencies = [
"async-trait",
"petri-driver",
@ -5540,7 +5540,7 @@ dependencies = [
[[package]]
name = "petri-steps"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?rev=a5906f6554b94f608ede3902daab106d5f8062c3#a5906f6554b94f608ede3902daab106d5f8062c3"
source = "git+https://github.com/lithoscomputer/petri.git?rev=c216cf226f2a6d072f9e82040d1fab5e8c7f3091#c216cf226f2a6d072f9e82040d1fab5e8c7f3091"
dependencies = [
"async-trait",
"petri-executor",
@ -5556,7 +5556,7 @@ dependencies = [
[[package]]
name = "petri-store"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?rev=a5906f6554b94f608ede3902daab106d5f8062c3#a5906f6554b94f608ede3902daab106d5f8062c3"
source = "git+https://github.com/lithoscomputer/petri.git?rev=c216cf226f2a6d072f9e82040d1fab5e8c7f3091#c216cf226f2a6d072f9e82040d1fab5e8c7f3091"
dependencies = [
"async-trait",
"getrandom 0.3.4",
@ -5571,7 +5571,7 @@ dependencies = [
[[package]]
name = "petri-testkit"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/petri.git?rev=a5906f6554b94f608ede3902daab106d5f8062c3#a5906f6554b94f608ede3902daab106d5f8062c3"
source = "git+https://github.com/lithoscomputer/petri.git?rev=c216cf226f2a6d072f9e82040d1fab5e8c7f3091#c216cf226f2a6d072f9e82040d1fab5e8c7f3091"
dependencies = [
"async-trait",
"petri-driver",
@ -5906,7 +5906,7 @@ dependencies = [
"once_cell",
"socket2",
"tracing",
"windows-sys 0.60.2",
"windows-sys 0.59.0",
]
[[package]]
@ -6354,7 +6354,7 @@ dependencies = [
"errno 0.3.14",
"libc",
"linux-raw-sys",
"windows-sys 0.61.2",
"windows-sys 0.59.0",
]
[[package]]
@ -6413,7 +6413,7 @@ dependencies = [
"security-framework",
"security-framework-sys",
"webpki-root-certs",
"windows-sys 0.61.2",
"windows-sys 0.59.0",
]
[[package]]
@ -7060,7 +7060,7 @@ version = "1.4.8"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c4db69cba1110affc0e9f7bcd48bbf87b3f4fc7c61fc9155afd4c469eb3d6c1b"
dependencies = [
"errno 0.2.8",
"errno 0.3.14",
"libc",
]
@ -7538,7 +7538,7 @@ dependencies = [
"getrandom 0.4.1",
"once_cell",
"rustix",
"windows-sys 0.61.2",
"windows-sys 0.59.0",
]
[[package]]
@ -7573,7 +7573,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "230a1b821ccbd75b185820a1f1ff7b14d21da1e442e22c0863ea5f08771a8874"
dependencies = [
"rustix",
"windows-sys 0.61.2",
"windows-sys 0.59.0",
]
[[package]]
@ -8620,7 +8620,7 @@ version = "0.1.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22"
dependencies = [
"windows-sys 0.61.2",
"windows-sys 0.59.0",
]
[[package]]

View file

@ -132,13 +132,13 @@ pebble-cli-core = { git = "https://github.com/lithoscomputer/pebble", rev = "a39
# lithos-llm and sandbox-driver revisions as this file, so the workspace links
# one copy of each. Only `fabro-petri` may depend on these packages; the keys
# carry the `petri_` prefix so the crate names say where they come from.
petri_runtime = { git = "https://github.com/lithoscomputer/petri.git", rev = "a5906f6554b94f608ede3902daab106d5f8062c3", package = "petri-runtime" }
petri_execution = { git = "https://github.com/lithoscomputer/petri.git", rev = "a5906f6554b94f608ede3902daab106d5f8062c3", package = "petri-execution" }
petri_store = { git = "https://github.com/lithoscomputer/petri.git", rev = "a5906f6554b94f608ede3902daab106d5f8062c3", package = "petri-store" }
petri_attractor_steps = { git = "https://github.com/lithoscomputer/petri.git", rev = "a5906f6554b94f608ede3902daab106d5f8062c3", package = "petri-attractor-steps" }
petri_frontend_attractor = { git = "https://github.com/lithoscomputer/petri.git", rev = "a5906f6554b94f608ede3902daab106d5f8062c3", package = "petri-frontend-attractor" }
petri_frontend_fabro = { git = "https://github.com/lithoscomputer/petri.git", rev = "a5906f6554b94f608ede3902daab106d5f8062c3", package = "petri-frontend-fabro" }
petri_testkit = { git = "https://github.com/lithoscomputer/petri.git", rev = "a5906f6554b94f608ede3902daab106d5f8062c3", package = "petri-testkit" }
petri_runtime = { git = "https://github.com/lithoscomputer/petri.git", rev = "c216cf226f2a6d072f9e82040d1fab5e8c7f3091", package = "petri-runtime" }
petri_execution = { git = "https://github.com/lithoscomputer/petri.git", rev = "c216cf226f2a6d072f9e82040d1fab5e8c7f3091", package = "petri-execution" }
petri_store = { git = "https://github.com/lithoscomputer/petri.git", rev = "c216cf226f2a6d072f9e82040d1fab5e8c7f3091", package = "petri-store" }
petri_attractor_steps = { git = "https://github.com/lithoscomputer/petri.git", rev = "c216cf226f2a6d072f9e82040d1fab5e8c7f3091", package = "petri-attractor-steps" }
petri_frontend_attractor = { git = "https://github.com/lithoscomputer/petri.git", rev = "c216cf226f2a6d072f9e82040d1fab5e8c7f3091", package = "petri-frontend-attractor" }
petri_frontend_fabro = { git = "https://github.com/lithoscomputer/petri.git", rev = "c216cf226f2a6d072f9e82040d1fab5e8c7f3091", package = "petri-frontend-fabro" }
petri_testkit = { git = "https://github.com/lithoscomputer/petri.git", rev = "c216cf226f2a6d072f9e82040d1fab5e8c7f3091", package = "petri-testkit" }
sentry = { version = "0.35", default-features = false, features = ["backtrace", "contexts", "ureq", "rustls"] }
fork = "0.2"
exec = "0.3"

View file

@ -12,14 +12,29 @@
//! Petri log, the last platform record consumed, and the delivery sequence
//! (`stream_seq`) it assigned to each item. The view therefore trails a
//! committed record and never leads one. No projection state of Petri's is
//! checkpointed: each pass replays the run through `replay_since`, which
//! rebuilds the engine and invocation state the derivation needs and
//! delivers only the events past the held positions.
//! checkpointed: a pass derives the events past the held positions from
//! the records alone, and every stored view equals a full replay
//! (`replay_run`) of the records it holds.
//!
//! A pass that finds new platform records committed between its read and
//! its write leaves the view alone and runs again, so the `runs` row never
//! moves backwards behind a concurrent lifecycle write.
//!
//! # The live run's cache
//!
//! A pass keeps in memory, per live run, Petri's replay of the run (a
//! `RunReplay`: the coordinator state, each execution's engine state, the
//! projection) and the view as the pass last committed it, so the next
//! pass reads and folds only the records past the ones the view holds and
//! costs the new records, not the run's length. The cache is never a
//! source of facts and never checkpointed: it is dropped when the run
//! records its finish, after ten idle minutes, when the stored view moves
//! under it, when the run is deleted, and with the process, and the first
//! pass after that rebuilds it by a full replay. A pass that commits
//! nothing (a platform record landed under it, or it failed before its
//! view transaction) keeps the events it derived for the next pass, so
//! nothing is derived twice or lost.
//!
//! # Where it runs
//!
//! In the server. [`Projector::signal`] schedules a pass for a run: the
@ -38,6 +53,8 @@
//! as incomplete with the replay's error; `inspect_run` decides
//! completeness once the run has recorded its finish.
mod cache;
use std::collections::{BTreeMap, BTreeSet, HashMap};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
@ -53,10 +70,11 @@ use petri_execution::{Access, CoordinatorEvent, RunKey, RunStore as _, inspect};
use petri_runtime::engine::Event;
use petri_store::StoreError;
use serde::{Deserialize, Serialize};
use tokio::sync::{Mutex as AsyncMutex, broadcast};
use tokio::sync::broadcast;
use tokio::time;
use tracing::{debug, info, warn};
use self::cache::{Caches, IDLE, RunCache};
use crate::SqliteRunStore;
use crate::projection::{self, FoldState, Item, RecordHealth, RunView};
@ -98,6 +116,11 @@ pub struct PassReport {
pub contended: bool,
pub petri_events: usize,
pub platform_records: usize,
/// How many of the run's records the pass fed through Petri's
/// derivation, before the held positions trimmed their events: the
/// pass's cost. The records past the cache for a live run, the whole
/// run for a pass that rebuilt it.
pub replayed_records: usize,
/// The last delivery sequence the view holds.
pub stream_seq: u64,
pub positions: Positions,
@ -129,6 +152,7 @@ pub enum ProjectError {
}
/// The stored view of a run, as the projection tables hold it.
#[derive(Clone)]
struct StoredView {
view: RunView,
positions: Positions,
@ -148,21 +172,22 @@ struct Slot {
/// the server both are the one database; a test may hand it the run
/// summary store's own pool for the views.
pub struct Projector {
records: DbPool,
pool: DbPool,
store: SqliteRunStore,
platform: PlatformRecordStore,
slots: Mutex<HashMap<RunId, Slot>>,
/// One pass at a time per run: a signalled pass and the startup pass
/// over the same run never interleave their reads and writes.
passes: Mutex<HashMap<RunId, Arc<AsyncMutex<()>>>>,
records: DbPool,
pool: DbPool,
store: SqliteRunStore,
platform: PlatformRecordStore,
slots: Mutex<HashMap<RunId, Slot>>,
/// One pass at a time per run (a signalled pass and the startup pass
/// over the same run never interleave their reads and writes), and the
/// cache each live run's passes continue from.
pub(crate) caches: Caches,
/// Test-only: stop the next pass after its reads, before its view
/// transaction, as a crash there would.
fault: AtomicBool,
fault: AtomicBool,
/// Sent after each committed pass that wrote stream rows: the run whose
/// stream grew. A wake-up for the stream's readers, never a source of
/// facts; a reader that lags re-reads from its cursor.
committed: broadcast::Sender<RunId>,
committed: broadcast::Sender<RunId>,
}
impl std::fmt::Debug for Projector {
@ -183,7 +208,7 @@ impl Projector {
records,
pool: views,
slots: Mutex::default(),
passes: Mutex::default(),
caches: Caches::default(),
fault: AtomicBool::new(false),
committed: broadcast::channel(COMMIT_SIGNAL_CAPACITY).0,
})
@ -214,6 +239,11 @@ impl Projector {
/// and its stream. The caller has ended the run's worker, so no writer
/// holds the lease.
pub async fn delete_run(&self, run_id: RunId) -> Result<(), ProjectError> {
// Under the run's pass lock: no pass reads the rows being deleted,
// and no cache outlives them.
let pass = self.caches.pass_of(run_id);
let mut cache = pass.lock().await;
*cache = None;
let id = run_id.to_string();
let mut views = self.pool.begin().await.map_err(ProjectError::Database)?;
for delete in [
@ -370,9 +400,34 @@ impl Projector {
/// One view pass for the run. Passes over one run run one at a time.
pub async fn project_run(&self, run_id: RunId) -> Result<PassReport, ProjectError> {
let pass = Arc::clone(lock(&self.passes).entry(run_id).or_default());
let _one_at_a_time = pass.lock().await;
let stored = self.load_view(&run_id).await?;
self.caches.sweep(IDLE);
let pass = self.caches.pass_of(run_id);
let mut slot = pass.lock().await;
// The view tables are the source of truth: a cache that no longer
// describes them (another projector committed a pass) is dropped.
let (positions, stream_seq) = stored_positions(&self.pool, run_id)
.await?
.unwrap_or_default();
let mut run = match slot.take() {
Some(cache) if cache.matches(&positions, stream_seq) => cache,
Some(_) => {
debug!(run_id = %run_id, "the stored view moved under the run's cache; rebuilding it");
RunCache::over(self.load_view(&run_id).await?)
}
None => RunCache::over(self.load_view(&run_id).await?),
};
let report = self.pass(run_id, &mut run).await;
// A finished run's records are complete: its cache is dropped, and
// the passes its late platform records take rebuild the view whole.
if !run.view.view.state.finished_run() {
*slot = Some(run);
}
report
}
/// The pass over the run's cache: read what is committed past the
/// positions the cache's view holds, fold it, and write the view.
async fn pass(&self, run_id: RunId, run: &mut RunCache) -> Result<PassReport, ProjectError> {
let key = RunKey::new(run_id.to_string());
let platform_head = self
.platform
@ -381,6 +436,7 @@ impl Projector {
.map_err(ProjectError::Store)?
.unwrap_or(0);
let petri_heads = self.petri_heads(&run_id).await?;
let stored = &run.view;
let at_head = platform_head == stored.positions.platform_seq
&& petri_heads.iter().all(|(log, head)| {
stored
@ -396,25 +452,35 @@ impl Projector {
contended: false,
petri_events: 0,
platform_records: 0,
replayed_records: 0,
stream_seq: stored.stream_seq,
positions: stored.positions,
health: stored.view.state.health,
positions: stored.positions.clone(),
health: stored.view.state.health.clone(),
});
}
let StoredView {
mut view,
mut positions,
mut stream_seq,
} = stored;
let platform_records = self
.platform
.read_after(&run_id, positions.platform_seq)
.read_after(&run_id, stored.positions.platform_seq)
.await
.map_err(ProjectError::Store)?;
let mut replayed_records = 0;
let (events, replay_failure) = match self.store.open(&key, Access::Read).await {
Ok(logs) => match events::replay_since(&*logs, &positions.held()).await {
Ok(events) => (events, None),
Ok(logs) => match run.replay.advance(&*logs).await {
Ok(new) => {
replayed_records = new.iter().filter(|event| event.id.index == 0).count();
// A rebuilt replay derives the run whole: only the events
// past the view's positions are new to it.
let held = run.view.positions.held();
let mut events = std::mem::take(&mut run.pending);
events.extend(new.into_iter().filter(|event| {
held.get(&event.id.source)
.is_none_or(|last| event.id > *last)
}));
(events, None)
}
// The replay stood still and is retried by the next pass;
// what it derived before stays pending.
Err(error) => {
let chain = collect_chain(&error).join(": ");
warn!(run_id = %run_id, error = %chain, "Petri run does not replay; the view holds");
@ -425,6 +491,9 @@ impl Projector {
Err(error) => return Err(ProjectError::Open(error)),
};
let mut view = run.view.view.clone();
let mut positions = run.view.positions.clone();
let mut stream_seq = run.view.stream_seq;
let run_finished = view.state.finished_run()
|| events.iter().any(|event| {
matches!(
@ -475,12 +544,85 @@ impl Projector {
};
rows.push(row);
}
drop(items);
view.state.health = self.health(&key, &view.state, replay_failure).await?;
if self.fault.swap(false, Ordering::SeqCst) {
run.pending = events;
return Err(ProjectError::Injected);
}
let written = self
.write_view(
run_id,
&view,
&positions,
stream_seq,
&rows,
platform_head_seen,
)
.await;
match written {
Ok(true) => {}
Ok(false) => {
debug!(run_id = %run_id, "platform records landed during the pass; running it again");
run.pending = events;
return Ok(PassReport {
run_id,
skipped: false,
contended: true,
petri_events: 0,
platform_records: 0,
replayed_records,
stream_seq: 0,
positions: Positions::default(),
health: RecordHealth::default(),
});
}
Err(error) => {
run.pending = events;
return Err(error);
}
}
debug!(
run_id = %run_id,
petri_events = events.len(),
platform_records = platform_records.len(),
replayed_records,
stream_seq,
"Petri projection pass committed"
);
let petri_events = events.len();
let health = view.state.health.clone();
run.committed(view, positions.clone(), stream_seq);
if !rows.is_empty() {
// No receiver is not an error: nobody follows the stream.
let _ = self.committed.send(run_id);
}
Ok(PassReport {
run_id,
skipped: false,
contended: false,
petri_events,
platform_records: platform_records.len(),
replayed_records,
stream_seq,
positions,
health,
})
}
/// The view transaction: the projection row, the stream rows and the
/// `runs` row, committed together, unless a platform record landed
/// since the pass read them (`false`: the view is left alone).
async fn write_view(
&self,
run_id: RunId,
view: &RunView,
positions: &Positions,
stream_seq: u64,
rows: &[StreamRow],
platform_head_seen: u64,
) -> Result<bool, ProjectError> {
let mut tx = self
.pool
.begin_with("BEGIN IMMEDIATE")
@ -494,23 +636,13 @@ impl Projector {
.await
.map_err(ProjectError::Database)?;
if u64::try_from(head_now).unwrap_or(0) != platform_head_seen {
debug!(run_id = %run_id, "platform records landed during the pass; running it again");
drop(tx);
return Ok(PassReport {
run_id,
skipped: false,
contended: true,
petri_events: 0,
platform_records: 0,
stream_seq: 0,
positions: Positions::default(),
health: RecordHealth::default(),
});
return Ok(false);
}
let projection_json =
serde_json::to_string(&view.projection).map_err(ProjectError::Encode)?;
let fold_json = serde_json::to_string(&view.state).map_err(ProjectError::Encode)?;
let positions_json = serde_json::to_string(&positions).map_err(ProjectError::Encode)?;
let positions_json = serde_json::to_string(positions).map_err(ProjectError::Encode)?;
sqlx::query(
"INSERT INTO petri_projection (run_id, projection_json, fold_json, positions_json, \
stream_seq, updated_at_ms) VALUES (?, ?, ?, ?, ?, ?) ON CONFLICT(run_id) DO UPDATE \
@ -527,7 +659,7 @@ impl Projector {
.execute(&mut *tx)
.await
.map_err(ProjectError::Database)?;
for row in &rows {
for row in rows {
sqlx::query(
"INSERT INTO petri_stream (run_id, stream_seq, item_kind, item_id, event_json) \
VALUES (?, ?, ?, ?, ?)",
@ -547,27 +679,7 @@ impl Projector {
.map_err(ProjectError::Store)?;
}
tx.commit().await.map_err(ProjectError::Database)?;
debug!(
run_id = %run_id,
petri_events = events.len(),
platform_records = platform_records.len(),
stream_seq,
"Petri projection pass committed"
);
if !rows.is_empty() {
// No receiver is not an error: nobody follows the stream.
let _ = self.committed.send(run_id);
}
Ok(PassReport {
run_id,
skipped: false,
contended: false,
petri_events: events.len(),
platform_records: platform_records.len(),
stream_seq,
positions,
health: view.state.health.clone(),
})
Ok(true)
}
/// The stored view of the run, or an empty one.
@ -729,6 +841,14 @@ impl petri_execution::RunLogs for SignallingLogs {
self.inner.read(log).await
}
async fn read_from(
&self,
log: &petri_execution::LogId,
seq: u64,
) -> Result<Vec<petri_execution::Record>, StoreError> {
self.inner.read_from(log, seq).await
}
async fn put_blob(&self, bytes: &[u8]) -> Result<petri_store::Digest, StoreError> {
self.inner.put_blob(bytes).await
}

View file

@ -0,0 +1,132 @@
//! The state a live run's passes continue from, kept in memory between
//! passes: Petri's replay of the run (the coordinator state, each
//! execution's engine state, the projection) and Fabro's view as the last
//! committed pass left it. With it a pass reads and folds only the records
//! past the ones the view holds, so its cost is the new records', not the
//! run's.
//!
//! The cache is never a source of facts. It is dropped when the run
//! records its finish, when it has not been used for [`IDLE`], when the
//! stored view moves under it (another projector committed a pass), when
//! the run is deleted, and with the process; the first pass after that
//! rebuilds it by a full replay, which is what every pass did before the
//! cache existed. A replay that fails (a torn tail) stands still and is
//! retried by the next pass. Every pass, cached or not, commits the same
//! rows: the rebuild test in `tests/projection.rs` compares the two.
use std::collections::HashMap;
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use std::time::{Duration, Instant};
use fabro_types::RunId;
use petri_execution::events::{RunEvent, RunReplay};
use tokio::sync::Mutex as AsyncMutex;
use super::{Positions, StoredView};
use crate::projection::RunView;
/// How long a run's cache is kept after its last pass. A run blocked on a
/// question for longer pays one full replay when its next record lands.
pub(super) const IDLE: Duration = Duration::from_mins(10);
/// One live run's cache.
pub(super) struct RunCache {
pub(super) replay: RunReplay,
/// Events derived by an earlier pass that committed nothing: a pass
/// that found a platform record landing under it, or that failed
/// before its view transaction. They lead the next pass's events.
pub(super) pending: Vec<RunEvent>,
/// The view as the last committed pass left it, with its positions.
pub(super) view: StoredView,
}
impl RunCache {
/// A cache over the stored view, with a replay that has consumed
/// nothing: the first advance replays the run whole.
pub(super) fn over(view: StoredView) -> Self {
Self {
replay: RunReplay::new(),
pending: Vec::new(),
view,
}
}
/// Whether the cache still describes the stored view: its positions and
/// delivery sequence are the ones the view tables hold.
pub(super) fn matches(&self, positions: &Positions, stream_seq: u64) -> bool {
self.view.stream_seq == stream_seq
&& self.view.positions.platform_seq == positions.platform_seq
&& self.view.positions.held() == positions.held()
}
/// The pass committed: the view moved on, and nothing is pending.
pub(super) fn committed(&mut self, view: RunView, positions: Positions, stream_seq: u64) {
self.view = StoredView {
view,
positions,
stream_seq,
};
self.pending.clear();
}
}
/// The caches of every run the projector passed over, each behind the
/// run's pass lock, so a pass and a sweep never race over one cache.
#[derive(Default)]
pub(crate) struct Caches {
runs: Mutex<HashMap<RunId, Entry>>,
}
struct Entry {
pass: Arc<AsyncMutex<Option<RunCache>>>,
touched: Instant,
}
impl Caches {
/// The run's pass lock, holding its cache if one is kept; the run counts
/// as used now.
pub(super) fn pass_of(&self, run_id: RunId) -> Arc<AsyncMutex<Option<RunCache>>> {
let mut runs = lock(&self.runs);
let entry = runs.entry(run_id).or_insert_with(|| Entry {
pass: Arc::default(),
touched: Instant::now(),
});
entry.touched = Instant::now();
Arc::clone(&entry.pass)
}
/// Drop the cache of every run not used for `idle`, and forget the runs
/// with no cache and no pass under way. A run whose pass is running is
/// in use and left alone. How many caches were dropped.
pub(crate) fn sweep(&self, idle: Duration) -> usize {
let mut runs = lock(&self.runs);
let mut dropped = 0;
runs.retain(|_, entry| {
if entry.touched.elapsed() < idle {
return true;
}
let Ok(mut cache) = entry.pass.try_lock() else {
return true;
};
if cache.take().is_some() {
dropped += 1;
}
drop(cache);
// An `Arc` held elsewhere is a pass about to take the lock: the
// entry stays so the run keeps one lock.
Arc::strong_count(&entry.pass) > 1
});
dropped
}
/// Whether a cache is kept for the run: a test's view of the cache.
pub(crate) fn holds(&self, run_id: RunId) -> bool {
let runs = lock(&self.runs);
runs.get(&run_id)
.is_some_and(|entry| entry.pass.try_lock().is_ok_and(|cache| cache.is_some()))
}
}
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex.lock().unwrap_or_else(PoisonError::into_inner)
}

View file

@ -513,11 +513,17 @@ impl RunLogs for SqliteRunLogs {
}
async fn read(&self, log: &LogId) -> Result<Vec<Record>, StoreError> {
self.read_from(log, 0).await
}
async fn read_from(&self, log: &LogId, seq: u64) -> 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",
"SELECT record_json FROM petri_records WHERE run_id = ? AND log = ? AND seq >= ? ORDER \
BY seq",
)
.bind(self.key.as_str())
.bind(log_id_text(log))
.bind(i64::try_from(seq).unwrap_or(i64::MAX))
.fetch_all(&self.shared.pool)
.await
.map_err(|cause| self.backend("read a log", cause))?;

View file

@ -6,6 +6,7 @@
use std::collections::HashMap;
use std::sync::{Mutex, MutexGuard, PoisonError};
use std::time::Duration;
use async_trait::async_trait;
use bytes::Bytes;
@ -16,6 +17,20 @@ pub use petri_testkit::run_store;
use crate::blobs::Blobs;
use crate::platform_records::{PlatformRecordError, PlatformRecords};
use crate::projector::Projector;
/// Whether the projector keeps a cache for the run: the replay and the
/// view its passes continue from.
#[must_use]
pub fn cache_held(projector: &Projector, run_id: RunId) -> bool {
projector.caches.holds(run_id)
}
/// Drop the projector's caches not used for `idle`, as its passes do
/// after the documented idle period; how many were dropped.
pub fn drop_idle_caches(projector: &Projector, idle: Duration) -> usize {
projector.caches.sweep(idle)
}
/// A blob table in memory.
#[derive(Debug, Default)]

View file

@ -3,7 +3,9 @@
//! missed its wake-ups catches up on the next signal; a crash between the
//! record commit and the view transaction is recovered by applying only the
//! missing suffix; two projectors over one store agree over nested child
//! executions; and a torn tail holds the view where it stands.
//! executions; a torn tail holds the view where it stands; and a pass over
//! a live run costs its new records, with the cache that makes it so
//! dropped at a restart, after the idle period and at the run's finish.
//!
//! Every run here takes its scope's environment through the sandbox-driver
//! host plugin, so the tests skip, and say why, when the executable is not
@ -25,13 +27,13 @@ use std::time::{Duration, Instant};
use fabro_db::DbPool;
use fabro_interview::ControlInterviewer;
use fabro_petri::SqliteRunStore;
use fabro_petri::blobs::{Blobs, RunBlobs};
use fabro_petri::check::Launch;
use fabro_petri::engine::{self, RunStatus as EngineRunStatus};
use fabro_petri::interview::{Approval, FabroInterviewer};
use fabro_petri::projector::{self, Projector};
use fabro_petri::runtime::RuntimeSpec;
use fabro_petri::{SqliteRunStore, test_support as petri_support};
use fabro_store::platform_records::{
PlatformRecord, PlatformRecordStore, RunCreatedRecord, RunLifecycleKind, RunLifecycleRecord,
};
@ -968,6 +970,181 @@ async fn a_torn_tail_holds_the_view_and_reports_the_run_incomplete() {
);
}
/// The run's records the projection reads (the coordinator log and the
/// execution logs; the sandbox ledger has no events) in the order they
/// were recorded: by `recorded_at`, the coordinator log first on a tie,
/// each log's own order kept. The ledger's records are copied to `staged`
/// first, since they are no part of any batch.
async fn in_recorded_order<'a>(
rows: &'a [(String, i64, i64, String)],
staged: &DbPool,
run_id: RunId,
) -> Vec<&'a (String, i64, i64, String)> {
let (ledger, projected): (Vec<_>, Vec<_>) = rows.iter().partition(|row| row.0 == "resources");
for row in ledger {
insert_petri_row(staged, run_id, row).await;
}
let mut ordered = projected;
ordered.sort_by_key(|row| (row.2, row.0 != "coordinator", row.0.clone(), row.1));
ordered
}
/// One committed pass over the run, run again while a platform record
/// contends it.
async fn committed_pass(projector: &Projector, run_id: RunId) -> projector::PassReport {
loop {
let report = projector
.project_run(run_id)
.await
.expect("the pass commits");
if !report.contended {
return report;
}
}
}
/// The records land in batches and a pass follows each: every pass feeds
/// only its batch through Petri's derivation, never the run so far, and
/// the view the batches build is the rebuild. The finished run's cache is
/// dropped, and a pass over it is skipped.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_pass_over_a_live_run_costs_its_new_records_not_the_run() {
const BATCH: usize = 7;
if host_plugin().is_none() {
return;
}
let scenario = parallel_scenario().await;
run_unobserved(&scenario).await;
let rows = petri_rows(&scenario.pool, scenario.run_id).await;
let staged = copy_run_without_records(&scenario.pool, scenario.run_id).await;
let ordered = in_recorded_order(&rows, &staged, scenario.run_id).await;
assert!(
ordered.len() > 4 * BATCH,
"enough records for several batches: {}",
ordered.len()
);
let projector = Projector::new(staged.clone(), staged.clone());
let mut replayed = Vec::new();
for batch in ordered.chunks(BATCH) {
for row in batch {
insert_petri_row(&staged, scenario.run_id, row).await;
}
let report = committed_pass(&projector, scenario.run_id).await;
assert!(!report.skipped, "a batch is folded: {report:?}");
assert!(
report.replayed_records <= batch.len(),
"pass {}: {} records replayed for a batch of {}",
replayed.len(),
report.replayed_records,
batch.len()
);
replayed.push(report.replayed_records);
}
assert_eq!(
replayed.iter().sum::<usize>(),
ordered.len(),
"every record was fed once: {replayed:?}"
);
assert_view_equals_rebuild(&staged, scenario.run_id).await;
assert!(
!petri_support::cache_held(&projector, scenario.run_id),
"a finished run's cache is dropped"
);
let again = committed_pass(&projector, scenario.run_id).await;
assert!(again.skipped, "nothing is left to fold: {again:?}");
assert!(again.health.complete, "{:?}", again.health.incomplete);
}
/// The cache is dropped with the process and after the idle period, and
/// rebuilt by one full replay: the first pass over new records after
/// either feeds the run so far through Petri's derivation, the next only
/// its new records. Nothing is checkpointed for it, and the view it
/// continues is the rebuild.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_restart_and_the_idle_period_drop_the_cache_and_one_full_replay_rebuilds_it() {
const BATCH: usize = 5;
if host_plugin().is_none() {
return;
}
let scenario = parallel_scenario().await;
run_unobserved(&scenario).await;
let rows = petri_rows(&scenario.pool, scenario.run_id).await;
let staged = copy_run_without_records(&scenario.pool, scenario.run_id).await;
let ordered = in_recorded_order(&rows, &staged, scenario.run_id).await;
let half = ordered.len() / 2;
assert!(half > 3 * BATCH, "enough records: {}", ordered.len());
let mut fed = 0;
let mut feed = |count: usize| {
let rows: Vec<_> = ordered[fed..(fed + count).min(ordered.len())].to_vec();
fed += rows.len();
rows
};
for row in feed(half) {
insert_petri_row(&staged, scenario.run_id, row).await;
}
let before = Projector::new(staged.clone(), staged.clone());
let first = committed_pass(&before, scenario.run_id).await;
assert_eq!(
first.replayed_records, half,
"the first pass replays the run so far"
);
assert!(!first.health.complete, "the run has not finished");
assert!(
petri_support::cache_held(&before, scenario.run_id),
"a live run's cache is kept"
);
drop(before);
// A restarted server builds a new projector: no cache, and the next
// pass replays the run whole once.
let after = Projector::new(staged.clone(), staged.clone());
assert!(
!petri_support::cache_held(&after, scenario.run_id),
"a restart holds no cache"
);
for row in feed(BATCH) {
insert_petri_row(&staged, scenario.run_id, row).await;
}
let rebuilt = committed_pass(&after, scenario.run_id).await;
assert_eq!(
rebuilt.replayed_records,
half + BATCH,
"the first pass after a restart replays the run so far"
);
assert!(petri_support::cache_held(&after, scenario.run_id));
for row in feed(BATCH) {
insert_petri_row(&staged, scenario.run_id, row).await;
}
let live = committed_pass(&after, scenario.run_id).await;
assert_eq!(
live.replayed_records, BATCH,
"the next pass replays its batch"
);
// The idle period passes: the cache is dropped, and rebuilt the same way.
assert_eq!(
petri_support::drop_idle_caches(&after, Duration::ZERO),
1,
"the run's cache was idle"
);
assert!(!petri_support::cache_held(&after, scenario.run_id));
for row in feed(BATCH) {
insert_petri_row(&staged, scenario.run_id, row).await;
}
let idle = committed_pass(&after, scenario.run_id).await;
assert_eq!(idle.replayed_records, half + 3 * BATCH);
let rest = feed(ordered.len());
let rest_len = rest.len();
for row in rest {
insert_petri_row(&staged, scenario.run_id, row).await;
}
let last = committed_pass(&after, scenario.run_id).await;
assert_eq!(last.replayed_records, rest_len);
assert!(last.health.complete, "{:?}", last.health.incomplete);
assert_view_equals_rebuild(&staged, scenario.run_id).await;
}
/// A gate scenario runs through the engine assembly with the interview
/// adapter, as a Fabro run does, over a store that signals the projector.
struct GateRun {