Merge pull request #886 from fabro-sh/petri-maintainability
Some checks failed
Rust / Format (push) Waiting to run
Rust / Clippy (push) Waiting to run
Rust / Rustdoc (push) Waiting to run
Rust / Generated Docs (push) Waiting to run
Rust / Test (Linux) (push) Waiting to run
Rust / Sandbox providers (Docker) (push) Waiting to run
Rust / Test (macOS) (push) Waiting to run
TypeScript / Typecheck (push) Has been cancelled
TypeScript / Test (push) Has been cancelled
TypeScript / Build (push) Has been cancelled

fabro-petri maintainability: ten refactors from the slopdetect review
This commit is contained in:
Bryan Helmkamp 2026-09-20 16:58:56 -04:00 • committed by GitHub
commit 40419cbd2b
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
31 changed files with 3998 additions and 3433 deletions

View file

@ -281,7 +281,7 @@ fn futures_lite_block_on<T>(future: impl std::future::Future<Output = T>) -> T {
fn restore_actions(server: &RunningServer, run_id: &str) -> Vec<String> {
let log = std::fs::read_to_string(server.worker_log(run_id)).unwrap_or_default();
log.lines()
.filter(|line| line.contains("sandbox workspace brought to its durable snapshot"))
.filter(|line| line.contains("workspace brought to its durable snapshot"))
.filter_map(|line| {
line.split_whitespace()
.find_map(|word| word.strip_prefix("action=").map(str::to_owned))

View file

@ -14,12 +14,13 @@
//! `RunOptions::run_key`.
use std::collections::HashMap;
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use std::sync::{Arc, Mutex};
use fabro_db::DbPool;
use fabro_petri::SqliteRunStore;
use fabro_petri::petri::{Access, OwnerId, RunKey, RunLogs, RunStore as _, StoreError};
use fabro_types::RunId;
use fabro_util::sync;
use tracing::debug;
pub(crate) struct PetriRuns {
@ -66,7 +67,7 @@ impl PetriRuns {
) -> Result<Arc<dyn RunLogs>, StoreError> {
let handle = self.store.open(&Self::key(&run_id), access.clone()).await?;
if let Some(owner) = access.owner() {
lock(&self.handles).insert((run_id, owner.clone()), Arc::clone(&handle));
sync::lock(&self.handles).insert((run_id, owner.clone()), Arc::clone(&handle));
}
Ok(handle)
}
@ -80,7 +81,7 @@ impl PetriRuns {
run_id: RunId,
owner: &OwnerId,
) -> Result<Arc<dyn RunLogs>, StoreError> {
if let Some(handle) = lock(&self.handles).get(&(run_id, owner.clone())) {
if let Some(handle) = sync::lock(&self.handles).get(&(run_id, owner.clone())) {
return Ok(Arc::clone(handle));
}
let holder = self.store.owner(&Self::key(&run_id)).await?;
@ -107,7 +108,7 @@ impl PetriRuns {
/// Drop the handle `owner` holds on the run: the worker's own release.
/// The store ends the lease when this was the owner's last handle.
pub(crate) fn release(&self, run_id: RunId, owner: &OwnerId) {
let handle = lock(&self.handles).remove(&(run_id, owner.clone()));
let handle = sync::lock(&self.handles).remove(&(run_id, owner.clone()));
debug!(
run_id = %run_id,
owner = %owner,
@ -132,7 +133,7 @@ impl PetriRuns {
/// releasing does not keep the lease.
pub(crate) fn worker_exited(&self, run_id: RunId) {
let dropped = {
let mut handles = lock(&self.handles);
let mut handles = sync::lock(&self.handles);
let owners: Vec<_> = handles
.keys()
.filter(|(held, _)| *held == run_id)
@ -154,10 +155,6 @@ impl PetriRuns {
}
}
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex.lock().unwrap_or_else(PoisonError::into_inner)
}
#[cfg(test)]
mod tests {
use std::collections::BTreeMap;
@ -226,14 +223,14 @@ mod tests {
}
fn launched_mode(&self) -> Option<&'static str> {
*lock(&self.mode)
*sync::lock(&self.mode)
}
}
#[async_trait::async_trait]
impl WorkerRuntime for HeldWorkerRuntime {
async fn start(&self, spec: WorkerLaunchSpec) -> anyhow::Result<StartedWorker> {
*lock(&self.mode) = Some(spec.mode);
*sync::lock(&self.mode) = Some(spec.mode);
self.running.store(true, Ordering::SeqCst);
let exit = Arc::clone(&self.exit);
let stderr: Pin<Box<dyn AsyncRead + Send + 'static>> = Box::pin(tokio::io::empty());

View file

@ -1201,6 +1201,19 @@ impl AppState {
&self.petri_projector
}
/// The status the server holds for a managed run, so a test can wait
/// for the run to settle in the server's own map (what the delete
/// precheck reads) and not only in the stored view, which can report
/// the run ended first.
#[cfg(any(test, feature = "test-support"))]
#[must_use]
pub fn test_managed_run_status(&self, run_id: &RunId) -> Option<RunStatus> {
self.runs
.lock()
.ok()
.and_then(|runs| runs.get(run_id).map(|managed_run| managed_run.status))
}
/// The pool the Petri view tables live in, so a test can read them.
#[cfg(any(test, feature = "test-support"))]
pub fn test_petri_view_pool(&self) -> DbPool {

View file

@ -28,7 +28,7 @@ use axum::body::Body;
use axum::http::{Request, StatusCode};
use fabro_petri::engine::{self, RunStatus};
use fabro_petri::petri::{Access, OwnerId, RunKey, RunStore as _};
use fabro_petri::{SqliteRunStore, projector};
use fabro_petri::{SqliteRunStore, test_support};
use fabro_server::server::AppState;
use fabro_server::test_support::{
TestAppStateBuilder, llm_overlay_with_provider_base_url, test_app_db_pool,
@ -237,7 +237,7 @@ pub(super) async fn settled_state(
/// How many items the run's projected stream holds.
async fn petri_stream_len(state: &AppState, run_id: &str) -> usize {
let id: RunId = run_id.parse().expect("the run id parses");
projector::stored_stream(&state.test_petri_view_pool(), id)
test_support::stored_stream(&state.test_petri_view_pool(), id)
.await
.expect("the stream reads")
.len()
@ -802,6 +802,9 @@ async fn deleting_a_run_prunes_its_host_workspace_through_petri() {
let store = state.test_petri_run_store();
let key = RunKey::new(run_id.clone());
wait_for_free_lease(store, &key).await;
// The view reports the run ended from Petri's own finish, before the
// server settles the managed run the delete precheck reads.
wait_for_managed_settle(&state, &run_id).await;
// A live handle on the run, as its worker holds one, refuses the
// delete: Petri will not prune under a lease someone holds.
@ -858,6 +861,21 @@ async fn deleting_a_run_prunes_its_host_workspace_through_petri() {
}
/// Wait until no owner holds the run's lease.
/// Wait until the server's own map holds the run as ended.
async fn wait_for_managed_settle(state: &AppState, run_id: &str) {
let run_id: RunId = run_id.parse().expect("a run id");
for _ in 0..500 {
if state
.test_managed_run_status(&run_id)
.is_none_or(fabro_types::RunStatus::is_terminal)
{
return;
}
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
panic!("the managed run did not settle");
}
async fn wait_for_free_lease(store: &SqliteRunStore, key: &RunKey) {
for _ in 0..500 {
if store.owner(key).await.expect("reads the lease").is_none() {

View file

@ -43,8 +43,8 @@ use std::time::Duration;
use fabro_checkpoint::author::GitAuthor;
use fabro_checkpoint::trailer::{self, Trailer};
use fabro_store::platform_records::{DecisionRef, OperationKey};
use fabro_types::DiffSummary;
use fabro_types::settings::run::RunCheckpointSettings;
use fabro_types::settings::run::{RunCheckpointSettings, RunNamespace};
use fabro_types::{DiffSummary, GitIdentitySource, SandboxProviderKind};
use petri_runtime::executor::{EnvError, ExecEnv, OutputMode, ProcessSpec, Sig};
use petri_runtime::ir::LogStream;
use tokio::process::Command;
@ -92,6 +92,54 @@ pub const EXCLUDE_DIRS: &[&str] = &[
".pytest_cache",
];
/// The settings a run's Git work runs under, as its namespace gives them:
/// who authors the checkpoint commits and where that identity came from,
/// the checkpoint settings, and whether the sandbox provider keeps the
/// workspaces on this host. The hooks and recovery both start from it.
#[derive(Clone, Debug)]
pub struct RunGitSettings {
pub author: GitAuthor,
pub identity_source: GitIdentitySource,
pub checkpoint: RunCheckpointSettings,
/// Whether the run's workspaces are on this host (the local sandbox
/// provider). A run elsewhere snapshots inside its sandboxes.
pub host_workspaces: bool,
}
impl From<&RunNamespace> for RunGitSettings {
fn from(settings: &RunNamespace) -> Self {
let author = settings
.git
.author
.as_ref()
.map(GitAuthor::from)
.unwrap_or_default();
let identity_source = if author.is_default() {
GitIdentitySource::Default
} else {
GitIdentitySource::Explicit
};
Self {
author,
identity_source,
checkpoint: settings.checkpoint.clone(),
host_workspaces: settings.environment.provider == SandboxProviderKind::LOCAL,
}
}
}
impl Default for RunGitSettings {
/// Fabro's default author and checkpoint settings, on this host.
fn default() -> Self {
Self {
author: GitAuthor::default(),
identity_source: GitIdentitySource::Default,
checkpoint: RunCheckpointSettings::default(),
host_workspaces: true,
}
}
}
/// The identity of one snapshot: the attempt whose files it holds.
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct CheckpointKey {
@ -178,6 +226,16 @@ impl CheckpointKey {
}
}
impl std::fmt::Display for CheckpointKey {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"execution {} firing {} attempt {}",
self.execution, self.firing, self.attempt
)
}
}
/// Why a snapshot could not be taken, found or restored.
#[derive(Debug, thiserror::Error)]
pub enum CheckpointError {
@ -340,53 +398,19 @@ impl RunWorkspaces {
.unwrap_or(false)
}
/// The host site of a workspace.
fn host(&self, workspace: &str) -> Site {
/// The host site of a workspace: where `git` runs for a workspace kept
/// on this host.
#[must_use]
pub fn host(&self, workspace: &str) -> Site {
Site::Host(self.workspace_path(workspace))
}
/// Commit the workspace's files on the run branch as the snapshot of
/// `key`, and publish it. An earlier commit of the same key that the
/// workspace still sits on, unchanged, is reused.
/// workspace still sits on, unchanged, is reused. On a sandbox site
/// `git` runs in the scope through its environment, and the commit
/// reaches the snapshot repository as a bundle.
pub async fn commit(
&self,
workspace: &str,
key: CheckpointKey,
node: &str,
status: &str,
) -> Result<Snapshot, CheckpointError> {
if !self.workspace_exists(workspace).await {
return Err(CheckpointError::WorkspaceMissing {
workspace: workspace.to_string(),
path: self.workspace_path(workspace),
});
}
self.commit_at(&self.host(workspace), workspace, key, node, status)
.await
}
/// [`commit`](Self::commit) for a workspace inside a sandbox: `git`
/// runs in the scope through `env`, and the commit reaches the
/// snapshot repository as a bundle.
pub async fn commit_in(
&self,
env: &Arc<dyn ExecEnv>,
workspace: &str,
key: CheckpointKey,
node: &str,
status: &str,
) -> Result<Snapshot, CheckpointError> {
self.commit_at(
&Site::Sandbox(Arc::clone(env)),
workspace,
key,
node,
status,
)
.await
}
async fn commit_at(
&self,
site: &Site,
workspace: &str,
@ -394,6 +418,14 @@ impl RunWorkspaces {
node: &str,
status: &str,
) -> Result<Snapshot, CheckpointError> {
if let Site::Host(path) = site {
if !fs::try_exists(path).await.unwrap_or(false) {
return Err(CheckpointError::WorkspaceMissing {
workspace: workspace.to_string(),
path: path.clone(),
});
}
}
let branched = self.ensure_repository(site).await?;
if let Some(existing) = self.published_sha(workspace, key).await? {
if self.head(site).await?.as_deref() == Some(existing.as_str())
@ -575,59 +607,20 @@ impl RunWorkspaces {
.is_some())
}
/// The workspace's `HEAD`, or `None` when it has no commit.
pub async fn workspace_head(&self, workspace: &str) -> Result<Option<String>, CheckpointError> {
self.head(&self.host(workspace)).await
}
/// [`workspace_head`](Self::workspace_head) for a workspace inside a
/// sandbox.
pub async fn workspace_head_in(
&self,
env: &Arc<dyn ExecEnv>,
) -> Result<Option<String>, CheckpointError> {
self.head(&Site::Sandbox(Arc::clone(env))).await
}
/// Whether the workspace sits on `sha` with nothing changed since.
pub async fn matches(&self, workspace: &str, sha: &str) -> Result<bool, CheckpointError> {
self.matches_at(&self.host(workspace), sha).await
}
/// [`matches`](Self::matches) for a workspace inside a sandbox.
pub async fn matches_in(
&self,
env: &Arc<dyn ExecEnv>,
sha: &str,
) -> Result<bool, CheckpointError> {
self.matches_at(&Site::Sandbox(Arc::clone(env)), sha).await
}
async fn matches_at(&self, site: &Site, sha: &str) -> Result<bool, CheckpointError> {
pub async fn matches(&self, site: &Site, sha: &str) -> Result<bool, CheckpointError> {
Ok(self.head(site).await?.as_deref() == Some(sha) && self.is_clean(site).await?)
}
/// Whether the host workspace's repository holds the commit `sha`, so a
/// reset can reach it; a directory that is no repository holds none.
pub async fn has_commit(&self, workspace: &str, sha: &str) -> Result<bool, CheckpointError> {
if !self.workspace_exists(workspace).await {
return Ok(false);
/// Whether the workspace's repository holds the commit `sha`, so a
/// reset can reach it (without a transfer, in a sandbox); a directory
/// that is gone or is no repository holds none.
pub async fn has_commit(&self, site: &Site, sha: &str) -> Result<bool, CheckpointError> {
if let Site::Host(path) = site {
if !fs::try_exists(path).await.unwrap_or(false) {
return Ok(false);
}
}
self.has_commit_at(&self.host(workspace), sha).await
}
/// Whether a sandbox workspace's repository holds the commit `sha`, so
/// a reset can reach it without a transfer.
pub async fn has_commit_in(
&self,
env: &Arc<dyn ExecEnv>,
sha: &str,
) -> Result<bool, CheckpointError> {
self.has_commit_at(&Site::Sandbox(Arc::clone(env)), sha)
.await
}
async fn has_commit_at(&self, site: &Site, sha: &str) -> Result<bool, CheckpointError> {
if self
.git_status(site, "rev-parse", &["rev-parse", "--git-dir"])
.await?
@ -647,16 +640,7 @@ impl RunWorkspaces {
/// Bring the workspace back to `sha`: tracked files reset, untracked
/// files removed, the excluded caches left alone.
pub async fn reset(&self, workspace: &str, sha: &str) -> Result<(), CheckpointError> {
self.reset_at(&self.host(workspace), sha).await
}
/// [`reset`](Self::reset) for a workspace inside a sandbox.
pub async fn reset_in(&self, env: &Arc<dyn ExecEnv>, sha: &str) -> Result<(), CheckpointError> {
self.reset_at(&Site::Sandbox(Arc::clone(env)), sha).await
}
async fn reset_at(&self, site: &Site, sha: &str) -> Result<(), CheckpointError> {
pub async fn reset(&self, site: &Site, sha: &str) -> Result<(), CheckpointError> {
self.git(site, "reset", &["reset", "-q", "--hard", sha])
.await?;
let mut clean = vec!["clean".to_string(), "-fdq".to_string()];
@ -673,21 +657,38 @@ impl RunWorkspaces {
}
/// Recreate a gone workspace from the published snapshot `key`, at
/// `sha`, on the run branch.
/// `sha`, on the run branch. On a host site the workspace directory is
/// created and the snapshot fetched from the repository beside it. Into
/// a sandbox the snapshot enters the scope as a bundle of the
/// checkpoint's ref, and the workspace, fresh or stale, is fetched from
/// it and forced onto the run branch at `sha`.
pub async fn restore(
&self,
site: &Site,
workspace: &str,
key: CheckpointKey,
sha: &str,
) -> Result<(), CheckpointError> {
let path = self.workspace_path(workspace);
fs::create_dir_all(&path)
match site {
Site::Host(path) => self.restore_on_host(path, workspace, key, sha).await,
Site::Sandbox(env) => self.restore_in_sandbox(env, workspace, key, sha).await,
}
}
async fn restore_on_host(
&self,
path: &Path,
workspace: &str,
key: CheckpointKey,
sha: &str,
) -> Result<(), CheckpointError> {
fs::create_dir_all(path)
.await
.map_err(|source| CheckpointError::Io {
path: path.clone(),
path: path.to_path_buf(),
source,
})?;
let site = Site::Host(path);
let site = Site::Host(path.to_path_buf());
self.git(&site, "init", &["init", "-q"]).await?;
let repository = self.snapshot_repository(workspace);
let repository = repository.to_string_lossy().into_owned();
@ -710,10 +711,7 @@ impl RunWorkspaces {
self.verify_restored(&site, sha).await
}
/// [`restore`](Self::restore) into a sandbox: the snapshot enters the
/// scope as a bundle of the checkpoint's ref, and the workspace, fresh
/// or stale, is fetched from it and forced onto the run branch at `sha`.
pub async fn restore_in(
async fn restore_in_sandbox(
&self,
env: &Arc<dyn ExecEnv>,
workspace: &str,
@ -766,7 +764,7 @@ impl RunWorkspaces {
"FETCH_HEAD",
])
.await?;
self.reset_at(&site, "HEAD").await?;
self.reset(&site, "HEAD").await?;
self.verify_restored(&site, sha).await
}
.await;
@ -1083,7 +1081,8 @@ impl RunWorkspaces {
.await
}
async fn head(&self, site: &Site) -> Result<Option<String>, CheckpointError> {
/// The workspace's `HEAD`, or `None` when it has no commit.
pub async fn head(&self, site: &Site) -> Result<Option<String>, CheckpointError> {
self.git_status(site, "rev-parse", &["rev-parse", "-q", "--verify", "HEAD"])
.await
}
@ -1360,7 +1359,13 @@ mod tests {
};
let first = workspaces
.commit(workspace, key, "build", "success")
.commit(
&workspaces.host(workspace),
workspace,
key,
"build",
"success",
)
.await
.expect("the commit");
assert!(!first.reused);
@ -1370,7 +1375,13 @@ mod tests {
"the first commit created the run branch in a fresh repository"
);
let again = workspaces
.commit(workspace, key, "build", "success")
.commit(
&workspaces.host(workspace),
workspace,
key,
"build",
"success",
)
.await
.expect("the second commit");
assert_eq!(again, Snapshot {
@ -1408,7 +1419,7 @@ mod tests {
);
assert!(
workspaces
.matches(workspace, &first.sha)
.matches(&workspaces.host(workspace), &first.sha)
.await
.expect("matches")
);
@ -1422,12 +1433,12 @@ mod tests {
.expect("an untracked file");
assert!(
!workspaces
.matches(workspace, &first.sha)
.matches(&workspaces.host(workspace), &first.sha)
.await
.expect("matches")
);
workspaces
.reset(workspace, &first.sha)
.reset(&workspaces.host(workspace), &first.sha)
.await
.expect("the reset");
assert_eq!(
@ -1449,7 +1460,7 @@ mod tests {
"the snapshot repository still knows the commit"
);
workspaces
.restore(workspace, key, &first.sha)
.restore(&workspaces.host(workspace), workspace, key, &first.sha)
.await
.expect("the restore");
assert_eq!(
@ -1460,7 +1471,7 @@ mod tests {
);
assert_eq!(
workspaces
.workspace_head(workspace)
.head(&workspaces.host(workspace))
.await
.expect("the head"),
Some(first.sha)
@ -1483,7 +1494,13 @@ mod tests {
attempt: 1,
};
let error = workspaces
.commit(workspace, key, "build", "success")
.commit(
&workspaces.host(workspace),
workspace,
key,
"build",
"success",
)
.await
.expect_err("the commit fails");
assert!(matches!(error, CheckpointError::Command { .. }), "{error}");
@ -1537,7 +1554,13 @@ mod tests {
attempt: 1,
};
let first = workspaces
.commit(workspace, first_key, "start", "success")
.commit(
&workspaces.host(workspace),
workspace,
first_key,
"start",
"success",
)
.await
.expect("the first commit");
assert_eq!(
@ -1555,7 +1578,13 @@ mod tests {
attempt: 1,
};
let second = workspaces
.commit(workspace, second_key, "write", "success")
.commit(
&workspaces.host(workspace),
workspace,
second_key,
"write",
"success",
)
.await
.expect("the second commit");
assert_eq!(second.branched, None);

View file

@ -509,16 +509,12 @@ pub async fn stage_labels(views: &DbPool, run_id: RunId) -> Result<StageLabels,
Ok(state
.stages
.iter()
.filter_map(|(key, stage)| {
let (execution, firing) = key.split_once(':')?;
Some((
(execution.parse().ok()?, firing.parse().ok()?),
StageLabel {
stage_id: stage.shown.then(|| stage.stage_id.to_string()),
node_name: stage.node_name.clone(),
visit: stage.visit,
},
))
.map(|(key, stage)| {
((key.execution, key.firing), StageLabel {
stage_id: stage.shown.then(|| stage.stage_id.to_string()),
node_name: stage.node_name.clone(),
visit: stage.visit,
})
})
.collect())
}

File diff suppressed because it is too large Load diff

View file

@ -62,13 +62,14 @@
use std::collections::HashMap;
use std::future::Future;
use std::sync::{Arc, Mutex, MutexGuard, PoisonError, Weak};
use std::sync::{Arc, Mutex, Weak};
use std::time::Duration;
use std::{fmt, mem, ptr};
use fabro_api::types::{PetriAccess, PetriAppendRequest, PetriOpenRequest, PetriRecord};
use fabro_client::{Client, api_failure_for};
use fabro_types::{BlobHash, RunId};
use fabro_util::sync;
use petri_store::{Access, Digest, LogId, OwnerId, Record, RunKey, RunLogs, RunStore, StoreError};
use serde_json::Value;
use tokio::runtime::Handle;
@ -148,7 +149,7 @@ impl HttpRunStore {
owner: OwnerId,
locator: String,
) -> Arc<HttpRunLogs> {
let mut live = lock(&self.shared.live);
let mut live = sync::lock(&self.shared.live);
let slot = (key.clone(), owner.clone());
if let Some(handle) = live.get(&slot).and_then(Weak::upgrade) {
return handle;
@ -185,7 +186,7 @@ impl Shared {
/// 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));
let pending = mem::take(&mut *sync::lock(&self.releases));
for release in pending {
// A release task never panics: it reports its own failure.
let _ = release.await;
@ -436,7 +437,7 @@ impl Drop for HttpRunLogs {
return;
};
{
let mut live = lock(&self.shared.live);
let mut live = sync::lock(&self.shared.live);
let slot = (self.key.clone(), owner.clone());
let this: *const Self = self;
if live
@ -454,7 +455,7 @@ impl Drop for HttpRunLogs {
let release = runtime.spawn(async move {
shared.release(&key, run_id, &owner).await;
});
lock(&self.shared.releases).push(release);
sync::lock(&self.shared.releases).push(release);
}
Err(_) => {
warn!(
@ -551,10 +552,6 @@ impl RunLogs for HttpRunLogs {
}
}
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex.lock().unwrap_or_else(PoisonError::into_inner)
}
#[cfg(test)]
mod tests {
use fabro_client::ApiFailure;

File diff suppressed because it is too large Load diff

View file

@ -0,0 +1,241 @@
//! Petri's coordinator events folded into the view: the run's start, its
//! invocations and executions, pause and unpause, its finish, and the
//! release of its sandbox (VIEWS.md "Run", "Sandbox").
use chrono::{DateTime, Utc};
use fabro_types::{
Conclusion, FailureCategory, FailureDetail, FailureReason, RunControlAction, RunFailure,
RunSandbox, RunStatus, RunTiming, StageOutcome, StartRecord, SuccessReason, usage_rollup,
};
use petri_execution::CoordinatorEvent;
use petri_execution::events::RunEvent;
use super::{FiringKey, InvocationRef, RunView, apply_status, settle_control};
impl RunView {
pub(super) fn fold_coordinator(
&mut self,
record: &CoordinatorEvent,
event: &RunEvent,
at: DateTime<Utc>,
) {
match record {
CoordinatorEvent::RunStarted {
root, forked_from, ..
} => {
self.state.root = Some(root.raw());
self.state.started_at = Some(event.recorded_at);
if let Some(projection) = self.projection.as_mut() {
// A fork's declaration names its source; a parse failure
// means the source was not a Fabro run, which the
// projection cannot show.
projection.forked_from = forked_from.as_ref().and_then(|origin| {
Some(fabro_types::ForkOrigin {
source_run_id: origin.source.as_str().parse().ok()?,
execution: origin.position.execution.raw(),
firing: origin.position.firing.raw(),
rerun_last: origin.rerun_last,
})
});
apply_status(projection, RunStatus::Running, at);
projection.start = Some(StartRecord {
start_time: at,
run_branch: self.state.run_branch.clone(),
base_sha: self.state.base_sha.clone(),
});
// The scope's sandbox is acquired next; `scope.acquired`
// or `scope.failed` settles it.
if let Some(sandbox) = projection.sandbox.take() {
projection.sandbox = Some(RunSandbox::initializing(sandbox.plan().clone()));
}
}
}
CoordinatorEvent::InvocationDeclared { invocation, .. } => {
let mut info = InvocationRef::default();
if let Some(parent) = &event.context.parent {
info.parent = Some((parent.execution.raw(), parent.firing.raw()));
if let Some((fork_firing, index)) = branch_slot(&parent.slot) {
let group = self
.state
.stages
.get(&FiringKey::new(parent.execution.raw(), fork_firing))
.map(|stage| stage.stage_id.clone());
if let Some(group) = group {
info.branch = Some((group, index));
}
}
}
self.state.invocations.insert(invocation.raw(), info);
}
CoordinatorEvent::ExecutionDeclared {
execution,
invocation,
..
} => {
self.state
.executions
.insert(execution.raw(), invocation.raw());
}
CoordinatorEvent::InvocationFinished { invocation, result } => {
let info = self.state.invocations.entry(invocation.raw()).or_default();
info.failure = result
.failure
.as_ref()
.map(|failure| failure.message.clone());
info.output = Some(result.output.clone());
}
CoordinatorEvent::RunPaused => {
if let Some(projection) = self.projection.as_mut() {
apply_status(projection, projection.status.paused(), at);
settle_control(projection, RunControlAction::Pause);
}
}
CoordinatorEvent::RunUnpaused => {
if let Some(projection) = self.projection.as_mut() {
apply_status(projection, projection.status.unpaused(), at);
settle_control(projection, RunControlAction::Unpause);
}
}
CoordinatorEvent::RunFinished { status } => {
self.state.finished = Some(status.to_string());
self.conclude(status.to_string().as_str(), at);
}
// ── Sandbox: the retention outcome (VIEWS.md "Sandbox") ─────────
// The instance stays on `Run.sandbox`: it names what ran, and
// `retained` says whether it still exists.
CoordinatorEvent::ScopeReleased {
invocation,
retained,
..
} => {
if Some(invocation.raw()) == self.state.root {
self.state.sandbox_retained = Some(*retained);
if let Some(sandbox) = self
.projection
.as_mut()
.and_then(|projection| projection.sandbox.as_mut())
{
sandbox.set_retained(*retained);
}
}
}
CoordinatorEvent::GraphRegistered { .. }
| CoordinatorEvent::ExecutionFinished { .. }
| CoordinatorEvent::InvocationCancelRequested { .. }
| CoordinatorEvent::RunNoteRecorded { .. } => {}
}
}
/// The run's conclusion, from its recorded finish and what the stages
/// The run's conclusion, from its recorded finish and what the stages
/// summed to.
fn conclude(&mut self, status: &str, at: DateTime<Utc>) {
let Some(projection) = self.projection.as_mut() else {
return;
};
let root = self
.state
.root
.and_then(|root| self.state.invocations.get(&root));
let failure_message = root.and_then(|root| root.failure.clone());
let (run_status, outcome, failure) = match status {
"success" => (
RunStatus::Succeeded {
reason: SuccessReason::Completed,
},
StageOutcome::Succeeded,
None,
),
"cancelled" => (
RunStatus::Failed {
reason: FailureReason::Cancelled,
},
StageOutcome::Failed {
retry_requested: false,
},
Some(RunFailure {
reason: FailureReason::Cancelled,
detail: FailureDetail::new(
failure_message
.clone()
.unwrap_or_else(|| "the run was cancelled".to_string()),
FailureCategory::Canceled,
),
}),
),
_ => (
RunStatus::Failed {
reason: FailureReason::WorkflowError,
},
StageOutcome::Failed {
retry_requested: false,
},
Some(RunFailure {
reason: FailureReason::WorkflowError,
detail: FailureDetail::new(
failure_message
.clone()
.unwrap_or_else(|| "the run failed".to_string()),
FailureCategory::Deterministic,
),
}),
),
};
apply_status(projection, run_status, at);
projection.pending_control = None;
projection.pending_interviews.clear();
let rollup = usage_rollup::usage_rollup_from_projection(projection);
let (stages, total_retries) = rollup.conclusion_stages(projection);
let wall_time_ms = self.state.started_at.map_or(0, |started| {
u64::try_from(at.timestamp_millis())
.unwrap_or(0)
.saturating_sub(started)
});
let timing = RunTiming::new(
wall_time_ms,
rollup.timing.inference_time_ms,
rollup.timing.tool_time_ms,
);
let last_checkpoint = projection.checkpoints.last();
projection.conclusion = Some(Conclusion {
timestamp: at,
status: outcome,
timing,
failure,
final_git_commit_sha: last_checkpoint
.and_then(|checkpoint| checkpoint.checkpoint.git_commit_sha.clone()),
stages,
usage: rollup.usage_if_present(),
total_retries,
diff: self
.state
.run_diff
.clone()
.or_else(|| last_checkpoint.map(|checkpoint| checkpoint.diff.clone()))
.unwrap_or_default(),
});
}
}
/// The fork firing and branch index a branch child's call slot names:
/// `branch:<fork>@<firing>:<index>:<target>`.
fn branch_slot(slot: &str) -> Option<(u64, u32)> {
let rest = slot.strip_prefix("branch:")?;
let mut parts = rest.splitn(3, ':');
let fork = parts.next()?;
let index = parts.next()?.parse::<u32>().ok()?;
let firing = fork.rsplit_once('@')?.1.parse::<u64>().ok()?;
Some((firing, index))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_branch_slot_names_the_fork_firing_and_the_index() {
assert_eq!(branch_slot("branch:fan@7:2:review"), Some((7, 2)));
assert_eq!(branch_slot("branch:fan@7:x:review"), None);
assert_eq!(branch_slot("child:0"), None);
}
}

View file

@ -0,0 +1,453 @@
//! Petri's engine and view events folded into the stages: a firing's
//! visit, its attempts and their outcomes, its wait states, and the scope
//! its sandbox was acquired in (VIEWS.md "Stages", "Sandbox").
use std::collections::BTreeMap;
use chrono::{DateTime, Utc};
use fabro_types::{
ModelUsage, ParallelBranchId, ParallelBranchResult, RunProjection, RunSandbox,
RunSandboxFailure, StageCompletion, StageHandler, StageId, StageOutcome, StageProjection,
StageState, StageTiming, first_event_seq, parse_blob_ref, timing,
};
use petri_execution::events::{Derived, RunEvent, Subject, ViewEvent, WaitState};
use petri_execution::{ExecutionId, InvocationId};
use petri_runtime::engine::{Admission, Event};
use petri_runtime::ir::{Metrics, Status};
use serde_json::Value;
use tracing::debug;
use super::model::{model_ref, usage_of};
use super::sandbox::{provider_kind, sandbox_instance, sandbox_plan_of};
use super::{FiringKey, RunView, StageRef, is_shown, node_meta_kind, stage_label, visit_of};
impl RunView {
pub(super) fn fold_engine(&mut self, engine: &Event, event: &RunEvent, at: DateTime<Utc>) {
let Some(execution) = event.context.execution else {
return;
};
match engine {
Event::AdmissionDecided { decision, .. } => {
if let Admission::Skip { outcome } = decision {
if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) {
stage.state = StageState::Skipped;
stage.completion = Some(StageCompletion {
outcome: StageOutcome::Skipped,
..completion(&outcome.status, at)
});
}
}
}
Event::StepStarted { attempt, .. } => {
if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) {
if attempt.raw() > 1 {
stage.clear_live_timing();
stage.output = None;
stage.output_bytes = None;
}
stage.state = StageState::Running;
stage.live_streaming = Some(true);
}
}
Event::StepProgressRecorded { ev, .. } => {
self.fold_progress(execution, event, ev, at);
}
Event::StepFinished {
firing,
attempt,
outcome,
} => {
self.state
.finished_firings
.insert(FiringKey::new(execution.raw(), firing.raw()));
let is_final = matches!(
event.derived,
Some(Derived::StepFinished { is_final: true, .. })
);
let node_name = event
.subject
.as_ref()
.map(|subject| subject.node.name.to_string());
if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) {
// The step's output: a string, or a command's `stdout`,
// either of which is a `blob://` reference when the
// step offloaded it. The reference stays as it is; the
// bytes it names are the live log's.
let output = outcome
.output
.as_str()
.or_else(|| outcome.output.get("stdout").and_then(Value::as_str));
if let Some(output) = output {
if parse_blob_ref(output).is_none() {
stage.output_bytes = Some(output.len() as u64);
}
stage.output = Some(output.to_string());
}
// A simulated step (a dry run) answers with its text.
let simulated = outcome
.output
.get("simulated")
.and_then(Value::as_bool)
.unwrap_or(false);
if simulated
&& matches!(
stage.handler,
Some(StageHandler::Prompt | StageHandler::Agent)
)
{
if let Some(text) = outcome.output.get("text").and_then(Value::as_str) {
stage.response = Some(text.to_string());
}
}
// An agent's answer: the `response.<node>` the step wrote
// into the run context, as the prompt step writes it.
if stage.handler == Some(StageHandler::Agent) {
let response = node_name
.as_deref()
.and_then(|name| {
outcome
.context_updates
.get(format!("response.{name}").as_str())
})
.and_then(Value::as_str)
.or_else(|| outcome.output.as_str());
if let Some(response) = response {
stage.response = Some(response.to_string());
}
}
stage.live_streaming = Some(false);
apply_metrics(stage, &outcome.metrics);
if is_final {
stage.completion = Some(completion(&outcome.status, at));
stage.termination = Some(match outcome.status {
Status::TimedOut => fabro_types::CommandTermination::TimedOut,
Status::Cancelled => fabro_types::CommandTermination::Cancelled,
Status::Success
| Status::PartialSuccess { .. }
| Status::Failure(_)
| Status::Skipped => fabro_types::CommandTermination::Exited,
});
} else {
stage.state = StageState::Retrying;
debug!(attempt = attempt.raw(), "attempt returned; a retry follows");
}
}
}
Event::ControlRequested { .. } => {
if let Some(Derived::ControlRequested {
deliverable: true,
answer: Some(answer),
}) = &event.derived
{
self.close_questions(
answer.question.as_deref(),
FiringKey::of_event(event),
at,
);
}
}
// ── Sandbox: the instance (VIEWS.md "Sandbox") ──────────────────
// The run's sandbox is the root invocation's scope. A child
// invocation's scope (a parallel branch) shares or owns another
// one and is not the run's; a re-acquisition (a resume, a
// replaced sandbox) names the current instance.
Event::ScopeAcquired {
sandbox,
duration_ms,
..
} => {
if let Some(projection) = self.root_scope_projection(event) {
let plan = sandbox_plan_of(projection);
projection.sandbox = Some(RunSandbox::ready(
plan.clone(),
sandbox_instance(&plan, sandbox, *duration_ms),
));
}
}
Event::ScopeFailed {
provider,
error,
causes,
duration_ms,
..
} => {
if let Some(projection) = self.root_scope_projection(event) {
let plan = sandbox_plan_of(projection);
let provider = provider
.as_deref()
.and_then(provider_kind)
.unwrap_or_else(|| plan.provider.clone());
projection.sandbox = Some(RunSandbox::failed(plan, RunSandboxFailure {
provider: provider.to_string(),
error: error.clone(),
causes: causes.clone(),
duration_ms: *duration_ms,
}));
}
}
Event::ExecutionStarted { .. }
| Event::TokenEmitted { .. }
| Event::RoutingResolved { .. }
| Event::RouteApplied { .. }
| Event::RetryElapsed { .. }
| Event::NodeExpanded { .. }
| Event::CancelRequested { .. }
| Event::KillRequested { .. } => {}
}
}
/// The projection, when `event` is a scope record of the root
/// The projection, when `event` is a scope record of the root
/// invocation: the run's own sandbox, not a child invocation's.
fn root_scope_projection(&mut self, event: &RunEvent) -> Option<&mut RunProjection> {
let root = self.state.root?;
if event.context.invocation.map(InvocationId::raw) != Some(root) {
return None;
}
self.projection.as_mut()
}
pub(super) fn fold_view(&mut self, view: &ViewEvent, event: &RunEvent, at: DateTime<Utc>) {
let Some(execution) = event.context.execution else {
return;
};
match view {
ViewEvent::VisitStarted { .. } => {
let Some(subject) = event.subject.as_ref() else {
return;
};
self.start_visit(execution, subject, at);
}
ViewEvent::WaitStateChanged { state } => {
if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) {
match state {
WaitState::AwaitingAdmission => {
if stage.state == StageState::Running {
stage.state = StageState::Pending;
}
}
WaitState::Running | WaitState::AwaitingAnswer | WaitState::Cancelling => {
stage.state = StageState::Running;
}
WaitState::AwaitingRetry => stage.state = StageState::Retrying,
}
}
}
ViewEvent::RetryScheduled { .. } => {
if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) {
stage.state = StageState::Retrying;
}
}
ViewEvent::VisitCompleted {
outcome,
executed,
attempts,
} => {
let Some(stage) = self.stage_of(execution, event.subject.as_ref()) else {
return;
};
stage.state = match outcome.status {
Status::Success => StageState::Succeeded,
Status::PartialSuccess { .. } => StageState::PartiallySucceeded,
Status::Failure(_) | Status::TimedOut => StageState::Failed,
Status::Skipped => StageState::Skipped,
Status::Cancelled => StageState::Cancelled,
};
if stage.completion.is_none() || !*executed {
stage.completion = Some(completion(&outcome.status, at));
}
if stage.timing.is_none() {
let wall = stage
.started_at
.map_or(0, |started| timing::elapsed_ms(started, at));
stage.set_authoritative_timing(StageTiming::new(wall, 0, 0));
}
debug!(attempts, "visit completed");
}
ViewEvent::ForkCompleted {
occurrence,
results,
..
} => {
let key = FiringKey::new(occurrence.execution.raw(), occurrence.firing.raw());
let Some(stage_id) = self
.state
.stages
.get(&key)
.map(|stage| stage.stage_id.clone())
else {
return;
};
let Some(projection) = self.projection.as_mut() else {
return;
};
if let Some(stage) = projection.stage_mut(&stage_id) {
stage.parallel_results = Some(
results
.iter()
.map(|result| ParallelBranchResult {
id: result.node.name.to_string(),
index: Some(result.branch.index as usize),
item_label: None,
status: stage_outcome(&result.status),
context_updates: BTreeMap::new(),
})
.collect(),
);
}
}
ViewEvent::ForkStarted { .. }
| ViewEvent::BranchCompleted { .. }
| ViewEvent::RunStalled { .. } => {}
}
}
/// A firing exists: register its stage and, when it is a logical stage,
/// A firing exists: register its stage and, when it is a logical stage,
/// show it.
pub(super) fn start_visit(
&mut self,
execution: ExecutionId,
subject: &Subject,
at: DateTime<Utc>,
) {
let Some(firing) = subject.firing else {
return;
};
let key = FiringKey::new(execution.raw(), firing.raw());
if self.state.stages.contains_key(&key) {
return;
}
let node_name = subject.node.name.to_string();
let visit = visit_of(subject);
let meta_kind = node_meta_kind(&subject.node);
let shown = is_shown(&subject.node);
// Only a shown stage takes a label: a lowering node (a branch's
// parent-side delegate shares its target's name) never competes with
// the stage it stands for.
let mut stage_id = StageId::new(node_name.clone(), visit);
if shown {
stage_id = stage_label(&node_name, visit, execution, &self.state.labels);
self.state.labels.insert(stage_id.to_string());
}
self.state.stages.insert(key, StageRef {
stage_id: stage_id.clone(),
shown,
node_name,
visit,
});
if !shown {
return;
}
let branch = self
.state
.executions
.get(&execution.raw())
.and_then(|invocation| self.state.invocations.get(invocation))
.and_then(|invocation| invocation.branch.clone());
let Some(projection) = self.projection.as_mut() else {
return;
};
let since_created = at
.signed_duration_since(projection.spec.run_id.created_at())
.num_milliseconds()
.max(0);
let ordinal = u32::try_from(since_created)
.unwrap_or(u32::MAX - 1)
.saturating_add(1);
let stage = projection.stage_entry(stage_id.node_id(), visit, first_event_seq(ordinal));
stage.handler = Some(StageHandler::from_handler_type(Some(meta_kind)));
stage.started_at = Some(at);
stage.graph_visit = Some(visit);
stage.state = StageState::Pending;
stage.parallel_branch_id = branch.map(|(group, index)| ParallelBranchId::new(group, index));
}
}
pub(super) fn stage_outcome(status: &Status) -> StageOutcome {
match status {
Status::Success => StageOutcome::Succeeded,
Status::PartialSuccess { .. } => StageOutcome::PartiallySucceeded,
Status::Failure(info) => StageOutcome::Failed {
retry_requested: info.class.as_str() == "retry_requested",
},
Status::Skipped => StageOutcome::Skipped,
Status::Cancelled | Status::TimedOut => StageOutcome::Failed {
retry_requested: false,
},
}
}
pub(super) fn failure_message(status: &Status) -> Option<String> {
match status {
Status::Failure(info)
| Status::PartialSuccess {
underlying: Some(info),
} => Some(info.message.clone()),
Status::TimedOut => Some("the step timed out".to_string()),
Status::Cancelled => Some("the step was cancelled".to_string()),
Status::Success | Status::PartialSuccess { underlying: None } | Status::Skipped => None,
}
}
/// A stage's completion from an attempt's status: its outcome and, for a
/// failure, the message.
fn completion(status: &Status, at: DateTime<Utc>) -> StageCompletion {
StageCompletion {
outcome: stage_outcome(status),
notes: None,
failure_reason: failure_message(status),
timestamp: at,
}
}
/// The finished attempt's metrics onto its stage: the timing and the usage
/// the backend reported.
fn apply_metrics(stage: &mut StageProjection, metrics: &Metrics) {
let custom = &metrics.custom;
let inference = custom
.get("pebble.inference_ms")
.and_then(Value::as_u64)
.unwrap_or(0);
let tool = custom
.get("pebble.tool_ms")
.and_then(Value::as_u64)
.unwrap_or(0);
let wall = metrics.duration_ms.unwrap_or(0);
let (inference, tool) = match stage.handler {
Some(StageHandler::Prompt) => (wall, 0),
Some(StageHandler::Command) => (0, wall),
_ => (inference, tool),
};
stage.set_authoritative_timing(StageTiming::new(wall, inference, tool).clamped_to_wall());
if let Some(usage) =
usage_of(custom.get("pebble.usage")).or_else(|| usage_of(custom.get("prompt.usage")))
{
stage.usage = usage;
}
if let Some(sessions) = custom
.get("pebble.subagents")
.and_then(|subagents| subagents.get("sessions"))
.and_then(Value::as_array)
{
let mut by_model: Vec<ModelUsage> = Vec::new();
for session in sessions {
let provider = session.get("provider").and_then(Value::as_str);
let model = session.get("model").and_then(Value::as_str);
let Some(usage) = usage_of(session.get("usage")) else {
continue;
};
let Some(model) = model.and_then(|model| model_ref(provider, model)) else {
continue;
};
if let Some(entry) = by_model.iter_mut().find(|entry| entry.model == model) {
entry.usage = entry.usage.saturating_add(usage);
} else {
by_model.push(ModelUsage::new(model, usage));
}
}
if !by_model.is_empty() {
stage.usage_by_model = by_model;
}
}
}

View file

@ -0,0 +1,438 @@
//! The projection of a Petri run: Petri's public events and Fabro's platform
//! records folded into the view Fabro's read side serves.
//!
//! The fold is pure. [`RunView`] holds the [`RunProjection`] the API serves
//! (`GET /runs/{id}/state`, the run list through its summary) and the
//! bookkeeping the fold needs between items ([`FoldState`]): which Petri
//! firing each stage is, which invocation each execution belongs to and
//! whether it is a parallel branch, which stage asked each open question.
//! Both halves are stored by the projector and reloaded for the next pass,
//! so a pass folds only the items past the committed positions.
//!
//! The mapping follows `VIEWS.md`, row by row. The stage key is `(execution,
//! firing)`; Fabro's `StageId` (`node@visit`) is the display label the
//! `RunProjection` keys stages by, and a label two firings would share (two
//! child invocations with the same node name and visit) is made unique by
//! naming the execution. What the matrix leaves default is left default
//! here and named in the crate's README.
//!
//! Every item the fold sees carries the delivery sequence the projector
//! assigned it (`stream_seq`), which a checkpoint keeps as its `seq`. A
//! stage's `first_event_seq`, the key the stage list sorts by, is not the
//! delivery sequence: two logs' records can be committed in an order that
//! differs from their recording times by a few positions, and the view
//! built live must equal the view rebuilt from the records alone. It is the
//! milliseconds from the run's creation to the stage's `visit.started`,
//! plus one, which is the same however the records were delivered.
mod coordinator;
mod engine;
mod model;
mod platform;
mod progress;
mod sandbox;
use std::collections::{BTreeMap, BTreeSet};
use std::fmt;
use std::str::FromStr;
use chrono::{DateTime, TimeZone as _, Utc};
use fabro_store::StagePosition;
use fabro_store::platform_records::StoredPlatformRecord;
use fabro_types::{
RunControlAction, RunDiff, RunId, RunProjection, RunStatus, StageId, StageProjection,
};
use petri_execution::ExecutionId;
use petri_execution::events::{NodeRef, RunEvent, Subject};
use serde::{Deserialize, Deserializer, Serialize, Serializer, de};
use serde_json::Value;
use tracing::debug;
/// One item the projector hands the fold, with its delivery sequence.
pub enum Item<'a> {
Petri(&'a RunEvent),
Platform(&'a StoredPlatformRecord),
}
/// A stage as the fold knows it: its label in the projection, and what it
/// learned about it.
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct StageRef {
pub stage_id: StageId,
/// Whether the stage is a logical one the projection shows, or a
/// lowering node it keeps off the list.
pub shown: bool,
/// The node's instance name and visit, for the collision rule.
pub node_name: String,
pub visit: u32,
}
/// What the fold knows about one invocation.
#[derive(Clone, Debug, Default, Serialize, Deserialize)]
pub struct InvocationRef {
/// The calling execution and firing, for a nested invocation.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parent: Option<(u64, u64)>,
/// The parallel group and branch index, for a branch child.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub branch: Option<(StageId, u32)>,
/// The result the invocation recorded, for the root.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub failure: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output: Option<Value>,
}
/// Whether the run's durable record is whole, as the projector last read it.
#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct RecordHealth {
pub complete: bool,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub incomplete: Vec<String>,
}
/// The fold's bookkeeping between items.
#[derive(Clone, Debug, Default, Serialize, Deserialize)]
pub struct FoldState {
/// Stages by firing.
#[serde(default)]
pub stages: BTreeMap<FiringKey, StageRef>,
/// Labels taken, so a second firing with the same name and visit gets
/// its own.
#[serde(default)]
pub labels: BTreeSet<String>,
#[serde(default)]
pub invocations: BTreeMap<u64, InvocationRef>,
/// Which invocation each execution belongs to.
#[serde(default)]
pub executions: BTreeMap<u64, u64>,
/// Open questions by id: the firing that asked.
#[serde(default)]
pub questions: BTreeMap<String, FiringKey>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub root: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub started_at: Option<u64>,
/// The run's recorded finish, when Petri recorded one.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub finished: Option<String>,
/// The run branch and base sha, when they arrive before `run.started`.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub run_branch: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub base_sha: Option<String>,
#[serde(default)]
pub checkpoints: u32,
/// The run's diff as its `run.diff` record gave it, whichever side of
/// the run's finish it arrived on.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub run_diff: Option<RunDiff>,
#[serde(default)]
pub health: RecordHealth,
/// Firings whose attempt has recorded a finish: what a position-keyed
/// platform record may be streamed behind.
#[serde(default)]
pub finished_firings: BTreeSet<FiringKey>,
/// Whether the run's sandbox still exists after its release
/// (`scope.released` `retained`): kept stopped, or deleted. Absent until
/// the root invocation's lease was released. The view carries the same
/// fact as `RunSandboxInstance.retained`.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub sandbox_retained: Option<bool>,
}
impl FoldState {
/// Whether Petri recorded the run's finish.
#[must_use]
pub fn finished_run(&self) -> bool {
self.finished.is_some()
}
}
/// The view of one run: what the API serves and what the fold keeps.
#[derive(Clone, Debug)]
pub struct RunView {
pub projection: Option<RunProjection>,
pub state: FoldState,
}
impl RunView {
#[must_use]
pub fn new() -> Self {
Self {
projection: None,
state: FoldState::default(),
}
}
/// Fold one item at its delivery sequence.
pub fn fold(&mut self, item: &Item<'_>, stream_seq: u64) {
match item {
Item::Platform(record) => self.fold_platform(record, stream_seq),
Item::Petri(event) => self.fold_petri(event),
}
}
/// The run's projection, once its `run.created` record was folded.
#[must_use]
pub fn projection(&self) -> Option<&RunProjection> {
self.projection.as_ref()
}
fn fold_petri(&mut self, event: &RunEvent) {
let at = millis(event.recorded_at);
if let Some(record) = event.coordinator() {
self.fold_coordinator(record, event, at);
} else if let Some(engine) = event.engine() {
self.fold_engine(engine, event, at);
} else if let Some(view) = event.view() {
self.fold_view(view, event, at);
}
if let Some(projection) = self.projection.as_mut() {
touch(projection, at);
}
}
/// The shown stage an event's subject firing belongs to.
fn stage_of(
&mut self,
execution: ExecutionId,
subject: Option<&Subject>,
) -> Option<&mut StageProjection> {
let firing = subject?.firing?;
let stage = self
.state
.stages
.get(&FiringKey::new(execution.raw(), firing.raw()))?;
if !stage.shown {
return None;
}
let stage_id = stage.stage_id.clone();
self.projection.as_mut()?.stage_mut(&stage_id)
}
}
impl Default for RunView {
fn default() -> Self {
Self::new()
}
}
// ── Shared by the folds ─────────────────────────────────────────────────
/// Apply a status transition; one the lifecycle refuses is logged and
/// skipped, since the view never fails the run.
fn apply_status(projection: &mut RunProjection, status: RunStatus, at: DateTime<Utc>) {
if let Err(error) = projection.try_apply_status(status, at) {
debug!(error = %error, "status transition not applied to the Petri projection");
}
}
fn touch(projection: &mut RunProjection, at: DateTime<Utc>) {
if at > projection.last_event_at {
projection.last_event_at = at;
}
}
/// A control the run acknowledged: the pending control is cleared when it
/// is the one that landed.
fn settle_control(projection: &mut RunProjection, action: RunControlAction) {
if projection.pending_control == Some(action) {
projection.pending_control = None;
}
}
/// The key of a stage: the execution and firing of the visit it shows. The
/// same fact a positioned platform record carries as its `StagePosition`.
/// It is written `<execution>:<firing>`, which is how the stored fold
/// state keys its maps.
#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct FiringKey {
pub execution: u64,
pub firing: u64,
}
impl FiringKey {
#[must_use]
pub fn new(execution: u64, firing: u64) -> Self {
Self { execution, firing }
}
/// The firing an event belongs to: its context's execution and its
/// subject's firing, when it has both.
#[must_use]
pub fn of_event(event: &RunEvent) -> Option<Self> {
let execution = event.context.execution?;
let firing = event.subject.as_ref()?.firing?;
Some(Self::new(execution.raw(), firing.raw()))
}
}
impl From<StagePosition> for FiringKey {
fn from(position: StagePosition) -> Self {
Self::new(position.execution, position.firing)
}
}
impl fmt::Display for FiringKey {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{}:{}", self.execution, self.firing)
}
}
/// A firing key that is not `<execution>:<firing>`.
#[derive(Debug, thiserror::Error)]
#[error("a firing key is `<execution>:<firing>`, not {0:?}")]
pub struct ParseFiringKeyError(String);
impl FromStr for FiringKey {
type Err = ParseFiringKeyError;
fn from_str(text: &str) -> Result<Self, Self::Err> {
let invalid = || ParseFiringKeyError(text.to_string());
let (execution, firing) = text.split_once(':').ok_or_else(invalid)?;
Ok(Self::new(
execution.parse().map_err(|_| invalid())?,
firing.parse().map_err(|_| invalid())?,
))
}
}
impl Serialize for FiringKey {
fn serialize<S: Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
serializer.collect_str(self)
}
}
impl<'de> Deserialize<'de> for FiringKey {
fn deserialize<D: Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
String::deserialize(deserializer)?
.parse()
.map_err(de::Error::custom)
}
}
/// Which firing of its node a subject is, 1-based.
#[must_use]
pub fn visit_of(subject: &Subject) -> u32 {
subject.visit.unwrap_or(1).max(1)
}
/// The role a frontend gave a node under `meta.kind`, or the empty string.
fn node_meta_kind(node: &NodeRef) -> &str {
node.meta.get("kind").and_then(Value::as_str).unwrap_or("")
}
/// Whether a node is a logical stage the projection shows, or a lowering
/// node it keeps off the list: one a frontend marked synthetic, or a
/// parallel branch's delegate.
#[must_use]
pub fn is_shown(node: &NodeRef) -> bool {
let synthetic = node
.meta
.get("synthetic")
.and_then(Value::as_bool)
.unwrap_or(false);
!synthetic && node_meta_kind(node) != "parallel.branch"
}
/// The label a shown firing takes, which is the stage id the projection
/// keys it by: `node@visit`, or `node/e<execution>@visit` when another
/// execution's firing already took that label. `taken` is every label given
/// so far; the caller adds the one returned. The interview adapter labels a
/// question's stage through this same rule, so the stage a question names
/// is the stage the projection shows.
#[must_use]
pub fn stage_label(
node_name: &str,
visit: u32,
execution: ExecutionId,
taken: &BTreeSet<String>,
) -> StageId {
let stage_id = StageId::new(node_name.to_string(), visit);
if taken.contains(&stage_id.to_string()) {
return StageId::new(format!("{node_name}/e{}", execution.raw()), visit);
}
stage_id
}
fn millis(recorded_at: u64) -> DateTime<Utc> {
Utc.timestamp_millis_opt(i64::try_from(recorded_at).unwrap_or(i64::MAX))
.single()
.unwrap_or_default()
}
/// The run id a Petri run key names.
#[must_use]
pub fn run_id_of(key: &str) -> Option<RunId> {
key.parse().ok()
}
#[cfg(test)]
mod tests {
use fabro_store::platform_records::{PlatformRecord, RunCreatedRecord};
use fabro_types::test_support as types_support;
use petri_runtime::driver::BranchRole;
use petri_runtime::ir::{FiringId, NodeId};
use super::*;
#[test]
fn a_taken_label_is_made_unique_by_the_execution() {
let mut view = RunView::new();
let created = StoredPlatformRecord {
seq: 1,
recorded_at: 1_000,
record: PlatformRecord::RunCreated(RunCreatedRecord {
spec: types_support::test_run_spec(),
title: Some("A run".to_string()),
parent_id: None,
retried_from: None,
web_url: None,
}),
position: None,
};
view.fold(&Item::Platform(&created), 1);
let subject = |name: &str| Subject {
node: NodeRef {
id: NodeId::new(1),
name: name.into(),
kind: "attractor/command".into(),
meta: serde_json::json!({ "kind": "command" }),
},
firing: Some(FiringId::new(4)),
visit: Some(1),
attempt: None,
generation: None,
branch: BranchRole::None,
};
view.start_visit(ExecutionId::new(1), &subject("build"), millis(2_000));
view.start_visit(ExecutionId::new(2), &subject("build"), millis(3_000));
let labels: Vec<String> = view
.projection()
.expect("the run was created")
.iter_stages()
.map(|(id, _)| id.to_string())
.collect();
assert_eq!(labels, vec!["build@1", "build/e2@1"]);
assert_eq!(view.state.stages.len(), 2);
}
#[test]
fn a_firing_key_is_stored_as_execution_colon_firing() {
let mut stages: BTreeMap<FiringKey, u32> = BTreeMap::new();
stages.insert(FiringKey::new(3, 7), 1);
let json = serde_json::to_string(&stages).expect("the map encodes");
assert_eq!(json, r#"{"3:7":1}"#);
let back: BTreeMap<FiringKey, u32> = serde_json::from_str(&json).expect("the map decodes");
assert_eq!(back, stages);
assert!("3-7".parse::<FiringKey>().is_err());
assert_eq!(
FiringKey::from(StagePosition {
execution: 3,
firing: 7,
}),
FiringKey::new(3, 7)
);
}
}

View file

@ -0,0 +1,41 @@
//! The model and usage facts the records carry, as the view names them.
use fabro_types::ModelRef;
use lithos_llm::catalog::{ModelId, ProviderId};
use lithos_llm::types::Usage;
use serde_json::Value;
pub(super) fn usage_of(value: Option<&Value>) -> Option<Usage> {
serde_json::from_value(value?.clone()).ok()
}
/// `provider/model` into its parts, or the model alone.
pub(super) fn split_model(model: &str) -> (Option<&str>, &str) {
match model.split_once('/') {
Some((provider, model)) if !provider.is_empty() && !model.is_empty() => {
(Some(provider), model)
}
_ => (None, model),
}
}
pub(super) fn model_ref(provider: Option<&str>, model: &str) -> Option<ModelRef> {
let provider = provider.filter(|provider| !provider.is_empty())?;
Some(ModelRef::new(
ProviderId::new(provider),
ModelId::new(model),
))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_model_selector_splits_into_provider_and_model() {
assert_eq!(split_model("openai/gpt-5.4"), (Some("openai"), "gpt-5.4"));
assert_eq!(split_model("gpt-5.4"), (None, "gpt-5.4"));
assert!(model_ref(None, "gpt-5.4").is_none());
assert!(model_ref(Some("openai"), "gpt-5.4").is_some());
}
}

View file

@ -0,0 +1,319 @@
//! The platform records folded into the view: the run's creation, its
//! lifecycle, its title and parent, the branch and identity the first
//! checkpoint recorded, every checkpoint and artifact, the run's diff, and
//! the pull request (VIEWS.md "Run", "Checkpoints", "Artifacts", "Pull
//! request").
use chrono::{DateTime, Utc};
use fabro_store::platform_records::{
PlatformRecord, RunLifecycleKind, RunLifecycleRecord, StoredPlatformRecord,
};
use fabro_types::{
CheckpointRecord as ViewCheckpoint, Conclusion, FailureCategory, FailureDetail, FailureReason,
PullRequestCreation, PullRequestCreationStatus, PullRequestLink, RunApproval, RunApprovalState,
RunArtifact, RunControlAction, RunDiff, RunFailure, RunProjection, RunSandbox, RunStatus,
StageOutcome, format_blob_ref,
};
use tracing::debug;
use super::sandbox::sandbox_plan;
use super::{FiringKey, RunView, apply_status, millis, settle_control, touch};
impl RunView {
pub(super) fn fold_platform(&mut self, stored: &StoredPlatformRecord, stream_seq: u64) {
let at = millis(stored.recorded_at);
if let PlatformRecord::RunCreated(created) = &stored.record {
let title = created
.title
.clone()
.unwrap_or_else(|| fabro_types::infer_run_title(created.spec.graph.goal()));
let mut projection = RunProjection::new(title, created.spec.clone(), at);
projection.parent_id = created.parent_id;
projection.retried_from = created.retried_from;
projection.web_url.clone_from(&created.web_url);
projection.sandbox = Some(RunSandbox::planned(sandbox_plan(
&projection.spec.settings.run.environment,
)));
self.projection = Some(projection);
return;
}
let Some(projection) = self.projection.as_mut() else {
debug!(
seq = stored.seq,
kind = %stored.record.kind(),
"platform record before run.created; not folded"
);
return;
};
touch(projection, at);
match &stored.record {
PlatformRecord::RunLifecycle(record) => fold_lifecycle(projection, record, at),
PlatformRecord::RunTitle(record) => projection.title.clone_from(&record.title),
PlatformRecord::RunParent(record) => projection.parent_id = record.parent_id,
PlatformRecord::RunArchived => projection.archived_at = Some(at),
PlatformRecord::RunUnarchived => projection.archived_at = None,
PlatformRecord::RunSuperseded(record) => {
projection.superseded_by = Some(record.new_run_id);
}
PlatformRecord::RunCreated(_)
| PlatformRecord::RunNotice(_)
| PlatformRecord::InterviewAnswered(_)
| PlatformRecord::NotificationSent(_)
| PlatformRecord::RunPaired(_) => {}
PlatformRecord::RunBranch(record) => {
self.state.run_branch.clone_from(&record.run_branch);
self.state.base_sha.clone_from(&record.base_sha);
if let Some(start) = projection.start.as_mut() {
start.run_branch.clone_from(&record.run_branch);
start.base_sha.clone_from(&record.base_sha);
}
}
PlatformRecord::GitIdentity(record) => {
projection.git_identity = Some(record.identity.clone());
}
PlatformRecord::Checkpoint(record) => {
self.state.checkpoints = self.state.checkpoints.saturating_add(1);
let stage = self
.state
.stages
.get(&FiringKey::new(record.execution, record.firing));
let current_node = stage.map_or_else(String::new, |stage| stage.node_name.clone());
let stage_id = stage
.filter(|stage| stage.shown)
.map(|stage| stage.stage_id.clone());
let checkpoint = fabro_types::Checkpoint {
timestamp: at,
current_node: current_node.clone(),
git_commit_sha: record.git_commit_sha.clone(),
};
// The patch stays in the blob table; the view carries its
// reference for a reader to resolve.
let patch = record.patch_blob.as_ref().map(format_blob_ref);
if let Some(stage) = stage_id.and_then(|stage_id| projection.stage_mut(&stage_id)) {
if patch.is_some() {
stage.diff.clone_from(&patch);
}
}
projection.checkpoints.push(ViewCheckpoint {
seq: u32::try_from(stream_seq).unwrap_or(u32::MAX),
checkpoint,
diff: RunDiff {
patch,
summary: record.diff_summary,
},
});
}
PlatformRecord::ArtifactCollected(record) => {
let stage = self
.state
.stages
.get(&FiringKey::new(record.execution, record.firing));
let Some(stage_id) = stage.map(|stage| stage.stage_id.clone()) else {
debug!(
seq = stored.seq,
path = record.path,
"artifact record for an unknown firing; not folded"
);
return;
};
projection.artifacts.push(RunArtifact {
stage_id,
retry: record.attempt,
relative_path: record.path.clone(),
size: record.bytes,
blob: record.blob,
});
}
PlatformRecord::RunDiff(record) => {
let diff = RunDiff {
patch: record.patch_blob.as_ref().map(format_blob_ref),
summary: record.diff_summary,
};
if let Some(conclusion) = projection.conclusion.as_mut() {
conclusion.diff = diff.clone();
}
self.state.run_diff = Some(diff);
}
PlatformRecord::PullRequestRequested(record) => {
projection.pull_request_creation = Some(PullRequestCreation {
id: record.creation_id,
status: PullRequestCreationStatus::Pending,
model: record.model.clone(),
force: record.force,
requested_at: at,
updated_at: at,
pull_request: None,
error: None,
});
}
PlatformRecord::PullRequestCreated(record) => {
let link = PullRequestLink {
owner: record.owner.clone(),
repo: record.repo.clone(),
number: record.number,
};
projection.pull_request = Some(link.clone());
if let Some(creation) = projection
.pull_request_creation
.as_mut()
.filter(|creation| creation.is_pending())
{
creation.succeed(link, at);
}
}
PlatformRecord::PullRequestFailed(record) => {
if let Some(creation) =
projection
.pull_request_creation
.as_mut()
.filter(|creation| {
creation.is_pending()
&& record
.creation_id
.is_none_or(|creation_id| creation_id == creation.id)
})
{
creation.fail(record.error.clone(), at);
}
}
PlatformRecord::PullRequestLinked(record) => {
let link = record.link();
projection.pull_request = Some(link.clone());
if let Some(creation) = projection
.pull_request_creation
.as_mut()
.filter(|creation| creation.is_pending())
{
creation.succeed(link, at);
}
}
PlatformRecord::PullRequestUnlinked(_) => {
projection.pull_request = None;
projection.pull_request_creation = None;
}
}
}
}
fn fold_lifecycle(projection: &mut RunProjection, record: &RunLifecycleRecord, at: DateTime<Utc>) {
use RunLifecycleKind as Kind;
match record.transition {
Kind::Submitted => apply_status(projection, RunStatus::Submitted, at),
Kind::StartRequested => {}
Kind::Unpaused => {
apply_status(projection, projection.status.unpaused(), at);
settle_control(projection, RunControlAction::Unpause);
}
Kind::Pending => {
if let Some(status) = record.status {
apply_status(projection, status, at);
}
projection.approval = Some(RunApproval {
state: RunApprovalState::Pending,
requested_at: at,
decided_at: None,
denial_reason: None,
});
}
Kind::Approved => {
if let Some(approval) = projection.approval.as_mut() {
approval.state = RunApprovalState::Approved;
approval.decided_at = Some(at);
}
}
Kind::Denied => {
if let Some(approval) = projection.approval.as_mut() {
approval.state = RunApprovalState::Denied;
approval.decided_at = Some(at);
approval.denial_reason.clone_from(&record.reason);
}
apply_status(
projection,
RunStatus::Failed {
reason: FailureReason::ApprovalDenied,
},
at,
);
}
Kind::Runnable => {
// A run left in flight by a restart goes back to the queue: the
// resume's `runnable` steps back from wherever the run stood.
let in_flight = matches!(
projection.status,
RunStatus::Starting
| RunStatus::Running
| RunStatus::Blocked { .. }
| RunStatus::Paused { .. }
);
if in_flight && record.status == Some(RunStatus::Runnable) {
projection.status = RunStatus::Runnable;
projection.status_updated_at = at;
} else if let Some(status) = record.status {
apply_status(projection, status, at);
}
}
Kind::Blocked => {
// A block that lands while the run is paused waits behind the
// pause: the unpause restores it.
match (projection.status, record.status) {
(RunStatus::Paused { .. }, Some(RunStatus::Blocked { blocked_reason })) => {
apply_status(
projection,
RunStatus::Paused {
prior_block: Some(blocked_reason),
},
at,
);
}
(_, Some(status)) => apply_status(projection, status, at),
(_, None) => {}
}
}
Kind::Unblocked => {
let status = match projection.status {
RunStatus::Paused { .. } => RunStatus::Paused { prior_block: None },
_ => record.status.unwrap_or(RunStatus::Running),
};
apply_status(projection, status, at);
}
Kind::Starting | Kind::Running | Kind::Removing | Kind::Dead => {
if let Some(status) = record.status {
apply_status(projection, status, at);
}
}
Kind::Paused => {
apply_status(projection, projection.status.paused(), at);
settle_control(projection, RunControlAction::Pause);
}
Kind::Succeeded | Kind::Failed => {
if let Some(status) = record.status {
apply_status(projection, status, at);
}
projection.pending_control = None;
if projection.conclusion.is_none() {
let (outcome, failure) = match record.status {
Some(RunStatus::Failed { reason }) => (
StageOutcome::Failed {
retry_requested: false,
},
Some(RunFailure {
reason,
detail: FailureDetail::new(
record
.reason
.clone()
.unwrap_or_else(|| "the run failed".to_string()),
FailureCategory::Deterministic,
),
}),
),
_ => (StageOutcome::Succeeded, None),
};
projection.conclusion = Some(Conclusion::outcome_only(at, outcome, failure));
}
}
Kind::CancelRequested => projection.pending_control = Some(RunControlAction::Cancel),
Kind::PauseRequested => projection.pending_control = Some(RunControlAction::Pause),
Kind::UnpauseRequested => projection.pending_control = Some(RunControlAction::Unpause),
}
}

View file

@ -0,0 +1,573 @@
//! A step's progress records folded into its stage: a command's log lines,
//! the payloads the Attractor steps emit (the prompt and its completion,
//! the fallback plan, the tools a session was offered, a parallel branch's
//! start), Pebble's coding-agent envelope, and a gate's question (VIEWS.md
//! "Agent activity", "Questions").
use chrono::{DateTime, Utc};
use fabro_types::{
BlockedReason, CodingAgentEvent, CodingEvent, InterviewOption, InterviewQuestionRecord,
PendingInterviewRecord, ReviewTarget, ReviewTargetKind, RunStatus, StageInferenceProjection,
StageModelUsage, StageProjection, ToolCategory, ToolSource, ToolSummary, timing,
};
use lithos_llm::types::{ReasoningEffort, Speed, Usage};
use petri_execution::ExecutionId;
use petri_execution::events::{Parsed, RunEvent};
use petri_runtime::ir::StepEvent;
use petri_runtime::steps::QuestionReference;
use serde::Deserialize;
use serde_json::Value;
use tracing::debug;
use super::model::{model_ref, split_model};
use super::{FiringKey, RunView, apply_status};
use crate::interview::question_type;
impl RunView {
pub(super) fn fold_progress(
&mut self,
execution: ExecutionId,
event: &RunEvent,
ev: &StepEvent,
at: DateTime<Utc>,
) {
match ev {
StepEvent::Log { line, .. } => {
if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) {
let output = stage.output.get_or_insert_default();
output.push_str(line);
output.push('\n');
stage.output_bytes = Some(output.len() as u64);
stage.live_streaming = Some(true);
}
}
StepEvent::Artifact { .. } => {}
StepEvent::Custom(payload) => {
if let Some(parsed) = event.parsed() {
self.fold_parsed(execution, event, parsed, at);
return;
}
let progress = match Progress::deserialize(payload) {
Ok(progress) => progress,
Err(error) => {
debug!(error = %error, "a progress payload did not decode; skipped");
return;
}
};
match progress {
Progress::Pebble { event: envelope } => {
self.fold_pebble(execution, event, &envelope, at);
}
Progress::Prompt { prompt, model } => {
if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) {
stage.prompt = prompt;
if let Some(model) = model.as_deref() {
let (provider, model_id) = split_model(model);
stage.provider_used = Some(StageModelUsage::new(
StageModelUsage::MODE_PROMPT,
provider.map(str::to_string),
Some(model_id.to_string()),
));
stage.model = model_ref(provider, model_id);
}
}
}
Progress::PromptCompleted { response, usage } => {
if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) {
stage.response = response;
if let Some(usage) = usage {
stage.usage = usage;
}
}
}
// `routes[0]` is the original route: what the stage was
// asked to run on, with its request controls.
Progress::FallbackPlan { routes } => {
if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) {
if let Some(route) = routes.into_iter().next() {
stage.model = model_ref(Some(&route.provider), &route.model);
stage.provider_used = Some(route.usage());
}
}
}
// The tools a native session was offered, once per
// session (VIEWS.md "Agent activity", tools available):
// the stage's list is the union over its sessions, by
// name, in the order the sessions listed them.
Progress::Tools { tools } => {
if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) {
for tool in tools {
if !stage
.agent_tools
.iter()
.any(|known| known.name == tool.name)
{
stage.agent_tools.push(tool.summary());
}
}
}
}
Progress::BranchStarted {
invocation,
index,
occurrence,
} => {
let group = self
.state
.stages
.get(&FiringKey::new(execution.raw(), occurrence.firing))
.map(|stage| stage.stage_id.clone());
if let Some(group) = group {
self.state.invocations.entry(invocation).or_default().branch =
Some((group, index));
}
}
Progress::Other => {}
}
}
}
}
fn fold_parsed(
&mut self,
execution: ExecutionId,
event: &RunEvent,
parsed: &Parsed,
at: DateTime<Utc>,
) {
match parsed {
Parsed::Question { question } => {
let Some(subject) = event.subject.as_ref() else {
return;
};
let Some(firing) = subject.firing else {
return;
};
let key = FiringKey::new(execution.raw(), firing.raw());
let label = self.state.stages.get(&key).map_or_else(
|| subject.node.name.to_string(),
|stage| stage.stage_id.to_string(),
);
self.state.questions.insert(question.id.clone(), key);
let Some(projection) = self.projection.as_mut() else {
return;
};
projection
.pending_interviews
.insert(question.id.clone(), PendingInterviewRecord {
question: InterviewQuestionRecord {
id: question.id.clone(),
text: question.text.clone(),
stage: label,
question_type: question_type(question),
options: question
.options
.iter()
.map(|option| InterviewOption {
key: option.key.clone(),
label: option.label.clone(),
description: option.description.clone(),
preview: option.preview.clone(),
})
.collect(),
allow_freeform: question.freeform,
timeout_seconds: question
.timeout_ms
.map(|timeout| timeout as f64 / 1000.0),
context_display: question.context.clone(),
review_target: question.reference.as_ref().and_then(review_target),
},
started_at: at,
});
apply_status(
projection,
RunStatus::Blocked {
blocked_reason: BlockedReason::HumanInputRequired,
},
at,
);
}
Parsed::QuestionExpired { expired } => {
self.close_questions(Some(expired.question.as_str()), None, at);
}
Parsed::Note { .. } => {}
}
}
/// Close one question by id, or every question of a firing, and unblock
/// Close one question by id, or every question of a firing, and unblock
/// the run when none is left.
pub(super) fn close_questions(
&mut self,
question: Option<&str>,
firing: Option<FiringKey>,
at: DateTime<Utc>,
) {
let closed: Vec<String> = match (question, firing) {
(Some(question), _) => vec![question.to_string()],
(None, Some(firing)) => self
.state
.questions
.iter()
.filter(|(_, asked_by)| **asked_by == firing)
.map(|(id, _)| id.clone())
.collect(),
(None, None) => Vec::new(),
};
for id in &closed {
self.state.questions.remove(id);
}
let Some(projection) = self.projection.as_mut() else {
return;
};
for id in &closed {
projection.pending_interviews.remove(id);
}
if projection.pending_interviews.is_empty()
&& matches!(projection.status, RunStatus::Blocked { .. })
{
apply_status(projection, RunStatus::Running, at);
}
}
fn fold_pebble(
&mut self,
execution: ExecutionId,
event: &RunEvent,
envelope: &CodingAgentEvent,
at: DateTime<Utc>,
) {
let Some(stage) = self.stage_of(execution, event.subject.as_ref()) else {
return;
};
let agent = stage.agent.get_or_insert_default();
agent.apply(envelope);
if stage.completion.is_none() {
stage.usage = agent.usage.saturating_add(agent.descendant_usage());
}
// A tool the stage's list names was called, by any of its sessions.
if let CodingEvent::ToolCallStarted { tool_name, .. } = &envelope.event {
if let Some(tool) = stage
.agent_tools
.iter_mut()
.find(|tool| tool.name == *tool_name)
{
tool.invoked = true;
}
}
let is_root = envelope.parent_session_id.is_none();
#[expect(
clippy::wildcard_enum_match_arm,
reason = "pebble's event vocabulary is non-exhaustive and only some events project"
)]
match &envelope.event {
CodingEvent::SessionStarted {
provider, model, ..
} if is_root => {
stage.provider_used = Some(StageModelUsage::new(
StageModelUsage::MODE_AGENT,
provider.clone(),
model.clone(),
));
if let Some(model) = model.as_deref() {
stage.model = model_ref(provider.as_deref(), model);
}
}
CodingEvent::LlmRequestStarted { requested_model } if is_root => {
stage.inference = Some(StageInferenceProjection {
session_id: envelope.session_id.clone(),
started_at: at,
requested_model: requested_model.clone(),
first_output_at: None,
first_output_kind: None,
retries: 0,
});
}
CodingEvent::LlmFirstOutput { kind } => {
if let Some(inference) = stage.inference.as_mut() {
if inference.session_id == envelope.session_id {
inference.first_output_at = Some(at);
inference.first_output_kind = Some(*kind);
}
}
}
CodingEvent::LlmRetry { .. } => {
if let Some(inference) = stage.inference.as_mut() {
if inference.session_id == envelope.session_id {
inference.retries = inference.retries.saturating_add(1);
inference.first_output_at = None;
inference.first_output_kind = None;
}
}
}
CodingEvent::AssistantMessage { model, .. } => {
if is_root {
if let Some(provider) = stage
.provider_used
.as_ref()
.and_then(|used| used.provider.as_deref())
{
stage.model = model_ref(Some(provider), model);
}
}
close_inference(stage, &envelope.session_id, at);
}
CodingEvent::Error { .. } | CodingEvent::RoundInterrupted { .. } => {
close_inference(stage, &envelope.session_id, at);
}
CodingEvent::SessionEnded => {
close_inference(stage, &envelope.session_id, at);
stage.close_tool_batch_for_session(&envelope.session_id, at);
}
CodingEvent::ToolCallStarted { tool_call_id, .. } if is_root => {
stage.open_tool_call(envelope.session_id.clone(), tool_call_id.clone(), at);
}
CodingEvent::ToolCallCompleted { tool_call_id, .. } if is_root => {
stage.close_tool_call(&envelope.session_id, tool_call_id, at);
}
_ => {}
}
}
}
/// The `StepEvent::Custom` payloads the fold reads, by their `kind`. The
/// kinds are Petri's, and a test holds each literal to the constant the
/// Attractor steps export, so a rename there fails here rather than
/// projecting nothing. A payload of another kind, or of no kind, is
/// `Other`.
#[derive(Deserialize)]
#[serde(tag = "kind")]
enum Progress {
/// Pebble's coding-agent envelope, forwarded by the agent step.
#[serde(rename = "pebble")]
Pebble { event: Box<CodingAgentEvent> },
/// The prompt step before its first model call: the prompt, and the
/// `provider/model` selector it runs on.
#[serde(rename = "attractor.prompt")]
Prompt {
#[serde(default)]
prompt: Option<String>,
#[serde(default)]
model: Option<String>,
},
/// The prompt step after its last model call.
#[serde(rename = "attractor.prompt.completed")]
PromptCompleted {
#[serde(default)]
response: Option<String>,
#[serde(default)]
usage: Option<Usage>,
},
/// A stage's fallback plan, once per stage: `routes[0]` is the
/// original route.
#[serde(rename = "attractor.fallback.plan")]
FallbackPlan {
#[serde(default)]
routes: Vec<PlannedRoute>,
},
/// The tools a native session was offered, once per session.
#[serde(rename = "attractor.tools")]
Tools {
#[serde(default)]
tools: Vec<OfferedTool>,
},
/// A parallel branch's child started: which fork visit it belongs to,
/// its index, and the child invocation.
#[serde(rename = "attractor.parallel.branch.started")]
BranchStarted {
invocation: u64,
index: u32,
occurrence: ForkOccurrence,
},
#[serde(other)]
Other,
}
/// One route of a fallback plan: the provider and model, with the request
/// controls the route carries.
#[derive(Deserialize)]
struct PlannedRoute {
provider: String,
model: String,
#[serde(default)]
reasoning_effort: Option<ReasoningEffort>,
#[serde(default)]
speed: Option<Speed>,
}
impl PlannedRoute {
/// The route as the stage's model usage: an agent route with its
/// controls.
fn usage(self) -> StageModelUsage {
StageModelUsage {
reasoning_effort: self.reasoning_effort,
speed: self.speed,
..StageModelUsage::new(
StageModelUsage::MODE_AGENT,
Some(self.provider),
Some(self.model),
)
}
}
}
/// The fork visit a branch belongs to: the fork step's firing in the
/// branch's execution.
#[derive(Deserialize)]
struct ForkOccurrence {
firing: u64,
}
/// One tool of an `attractor.tools` payload: the name and description as
/// recorded, Pebble's `source` as it is, and Petri's origin category
/// (`builtin`, `mcp`, `host`, `question`, `subagent`).
#[derive(Deserialize)]
struct OfferedTool {
name: String,
#[serde(default)]
description: String,
/// Left as recorded: a source Pebble adds later still lists the tool,
/// under the default source.
#[serde(default)]
source: Value,
#[serde(default)]
category: Option<String>,
}
impl OfferedTool {
/// The tool as the stage's list carries it. Pebble's behavioural
/// category is kept where Petri's says which (a sub-agent tool); every
/// other tool is `other`, because the payload carries Petri's origin
/// category, not Pebble's permission class. `invoked` starts false and
/// flips on the session's `ToolCallStarted`.
fn summary(self) -> ToolSummary {
let category = match self.category.as_deref() {
Some("subagent") => ToolCategory::Subagent,
_ => ToolCategory::Other,
};
ToolSummary {
name: self.name,
description: self.description,
source: serde_json::from_value::<ToolSource>(self.source).unwrap_or_default(),
category,
invoked: false,
}
}
}
/// The question's `reference` as Fabro's review target, when it is one
/// Fabro's validation admits (a `document`, or a reference without a kind,
/// with a label and an absolute HTTP URL within Fabro's limits).
fn review_target(reference: &QuestionReference) -> Option<ReviewTarget> {
let kind = match reference.kind.as_deref() {
Some("document") | None => ReviewTargetKind::Document,
Some(_) => return None,
};
ReviewTarget::new(reference.label.clone(), reference.url.clone(), kind).ok()
}
fn close_inference(stage: &mut StageProjection, session_id: &str, at: DateTime<Utc>) {
let open = stage
.inference
.as_ref()
.is_some_and(|inference| inference.session_id == session_id);
if !open {
return;
}
if let Some(inference) = stage.inference.take() {
stage.accumulate_inference_ms(timing::elapsed_ms(inference.started_at, at));
}
}
#[cfg(test)]
mod tests {
use petri_attractor_steps::{fallback, parallel, pebble, prompt};
use serde_json::json;
use super::*;
/// The kinds the fold matches are the ones the Attractor steps emit.
#[test]
fn the_progress_kinds_are_petris() {
let kind_of = |value: Value| -> &'static str {
match Progress::deserialize(&value).expect("a known kind decodes") {
Progress::Pebble { .. } => "pebble",
Progress::Prompt { .. } => "prompt",
Progress::PromptCompleted { .. } => "prompt.completed",
Progress::FallbackPlan { .. } => "fallback.plan",
Progress::Tools { .. } => "tools",
Progress::BranchStarted { .. } => "branch.started",
Progress::Other => "other",
}
};
assert_eq!(kind_of(json!({ "kind": prompt::PROMPT_EVENT })), "prompt");
assert_eq!(
kind_of(json!({ "kind": prompt::COMPLETED_EVENT })),
"prompt.completed"
);
assert_eq!(
kind_of(json!({ "kind": fallback::PLAN_EVENT })),
"fallback.plan"
);
assert_eq!(kind_of(json!({ "kind": pebble::tools::EVENT })), "tools");
assert_eq!(
kind_of(json!({
"kind": parallel::BRANCH_STARTED_EVENT,
"invocation": 3,
"index": 1,
"occurrence": { "fork": "fan", "firing": 7 },
})),
"branch.started"
);
assert_eq!(
kind_of(json!({ "kind": parallel::BRANCH_COMPLETED_EVENT })),
"other"
);
}
/// A plan's original route carries its request controls onto the
/// stage's model usage.
#[test]
fn a_fallback_plan_route_keeps_its_controls() {
let progress = Progress::deserialize(&json!({
"kind": fallback::PLAN_EVENT,
"routes": [
{ "position": 0, "provider": "openai", "model": "gpt-5.4",
"reasoning_effort": "high", "speed": null },
{ "position": 1, "provider": "anthropic", "model": "claude" },
],
}))
.expect("the plan decodes");
let Progress::FallbackPlan { routes } = progress else {
panic!("not a plan");
};
let usage = routes.into_iter().next().expect("a route").usage();
assert_eq!(usage.mode, StageModelUsage::MODE_AGENT);
assert_eq!(usage.provider.as_deref(), Some("openai"));
assert_eq!(usage.model.as_deref(), Some("gpt-5.4"));
assert_eq!(usage.reasoning_effort, Some(ReasoningEffort::High));
assert_eq!(usage.speed, None);
}
/// A tool lists under Petri's category, with its source as recorded.
#[test]
fn an_offered_tool_maps_to_its_summary() {
let progress = Progress::deserialize(&json!({
"kind": pebble::tools::EVENT,
"tools": [
{ "name": "spawn_agent", "description": "a child", "source": "native",
"category": "subagent" },
{ "name": "Read", "category": "builtin" },
],
}))
.expect("the tools decode");
let Progress::Tools { tools } = progress else {
panic!("not a tool list");
};
let summaries: Vec<ToolSummary> = tools.into_iter().map(OfferedTool::summary).collect();
assert_eq!(summaries[0].category, ToolCategory::Subagent);
assert_eq!(summaries[1].category, ToolCategory::Other);
assert_eq!(summaries[1].description, "");
assert!(!summaries[0].invoked);
}
}

View file

@ -0,0 +1,76 @@
//! The run's sandbox as the view shows it: the plan its settings give, and
//! the instance Petri's scope records name (VIEWS.md "Sandbox").
use fabro_types::settings::run::RunEnvironmentSettings;
use fabro_types::{
RunProjection, RunSandboxInstance, RunSandboxPlan, RunSandboxRuntime, SandboxProviderKind,
};
use petri_runtime::ir::SandboxInstance;
pub(super) fn sandbox_plan(settings: &RunEnvironmentSettings) -> RunSandboxPlan {
RunSandboxPlan {
provider: settings.provider.clone(),
image: (settings.provider == SandboxProviderKind::DOCKER)
.then(|| settings.image.docker.clone())
.flatten()
.filter(|image| !image.is_empty()),
snapshot: None,
}
}
/// The plan the projection's sandbox carries, or the one its environment
/// settings give when no sandbox was projected yet.
pub(super) fn sandbox_plan_of(projection: &RunProjection) -> RunSandboxPlan {
projection.sandbox.as_ref().map_or_else(
|| sandbox_plan(&projection.spec.settings.run.environment),
|sandbox| sandbox.plan().clone(),
)
}
/// Fabro's name for the provider Petri's `scope.acquired` names: Petri's
/// `host` is Fabro's `local`; every other kind is spelled the same. `None`
/// for a name that is no provider kind.
pub(super) fn provider_kind(provider: &str) -> Option<SandboxProviderKind> {
if provider == "host" {
return Some(SandboxProviderKind::LOCAL);
}
SandboxProviderKind::try_new(provider).ok()
}
/// The run's sandbox instance from Petri's record of the scope's
/// acquisition: the provider, the provider's id for the sandbox (what a
/// reconnect attaches by), its image and snapshot when the provider knows
/// them, the working directory, and how long the acquisition took. The
/// clone fields stay unset: Petri's checkout copies the bound repository
/// into the workspace and is not a clone Fabro made, and the workspace
/// roots are the provider's own layout, read live. `retained` waits for
/// the scope's release.
pub(super) fn sandbox_instance(
plan: &RunSandboxPlan,
sandbox: &SandboxInstance,
ready_duration_ms: u64,
) -> RunSandboxInstance {
RunSandboxInstance {
provider: provider_kind(&sandbox.provider)
.unwrap_or_else(|| plan.provider.clone()),
image: sandbox
.image
.as_ref()
.map(ToString::to_string)
.or_else(|| plan.image.clone()),
snapshot: sandbox.snapshot.as_ref().map(ToString::to_string),
runtime: RunSandboxRuntime {
id: sandbox.instance.to_string(),
working_directory: sandbox.working_directory.to_string(),
repo_cloned: None,
clone_origin_url: None,
clone_branch: None,
workspace_root: None,
repos_root: None,
primary_repo_path: None,
primary_repo_link: None,
},
ready_duration_ms: Some(ready_duration_ms),
retained: None,
}
}

View file

@ -15,10 +15,11 @@
//! rows: the rebuild test in `tests/projection.rs` compares the two.
use std::collections::HashMap;
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use fabro_types::RunId;
use fabro_util::sync;
use petri_execution::events::{RunEvent, RunReplay};
use tokio::sync::Mutex as AsyncMutex;
@ -86,7 +87,7 @@ 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 mut runs = sync::lock(&self.runs);
let entry = runs.entry(run_id).or_insert_with(|| Entry {
pass: Arc::default(),
touched: Instant::now(),
@ -99,7 +100,7 @@ impl Caches {
/// 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 runs = sync::lock(&self.runs);
let mut dropped = 0;
runs.retain(|_, entry| {
if entry.touched.elapsed() < idle {
@ -120,13 +121,10 @@ impl Caches {
}
/// Whether a cache is kept for the run: a test's view of the cache.
#[cfg(any(test, feature = "test-support"))]
pub(crate) fn holds(&self, run_id: RunId) -> bool {
let runs = lock(&self.runs);
let runs = sync::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

@ -0,0 +1,343 @@
//! The order one pass streams its new items in.
use std::collections::BTreeSet;
use fabro_store::platform_records::StoredPlatformRecord;
use petri_execution::events::{EventSource, RunEvent};
use petri_runtime::engine::Event;
use crate::projection::{FiringKey, Item};
/// The order one pass streams its new items in, and how many platform
/// records it holds back for a later pass.
///
/// Every item is first ordered by `recorded_at` (stable: the coordinator
/// log before an execution log before a platform record on a tie, and each
/// log's own order kept). A platform record that carries a Petri position
/// (a checkpoint, keyed on `(execution, firing)`) is then placed by that
/// position, not by its clock, because the server stamps the record and the
/// worker stamps Petri's records and the two clocks can tie or invert:
///
/// - before the firing's first `routing.resolved` event in the pass, which is
/// right after the firing's finish (its `step.finished` and the
/// `visit.completed` attached to it) and before the next firing's
/// `visit.started`, which is attached to that routing record;
/// - else after the last event of the firing in the pass;
/// - else, when the firing finished in an earlier pass, before the first event
/// of a later firing (a larger firing id) in the same execution, or where its
/// `recorded_at` put it;
/// - else the record is held back, with every platform record after it, and the
/// pass consumes platform records only up to it. The hook that writes a
/// checkpoint record runs after the driver appended the attempt's finish, but
/// the driver's store writer flushes that record on its own schedule, so the
/// platform record can be committed before its firing's `step.finished`;
/// holding it keeps the stream's order the same live and on a rebuild.
/// Nothing is held once the run has recorded its finish.
///
/// The rule reads only the pass's own items and the firings already
/// finished, so a record is never streamed before its firing's finish and
/// never after the firing's routes.
pub(crate) fn order_items<'a>(
events: &'a [RunEvent],
platform_records: &'a [StoredPlatformRecord],
finished_before: &BTreeSet<FiringKey>,
run_finished: bool,
) -> (Vec<Item<'a>>, usize) {
let finished_in_pass = |at: FiringKey| {
events.iter().any(|event| {
FiringKey::of_event(event) == Some(at)
&& matches!(event.engine(), Some(Event::StepFinished { .. }))
})
};
let finished = |at: FiringKey| finished_in_pass(at) || finished_before.contains(&at);
// Platform records are consumed in seq order: the first one whose firing
// has not finished holds itself and everything after it.
let consumed = if run_finished {
platform_records.len()
} else {
platform_records
.iter()
.position(|record| {
record
.position
.is_some_and(|position| !finished(FiringKey::from(position)))
})
.unwrap_or(platform_records.len())
};
let held = platform_records.len() - consumed;
let platform_records = &platform_records[..consumed];
let mut items: Vec<(u64, u8, Item<'a>)> =
Vec::with_capacity(events.len() + platform_records.len());
for event in events {
let rank = match event.id.source {
EventSource::Coordinator => 0,
EventSource::Execution { .. } => 1,
};
items.push((event.recorded_at, rank, Item::Petri(event)));
}
for record in platform_records {
items.push((record.recorded_at, 2, Item::Platform(record)));
}
items.sort_by_key(|(recorded_at, rank, _)| (*recorded_at, *rank));
let item_firing = |item: &Item<'a>| match item {
Item::Petri(event) => FiringKey::of_event(event),
Item::Platform(_) => None,
};
let is_routing = |item: &Item<'a>| {
matches!(
item,
Item::Petri(event) if matches!(event.engine(), Some(Event::RoutingResolved { .. }))
)
};
// The key of each item: its index in clock order, and whether it sits
// before (0), at (1) or after (2) that index.
let mut keys: Vec<(usize, u8)> = (0..items.len()).map(|index| (index, 1)).collect();
for (index, (_, _, item)) in items.iter().enumerate() {
let Item::Platform(record) = item else {
continue;
};
let Some(position) = record.position else {
continue;
};
let at = FiringKey::from(position);
let first_routing = items
.iter()
.position(|(_, _, other)| item_firing(other) == Some(at) && is_routing(other));
let last_of_firing = items
.iter()
.rposition(|(_, _, other)| item_firing(other) == Some(at));
let first_later = items.iter().position(|(_, _, other)| {
item_firing(other)
.is_some_and(|key| key.execution == at.execution && key.firing > at.firing)
});
keys[index] = if let Some(before) = first_routing {
(before, 0)
} else if let Some(after) = last_of_firing {
(after, 2)
} else if let Some(before) = first_later {
(before, 0)
} else {
(index, 1)
};
}
let mut order: Vec<usize> = (0..items.len()).collect();
order.sort_by_key(|index| keys[*index]);
let mut ordered: Vec<Option<Item<'a>>> =
items.into_iter().map(|(_, _, item)| Some(item)).collect();
let items = order
.into_iter()
.map(|index| ordered[index].take().expect("each item is placed once"))
.collect();
(items, held)
}
#[cfg(test)]
mod tests {
use fabro_store::PlatformRecord;
use fabro_store::platform_records::{CheckpointRecord, StagePosition};
use petri_execution::events::{Context, EventId, NodeRef, Record, RecordOrigin, Subject};
use petri_execution::{ExecutionId, StoredEngineRecord};
use petri_runtime::driver::BranchRole;
use petri_runtime::engine::{DecisionId, EventOrigin, RouteApplied};
use petri_runtime::ir::{Attempt, FiringId, NodeId, Outcome, Status};
use super::*;
use crate::projector::stream;
/// A firing's engine event at `seq`, recorded at `at`.
fn engine_event(seq: u64, firing: u64, at: u64, body: Event) -> RunEvent {
RunEvent {
id: EventId {
source: EventSource::Execution {
execution: ExecutionId::new(0),
},
seq,
index: 0,
},
origin: RecordOrigin::External,
context: Context {
invocation: None,
execution: Some(ExecutionId::new(0)),
parent: None,
},
subject: Some(Subject {
node: NodeRef {
id: NodeId::new(1),
name: format!("n{firing}").into(),
kind: "attractor/command".into(),
meta: serde_json::Value::Null,
},
firing: Some(FiringId::new(firing)),
visit: Some(1),
attempt: Some(Attempt::FIRST),
generation: None,
branch: BranchRole::None,
}),
observed_at: None,
recorded_at: at,
record: Some(Record::Engine(StoredEngineRecord {
seq,
origin: EventOrigin::External,
recorded_at: at,
body,
})),
derived: None,
}
}
fn finished(seq: u64, firing: u64, at: u64) -> RunEvent {
engine_event(seq, firing, at, Event::StepFinished {
firing: FiringId::new(firing),
attempt: Attempt::FIRST,
outcome: Outcome::new(Status::Success, serde_json::Value::Null),
})
}
fn routing(seq: u64, firing: u64, at: u64) -> RunEvent {
engine_event(seq, firing, at, Event::RoutingResolved {
decision_id: DecisionId::route(FiringId::new(firing), Attempt::FIRST),
groups: Vec::new(),
})
}
fn applied(seq: u64, firing: u64, at: u64) -> RunEvent {
engine_event(seq, firing, at, Event::RouteApplied {
applied: RouteApplied::None {
firing: FiringId::new(firing),
group: 0,
},
})
}
fn started(seq: u64, firing: u64, at: u64) -> RunEvent {
engine_event(seq, firing, at, Event::StepStarted {
firing: FiringId::new(firing),
attempt: Attempt::FIRST,
})
}
fn checkpoint(seq: u64, firing: u64, at: u64) -> StoredPlatformRecord {
StoredPlatformRecord {
seq,
recorded_at: at,
record: PlatformRecord::Checkpoint(CheckpointRecord {
execution: 0,
firing,
attempt: Some(1),
workspace: None,
git_commit_sha: Some("abc".to_string()),
diff_summary: None,
patch_blob: None,
operation: None,
}),
position: Some(StagePosition {
execution: 0,
firing,
}),
}
}
fn names(items: &[Item<'_>]) -> Vec<String> {
items
.iter()
.map(|item| match item {
Item::Petri(event) => stream::event_id_text(&event.id),
Item::Platform(record) => format!("platform {}", record.seq),
})
.collect()
}
/// A firing's events, then a later firing's events, then a checkpoint
/// for the first firing stamped later than all of them: the stream puts
/// the checkpoint right after the first firing's finish, before its
/// routes and before the later firing.
#[test]
fn a_positioned_record_follows_its_firings_finish_whatever_its_clock_says() {
let events = vec![
finished(10, 1, 100),
routing(11, 1, 101),
applied(12, 1, 102),
started(13, 2, 103),
finished(14, 2, 104),
];
let records = vec![checkpoint(1, 1, 250)];
let (items, held) = order_items(&events, &records, &BTreeSet::new(), false);
assert_eq!(held, 0);
assert_eq!(names(&items), vec![
"execution 0/10/0",
"platform 1",
"execution 0/11/0",
"execution 0/12/0",
"execution 0/13/0",
"execution 0/14/0",
]);
}
/// With the firing finished in an earlier pass, the record goes before
/// the first event of a later firing; a record with no position keeps
/// its clock order.
#[test]
fn a_positioned_record_precedes_later_firings_and_an_unpositioned_one_keeps_its_clock() {
let events = vec![started(13, 2, 103), finished(14, 2, 104)];
let records = vec![checkpoint(1, 1, 250)];
let finished_before: BTreeSet<FiringKey> = [FiringKey::new(0, 1)].into_iter().collect();
let (items, held) = order_items(&events, &records, &finished_before, false);
assert_eq!(held, 0);
assert_eq!(names(&items), vec![
"platform 1",
"execution 0/13/0",
"execution 0/14/0",
]);
let unpositioned = StoredPlatformRecord {
position: None,
..checkpoint(2, 1, 250)
};
let unpositioned = [unpositioned];
let (items, held) = order_items(&events, &unpositioned, &BTreeSet::new(), false);
assert_eq!(held, 0);
assert_eq!(names(&items), vec![
"execution 0/13/0",
"execution 0/14/0",
"platform 2",
]);
}
/// A record whose firing has no finish yet, in the stream or in the pass,
/// is held back with everything after it until the finish arrives, or
/// until the run has finished.
#[test]
fn a_positioned_record_is_held_until_its_firings_finish_is_in_the_stream() {
let events = vec![started(13, 2, 103), finished(14, 2, 104)];
let records = vec![
checkpoint(1, 2, 50),
checkpoint(2, 3, 60),
checkpoint(3, 2, 70),
];
let (items, held) = order_items(&events, &records, &BTreeSet::new(), false);
assert_eq!(
held, 2,
"the record for firing 3 holds itself and the one after it"
);
assert_eq!(names(&items), vec![
"execution 0/13/0",
"execution 0/14/0",
"platform 1",
]);
// Once the run finished, a record for a firing that never finished
// keeps its clock order; the firing's own records still follow its
// finish, in their seq order.
let (items, held) = order_items(&events, &records, &BTreeSet::new(), true);
assert_eq!(held, 0, "nothing is held once the run finished");
assert_eq!(names(&items), vec![
"platform 2",
"execution 0/13/0",
"execution 0/14/0",
"platform 1",
"platform 3",
]);
}
}

View file

@ -0,0 +1,99 @@
//! The in-process wake-up: a run store whose appends signal the projector,
//! for a run that executes in the same process as the projector, over the
//! SQLite store directly, where no append endpoint is there to signal.
use std::sync::Arc;
use fabro_types::RunId;
use petri_execution::{Access, RunKey};
use petri_store::StoreError;
use super::Projector;
use crate::projection;
impl Projector {
/// A run store whose appends signal this projector: for a run that
/// executes in the same process as the projector, over the SQLite store
/// directly, where no append endpoint is there to signal. The signal is
/// sent after the store's append returned, so the records it covers are
/// durable before the view sees them.
pub fn observe_store(
self: &Arc<Self>,
inner: Arc<dyn petri_execution::RunStore>,
) -> Arc<dyn petri_execution::RunStore> {
Arc::new(SignallingStore {
inner,
projector: Arc::clone(self),
})
}
}
/// A run store that signals a projector after each append.
struct SignallingStore {
inner: Arc<dyn petri_execution::RunStore>,
projector: Arc<Projector>,
}
#[async_trait::async_trait]
impl petri_execution::RunStore for SignallingStore {
async fn open(
&self,
key: &RunKey,
access: Access,
) -> Result<Arc<dyn petri_execution::RunLogs>, StoreError> {
let logs = self.inner.open(key, access).await?;
Ok(Arc::new(SignallingLogs {
inner: logs,
run_id: projection::run_id_of(key.as_str()),
projector: Arc::clone(&self.projector),
}))
}
}
struct SignallingLogs {
inner: Arc<dyn petri_execution::RunLogs>,
run_id: Option<RunId>,
projector: Arc<Projector>,
}
#[async_trait::async_trait]
impl petri_execution::RunLogs for SignallingLogs {
fn locator(&self) -> String {
self.inner.locator()
}
async fn append(
&self,
log: &petri_execution::LogId,
records: &[petri_execution::Record],
) -> Result<(), StoreError> {
self.inner.append(log, records).await?;
if let Some(run_id) = self.run_id {
self.projector.signal(run_id);
}
Ok(())
}
async fn read(
&self,
log: &petri_execution::LogId,
) -> Result<Vec<petri_execution::Record>, StoreError> {
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
}
async fn get_blob(&self, digest: petri_store::Digest) -> Result<Option<Vec<u8>>, StoreError> {
self.inner.get_blob(digest).await
}
}

View file

@ -0,0 +1,113 @@
//! The run's stream as the view tables hold it: one row per folded item,
//! written by a pass and read back in `stream_seq` order.
use fabro_db::DbPool;
use fabro_types::{RunId, RunStreamItem, RunStreamItemKind};
use petri_execution::events::{EventId, EventSource};
use super::{Positions, ProjectError};
use crate::projection::{Item, RunView};
pub(crate) struct StreamRow {
pub(super) stream_seq: u64,
pub(super) item_kind: &'static str,
pub(super) item_id: String,
pub(super) event_json: String,
}
/// Fold the items into the view in order, each at the next delivery
/// sequence, advancing the positions with each, and produce its stream
/// row.
pub(crate) fn stream_rows(
items: &[Item<'_>],
view: &mut RunView,
positions: &mut Positions,
stream_seq: &mut u64,
) -> Result<Vec<StreamRow>, ProjectError> {
let mut rows = Vec::with_capacity(items.len());
for item in items {
*stream_seq += 1;
view.fold(item, *stream_seq);
let row = match item {
Item::Petri(event) => {
positions.advance(event.id);
StreamRow {
stream_seq: *stream_seq,
item_kind: "petri",
item_id: event_id_text(&event.id),
event_json: serde_json::to_string(event).map_err(ProjectError::Encode)?,
}
}
Item::Platform(record) => {
positions.platform_seq = record.seq;
StreamRow {
stream_seq: *stream_seq,
item_kind: "platform",
item_id: record.seq.to_string(),
event_json: serde_json::to_string(record).map_err(ProjectError::Encode)?,
}
}
};
rows.push(row);
}
Ok(rows)
}
/// The run's stream past the cursor, read from the view tables: up to
/// `limit` rows with `stream_seq > after`, in order, in Fabro's envelope.
pub(super) async fn stream_after(
views: &DbPool,
run_id: RunId,
after: u64,
limit: usize,
) -> Result<Vec<RunStreamItem>, ProjectError> {
let rows: Vec<(i64, String, String, String)> = sqlx::query_as(
"SELECT stream_seq, item_kind, item_id, event_json FROM petri_stream WHERE run_id = ? AND \
stream_seq > ? ORDER BY stream_seq LIMIT ?",
)
.bind(run_id.to_string())
.bind(column(after))
.bind(i64::try_from(limit).unwrap_or(i64::MAX))
.fetch_all(views)
.await
.map_err(ProjectError::Database)?;
rows.into_iter()
.map(|(stream_seq, item_kind, item_id, event_json)| {
let item: serde_json::Value =
serde_json::from_str(&event_json).map_err(ProjectError::Encode)?;
let kind = match item_kind.as_str() {
"platform" => RunStreamItemKind::Platform,
_ => RunStreamItemKind::Petri,
};
let recorded_at = item
.get("recorded_at")
.and_then(serde_json::Value::as_u64)
.unwrap_or(0);
Ok(RunStreamItem {
run_id,
stream_seq: u64::try_from(stream_seq).unwrap_or(0),
kind,
id: item_id,
recorded_at,
item,
})
})
.collect()
}
/// A Petri event id as the stream names it: `<log>/<seq>/<index>`.
#[must_use]
pub(crate) fn event_id_text(id: &EventId) -> String {
format!("{}/{}/{}", log_text(&id.source), id.seq, id.index)
}
pub(super) fn log_text(source: &EventSource) -> String {
match source {
EventSource::Coordinator => "coordinator".to_string(),
EventSource::Execution { execution } => format!("execution {execution}"),
}
}
pub(super) fn column(value: u64) -> i64 {
i64::try_from(value).unwrap_or(i64::MAX)
}

View file

@ -31,45 +31,41 @@
//! the worker's run reaches, so its target is deferred, and the worker's
//! hooks read the same plan and apply it through the scope's environment
//! at `scope_acquired`, before the first attempt runs there
//! ([`bring_sandbox_to`]).
//! ([`bring_to`]).
use std::collections::BTreeMap;
use std::path::PathBuf;
use std::sync::Arc;
use fabro_checkpoint::author::GitAuthor;
use fabro_store::platform_records::CheckpointRecord;
use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition};
use fabro_types::settings::run::{RunCheckpointSettings, RunNamespace};
use fabro_types::{RunId, SandboxProviderKind};
use fabro_types::RunId;
use fabro_types::settings::run::RunNamespace;
use petri_execution::host::{self, HostError};
use petri_execution::inspect::{self, ExecutionInspection, InspectError};
use petri_execution::{Access, InvocationId, RunKey, RunStore};
use petri_runtime::executor::ExecEnv;
use petri_store::StoreError;
use tracing::info;
use crate::checkpoint::{CHECKPOINT_FAILED_CLASS, CheckpointError, CheckpointKey, RunWorkspaces};
use crate::checkpoint::{
CHECKPOINT_FAILED_CLASS, CheckpointError, CheckpointKey, RunGitSettings, RunWorkspaces, Site,
};
use crate::platform_records::{PlatformRecordError, PlatformRecords};
use crate::workspace::{WorkspaceLookup, WorkspaceLookupError};
/// What recovery needs: the run, where its workspaces are, its records.
/// What recovery needs: the run, where its workspaces are, its records,
/// and its Git settings.
pub struct RecoveryRequest {
pub run_id: RunId,
pub run_id: RunId,
/// The run directory Petri ran under (the run's `petri` scratch).
pub run_dir: PathBuf,
pub store: Arc<dyn RunStore>,
pub records: Arc<dyn PlatformRecords>,
pub author: GitAuthor,
pub checkpoint: RunCheckpointSettings,
/// Whether the run's workspaces are on this host.
pub host_workspaces: bool,
pub run_dir: PathBuf,
pub store: Arc<dyn RunStore>,
pub records: Arc<dyn PlatformRecords>,
pub git: RunGitSettings,
}
impl RecoveryRequest {
/// The request a run's settings give: its Git author, its checkpoint
/// settings, and whether its sandbox provider keeps workspaces on this
/// host.
/// The request a run's settings give.
#[must_use]
pub fn for_run(
run_id: RunId,
@ -83,14 +79,7 @@ impl RecoveryRequest {
run_dir,
store,
records,
author: settings
.git
.author
.as_ref()
.map(GitAuthor::from)
.unwrap_or_default(),
checkpoint: settings.checkpoint.clone(),
host_workspaces: settings.environment.provider == SandboxProviderKind::LOCAL,
git: RunGitSettings::from(settings),
}
}
}
@ -172,12 +161,6 @@ pub enum RecoveryError {
},
}
/// One live execution's last durable finish and the snapshot it names.
struct Target {
execution: u64,
key: CheckpointKey,
}
/// Decide how the run continues: the snapshot every live workspace must sit
/// on, from the records and the snapshot repository, with a lost record
/// reconciled from the repository. Nothing is touched.
@ -220,80 +203,14 @@ pub async fn plan(
if let Some(failed) = checkpoint_failure(&inspection.executions) {
return Ok(Plan::Failed { reason: failed });
}
let lookup = WorkspaceLookup::new(Arc::clone(&store), key);
let recorded = recorded_checkpoints(records, run_id).await?;
// The snapshot each live execution's workspace must sit on. A live
// execution is one whose log records no exit: `inspect_run` reports it
// as incomplete.
let mut candidates: BTreeMap<String, Vec<(Target, String)>> = BTreeMap::new();
for execution in inspection
.executions
.iter()
.filter(|execution| execution.status == "incomplete")
{
let Some(target) = last_finish(execution) else {
continue;
};
let owned = lookup
.of_invocation(execution.invocation)
.await
.map_err(RecoveryError::Lookup)?;
if owned.is_empty() {
continue;
}
let mut found = false;
for workspace in owned {
let sha = match recorded.get(&target.key) {
Some((recorded_workspace, sha))
if recorded_workspace.as_deref().is_none_or(|w| w == workspace) =>
{
Some(sha.clone())
}
_ => {
let sha = workspaces
.find(&workspace, target.key)
.await
.map_err(|source| RecoveryError::Workspace {
workspace: workspace.clone(),
source,
})?;
if let Some(sha) = &sha {
reconcile_record(records, run_id, target.key, &workspace, sha).await?;
}
sha
}
};
if let Some(sha) = sha {
found = true;
candidates
.entry(workspace)
.or_default()
.push((Target { ..target }, sha));
}
}
if !found {
return Ok(Plan::Failed {
reason: format!(
"no checkpoint snapshot exists for the last durable finish of execution {} \
(firing {} attempt {}); the run cannot resume on stale files",
target.execution, target.key.firing, target.key.attempt
),
});
}
}
let mut targets = BTreeMap::new();
for (workspace, candidates) in candidates {
let sha = newest(workspaces, &workspace, &candidates).await?;
let key = candidates
.iter()
.find(|(_, candidate)| *candidate == sha)
.map_or(candidates[0].0.key, |(target, _)| target.key);
targets.insert(workspace, RestoreTarget { key, sha });
}
Ok(Plan::Resume { targets })
let planner = Planner {
records,
run_id,
workspaces,
lookup: WorkspaceLookup::new(store, key),
recorded: recorded_checkpoints(records, run_id).await?,
};
planner.targets(&inspection.executions).await
}
/// Decide how the run continues, and bring its host workspaces to their
@ -302,8 +219,8 @@ pub async fn recover(request: RecoveryRequest) -> Result<Recovery, RecoveryError
let workspaces = RunWorkspaces::new(
request.run_dir.clone(),
request.run_id.to_string(),
request.author.clone(),
&request.checkpoint,
request.git.author.clone(),
&request.git.checkpoint,
);
let targets = match plan(
Arc::clone(&request.store),
@ -320,8 +237,14 @@ pub async fn recover(request: RecoveryRequest) -> Result<Recovery, RecoveryError
let mut recovered = Vec::new();
for (workspace, target) in targets {
let action = if request.host_workspaces {
bring_host_to(&workspaces, &workspace, &target).await?
let action = if request.git.host_workspaces {
bring_to(
&workspaces,
&workspaces.host(&workspace),
&workspace,
&target,
)
.await?
} else {
WorkspaceAction::Deferred
};
@ -343,6 +266,178 @@ pub async fn recover(request: RecoveryRequest) -> Result<Recovery, RecoveryError
})
}
/// The snapshot one live execution's last durable finish names on a
/// workspace.
struct Candidate {
key: CheckpointKey,
sha: String,
}
/// The decision over one run's records and snapshot repository.
struct Planner<'a> {
records: &'a dyn PlatformRecords,
run_id: &'a RunId,
workspaces: &'a RunWorkspaces,
lookup: WorkspaceLookup,
/// The run's checkpoint records by key: the workspace they name and
/// the commit.
recorded: BTreeMap<CheckpointKey, (Option<String>, String)>,
}
impl Planner<'_> {
/// The snapshot each live execution's workspace must sit on. A live
/// execution is one whose log records no exit: `inspect_run` reports
/// it as incomplete. A workspace several live executions share is
/// brought to the newest of their snapshots.
async fn targets(&self, executions: &[ExecutionInspection]) -> Result<Plan, RecoveryError> {
let mut candidates: BTreeMap<String, Vec<Candidate>> = BTreeMap::new();
for execution in executions
.iter()
.filter(|execution| execution.status == "incomplete")
{
let Some(key) = last_finish(execution) else {
continue;
};
let owned = self
.lookup
.of_invocation(execution.invocation)
.await
.map_err(RecoveryError::Lookup)?;
if owned.is_empty() {
continue;
}
let mut found = false;
for workspace in owned {
if let Some(sha) = self.snapshot_of(&workspace, key).await? {
found = true;
candidates
.entry(workspace)
.or_default()
.push(Candidate { key, sha });
}
}
if !found {
return Ok(Plan::Failed {
reason: format!(
"no checkpoint snapshot exists for the last durable finish of {key}; the \
run cannot resume on stale files"
),
});
}
}
let mut targets = BTreeMap::new();
for (workspace, candidates) in candidates {
let target = self.newest(&workspace, &candidates).await?;
targets.insert(workspace, target);
}
Ok(Plan::Resume { targets })
}
/// The snapshot of `key` in `workspace`: the commit its record names,
/// when the record names this workspace or none; else the commit found
/// by its key in the workspace's snapshot repository or history, which
/// is then recorded again for the record the crash lost. `None` when no
/// snapshot exists.
async fn snapshot_of(
&self,
workspace: &str,
key: CheckpointKey,
) -> Result<Option<String>, RecoveryError> {
if let Some((recorded_workspace, sha)) = self.recorded.get(&key) {
if recorded_workspace
.as_deref()
.is_none_or(|recorded| recorded == workspace)
{
return Ok(Some(sha.clone()));
}
}
let found = self
.workspaces
.find(workspace, key)
.await
.map_err(|source| RecoveryError::Workspace {
workspace: workspace.to_string(),
source,
})?;
if let Some(sha) = &found {
self.reconcile_record(key, workspace, sha).await?;
}
Ok(found)
}
/// Write the record a crash lost, from the commit found by its key.
async fn reconcile_record(
&self,
key: CheckpointKey,
workspace: &str,
sha: &str,
) -> Result<(), RecoveryError> {
info!(
run_id = %self.run_id,
execution = key.execution,
firing = key.firing,
attempt = key.attempt,
sha,
"checkpoint record reconciled from the run branch"
);
let record = PlatformRecord::Checkpoint(CheckpointRecord {
execution: key.execution,
firing: key.firing,
attempt: Some(key.attempt),
workspace: Some(workspace.to_string()),
git_commit_sha: Some(sha.to_string()),
diff_summary: None,
patch_blob: None,
operation: Some(key.operation()),
});
self.records
.append(
self.run_id,
&record,
Some(StagePosition {
execution: key.execution,
firing: key.firing,
}),
)
.await
.map_err(RecoveryError::Records)?;
Ok(())
}
/// Of the snapshots live executions name on one workspace, the one
/// every other descends from, else the last named. Two executions that
/// name the same commit share it under the first one's key.
async fn newest(
&self,
workspace: &str,
candidates: &[Candidate],
) -> Result<RestoreTarget, RecoveryError> {
let mut chosen = &candidates[0];
for candidate in &candidates[1..] {
if self
.workspaces
.is_ancestor(workspace, &chosen.sha, &candidate.sha)
.await
.map_err(|source| RecoveryError::Workspace {
workspace: workspace.to_string(),
source,
})?
{
chosen = candidate;
}
}
let chosen = candidates
.iter()
.find(|candidate| candidate.sha == chosen.sha)
.unwrap_or(chosen);
Ok(RestoreTarget {
key: chosen.key,
sha: chosen.sha.clone(),
})
}
}
/// The reason a run with a failed checkpoint is reported failed, when it
/// has one.
fn checkpoint_failure(executions: &[ExecutionInspection]) -> Option<String> {
@ -368,16 +463,14 @@ fn checkpoint_failure(executions: &[ExecutionInspection]) -> Option<String> {
})
}
/// The last `StepFinished` of an execution's log.
fn last_finish(execution: &ExecutionInspection) -> Option<Target> {
/// The last `StepFinished` of an execution's log, as the key of its
/// snapshot.
fn last_finish(execution: &ExecutionInspection) -> Option<CheckpointKey> {
let attempt = execution.engine.as_ref()?.attempts.last()?;
Some(Target {
Some(CheckpointKey {
execution: execution.execution.raw(),
key: CheckpointKey {
execution: execution.execution.raw(),
firing: attempt.firing,
attempt: attempt.attempt,
},
firing: attempt.firing,
attempt: attempt.attempt,
})
}
@ -407,75 +500,14 @@ async fn recorded_checkpoints(
Ok(recorded)
}
/// Write the record a crash lost, from the commit found by its key.
async fn reconcile_record(
records: &dyn PlatformRecords,
run_id: &RunId,
key: CheckpointKey,
workspace: &str,
sha: &str,
) -> Result<(), RecoveryError> {
info!(
run_id = %run_id,
execution = key.execution,
firing = key.firing,
attempt = key.attempt,
sha,
"checkpoint record reconciled from the run branch"
);
let record = PlatformRecord::Checkpoint(CheckpointRecord {
execution: key.execution,
firing: key.firing,
attempt: Some(key.attempt),
workspace: Some(workspace.to_string()),
git_commit_sha: Some(sha.to_string()),
diff_summary: None,
patch_blob: None,
operation: Some(key.operation()),
});
records
.append(
run_id,
&record,
Some(StagePosition {
execution: key.execution,
firing: key.firing,
}),
)
.await
.map_err(RecoveryError::Records)?;
Ok(())
}
/// Of the snapshots live executions name on one workspace, the one every
/// other descends from, else the last named.
async fn newest(
workspaces: &RunWorkspaces,
workspace: &str,
targets: &[(Target, String)],
) -> Result<String, RecoveryError> {
let mut chosen = &targets[0].1;
for (_, sha) in &targets[1..] {
if workspaces
.is_ancestor(workspace, chosen, sha)
.await
.map_err(|source| RecoveryError::Workspace {
workspace: workspace.to_string(),
source,
})?
{
chosen = sha;
}
}
Ok(chosen.clone())
}
/// Verify, reset or restore the host workspace onto its target: a
/// Verify, reset or restore the workspace at `site` onto its target: a
/// workspace that still holds the commit is verified or reset in place; a
/// gone one, or a fresh directory with no history (a fork's first
/// acquisition), is restored from the snapshot repository.
pub async fn bring_host_to(
/// gone one, a fresh directory with no history (a fork's first
/// acquisition), or a sandbox whose repository lost the commit, is restored
/// from the snapshot repository (a bundle of the snapshot, into a sandbox).
pub async fn bring_to(
workspaces: &RunWorkspaces,
site: &Site,
workspace: &str,
target: &RestoreTarget,
) -> Result<WorkspaceAction, RecoveryError> {
@ -484,64 +516,22 @@ pub async fn bring_host_to(
source,
};
if workspaces
.has_commit(workspace, &target.sha)
.has_commit(site, &target.sha)
.await
.map_err(failed)?
{
if workspaces
.matches(workspace, &target.sha)
.matches(site, &target.sha)
.await
.map_err(failed)?
{
return Ok(WorkspaceAction::Verified);
}
workspaces
.reset(workspace, &target.sha)
.await
.map_err(failed)?;
workspaces.reset(site, &target.sha).await.map_err(failed)?;
return Ok(WorkspaceAction::Reset);
}
workspaces
.restore(workspace, target.key, &target.sha)
.await
.map_err(failed)?;
Ok(WorkspaceAction::Restored)
}
/// Verify, reset or restore a sandbox workspace onto its target, through
/// the scope's environment: a retained sandbox that still holds the commit
/// is verified or reset in place; a fresh one, or one whose repository
/// lost the commit, is restored from a bundle of the snapshot.
pub async fn bring_sandbox_to(
workspaces: &RunWorkspaces,
env: &Arc<dyn ExecEnv>,
workspace: &str,
target: &RestoreTarget,
) -> Result<WorkspaceAction, RecoveryError> {
let failed = |source| RecoveryError::Workspace {
workspace: workspace.to_string(),
source,
};
if workspaces
.has_commit_in(env, &target.sha)
.await
.map_err(failed)?
{
if workspaces
.matches_in(env, &target.sha)
.await
.map_err(failed)?
{
return Ok(WorkspaceAction::Verified);
}
workspaces
.reset_in(env, &target.sha)
.await
.map_err(failed)?;
return Ok(WorkspaceAction::Reset);
}
workspaces
.restore_in(env, workspace, target.key, &target.sha)
.restore(site, workspace, target.key, &target.sha)
.await
.map_err(failed)?;
Ok(WorkspaceAction::Restored)

View file

@ -55,13 +55,14 @@
use std::collections::HashMap;
use std::error::Error;
use std::sync::{Arc, Mutex, MutexGuard, PoisonError, Weak};
use std::sync::{Arc, Mutex, 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 fabro_util::sync;
use petri_store::{
Access, Digest, ExecutionId, LogId, OwnerId, Record, RunKey, RunLogs, RunStore, StoreError,
};
@ -150,7 +151,7 @@ impl SqliteRunStore {
if result.rows_affected() == 0 {
return Err(self.shared.not_found(key));
}
lock(&self.shared.live).remove(key);
sync::lock(&self.shared.live).remove(key);
debug!(run_id = %key, "Petri run lease released from outside");
Ok(())
}
@ -174,7 +175,7 @@ impl SqliteRunStore {
/// 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);
let mut live = sync::lock(&self.shared.live);
if let Some(handle) = live.get(key).and_then(Weak::upgrade) {
if handle.owner.as_ref() == Some(&owner) {
return handle;
@ -214,7 +215,7 @@ impl Shared {
/// 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));
let pending = mem::take(&mut *sync::lock(&self.releases));
for release in pending {
// A release task never panics: it reports its own failure.
let _ = release.await;
@ -409,7 +410,7 @@ impl Drop for SqliteRunLogs {
return;
};
{
let mut live = lock(&self.shared.live);
let mut live = sync::lock(&self.shared.live);
let this: *const Self = self;
if live
.get(&self.key)
@ -425,7 +426,7 @@ impl Drop for SqliteRunLogs {
let release = runtime.spawn(async move {
shared.release_owner(&key, &owner).await;
});
lock(&self.shared.releases).push(release);
sync::lock(&self.shared.releases).push(release);
}
Err(_) => {
warn!(
@ -558,10 +559,6 @@ impl RunLogs for SqliteRunLogs {
}
}
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()

View file

@ -1,23 +1,35 @@
//! Petri's test kit, for Fabro crates that check a store implementation
//! against Petri's contract from their own tests, an in-memory platform
//! against Petri's contract from their own tests; an in-memory platform
//! record store and an in-memory blob table for tests of the hooks and
//! recovery. Compiled only with the `test-support` feature, which a
//! recovery; and the readers over the view tables a test compares a live
//! view with. Compiled only with the `test-support` feature, which a
//! dev-dependency turns on.
use std::collections::HashMap;
use std::sync::{Mutex, MutexGuard, PoisonError};
use std::collections::{BTreeSet, HashMap};
use std::sync::Mutex;
use std::time::Duration;
use async_trait::async_trait;
use bytes::Bytes;
use fabro_store::platform_records::now_ms;
use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition, StoredPlatformRecord};
use fabro_db::DbPool;
use fabro_store::platform_records::{PlatformRecordStore, now_ms};
use fabro_store::{
PlatformRecord, PlatformRecordKind, RunProjection, StagePosition, StoredPlatformRecord,
};
use fabro_types::{BlobHash, RunId};
use fabro_util::error::collect_chain;
use fabro_util::sync;
use petri_execution::events::{self, RunEvent};
use petri_execution::{Access, CoordinatorEvent, RunKey, RunStore as _};
use petri_store::StoreError;
pub use petri_testkit::run_store;
use tracing::warn;
use crate::SqliteRunStore;
use crate::blobs::Blobs;
use crate::platform_records::{PlatformRecordError, PlatformRecords};
use crate::projector::Projector;
use crate::projection::RunView;
use crate::projector::{self, Positions, ProjectError, Projector, order, stream};
/// Whether the projector keeps a cache for the run: the replay and the
/// view its passes continue from.
@ -47,7 +59,7 @@ impl MemoryBlobs {
/// How many blobs the table holds.
#[must_use]
pub fn len(&self) -> usize {
lock(&self.rows).len()
sync::lock(&self.rows).len()
}
#[must_use]
@ -60,12 +72,12 @@ impl MemoryBlobs {
impl Blobs for MemoryBlobs {
async fn write(&self, bytes: &[u8]) -> anyhow::Result<BlobHash> {
let hash = BlobHash::new(bytes);
lock(&self.rows).insert(hash, bytes.to_vec());
sync::lock(&self.rows).insert(hash, bytes.to_vec());
Ok(hash)
}
async fn read(&self, hash: &BlobHash) -> anyhow::Result<Option<Bytes>> {
Ok(lock(&self.rows)
Ok(sync::lock(&self.rows)
.get(hash)
.map(|bytes| Bytes::copy_from_slice(bytes)))
}
@ -86,14 +98,13 @@ impl MemoryPlatformRecords {
/// Every record of the run, in seq order.
#[must_use]
pub fn records(&self, run_id: &RunId) -> Vec<StoredPlatformRecord> {
lock(&self.runs).get(run_id).cloned().unwrap_or_default()
sync::lock(&self.runs)
.get(run_id)
.cloned()
.unwrap_or_default()
}
}
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex.lock().unwrap_or_else(PoisonError::into_inner)
}
#[async_trait]
impl PlatformRecords for MemoryPlatformRecords {
async fn append(
@ -102,7 +113,7 @@ impl PlatformRecords for MemoryPlatformRecords {
record: &PlatformRecord,
position: Option<StagePosition>,
) -> Result<StoredPlatformRecord, PlatformRecordError> {
let mut runs = lock(&self.runs);
let mut runs = sync::lock(&self.runs);
let records = runs.entry(*run_id).or_default();
let stored = StoredPlatformRecord {
seq: records.len() as u64 + 1,
@ -126,3 +137,100 @@ impl PlatformRecords for MemoryPlatformRecords {
.collect())
}
}
/// The run's projection rebuilt from its records alone, with nothing
/// stored: what a fresh projector would commit over the same records. A test
/// compares it with the live view. `records` and `views` are the two pools
/// [`Projector::new`] takes.
pub async fn rebuild(
records: &DbPool,
views: &DbPool,
run_id: RunId,
) -> Result<(Option<RunProjection>, Positions, u64), ProjectError> {
let store = SqliteRunStore::new(records.clone());
let platform = PlatformRecordStore::new(views.clone());
let key = RunKey::new(run_id.to_string());
let platform_records = platform.read(&run_id).await.map_err(ProjectError::Store)?;
let events = match store.open(&key, Access::Read).await {
Ok(logs) => events::replay_run(&*logs)
.await
.inspect_err(|error| {
warn!(error = %collect_chain(error).join(": "), "rebuild: the run does not replay");
})
.unwrap_or_default(),
Err(StoreError::NotFound { .. }) => Vec::new(),
Err(error) => return Err(ProjectError::Open(error)),
};
let run_finished = events.iter().any(|event| {
matches!(
event.coordinator(),
Some(CoordinatorEvent::RunFinished { .. })
)
});
let (items, _held) =
order::order_items(&events, &platform_records, &BTreeSet::new(), run_finished);
let mut view = RunView::new();
let mut positions = Positions::default();
let mut stream_seq = 0;
stream::stream_rows(&items, &mut view, &mut positions, &mut stream_seq)?;
Ok((view.projection, positions, stream_seq))
}
/// The stored view's positions and stream sequence; `views` is the pool
/// the view tables live in.
pub async fn stored_positions(
views: &DbPool,
run_id: RunId,
) -> Result<Option<(Positions, u64)>, ProjectError> {
projector::stored_positions(views, run_id).await
}
/// The stored view's projection, for a test or a reader outside the store.
pub async fn stored_projection(
views: &DbPool,
run_id: RunId,
) -> Result<Option<RunProjection>, ProjectError> {
let json: Option<String> =
sqlx::query_scalar("SELECT projection_json FROM petri_projection WHERE run_id = ?")
.bind(run_id.to_string())
.fetch_optional(views)
.await
.map_err(ProjectError::Database)?;
json.map(|json| serde_json::from_str(&json).map_err(ProjectError::Encode))
.transpose()
}
/// The stream rows of a run: `(stream_seq, item_kind, item_id)`, in order.
pub async fn stored_stream(
views: &DbPool,
run_id: RunId,
) -> Result<Vec<(u64, String, String)>, ProjectError> {
let rows: Vec<(i64, String, String)> = sqlx::query_as(
"SELECT stream_seq, item_kind, item_id FROM petri_stream WHERE run_id = ? ORDER BY stream_seq",
)
.bind(run_id.to_string())
.fetch_all(views)
.await
.map_err(ProjectError::Database)?;
Ok(rows
.into_iter()
.map(|(seq, kind, id)| (u64::try_from(seq).unwrap_or(0), kind, id))
.collect())
}
/// Every stored platform record of a run, for a reader outside the store.
pub async fn stored_platform_records(
views: &DbPool,
run_id: RunId,
) -> Result<Vec<StoredPlatformRecord>, ProjectError> {
PlatformRecordStore::new(views.clone())
.read(&run_id)
.await
.map_err(ProjectError::Store)
}
/// A recorded event's projection is what `RunEvent` serializes to.
#[must_use]
pub fn event_json(event: &RunEvent) -> serde_json::Value {
serde_json::to_value(event).unwrap_or_default()
}

View file

@ -24,7 +24,9 @@ use fabro_checkpoint::author::GitAuthor;
use fabro_petri::admission::AdmittedGraphs;
use fabro_petri::blobs::Blobs;
use fabro_petri::check::{self, Bundle, CheckRequest, Launch};
use fabro_petri::checkpoint::{CHECKPOINT_FAILED_CLASS, CheckpointKey, RunWorkspaces};
use fabro_petri::checkpoint::{
CHECKPOINT_FAILED_CLASS, CheckpointKey, RunGitSettings, RunWorkspaces,
};
use fabro_petri::controls::RunControls;
use fabro_petri::engine::{self, Execution, RunRequest, RunStatus};
use fabro_petri::hooks::HooksSpec;
@ -166,13 +168,13 @@ impl Harness {
fn hooks(&self, provider: &SandboxProviderKind) -> HooksSpec {
HooksSpec {
records: Arc::clone(&self.records) as Arc<dyn PlatformRecords>,
author: GitAuthor::default(),
identity_source: GitIdentitySource::Default,
checkpoint: RunCheckpointSettings::default(),
artifacts: self.artifacts.clone(),
host_workspaces: *provider == SandboxProviderKind::LOCAL,
test_gates: None,
records: Arc::clone(&self.records) as Arc<dyn PlatformRecords>,
git: RunGitSettings {
host_workspaces: *provider == SandboxProviderKind::LOCAL,
..RunGitSettings::default()
},
artifacts: self.artifacts.clone(),
test_gates: None,
}
}
@ -285,13 +287,11 @@ impl Harness {
async fn recover(&self) -> Recovery {
recovery::recover(RecoveryRequest {
run_id: self.run_id,
run_dir: self.run_dir.clone(),
store: Arc::clone(&self.store) as Arc<dyn petri_store::RunStore>,
records: Arc::clone(&self.records) as Arc<dyn PlatformRecords>,
author: GitAuthor::default(),
checkpoint: RunCheckpointSettings::default(),
host_workspaces: true,
run_id: self.run_id,
run_dir: self.run_dir.clone(),
store: Arc::clone(&self.store) as Arc<dyn petri_store::RunStore>,
records: Arc::clone(&self.records) as Arc<dyn PlatformRecords>,
git: RunGitSettings::default(),
})
.await
.expect("recovery decides")

View file

@ -418,15 +418,15 @@ fn json<T: serde::Serialize>(value: &T) -> serde_json::Value {
/// The stored view equals the view rebuilt from the records alone: the
/// projection, the positions and the delivery sequence.
async fn assert_view_equals_rebuild(pool: &DbPool, run_id: RunId) {
let stored = projector::stored_projection(pool, run_id)
let stored = petri_support::stored_projection(pool, run_id)
.await
.expect("the stored projection reads")
.expect("the run has a stored projection");
let (stored_positions, stored_stream_seq) = projector::stored_positions(pool, run_id)
let (stored_positions, stored_stream_seq) = petri_support::stored_positions(pool, run_id)
.await
.expect("the positions read")
.expect("the run has positions");
let (rebuilt, positions, stream_seq) = projector::rebuild(pool, pool, run_id)
let (rebuilt, positions, stream_seq) = petri_support::rebuild(pool, pool, run_id)
.await
.expect("the run rebuilds");
let rebuilt = rebuilt.expect("the rebuild has a projection");
@ -442,7 +442,7 @@ async fn assert_view_equals_rebuild(pool: &DbPool, run_id: RunId) {
positions.petri.sort();
assert_eq!(stored_positions, positions);
assert_eq!(stored_stream_seq, stream_seq);
let stream = projector::stored_stream(pool, run_id)
let stream = petri_support::stored_stream(pool, run_id)
.await
.expect("the stream reads");
let seqs: Vec<u64> = stream.iter().map(|(seq, _, _)| *seq).collect();
@ -454,7 +454,7 @@ async fn assert_view_equals_rebuild(pool: &DbPool, run_id: RunId) {
}
async fn stage_states(pool: &DbPool, run_id: RunId) -> Vec<(String, StageState)> {
let stored = projector::stored_projection(pool, run_id)
let stored = petri_support::stored_projection(pool, run_id)
.await
.expect("the stored projection reads")
.expect("the run has a stored projection");
@ -472,7 +472,7 @@ async fn the_hello_bundle_projects_live_as_it_rebuilds() {
let scenario = hello_scenario().await;
run_live(&scenario).await;
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
let stored = projector::stored_projection(&scenario.pool, scenario.run_id)
let stored = petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
@ -501,7 +501,7 @@ async fn a_large_output_projects_as_its_blob_reference() {
let scenario = large_output_scenario().await;
run_live(&scenario).await;
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
let stored = projector::stored_projection(&scenario.pool, scenario.run_id)
let stored = petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
@ -525,7 +525,7 @@ async fn a_command_workflow_projects_live_as_it_rebuilds() {
let scenario = command_scenario().await;
run_live(&scenario).await;
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
let stored = projector::stored_projection(&scenario.pool, scenario.run_id)
let stored = petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
@ -554,7 +554,7 @@ async fn a_parallel_workflow_projects_its_branches_as_child_executions() {
let scenario = parallel_scenario().await;
run_live(&scenario).await;
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
let stored = projector::stored_projection(&scenario.pool, scenario.run_id)
let stored = petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
@ -593,7 +593,7 @@ async fn dropped_wake_ups_are_caught_up_by_the_next_signal() {
let scenario = command_scenario().await;
run_unobserved(&scenario).await;
assert!(
projector::stored_projection(&scenario.pool, scenario.run_id)
petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.is_none(),
@ -780,12 +780,13 @@ async fn a_crash_between_the_record_commit_and_the_view_applies_only_the_suffix(
.expect("the first pass commits");
assert!(!pass.skipped);
assert!(!pass.health.complete, "the run has not finished");
let (positions_before, stream_before) = projector::stored_positions(&replayed, scenario.run_id)
.await
.expect("reads")
.expect("positions");
let (positions_before, stream_before) =
petri_support::stored_positions(&replayed, scenario.run_id)
.await
.expect("reads")
.expect("positions");
assert_eq!(pass.stream_seq, stream_before);
let stream_rows_before = projector::stored_stream(&replayed, scenario.run_id)
let stream_rows_before = petri_support::stored_stream(&replayed, scenario.run_id)
.await
.expect("reads")
.len();
@ -801,7 +802,7 @@ async fn a_crash_between_the_record_commit_and_the_view_applies_only_the_suffix(
"{crashed:?}"
);
assert_eq!(
projector::stored_positions(&replayed, scenario.run_id)
petri_support::stored_positions(&replayed, scenario.run_id)
.await
.expect("reads")
.expect("positions"),
@ -813,11 +814,12 @@ async fn a_crash_between_the_record_commit_and_the_view_applies_only_the_suffix(
let after = Projector::new(replayed.clone(), replayed.clone());
let report = after.startup_pass().await.expect("the restart catches up");
assert_eq!((report.runs, report.projected), (1, 1));
let (positions_after, stream_after) = projector::stored_positions(&replayed, scenario.run_id)
.await
.expect("reads")
.expect("positions");
let stream_rows_after = projector::stored_stream(&replayed, scenario.run_id)
let (positions_after, stream_after) =
petri_support::stored_positions(&replayed, scenario.run_id)
.await
.expect("reads")
.expect("positions");
let stream_rows_after = petri_support::stored_stream(&replayed, scenario.run_id)
.await
.expect("reads");
// Only the suffix was applied: the stream grew by the suffix's events,
@ -845,11 +847,11 @@ async fn a_crash_between_the_record_commit_and_the_view_applies_only_the_suffix(
// And the copy agrees with the run projected in one go over the source.
let source = Projector::new(scenario.pool.clone(), scenario.pool.clone());
source.startup_pass().await.expect("the source projects");
let whole = projector::stored_projection(&scenario.pool, scenario.run_id)
let whole = petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
let pieced = projector::stored_projection(&replayed, scenario.run_id)
let pieced = petri_support::stored_projection(&replayed, scenario.run_id)
.await
.expect("reads")
.expect("stored");
@ -904,11 +906,11 @@ async fn a_restarted_projector_agrees_over_nested_child_executions() {
let whole = Projector::new(scenario.pool.clone(), scenario.pool.clone());
whole.startup_pass().await.expect("the source projects");
let one_go = projector::stored_projection(&scenario.pool, scenario.run_id)
let one_go = petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
let restarted = projector::stored_projection(&staged, scenario.run_id)
let restarted = petri_support::stored_projection(&staged, scenario.run_id)
.await
.expect("reads")
.expect("stored");
@ -937,11 +939,11 @@ async fn a_torn_tail_holds_the_view_and_reports_the_run_incomplete() {
.await
.expect("the clean pass commits");
assert!(clean.health.complete, "{:?}", clean.health.incomplete);
let (positions, stream_seq) = projector::stored_positions(&scenario.pool, scenario.run_id)
let (positions, stream_seq) = petri_support::stored_positions(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("positions");
let before = projector::stored_projection(&scenario.pool, scenario.run_id)
let before = petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
@ -973,7 +975,7 @@ async fn a_torn_tail_holds_the_view_and_reports_the_run_incomplete() {
held.health
);
let (positions_after, stream_after) =
projector::stored_positions(&scenario.pool, scenario.run_id)
petri_support::stored_positions(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("positions");
@ -982,7 +984,7 @@ async fn a_torn_tail_holds_the_view_and_reports_the_run_incomplete() {
"the view did not advance past the tear"
);
assert_eq!(stream_after, stream_seq);
let after = projector::stored_projection(&scenario.pool, scenario.run_id)
let after = petri_support::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
@ -1240,9 +1242,10 @@ impl GateRun {
async fn pending(&self) -> fabro_types::RunProjection {
let deadline = Instant::now() + Duration::from_secs(30);
loop {
let stored = projector::stored_projection(&self.scenario.pool, self.scenario.run_id)
.await
.expect("the stored projection reads");
let stored =
petri_support::stored_projection(&self.scenario.pool, self.scenario.run_id)
.await
.expect("the stored projection reads");
if let Some(stored) = stored.filter(|stored| !stored.pending_interviews.is_empty()) {
return stored;
}
@ -1255,7 +1258,7 @@ impl GateRun {
}
async fn stored(&self) -> fabro_types::RunProjection {
projector::stored_projection(&self.scenario.pool, self.scenario.run_id)
petri_support::stored_projection(&self.scenario.pool, self.scenario.run_id)
.await
.expect("the stored projection reads")
.expect("the run has a stored projection")
@ -1282,7 +1285,7 @@ async fn an_expired_question_is_pending_while_the_gate_waits_and_closes_on_the_e
let run_id = gate.scenario.run_id;
let deadline = Instant::now() + Duration::from_secs(30);
loop {
let stored = projector::stored_projection(&pool, run_id)
let stored = petri_support::stored_projection(&pool, run_id)
.await
.expect("the stored projection reads");
if let Some(stored) = stored.filter(|stored| !stored.pending_interviews.is_empty()) {

View file

@ -40,3 +40,27 @@ pub struct Conclusion {
#[serde(default)]
pub diff: RunDiff,
}
impl Conclusion {
/// A conclusion that records only how the run ended: no timing, stages,
/// usage or diff. What a terminal lifecycle record gives when the
/// engine recorded no finish of its own.
#[must_use]
pub fn outcome_only(
timestamp: DateTime<Utc>,
status: StageOutcome,
failure: Option<RunFailure>,
) -> Self {
Self {
timestamp,
status,
timing: RunTiming::default(),
failure,
final_git_commit_sha: None,
stages: Vec::new(),
usage: None,
total_retries: 0,
diff: RunDiff::default(),
}
}
}

View file

@ -122,6 +122,19 @@ impl StageModelUsage {
pub const MODE_AGENT: &'static str = "agent";
pub const MODE_ACP: &'static str = "acp";
/// The usage record of a stage that named its provider and model, with
/// no request controls.
#[must_use]
pub fn new(mode: &str, provider: Option<String>, model: Option<String>) -> Self {
Self {
mode: mode.to_string(),
provider,
model,
reasoning_effort: None,
speed: None,
}
}
/// Build the usage record from a `stage.prompt` event, returning `None`
/// when the event carried no model metadata.
#[must_use]

View file

@ -120,6 +120,27 @@ impl RunStatus {
}
}
/// The status a pause takes the run to: paused, remembering the block
/// the run was under so the unpause can restore it.
#[must_use]
pub fn paused(self) -> Self {
Self::Paused {
prior_block: self.blocked_reason(),
}
}
/// The status an unpause takes the run to: back to the block the pause
/// remembered, else running.
#[must_use]
pub fn unpaused(self) -> Self {
match self {
Self::Paused {
prior_block: Some(blocked_reason),
} => Self::Blocked { blocked_reason },
_ => Self::Running,
}
}
pub fn terminal_status(self) -> Option<TerminalStatus> {
match self {
Self::Succeeded { reason } => Some(TerminalStatus::Succeeded { reason }),

View file

@ -12,6 +12,7 @@ pub mod printer;
pub mod run_log;
pub mod session_secret;
pub mod shell;
pub mod sync;
pub mod terminal;
pub mod text;
pub mod time;

View file

@ -0,0 +1,11 @@
//! Locking helpers shared across crates.
use std::sync::{Mutex, MutexGuard, PoisonError};
/// Lock a mutex, recovering the guard when another holder panicked. The
/// state such a mutex guards is bookkeeping (a cache, a set of ids, a
/// counter) that stays usable after a panic elsewhere, so the poison is
/// cleared rather than propagated.
pub fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex.lock().unwrap_or_else(PoisonError::into_inner)
}