Checkpoint a Petri run's stages and recover its workspaces on restart

Fabro's hooks on a Petri run wrap the hooks the runtime installed for
`[[run.hooks]]` and forward every point. In `prepare_result`, before the
finish is recorded, they commit the stage's files on the run branch of
its host workspace with Fabro's author identity and the run, execution,
firing and attempt as trailers, and publish the commit to a snapshot
repository beside the run's workspaces under a ref per checkpoint. A
stage that failed on its own terms is committed like a successful one; a
commit that fails is fatal: the outcome becomes a `checkpoint_failed`
failure, the run is cancelled through the coordinator handle, and the
transition refuses the firing's routes. In `transition` they write the
platform checkpoint record, keyed on the Petri position and the
checkpoint's operation identity, and a failed write is a recorded
problem.

On restart the server runs the recovery protocol before it relaunches a
worker: a run with a failed checkpoint is reported failed; otherwise
every live execution's last durable finish names the snapshot its
workspace is verified against, reset to, or restored from, with a lost
record reconciled from the snapshot repository, and a finish with no
snapshot fails the run rather than resume it on stale files.

The worker reaches the platform records over two new worker-scoped
endpoints; the server reaches the table directly. A test gate directory
lets the CLI scenarios hold a checkpoint at a named point.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-18 00:12:14 -04:00
parent abbc7ca11d
commit 01beea0a6c
No known key found for this signature in database
20 changed files with 3197 additions and 37 deletions

3
Cargo.lock generated
View file

@ -2893,12 +2893,15 @@ dependencies = [
"bytes",
"fabro-api",
"fabro-auth",
"fabro-checkpoint",
"fabro-client",
"fabro-db",
"fabro-http",
"fabro-llm",
"fabro-petri",
"fabro-store",
"fabro-types",
"fabro-util",
"lithos-llm",
"petri-attractor-steps",
"petri-execution",

View file

@ -3418,6 +3418,59 @@ paths:
schema:
$ref: "#/components/schemas/ErrorResponse"
/api/v1/runs/{id}/petri/platform-records:
get:
operationId: listPetriPlatformRecords
tags: [Run Internals]
summary: List Petri Platform Records
description: |
The run's platform records (Fabro's own facts about a Petri run: a
checkpoint commit, a pull request, a notification), in `seq` order,
optionally of one kind. What a run's worker reads to find an effect
it already performed before performing it again.
parameters:
- $ref: "#/components/parameters/RunId"
- $ref: "#/components/parameters/PetriPlatformRecordKind"
responses:
"200":
description: The run's platform records
content:
application/json:
schema:
$ref: "#/components/schemas/PetriPlatformRecordList"
post:
operationId: appendPetriPlatformRecord
tags: [Run Internals]
summary: Append Petri Platform Record
description: |
Stores one platform record at the run's next `seq`, tied to the
Petri stage named by `execution` and `firing` when it belongs to one.
The record is the JSON of a Fabro platform record, tagged by `kind`.
parameters:
- $ref: "#/components/parameters/RunId"
requestBody:
required: true
content:
application/json:
schema:
$ref: "#/components/schemas/PetriPlatformRecordAppendRequest"
responses:
"200":
description: The record as stored
content:
application/json:
schema:
$ref: "#/components/schemas/PetriPlatformRecord"
"400":
description: The record is not a platform record
headers:
x-request-id:
$ref: "#/components/headers/XRequestId"
content:
application/json:
schema:
$ref: "#/components/schemas/ErrorResponse"
/api/v1/runs/{id}/stages/{stageId}/logs/output:
get:
operationId: getRunStageCommandLog
@ -6222,6 +6275,17 @@ components:
type: string
example: 18f3c2a9e1b4-42017-0-9f3a1c7e2b5d
PetriPlatformRecordKind:
name: kind
in: query
required: false
description: >-
Only the platform records of this kind, as its `kind` tag spells it
(`checkpoint`, `pull_request.created`, ...).
schema:
type: string
example: checkpoint
ArtifactFilename:
name: filename
in: query
@ -11000,6 +11064,71 @@ components:
items:
$ref: "#/components/schemas/PetriRecord"
PetriPlatformRecord:
description: >-
One of Fabro's platform records of a Petri run, as stored: the
record's JSON tagged by `kind`, its position in the run's platform
record sequence, and the Petri stage it belongs to when it belongs
to one.
type: object
required:
- seq
- recorded_at
- record
properties:
seq:
type: integer
format: uint64
description: The record's position in the run's platform records, from 1.
example: 4
recorded_at:
type: integer
format: uint64
description: Milliseconds since the Unix epoch when the record was stored.
example: 1758067200123
record:
type: object
additionalProperties: true
description: The platform record itself, tagged by `kind`.
execution:
type: integer
format: uint64
description: The Petri execution the record belongs to, with `firing`.
firing:
type: integer
format: uint64
description: The Petri firing the record belongs to, with `execution`.
PetriPlatformRecordAppendRequest:
description: One platform record to store for the run.
type: object
required:
- record
properties:
record:
type: object
additionalProperties: true
description: The platform record, tagged by `kind`.
execution:
type: integer
format: uint64
description: The Petri execution the record belongs to, with `firing`.
firing:
type: integer
format: uint64
description: The Petri firing the record belongs to, with `execution`.
PetriPlatformRecordList:
description: The run's platform records, in `seq` order.
type: object
required:
- records
properties:
records:
type: array
items:
$ref: "#/components/schemas/PetriPlatformRecord"
CommandTermination:
description: Terminal state for a command execution.
type: string

View file

@ -24,6 +24,10 @@
//! is lost for good cancels the run the same way, and the worker exits with
//! that loss as its error once the run has settled.
//!
//! Fabro's hooks ride the run with their platform records over the same
//! client: the checkpoint commit in the run's host workspace before every
//! durable finish, and its record after every route.
//!
//! The runtime's settings layer is left empty here: the run's graphs were
//! lowered and admitted at create time with the server's layer, and nothing
//! lowers again at execution. The model client is built from the worker's
@ -40,9 +44,12 @@ use fabro_client::{Client, ServerTarget};
use fabro_interview::ControlInterviewer;
use fabro_llm::credentials::{CredentialProvider, readiness};
use fabro_petri::engine::{self, Conclusion, Execution, RunRequest};
use fabro_petri::hooks::HooksSpec;
use fabro_petri::petri::OwnerId;
use fabro_petri::platform_records::HttpPlatformRecords;
use fabro_petri::runtime::{self, RuntimeSpec};
use fabro_petri::{HttpRunStore, admission};
use fabro_static::EnvVars;
use fabro_store::RunProjection;
use fabro_types::settings::run::RunMode;
use fabro_types::{FailureReason, RunId, RunTiming, StageOutcome, SuccessReason};
@ -141,6 +148,11 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
}
runner::set_worker_title(&run_id, WorkerTitlePhase::Running);
let hooks = HooksSpec::for_run(
Arc::new(HttpPlatformRecords::new(worker.client.clone_for_reuse())),
&worker.run_state.spec.settings.run,
)
.with_test_gates(test_checkpoint_gates());
let request = RunRequest {
run_id: run_id.to_string(),
run_dir: worker.run_dir.join("petri"),
@ -156,6 +168,7 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
.provider
.clone(),
cancel: cancel_token.clone(),
hooks: Some(hooks),
};
let run = Box::pin(engine::run(request));
tokio::pin!(run);
@ -227,6 +240,15 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
}
}
/// A test's checkpoint gate directory, when the server forwarded one.
#[expect(
clippy::disallowed_methods,
reason = "the gate directory is a test-only process-env facade the server forwards by name"
)]
fn test_checkpoint_gates() -> Option<PathBuf> {
std::env::var_os(EnvVars::FABRO_TEST_CHECKPOINT_GATES).map(PathBuf::from)
}
/// The runtime the worker hands Petri: no settings layer (nothing lowers
/// at execution), the model client over the worker's catalog and vault for
/// the providers whose credentials resolve, and the run's mode.

View file

@ -18,11 +18,13 @@ use std::sync::Arc;
use axum::extract::DefaultBodyLimit;
use axum::routing::{get, post};
use fabro_api::types::{
PetriAccess, PetriAppendRequest, PetriOpenRequest, PetriOpenResponse, PetriRecord,
PetriRecordList, PetriReleaseRequest, WriteBlobResponse,
PetriAccess, PetriAppendRequest, PetriOpenRequest, PetriOpenResponse, PetriPlatformRecord,
PetriPlatformRecordAppendRequest, PetriPlatformRecordList, PetriRecord, PetriRecordList,
PetriReleaseRequest, WriteBlobResponse,
};
use fabro_petri::petri::{Access, Digest, OwnerId, Record, StoreError};
use fabro_petri::run_store::{log_id_text, parse_log_id};
use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition, StoredPlatformRecord};
use fabro_types::BlobHash;
use fabro_util::error::collect_chain;
use serde_json::{Map, Value, json};
@ -51,6 +53,10 @@ pub(super) fn routes() -> Router<Arc<AppState>> {
post(write_blob).layer(DefaultBodyLimit::disable()),
)
.route("/runs/{id}/petri/blobs/{blobHash}", get(read_blob))
.route(
"/runs/{id}/petri/platform-records",
get(list_platform_records).post(append_platform_record),
)
}
#[derive(serde::Deserialize)]
@ -58,6 +64,11 @@ struct OwnerQuery {
owner: String,
}
#[derive(serde::Deserialize)]
struct KindQuery {
kind: Option<String>,
}
async fn open_run(
RequireWorkerRunScoped(id): RequireWorkerRunScoped,
State(state): State<Arc<AppState>>,
@ -207,6 +218,115 @@ async fn read_blob(
}
}
/// The run's platform records, of one kind when the query names it.
async fn list_platform_records(
RequireWorkerRunScoped(id): RequireWorkerRunScoped,
State(state): State<Arc<AppState>>,
Query(query): Query<KindQuery>,
) -> Response {
let store = state.stores.run_summaries.platform_records();
let records = match query.kind.as_deref() {
Some(kind) => match kind.parse::<PlatformRecordKind>() {
Ok(kind) => store.read_kind(&id, kind).await,
Err(_) => {
return ApiError::bad_request(format!("`{kind}` is not a platform record kind."))
.into_response();
}
},
None => store.read(&id).await,
};
match records {
Ok(records) => match records
.into_iter()
.map(|stored| wire_platform_record(&stored))
.collect::<Result<Vec<_>, _>>()
{
Ok(records) => Json(PetriPlatformRecordList { records }).into_response(),
Err(err) => err.into_response(),
},
Err(err) => platform_store_error_response(id, &err),
}
}
/// Store one platform record for the run and wake its projector.
async fn append_platform_record(
RequireWorkerRunScoped(id): RequireWorkerRunScoped,
State(state): State<Arc<AppState>>,
Json(request): Json<PetriPlatformRecordAppendRequest>,
) -> Response {
let record: PlatformRecord = match serde_json::from_value(Value::Object(request.record)) {
Ok(record) => record,
Err(err) => {
return ApiError::bad_request(format!("Invalid platform record: {err}"))
.into_response();
}
};
let position = match (request.execution, request.firing) {
(Some(execution), Some(firing)) => Some(StagePosition { execution, firing }),
_ => None,
};
let summaries = &state.stores.run_summaries;
match summaries
.platform_records()
.append(&id, &record, position)
.await
{
Ok(stored) => {
summaries.notify_platform_record(id);
match wire_platform_record(&stored) {
Ok(record) => Json(record).into_response(),
Err(err) => err.into_response(),
}
}
Err(err) => platform_store_error_response(id, &err),
}
}
/// A stored platform record as the wire carries it.
fn wire_platform_record(stored: &StoredPlatformRecord) -> Result<PetriPlatformRecord, ApiError> {
let record = serde_json::to_value(&stored.record).map_err(|err| {
ApiError::with_code(
StatusCode::INTERNAL_SERVER_ERROR,
format!(
"The stored platform record at seq {} does not encode: {err}",
stored.seq
),
"petri_store_failed",
)
})?;
let Value::Object(record) = record else {
return Err(ApiError::with_code(
StatusCode::INTERNAL_SERVER_ERROR,
format!(
"The stored platform record at seq {} is not a JSON object.",
stored.seq
),
"petri_store_failed",
));
};
Ok(PetriPlatformRecord {
seq: stored.seq,
recorded_at: stored.recorded_at,
record,
execution: stored.position.map(|position| position.execution),
firing: stored.position.map(|position| position.firing),
})
}
fn platform_store_error_response(run_id: RunId, err: &fabro_store::Error) -> Response {
tracing::error!(
run_id = %run_id,
error = %collect_chain(err).join(": "),
"platform record store failed"
);
ApiError::with_code(
StatusCode::INTERNAL_SERVER_ERROR,
"The platform record store failed; see the server log.",
"petri_store_failed",
)
.into_response()
}
fn unknown_log(log: &str) -> Response {
ApiError::bad_request(format!(
"`{log}` is not a Petri log: expected `coordinator`, `resources` or `execution <n>`."

View file

@ -20,7 +20,10 @@
//! the read-side item that follows.
//!
//! After a server restart, [`reconcile_on_startup`] hands a Petri run the
//! previous server left in flight back to a worker in resume mode.
//! previous server left in flight back to a worker in resume mode, once the
//! recovery protocol (`fabro_petri::recovery`) has brought every live
//! workspace to the snapshot its durable state names, or reports the run
//! failed when it cannot.
use std::collections::{BTreeMap, HashSet};
use std::sync::Arc;
@ -30,7 +33,10 @@ use fabro_config::{SettingsLayer, Storage};
use fabro_llm::selection;
use fabro_petri::check::{self, Bundle, CheckError, CheckRequest, Diagnostic, Launch};
use fabro_petri::engine::{self, Conclusion, Execution, RunRequest};
use fabro_petri::hooks::HooksSpec;
use fabro_petri::petri::StoreError;
use fabro_petri::platform_records::SqlitePlatformRecords;
use fabro_petri::recovery::{self, Recovery, RecoveryRequest};
use fabro_petri::runtime::{self, RuntimeSpec};
use fabro_petri::{SqliteRunStore, admission};
use fabro_types::settings::run::RunMode;
@ -337,6 +343,12 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
}
let (_, eligible) = state.resolve_llm_client_with_ready_ids().await;
let dry_run = run_state.spec.settings.run.execution.mode == RunMode::DryRun;
let hooks = HooksSpec::for_run(
Arc::new(SqlitePlatformRecords::new(Arc::clone(
&state.stores.run_summaries,
))),
&run_state.spec.settings.run,
);
let request = RunRequest {
run_id: run_id.to_string(),
run_dir: run_dir.join("petri"),
@ -345,6 +357,7 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
runtime: runtime_spec(&state, &eligible, dry_run),
provider: run_state.spec.settings.run.environment.provider.clone(),
cancel,
hooks: Some(hooks),
};
let result = Box::pin(engine::run(request)).await;
let timing = RunTiming {
@ -383,21 +396,22 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
}
/// Bring a Petri run the server left in flight back to its worker after a
/// restart: the run continues from its records, as Petri's own resume does.
/// restart: the run continues from its records, as Petri's own resume does,
/// on workspaces that match them.
///
/// The lease the previous worker held is released from outside, which
/// fences that worker should it still be alive; then the run is asked to
/// start again as a resume (`run.start_requested` with `resume`, then
/// `run.runnable`, the same pair the API's resume appends), and a managed
/// run is registered for the scheduler in resume mode when Petri's store
/// holds the run, else in start mode: a worker that died before it created
/// the run's record left nothing to continue from, so the run starts from
/// its admitted graphs.
///
/// Full recovery, where the workspace a resumed stage sees is restored to
/// the snapshot its durable state names, is the integration plan's F3.5.
/// Until it lands, a retained workspace is used as the previous worker left
/// it.
/// fences that worker should it still be alive. Then the recovery protocol
/// reads the run's durable execution state: a run with a failed checkpoint
/// is reported failed here and never resumed; otherwise every live
/// workspace on this host is verified against, reset to, or restored from
/// the snapshot its last durable finish names, and a finish with no
/// snapshot fails the run rather than resume it on stale files. The run is
/// then asked to start again as a resume (`run.start_requested` with
/// `resume`, then `run.runnable`, the same pair the API's resume appends),
/// and a managed run is registered for the scheduler in resume mode when
/// Petri's store holds the run, else in start mode: a worker that died
/// before it created the run's record left nothing to continue from, so the
/// run starts from its admitted graphs.
pub(crate) async fn reconcile_on_startup(
state: &Arc<AppState>,
run_id: RunId,
@ -412,8 +426,46 @@ pub(crate) async fn reconcile_on_startup(
return Err(anyhow::Error::new(err).context("releasing the Petri run's lease"));
}
};
let run_dir = Storage::new(state.server_storage_dir())
.run_scratch(&run_id)
.root()
.to_path_buf();
let mode = if held {
RunExecutionMode::Resume
let request = RecoveryRequest::for_run(
run_id,
run_dir.join("petri"),
Arc::new(SqliteRunStore::new(state.db_pool.clone())),
Arc::new(SqlitePlatformRecords::new(Arc::clone(
&state.stores.run_summaries,
))),
&run_state.spec.settings.run,
);
match recovery::recover(request)
.await
.map_err(|err| anyhow::Error::new(err).context("recovering the Petri run"))?
{
Recovery::Start => RunExecutionMode::Start,
Recovery::Resume { workspaces } => {
info!(
run_id = %run_id,
workspaces = workspaces.len(),
"Petri run's workspaces match its durable state"
);
RunExecutionMode::Resume
}
Recovery::Failed { reason } => {
warn!(
run_id = %run_id,
petri_key = %key,
error = %reason,
"Petri run left in flight by the previous server cannot resume; reporting it failed"
);
let (_, _, event) =
failed(FailureReason::WorkflowError, reason, RunTiming::default());
workflow_event::append_event(run_store, &run_id, &event).await?;
return Ok(());
}
}
} else {
RunExecutionMode::Start
};
@ -435,10 +487,6 @@ pub(crate) async fn reconcile_on_startup(
] {
workflow_event::append_event(run_store, &run_id, &event).await?;
}
let run_dir = Storage::new(state.server_storage_dir())
.run_scratch(&run_id)
.root()
.to_path_buf();
let mut runs = state.runs.lock().expect("runs lock poisoned");
runs.insert(
run_id,

View file

@ -62,6 +62,9 @@ const WORKER_ENV_ALLOWLIST: &[&str] = &[
EnvVars::PETRI_SANDBOX_PLUGIN_DEV,
EnvVars::PETRI_SANDBOX_DOCKER_HOST_ADDRESS,
EnvVars::PETRI_SANDBOX_ACTION_HOST_IMAGE,
// A test's checkpoint gates: the worker's hooks hold at a named point
// until the test releases them, so a crash can be placed there.
EnvVars::FABRO_TEST_CHECKPOINT_GATES,
];
const RENDER_GRAPH_ENV_ALLOWLIST: &[&str] = &[EnvVars::PATH, EnvVars::HOME, EnvVars::TMPDIR];

View file

@ -25,6 +25,8 @@ fabro-db = { path = "../../foundation/fabro-db" }
fabro-http.workspace = true
fabro-store = { path = "../fabro-store" }
fabro-types = { path = "../../foundation/fabro-types" }
fabro-checkpoint = { path = "../fabro-checkpoint" }
fabro-util = { path = "../../foundation/fabro-util" }
petri_runtime.workspace = true
petri_execution.workspace = true
petri_store.workspace = true
@ -46,6 +48,7 @@ tokio-util.workspace = true
tracing.workspace = true
[dev-dependencies]
fabro-petri = { path = ".", features = ["test-support"] }
fabro-auth = { path = "../../foundation/fabro-auth", features = ["test-support"] }
fabro-llm = { path = "../fabro-llm", features = ["test-support"] }
fabro-store = { path = "../fabro-store", features = ["test-support"] }

View file

@ -0,0 +1,864 @@
//! Git snapshots of a Petri run's workspaces: the checkpoint commit and
//! what recovery does with it.
//!
//! A Fabro stage's files are committed on the run branch of its workspace
//! before the stage's finish is recorded, so a durable finish implies a
//! durable snapshot (the integration plan's F3.1). The commit message
//! carries the snapshot's identity as trailers, the run key, execution,
//! firing and attempt, so a restart reconciles a missing platform record
//! from the branch alone.
//!
//! # Where the workspace is
//!
//! Petri's host backend keeps a scope's workspace under the run directory
//! at `scopes/<workspace id>/work`, the layout `HostExecutor::workspace_for`
//! names. This module reaches it there and runs `git` on the host, which
//! is where the worker, and the server at recovery, run. A Docker or
//! Daytona workspace lives inside its sandbox, out of reach of this module:
//! the hooks record that no snapshot was taken and recovery resumes such a
//! run on the retained sandbox as it was left.
//!
//! # The snapshot repository
//!
//! Every checkpoint commit is also pushed to a bare repository beside the
//! run's workspaces, `snapshots/<workspace id>.git`, under an immutable ref
//! per checkpoint (`refs/checkpoints/<execution>/<firing>/<attempt>`). A
//! workspace that is gone at recovery is restored from it, and the refs
//! are what recovery reconciles a missing record from.
use std::path::{Path, PathBuf};
use std::process::Stdio;
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::settings::run::RunCheckpointSettings;
use tokio::process::Command;
use tokio::{fs, time};
/// The failure class of a stage whose checkpoint commit failed: fatal to
/// the run, and terminal for a restart.
pub const CHECKPOINT_FAILED_CLASS: &str = "checkpoint_failed";
/// The effect kind of a checkpoint in its operation identity.
pub const CHECKPOINT_EFFECT: &str = "checkpoint";
pub const RUN_TRAILER: &str = "Fabro-Run";
pub const EXECUTION_TRAILER: &str = "Fabro-Execution";
pub const FIRING_TRAILER: &str = "Fabro-Firing";
pub const ATTEMPT_TRAILER: &str = "Fabro-Attempt";
const FOOTER: &str = "\u{2692}\u{fe0f} Generated with [Fabro](https://fabro.sh)";
const REFS_PREFIX: &str = "refs/checkpoints/";
/// Directories never committed, the legacy executor's list: build output
/// and dependency caches a stage regenerates.
pub const EXCLUDE_DIRS: &[&str] = &[
".git",
"node_modules",
".pnpm-store",
".npm",
"target",
".next",
"__pycache__",
".venv",
"venv",
".cache",
".tox",
".pytest_cache",
];
/// The identity of one snapshot: the attempt whose files it holds.
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct CheckpointKey {
pub execution: u64,
pub firing: u64,
pub attempt: u32,
}
impl CheckpointKey {
/// The immutable ref the snapshot is published under.
#[must_use]
pub fn snapshot_ref(self) -> String {
format!(
"{REFS_PREFIX}{}/{}/{}",
self.execution, self.firing, self.attempt
)
}
/// The operation identity of the checkpoint effect: the attempt's
/// decision in its execution, effect kind `checkpoint`.
#[must_use]
pub fn operation(self) -> OperationKey {
OperationKey {
execution: self.execution,
decision: DecisionRef::AttemptStart {
firing: self.firing,
attempt: self.attempt,
},
effect: CHECKPOINT_EFFECT.to_string(),
}
}
/// The key an operation identity names, when it is a checkpoint's.
#[must_use]
pub fn from_operation(operation: &OperationKey) -> Option<Self> {
match operation.decision {
DecisionRef::AttemptStart { firing, attempt }
if operation.effect == CHECKPOINT_EFFECT =>
{
Some(Self {
execution: operation.execution,
firing,
attempt,
})
}
DecisionRef::AttemptStart { .. }
| DecisionRef::ExecutionStart
| DecisionRef::Route { .. } => None,
}
}
fn from_ref(name: &str) -> Option<Self> {
let mut parts = name.strip_prefix(REFS_PREFIX)?.split('/');
let execution = parts.next()?.parse().ok()?;
let firing = parts.next()?.parse().ok()?;
let attempt = parts.next()?.parse().ok()?;
parts.next().is_none().then_some(Self {
execution,
firing,
attempt,
})
}
/// The key a checkpoint commit's message carries in its trailers.
#[must_use]
pub fn from_message(message: &str) -> Option<Self> {
Some(Self {
execution: trailer::parse(message, EXECUTION_TRAILER)?.parse().ok()?,
firing: trailer::parse(message, FIRING_TRAILER)?.parse().ok()?,
attempt: trailer::parse(message, ATTEMPT_TRAILER)?.parse().ok()?,
})
}
}
/// Why a snapshot could not be taken, found or restored.
#[derive(Debug, thiserror::Error)]
pub enum CheckpointError {
#[error("the workspace `{workspace}` does not exist at {}", path.display())]
WorkspaceMissing {
workspace: String,
path: PathBuf,
},
#[error("git {action} failed ({status}): {detail}")]
Command {
action: String,
status: String,
detail: String,
},
#[error("git {action} could not run")]
Spawn {
action: String,
#[source]
source: std::io::Error,
},
#[error("git {action} did not finish within {timeout:?}")]
TimedOut { action: String, timeout: Duration },
#[error("the workspace could not be prepared at {}", path.display())]
Io {
path: PathBuf,
#[source]
source: std::io::Error,
},
#[error("the restored workspace is at {actual}, not the snapshot {expected}")]
RestoreMismatch { expected: String, actual: String },
}
/// A checkpoint commit: the commit, and whether an earlier attempt of the
/// same operation had already made it.
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Snapshot {
pub sha: String,
pub reused: bool,
}
/// One published snapshot of a workspace.
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct PublishedSnapshot {
pub key: CheckpointKey,
pub sha: String,
}
/// The workspaces of one run on this host, and the Git operations Fabro
/// performs on them.
#[derive(Clone, Debug)]
pub struct RunWorkspaces {
run_dir: PathBuf,
run_id: String,
author: GitAuthor,
exclude_globs: Vec<String>,
timeout: Duration,
}
impl RunWorkspaces {
#[must_use]
pub fn new(
run_dir: PathBuf,
run_id: String,
author: GitAuthor,
settings: &RunCheckpointSettings,
) -> Self {
Self {
run_dir,
run_id,
author,
exclude_globs: settings.exclude_globs.clone(),
timeout: Duration::from_millis(settings.commit_timeout_ms.max(1)),
}
}
/// The run branch every workspace of the run commits on.
#[must_use]
pub fn run_branch(&self) -> String {
format!("fabro/run/{}", self.run_id)
}
/// Where the host backend keeps the workspace: `scopes/<id>/work` under
/// the run directory.
#[must_use]
pub fn workspace_path(&self, workspace: &str) -> PathBuf {
self.run_dir.join("scopes").join(workspace).join("work")
}
/// The bare repository the workspace's snapshots are published to.
#[must_use]
pub fn snapshot_repository(&self, workspace: &str) -> PathBuf {
self.run_dir
.join("snapshots")
.join(format!("{workspace}.git"))
}
/// Whether the workspace exists on this host.
pub async fn workspace_exists(&self, workspace: &str) -> bool {
fs::try_exists(self.workspace_path(workspace))
.await
.unwrap_or(false)
}
/// 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.
pub async fn commit(
&self,
workspace: &str,
key: CheckpointKey,
node: &str,
status: &str,
) -> Result<Snapshot, CheckpointError> {
let path = self.workspace_path(workspace);
if !self.workspace_exists(workspace).await {
return Err(CheckpointError::WorkspaceMissing {
workspace: workspace.to_string(),
path,
});
}
self.ensure_repository(&path).await?;
if let Some(existing) = self.published_sha(workspace, key).await? {
if self.head(&path).await?.as_deref() == Some(existing.as_str())
&& self.is_clean(&path).await?
{
return Ok(Snapshot {
sha: existing,
reused: true,
});
}
}
let mut add = vec![
"add".to_string(),
"-A".to_string(),
"--".to_string(),
".".to_string(),
];
add.extend(
EXCLUDE_DIRS
.iter()
.map(|dir| format!(":(glob,exclude)**/{dir}/**")),
);
add.extend(
self.exclude_globs
.iter()
.map(|glob| format!(":(glob,exclude){glob}")),
);
self.git(&path, "add", &add).await?;
let message = self.message(key, node, status);
let user_name = format!("user.name={}", self.author.name);
let user_email = format!("user.email={}", self.author.email);
self.git(&path, "commit", &[
"-c",
&user_name,
"-c",
&user_email,
"commit",
"-q",
"--allow-empty",
"-m",
&message,
])
.await?;
let sha = self.git(&path, "rev-parse", &["rev-parse", "HEAD"]).await?;
self.publish(workspace, &path, key, &sha).await?;
Ok(Snapshot { sha, reused: false })
}
/// The commit of `key`, from the snapshot repository first, else from
/// the workspace's own history by the trailers.
pub async fn find(
&self,
workspace: &str,
key: CheckpointKey,
) -> Result<Option<String>, CheckpointError> {
if let Some(sha) = self.published_sha(workspace, key).await? {
return Ok(Some(sha));
}
let path = self.workspace_path(workspace);
if !self.workspace_exists(workspace).await || self.head(&path).await?.is_none() {
return Ok(None);
}
let listed = self
.git(&path, "log", &[
"log",
"--format=%H",
"--extended-regexp",
&format!("--grep=^{EXECUTION_TRAILER}: {}$", key.execution),
&format!("--grep=^{FIRING_TRAILER}: {}$", key.firing),
&format!("--grep=^{ATTEMPT_TRAILER}: {}$", key.attempt),
"--all-match",
"HEAD",
])
.await?;
Ok(listed.lines().next().map(str::to_owned))
}
/// Every snapshot published for the workspace.
pub async fn published(
&self,
workspace: &str,
) -> Result<Vec<PublishedSnapshot>, CheckpointError> {
let repository = self.snapshot_repository(workspace);
if !fs::try_exists(&repository).await.unwrap_or(false) {
return Ok(Vec::new());
}
let listed = self
.git(&repository, "for-each-ref", &[
"for-each-ref",
"--format=%(refname) %(objectname)",
REFS_PREFIX,
])
.await?;
Ok(listed
.lines()
.filter_map(|line| {
let (name, sha) = line.split_once(' ')?;
Some(PublishedSnapshot {
key: CheckpointKey::from_ref(name)?,
sha: sha.to_string(),
})
})
.collect())
}
/// Whether `ancestor` is reachable from `descendant` in the workspace's
/// published history.
pub async fn is_ancestor(
&self,
workspace: &str,
ancestor: &str,
descendant: &str,
) -> Result<bool, CheckpointError> {
let repository = self.snapshot_repository(workspace);
Ok(self
.git_status(&repository, "merge-base", &[
"merge-base",
"--is-ancestor",
ancestor,
descendant,
])
.await?
.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> {
let path = self.workspace_path(workspace);
self.head(&path).await
}
/// Whether the workspace sits on `sha` with nothing changed since.
pub async fn matches(&self, workspace: &str, sha: &str) -> Result<bool, CheckpointError> {
let path = self.workspace_path(workspace);
Ok(self.head(&path).await?.as_deref() == Some(sha) && self.is_clean(&path).await?)
}
/// 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> {
let path = self.workspace_path(workspace);
self.git(&path, "reset", &["reset", "-q", "--hard", sha])
.await?;
let mut clean = vec!["clean".to_string(), "-fdq".to_string()];
for dir in EXCLUDE_DIRS {
clean.push("-e".to_string());
clean.push((*dir).to_string());
}
for glob in &self.exclude_globs {
clean.push("-e".to_string());
clean.push(glob.clone());
}
self.git(&path, "clean", &clean).await?;
Ok(())
}
/// Recreate a gone workspace from the published snapshot `key`, at
/// `sha`, on the run branch.
pub async fn restore(
&self,
workspace: &str,
key: CheckpointKey,
sha: &str,
) -> Result<(), CheckpointError> {
let path = self.workspace_path(workspace);
fs::create_dir_all(&path)
.await
.map_err(|source| CheckpointError::Io {
path: path.clone(),
source,
})?;
self.git(&path, "init", &["init", "-q"]).await?;
let repository = self.snapshot_repository(workspace);
let repository = repository.to_string_lossy().into_owned();
self.git(&path, "fetch", &[
"fetch",
"-q",
&repository,
&key.snapshot_ref(),
])
.await?;
let branch = self.run_branch();
self.git(&path, "checkout", &[
"checkout",
"-q",
"-B",
&branch,
"FETCH_HEAD",
])
.await?;
let actual = self.git(&path, "rev-parse", &["rev-parse", "HEAD"]).await?;
if actual != sha {
return Err(CheckpointError::RestoreMismatch {
expected: sha.to_string(),
actual,
});
}
Ok(())
}
/// The commit message: Fabro's subject, the footer, and the identity
/// trailers last, so `git interpret-trailers` and
/// [`CheckpointKey::from_message`] both read them.
fn message(&self, key: CheckpointKey, node: &str, status: &str) -> String {
let subject = format!("fabro({}): {node} ({status})", self.run_id);
let execution = key.execution.to_string();
let firing = key.firing.to_string();
let attempt = key.attempt.to_string();
let mut trailers = vec![
Trailer {
key: RUN_TRAILER,
value: &self.run_id,
},
Trailer {
key: EXECUTION_TRAILER,
value: &execution,
},
Trailer {
key: FIRING_TRAILER,
value: &firing,
},
Trailer {
key: ATTEMPT_TRAILER,
value: &attempt,
},
];
let defaults = GitAuthor::default();
let co_author = format!("{} <{}>", defaults.name, defaults.email);
if !self.author.is_default() {
trailers.push(Trailer {
key: "Co-Authored-By",
value: &co_author,
});
}
trailer::format_message(&subject, FOOTER, &trailers)
}
/// A repository on the run branch, initialised when the workspace has
/// none.
async fn ensure_repository(&self, path: &Path) -> Result<(), CheckpointError> {
if self
.git_status(path, "rev-parse", &["rev-parse", "--git-dir"])
.await?
.is_none()
{
self.git(path, "init", &["init", "-q"]).await?;
}
let branch = self.run_branch();
let current = self
.git_status(path, "symbolic-ref", &[
"symbolic-ref",
"-q",
"--short",
"HEAD",
])
.await?;
if current.as_deref() != Some(branch.as_str()) {
self.git(path, "checkout", &["checkout", "-q", "-B", &branch])
.await?;
}
Ok(())
}
async fn publish(
&self,
workspace: &str,
path: &Path,
key: CheckpointKey,
sha: &str,
) -> Result<(), CheckpointError> {
let repository = self.snapshot_repository(workspace);
if !fs::try_exists(&repository).await.unwrap_or(false) {
fs::create_dir_all(&repository)
.await
.map_err(|source| CheckpointError::Io {
path: repository.clone(),
source,
})?;
self.git(&repository, "init --bare", &["init", "-q", "--bare"])
.await?;
}
let refspec = format!("{sha}:{}", key.snapshot_ref());
let repository = repository.to_string_lossy().into_owned();
self.git(path, "push", &[
"push",
"-q",
"--force",
&repository,
&refspec,
])
.await?;
Ok(())
}
async fn published_sha(
&self,
workspace: &str,
key: CheckpointKey,
) -> Result<Option<String>, CheckpointError> {
let repository = self.snapshot_repository(workspace);
if !fs::try_exists(&repository).await.unwrap_or(false) {
return Ok(None);
}
self.git_status(&repository, "rev-parse", &[
"rev-parse",
"-q",
"--verify",
&key.snapshot_ref(),
])
.await
}
async fn head(&self, path: &Path) -> Result<Option<String>, CheckpointError> {
self.git_status(path, "rev-parse", &["rev-parse", "-q", "--verify", "HEAD"])
.await
}
async fn is_clean(&self, path: &Path) -> Result<bool, CheckpointError> {
let status = self.git(path, "status", &["status", "--porcelain"]).await?;
Ok(status.trim().is_empty())
}
/// Run `git` in `cwd`; a non-zero exit is the error.
async fn git<S: AsRef<str>>(
&self,
cwd: &Path,
action: &str,
args: &[S],
) -> Result<String, CheckpointError> {
let output = self.run(cwd, action, args).await?;
if output.status.success() {
Ok(String::from_utf8_lossy(&output.stdout).trim().to_owned())
} else {
Err(CheckpointError::Command {
action: action.to_string(),
status: output.status.to_string(),
detail: detail(&output.stderr),
})
}
}
/// Run `git` in `cwd`; a non-zero exit is `None`, for the queries whose
/// answer it is (an unborn `HEAD`, a missing ref, no repository).
async fn git_status<S: AsRef<str>>(
&self,
cwd: &Path,
action: &str,
args: &[S],
) -> Result<Option<String>, CheckpointError> {
let output = self.run(cwd, action, args).await?;
Ok(output
.status
.success()
.then(|| String::from_utf8_lossy(&output.stdout).trim().to_owned()))
}
async fn run<S: AsRef<str>>(
&self,
cwd: &Path,
action: &str,
args: &[S],
) -> Result<std::process::Output, CheckpointError> {
let mut command = Command::new("git");
command
.args([
"-c",
"core.hooksPath=/dev/null",
"-c",
"commit.gpgsign=false",
"-c",
"gc.auto=0",
"-c",
"advice.detachedHead=false",
"-c",
"init.defaultBranch=main",
])
.args(args.iter().map(AsRef::as_ref))
.current_dir(cwd)
.env("GIT_TERMINAL_PROMPT", "0")
.env_remove("GIT_DIR")
.env_remove("GIT_WORK_TREE")
.env_remove("GIT_INDEX_FILE")
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.kill_on_drop(true);
match time::timeout(self.timeout, command.output()).await {
Ok(Ok(output)) => Ok(output),
Ok(Err(source)) => Err(CheckpointError::Spawn {
action: action.to_string(),
source,
}),
Err(_) => Err(CheckpointError::TimedOut {
action: action.to_string(),
timeout: self.timeout,
}),
}
}
}
/// The tail of git's stderr for an error message: what the run's record
/// carries about the failure, bounded.
fn detail(stderr: &[u8]) -> String {
const LIMIT: usize = 512;
let text = String::from_utf8_lossy(stderr);
let text = text.trim();
if text.is_empty() {
return "no output".to_string();
}
let start = text.len().saturating_sub(LIMIT);
let start = text
.char_indices()
.map(|(index, _)| index)
.find(|index| *index >= start)
.unwrap_or(0);
text[start..].to_string()
}
#[cfg(test)]
mod tests {
use super::*;
fn workspaces(dir: &Path) -> RunWorkspaces {
RunWorkspaces::new(
dir.to_path_buf(),
"run-1".to_string(),
GitAuthor::default(),
&RunCheckpointSettings::default(),
)
}
#[test]
fn a_key_round_trips_through_its_ref_and_its_operation() {
let key = CheckpointKey {
execution: 3,
firing: 17,
attempt: 2,
};
assert_eq!(key.snapshot_ref(), "refs/checkpoints/3/17/2");
assert_eq!(CheckpointKey::from_ref(&key.snapshot_ref()), Some(key));
assert_eq!(CheckpointKey::from_ref("refs/heads/main"), None);
assert_eq!(CheckpointKey::from_operation(&key.operation()), Some(key));
assert_eq!(
CheckpointKey::from_operation(&OperationKey {
execution: 3,
decision: DecisionRef::Route {
firing: 17,
attempt: 2,
},
effect: CHECKPOINT_EFFECT.to_string(),
}),
None
);
}
#[test]
fn the_message_carries_the_identity_as_trailers_last() {
let dir = tempfile::tempdir().expect("a temp dir");
let key = CheckpointKey {
execution: 0,
firing: 4,
attempt: 1,
};
let message = workspaces(dir.path()).message(key, "build", "success");
assert!(message.starts_with("fabro(run-1): build (success)\n\n"));
assert_eq!(CheckpointKey::from_message(&message), Some(key));
assert_eq!(trailer::parse(&message, RUN_TRAILER), Some("run-1"));
}
#[tokio::test]
async fn a_commit_is_published_found_and_restored() {
let dir = tempfile::tempdir().expect("a temp dir");
let workspaces = workspaces(dir.path());
let workspace = "invocation-0-scope-0";
let path = workspaces.workspace_path(workspace);
fs::create_dir_all(&path).await.expect("the workspace");
fs::write(path.join("out.txt"), "one\n")
.await
.expect("a file");
let key = CheckpointKey {
execution: 0,
firing: 2,
attempt: 1,
};
let first = workspaces
.commit(workspace, key, "build", "success")
.await
.expect("the commit");
assert!(!first.reused);
let again = workspaces
.commit(workspace, key, "build", "success")
.await
.expect("the second commit");
assert_eq!(again, Snapshot {
sha: first.sha.clone(),
reused: true,
});
assert_eq!(
workspaces.find(workspace, key).await.expect("the lookup"),
Some(first.sha.clone())
);
assert_eq!(
workspaces.published(workspace).await.expect("the listing"),
vec![PublishedSnapshot {
key,
sha: first.sha.clone(),
}]
);
assert!(
workspaces
.matches(workspace, &first.sha)
.await
.expect("matches")
);
// The stage goes on, then the workspace is lost.
fs::write(path.join("out.txt"), "two\n")
.await
.expect("a change");
fs::write(path.join("scratch.txt"), "junk\n")
.await
.expect("an untracked file");
assert!(
!workspaces
.matches(workspace, &first.sha)
.await
.expect("matches")
);
workspaces
.reset(workspace, &first.sha)
.await
.expect("the reset");
assert_eq!(
fs::read_to_string(path.join("out.txt"))
.await
.expect("the file"),
"one\n"
);
assert!(
!fs::try_exists(path.join("scratch.txt"))
.await
.expect("exists")
);
fs::remove_dir_all(&path).await.expect("the workspace goes");
assert_eq!(
workspaces.find(workspace, key).await.expect("the lookup"),
Some(first.sha.clone()),
"the snapshot repository still knows the commit"
);
workspaces
.restore(workspace, key, &first.sha)
.await
.expect("the restore");
assert_eq!(
fs::read_to_string(path.join("out.txt"))
.await
.expect("the restored file"),
"one\n"
);
assert_eq!(
workspaces
.workspace_head(workspace)
.await
.expect("the head"),
Some(first.sha)
);
}
#[tokio::test]
async fn an_unusable_repository_fails_the_commit() {
let dir = tempfile::tempdir().expect("a temp dir");
let workspaces = workspaces(dir.path());
let workspace = "invocation-0-scope-0";
let path = workspaces.workspace_path(workspace);
fs::create_dir_all(&path).await.expect("the workspace");
fs::write(path.join(".git"), "garbage\n")
.await
.expect("a broken gitfile");
let key = CheckpointKey {
execution: 0,
firing: 2,
attempt: 1,
};
let error = workspaces
.commit(workspace, key, "build", "success")
.await
.expect_err("the commit fails");
assert!(matches!(error, CheckpointError::Command { .. }), "{error}");
}
#[test]
fn detail_keeps_the_tail_of_long_output() {
let long = "x".repeat(600);
assert_eq!(detail(long.as_bytes()).len(), 512);
assert_eq!(detail(b""), "no output");
}
}

View file

@ -16,16 +16,19 @@
//! reports is what the durable record says.
//!
//! What the standalone runner's defaults give the run: Petri's local hook
//! service for `[[run.hooks]]`, no `ExecutionHooks` of Fabro's own, the
//! [`Unattended`] interviewer that fails any question, no host tools, and
//! `Retention::Always` for every workspace, Fabro's default. Cancellation
//! service for `[[run.hooks]]`, the [`Unattended`] interviewer that fails
//! any question, no host tools, and `Retention::Always` for every
//! workspace, Fabro's default. With a [`HooksSpec`], Fabro's own
//! [`FabroHooks`] wrap the local service: the checkpoint commit before every
//! durable finish and its platform record after every route, with a failed
//! commit ending the run as a `checkpoint_failed` failure. Cancellation
//! rides the caller's token: when it fires, the root invocation is cancelled
//! politely and Petri records why.
//!
//! A resume here is Petri's own: the run continues from its records, and
//! sandbox leases are reconciled by label. Full recovery, where the
//! workspace a resumed stage sees is restored to the snapshot its durable
//! state names, is the integration plan's F3.5 and lands after this.
//! sandbox leases are reconciled by label. What the workspaces look like
//! when it does is the server's business before it relaunches the worker
//! ([`recovery`](crate::recovery)).
//!
//! No stage or agent event is projected into Fabro's tables here; the
//! caller appends only the run lifecycle events Fabro's read side needs to
@ -35,12 +38,13 @@
use std::path::PathBuf;
use std::sync::Arc;
use fabro_types::{FailureReason, SandboxProviderKind};
use fabro_types::{FailureReason, RunId, SandboxProviderKind};
use petri_execution::host::{self, HostError, HostRun};
use petri_execution::inspect::{self, InspectError, RunInspection};
use petri_execution::{
Access, CancelReason, InterviewDispatcher, InvocationId, RECEIPT_FILE, RunKey, RunStore,
};
use petri_runtime::driver::lifecycle::ExecutionHooks;
use petri_runtime::executor::Retention;
use petri_runtime::{RunOptions, SandboxBackend};
use tokio::fs;
@ -48,6 +52,7 @@ use tokio_util::sync::CancellationToken;
use tracing::{debug, info, warn};
use crate::admission::AdmittedGraphs;
use crate::hooks::{FabroHooks, HooksSpec};
use crate::interviewer::Unattended;
use crate::runtime::RuntimeSpec;
@ -78,6 +83,9 @@ pub struct RunRequest {
pub provider: SandboxProviderKind,
/// Fires to cancel the run.
pub cancel: CancellationToken,
/// Fabro's hooks: the checkpoint commit and its record. `None` runs
/// with Petri's local hook service alone.
pub hooks: Option<HooksSpec>,
}
/// The recorded status of a finished run.
@ -140,17 +148,37 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
options.run_key = Some(key.clone());
options.retention = Retention::Always;
options.sandbox.backend = backend;
let runtime = request
let mut runtime = request
.runtime
.runtime(true)
.store(Arc::clone(&request.store))
.options(options);
let fabro_hooks = request.hooks.map(|spec| {
let inner = runtime
.installed_hooks()
.unwrap_or_else(|| Arc::new(NoHooks));
let run_id = spec_run_id(&request.run_id);
Arc::new(FabroHooks::new(
spec,
inner,
run_id,
key.clone(),
request.run_dir.clone(),
Arc::clone(&request.store),
))
});
if let Some(hooks) = &fabro_hooks {
runtime = runtime.hooks(Arc::clone(hooks) as Arc<dyn ExecutionHooks>);
}
let dispatcher = InterviewDispatcher::new(Arc::new(Unattended));
let cancel = request.cancel.clone();
let mut cancel_task = None;
let with_handle = |handle: petri_execution::CoordinatorHandle, secrets| {
dispatcher.wire(handle.clone(), secrets);
if let Some(hooks) = &fabro_hooks {
hooks.attach(handle.clone());
}
cancel_task = Some(tokio::spawn(async move {
cancel.cancelled().await;
info!("cancelling the Petri run");
@ -187,9 +215,37 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
Err(error) => warn!(error = %error, "Petri run ended with a host error"),
}
let inspection = inspect(request.store.as_ref(), &key).await?;
outcome(inspection, result.err())
let mut outcome = outcome(inspection, result.err())?;
// A failed checkpoint cancelled the run; what Fabro reports is the
// checkpoint failure, not a cancellation.
if let Some(failure) = fabro_hooks
.as_ref()
.and_then(|hooks| hooks.checkpoint_failure())
{
outcome.status = RunStatus::Failed;
outcome.failure = Some(failure);
}
Ok(outcome)
}
/// The Fabro run id the run key names. A key that is not one (a test's
/// bare key) still gets hooks, under a fresh id for its platform records.
fn spec_run_id(run_id: &str) -> RunId {
run_id.parse().unwrap_or_else(|_| {
warn!(
run_id,
"the Petri run key is not a Fabro run id; platform records use a fresh id"
);
RunId::new()
})
}
/// No host hooks at all: what Fabro's hooks wrap when the runtime installed
/// none.
struct NoHooks;
impl ExecutionHooks for NoHooks {}
/// What the run's record says, read through a handle that holds no lease:
/// the same derivation [`run`] ends with, for a caller that only holds the
/// store, such as a test checking a finished run.

View file

@ -0,0 +1,566 @@
//! Fabro's awaited extension points on a Petri run: the checkpoint commit,
//! its platform record, and the run-level ends, wrapped around Petri's own
//! hook service so `[[run.hooks]]` keep running.
//!
//! [`FabroHooks`] implements Petri's `ExecutionHooks` and is installed with
//! `Runtime::hooks` by [`engine::run`](crate::engine::run). It holds the
//! hooks the runtime installed before it (Petri's local hook service behind
//! its adapter, which serves `[[run.hooks]]`) and forwards every point to
//! them, `run_finished` and `scope_released` included, the way Petri's
//! embedding host does. Its own work, at each point:
//!
//! - `prepare_result`: the checkpoint commit, before the `StepFinished` record
//! is appended, so a durable finish implies a durable snapshot. A stage that
//! failed on its own terms is committed like a successful one; only a
//! cancelled attempt is not. A failed commit is fatal to the run: the outcome
//! becomes a failure of class `checkpoint_failed`, the run is cancelled
//! through the coordinator handle, and `transition` refuses the firing's
//! routes, so no route is taken.
//! - `transition`: the platform checkpoint record, keyed on the Petri position
//! and the checkpoint's operation identity. A failed write is a recorded
//! problem on the transition, never a blocked route.
//! - `run_finished` and `scope_released`: forwarded, so the local service runs
//! `run_complete`, `run_failed` and `sandbox_cleanup` with the sandbox in
//! place. Fabro's own end-of-run work (the terminal lifecycle event,
//! notifications on it) is the run lifecycle path's, on the worker's and
//! server's side of the engine, and the workspace's retention is Petri's
//! (`Retention::Always`).
//!
//! # Operation identities
//!
//! Every external effect here is keyed on `(run key, execution, DecisionId,
//! effect kind)` from the hook context and deduplicated on retry: the
//! checkpoint's key is the attempt's decision in its execution, effect
//! `checkpoint`. A re-dispatched attempt whose commit already landed
//! reuses it when the workspace still sits on it unchanged (see
//! [`RunWorkspaces::commit`]); a reissued routing decision finds the
//! record, or the commit by its trailers, and writes nothing twice.
//!
//! # Where the workspace is
//!
//! The commit runs on the host, in the workspace Petri's host backend keeps
//! under the run directory (`crate::checkpoint`). A run on Docker or
//! Daytona has no workspace this process can reach; its hooks record that
//! no snapshot was taken and leave the run to continue as before.
use std::collections::{HashMap, HashSet};
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, MutexGuard, OnceLock, PoisonError};
use std::time::Duration;
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_util::error::collect_chain;
use petri_execution::{CancelReason, CoordinatorHandle, InvocationId, RunKey, RunStore};
use petri_runtime::driver::lifecycle::{
AdmitAttempt, AttemptDecision, ExecutionHooks, HookContext, Note, PrepareError, PrepareResult,
Prepared, Recorded, ResultOrigin, RunFinished, ScopeReleased, Transition, TransitionError,
TransitionReport,
};
use petri_runtime::ir::{FailureInfo, ScopeId, Status};
use serde_json::json;
use tokio::sync::OnceCell;
use tokio::{fs, time};
use tracing::{debug, info, warn};
use crate::checkpoint::{CHECKPOINT_FAILED_CLASS, CheckpointKey, RunWorkspaces};
use crate::platform_records::PlatformRecords;
use crate::workspace::{self, WorkspaceLookup};
/// The note kind the hooks record on a firing about its checkpoint.
pub const CHECKPOINT_NOTE: &str = "fabro.checkpoint";
/// How often a held checkpoint polls its test gate.
const GATE_POLL: Duration = Duration::from_millis(50);
/// What Fabro's hooks need beside the run: where the platform records go,
/// who authors the commits, and the checkpoint settings.
pub struct HooksSpec {
pub records: Arc<dyn PlatformRecords>,
pub author: GitAuthor,
pub checkpoint: RunCheckpointSettings,
/// Whether the run's workspaces are on this host (the local sandbox
/// provider). A run elsewhere takes no snapshot.
pub host_workspaces: bool,
/// A test's gate directory: a checkpoint point named by a `.hold` file
/// there waits for its `.release` file. `None` outside tests.
pub test_gates: Option<PathBuf>,
}
impl HooksSpec {
/// The spec a run's settings give: its Git author, its checkpoint
/// settings, and whether its sandbox provider keeps workspaces on this
/// host.
#[must_use]
pub fn for_run(records: Arc<dyn PlatformRecords>, settings: &RunNamespace) -> Self {
Self {
records,
author: settings
.git
.author
.as_ref()
.map(GitAuthor::from)
.unwrap_or_default(),
checkpoint: settings.checkpoint.clone(),
host_workspaces: settings.environment.provider == SandboxProviderKind::LOCAL,
test_gates: None,
}
}
#[must_use]
pub fn with_test_gates(mut self, gates: Option<PathBuf>) -> Self {
self.test_gates = gates;
self
}
}
/// Fabro's `ExecutionHooks`, around the hooks the runtime installed.
pub struct FabroHooks {
inner: Arc<dyn ExecutionHooks>,
run_id: RunId,
records: Arc<dyn PlatformRecords>,
workspaces: RunWorkspaces,
lookup: WorkspaceLookup,
host_workspaces: bool,
test_gates: Option<PathBuf>,
handle: OnceLock<CoordinatorHandle>,
/// The workspace and commit of every checkpoint this process made.
committed: Mutex<HashMap<CheckpointKey, (String, String)>>,
/// Which checkpoints have their platform record, loaded from the store
/// once and kept up to date with every append.
recorded: Mutex<HashSet<CheckpointKey>>,
recorded_loaded: OnceCell<()>,
/// Inherited workspaces resolved through the run's records.
inherited: Mutex<HashMap<InvocationId, Option<String>>>,
/// The checkpoint failure that ended the run, when one did.
failure: Mutex<Option<String>>,
unreachable_noted: AtomicBool,
}
impl FabroHooks {
/// Wrap `inner` (the hooks `Runtime::installed_hooks` returned) for the
/// run whose records are in `store` under `run_key`, with its
/// workspaces under `run_dir`.
#[must_use]
pub fn new(
spec: HooksSpec,
inner: Arc<dyn ExecutionHooks>,
run_id: RunId,
run_key: RunKey,
run_dir: PathBuf,
store: Arc<dyn RunStore>,
) -> Self {
let workspaces =
RunWorkspaces::new(run_dir, run_id.to_string(), spec.author, &spec.checkpoint);
Self {
inner,
run_id,
records: spec.records,
workspaces,
lookup: WorkspaceLookup::new(store, run_key),
host_workspaces: spec.host_workspaces,
test_gates: spec.test_gates,
handle: OnceLock::new(),
committed: Mutex::default(),
recorded: Mutex::default(),
recorded_loaded: OnceCell::new(),
inherited: Mutex::default(),
failure: Mutex::default(),
unreachable_noted: AtomicBool::new(false),
}
}
/// Hand the hooks the running coordinator, so a fatal checkpoint can
/// cancel the run. Called once, from the host's handle callback.
pub fn attach(&self, handle: CoordinatorHandle) {
if self.handle.set(handle).is_err() {
debug!("the coordinator handle was already attached to the hooks");
}
}
/// The checkpoint failure that ended the run, when one did: what the
/// engine reports the run failed with.
#[must_use]
pub fn checkpoint_failure(&self) -> Option<String> {
lock(&self.failure).clone()
}
/// The run's workspaces on this host, as the hooks reach them.
#[must_use]
pub fn workspaces(&self) -> &RunWorkspaces {
&self.workspaces
}
fn fail_run(&self, message: &str) {
let mut failure = lock(&self.failure);
if failure.is_none() {
*failure = Some(message.to_string());
}
drop(failure);
if let Some(handle) = self.handle.get() {
info!(run_id = %self.run_id, "cancelling the Petri run after a failed checkpoint");
handle.cancel_root_for(CancelReason::Control);
} else {
warn!(
run_id = %self.run_id,
"no coordinator handle is attached; the failed checkpoint cannot cancel the run"
);
}
}
/// The workspace id of `scope` in the context's invocation: the
/// isolated name when its workspace exists, else the inherited one the
/// records name, else the isolated name for the caller to report.
async fn workspace_of(&self, context: &HookContext, scope: ScopeId) -> Result<String, String> {
let isolated = workspace::isolated_workspace(context.invocation, scope);
if self.workspaces.workspace_exists(&isolated).await {
return Ok(isolated);
}
let cached = lock(&self.inherited).get(&context.invocation).cloned();
let inherited = if let Some(inherited) = cached {
inherited
} else {
let inherited = self
.lookup
.inherited(context.invocation)
.await
.map_err(|error| {
format!(
"the workspace of scope {scope} in invocation {} could not be found: {}",
context.invocation,
collect_chain(&error).join(": ")
)
})?;
lock(&self.inherited).insert(context.invocation, inherited.clone());
inherited
};
Ok(inherited.unwrap_or(isolated))
}
/// The checkpoint commit for one attempt's result. `Ok(Some)` is the
/// note to record, `Ok(None)` nothing to record, `Err` the fatal
/// failure message.
async fn snapshot(
&self,
context: &HookContext,
scope: ScopeId,
key: CheckpointKey,
node: &str,
status: &Status,
origin: ResultOrigin,
) -> Result<Option<Note>, String> {
if !self.host_workspaces {
if !self.unreachable_noted.swap(true, Ordering::SeqCst) {
warn!(
run_id = %self.run_id,
"the run's workspaces are not on this host; no checkpoint snapshot is taken"
);
}
return Ok(Some(Note::new(
CHECKPOINT_NOTE,
json!({
"execution": key.execution,
"firing": key.firing,
"attempt": key.attempt,
"skipped": "the workspace is not on this host",
}),
)));
}
let workspace = self.workspace_of(context, scope).await?;
if !self.workspaces.workspace_exists(&workspace).await {
// A skipped node or a driver-made outcome may precede the scope's
// environment; nothing of the stage's is on disk to snapshot.
if origin == ResultOrigin::Driver || matches!(status, Status::Skipped) {
return Ok(Some(Note::new(
CHECKPOINT_NOTE,
json!({
"execution": key.execution,
"firing": key.firing,
"attempt": key.attempt,
"workspace": workspace,
"skipped": "the workspace does not exist yet",
}),
)));
}
return Err(format!(
"the workspace `{workspace}` of scope {scope} does not exist at {}",
self.workspaces.workspace_path(&workspace).display()
));
}
self.gate("commit", node).await;
match self
.workspaces
.commit(&workspace, key, node, status.tag())
.await
{
Ok(snapshot) => {
debug!(
run_id = %self.run_id,
node,
execution = key.execution,
firing = key.firing,
attempt = key.attempt,
reused = snapshot.reused,
"checkpoint committed"
);
lock(&self.committed).insert(key, (workspace.clone(), snapshot.sha.clone()));
Ok(Some(Note::new(
CHECKPOINT_NOTE,
json!({
"execution": key.execution,
"firing": key.firing,
"attempt": key.attempt,
"workspace": workspace,
"git_commit_sha": snapshot.sha,
"reused": snapshot.reused,
}),
)))
}
Err(error) => Err(format!(
"the checkpoint commit of `{node}` failed: {}",
collect_chain(&error).join(": ")
)),
}
}
/// The checkpoint's platform record, once per operation identity.
async fn record(
&self,
context: &HookContext,
scope: ScopeId,
key: CheckpointKey,
) -> Result<(), String> {
self.recorded_loaded
.get_or_try_init(|| self.load_recorded())
.await?;
if lock(&self.recorded).contains(&key) {
return Ok(());
}
let committed = lock(&self.committed).get(&key).cloned();
let (workspace, sha) = if let Some(committed) = committed {
committed
} else {
let workspace = self.workspace_of(context, scope).await?;
let sha = self
.workspaces
.find(&workspace, key)
.await
.map_err(|error| {
format!(
"the checkpoint commit could not be looked up: {}",
collect_chain(&error).join(": ")
)
})?
.ok_or_else(|| {
format!(
"no checkpoint commit exists for execution {} firing {} attempt {}",
key.execution, key.firing, key.attempt
)
})?;
(workspace, sha)
};
let record = PlatformRecord::Checkpoint(CheckpointRecord {
execution: key.execution,
firing: key.firing,
attempt: Some(key.attempt),
workspace: Some(workspace),
git_commit_sha: Some(sha),
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(|error| {
format!(
"the checkpoint record could not be written: {}",
collect_chain(&error).join(": ")
)
})?;
lock(&self.recorded).insert(key);
Ok(())
}
/// The checkpoints already recorded for the run, read once: what a
/// resume's reissued routing decisions must not record again.
async fn load_recorded(&self) -> Result<(), String> {
let stored = self
.records
.read_kind(&self.run_id, PlatformRecordKind::Checkpoint)
.await
.map_err(|error| {
format!(
"the run's checkpoint records could not be read: {}",
collect_chain(&error).join(": ")
)
})?;
let mut recorded = lock(&self.recorded);
for record in stored {
let PlatformRecord::Checkpoint(checkpoint) = &record.record else {
continue;
};
if let Some(key) = checkpoint
.operation
.as_ref()
.and_then(CheckpointKey::from_operation)
{
recorded.insert(key);
}
}
Ok(())
}
/// Hold at a test gate when one is set for this point and node.
async fn gate(&self, point: &str, node: &str) {
let Some(dir) = &self.test_gates else {
return;
};
let hold = dir.join(format!("{point}.{node}.hold"));
if !fs::try_exists(&hold).await.unwrap_or(false) {
return;
}
let release = dir.join(format!("{point}.{node}.release"));
info!(point, node, "checkpoint held at a test gate");
while !fs::try_exists(&release).await.unwrap_or(false) {
time::sleep(GATE_POLL).await;
}
info!(point, node, "checkpoint released by its test gate");
}
}
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex.lock().unwrap_or_else(PoisonError::into_inner)
}
fn is_checkpoint_failure(status: &Status) -> bool {
matches!(status, Status::Failure(info) if info.class.as_str() == CHECKPOINT_FAILED_CLASS)
}
#[async_trait::async_trait]
impl ExecutionHooks for FabroHooks {
async fn before_attempt(
&self,
context: &HookContext,
request: AdmitAttempt,
) -> AttemptDecision {
self.inner.before_attempt(context, request).await
}
async fn prepare_result(
&self,
context: &HookContext,
request: PrepareResult,
) -> Result<Prepared, PrepareError> {
let node = request.view.node_name().to_owned();
let scope = request.view.scope;
let key = CheckpointKey {
execution: context.execution.raw(),
firing: request.view.firing.raw(),
attempt: request.view.attempt.raw(),
};
let original = request.outcome.status.clone();
let origin = request.origin;
let mut prepared = self.inner.prepare_result(context, request).await?;
let effective = prepared.adjustment.status.clone().unwrap_or(original);
if matches!(effective, Status::Cancelled) {
return Ok(prepared);
}
match self
.snapshot(context, scope, key, &node, &effective, origin)
.await
{
Ok(Some(note)) => prepared.notes.push(note),
Ok(None) => {}
Err(message) => {
warn!(
run_id = %self.run_id,
node,
execution = key.execution,
firing = key.firing,
attempt = key.attempt,
error = %message,
"checkpoint failed; the run ends"
);
self.fail_run(&message);
prepared.adjustment.status = Some(Status::Failure(
FailureInfo::new(message.clone()).with_class(CHECKPOINT_FAILED_CLASS),
));
prepared.adjustment.reason = Some(message);
}
}
Ok(prepared)
}
async fn after_record(&self, context: &HookContext, recorded: Recorded) -> Vec<Note> {
self.inner.after_record(context, recorded).await
}
async fn transition(
&self,
context: &HookContext,
transition: Transition,
) -> Result<TransitionReport, TransitionError> {
if is_checkpoint_failure(&transition.outcome.status) {
return Err(TransitionError::new(
"the stage's checkpoint commit failed; no route is taken",
));
}
let node = transition.view.node_name().to_owned();
let scope = transition.view.scope;
let key = CheckpointKey {
execution: context.execution.raw(),
firing: transition.view.firing.raw(),
attempt: transition.view.attempt.raw(),
};
let mut problems = Vec::new();
if self.host_workspaces {
self.gate("record", &node).await;
if let Err(problem) = self.record(context, scope, key).await {
warn!(
run_id = %self.run_id,
node,
execution = key.execution,
firing = key.firing,
error = %problem,
"the checkpoint record was not written"
);
problems.push(problem);
}
}
let mut report = self.inner.transition(context, transition).await?;
report.problems.extend(problems);
Ok(report)
}
async fn run_finished(&self, context: &HookContext, finished: RunFinished) -> Vec<Note> {
info!(
run_id = %self.run_id,
status = ?finished.status,
failure = finished.failure.as_deref().unwrap_or(""),
"Petri run finished; running the run-end hooks"
);
self.inner.run_finished(context, finished).await
}
async fn scope_released(&self, context: &HookContext, released: ScopeReleased) -> Vec<Note> {
debug!(
run_id = %self.run_id,
scope = %released.scope,
outcome = ?released.outcome,
"scope released; running the sandbox cleanup hooks"
);
self.inner.scope_released(context, released).await
}
}

View file

@ -23,22 +23,37 @@
//! - [`HttpRunStore`]: the same store as a run's worker process reaches it,
//! over the server's API with the worker's token and its launch id as the
//! lease owner;
//! - the platform adapters still to come: hooks, interviews over Fabro's API,
//! secrets, output storage, the run tools, the event projection.
//! - [`hooks`]: Fabro's `ExecutionHooks`, the checkpoint commit in
//! `prepare_result` and its platform record in `transition`, around Petri's
//! own hook service for `[[run.hooks]]`;
//! - [`checkpoint`]: the Git snapshots of a run's host workspaces and the
//! snapshot repository they are published to;
//! - [`recovery`]: the resume-on-restart protocol, which brings every live
//! workspace to the snapshot its durable state names before the run goes back
//! to a worker;
//! - [`platform_records`]: Fabro's platform records as the adapters reach them,
//! in the server's database or over its API from a worker;
//! - the platform adapters still to come: interviews over Fabro's API, secrets,
//! output storage, the run tools, the event projection.
//!
//! The Petri packages are pinned by revision in the workspace `Cargo.toml`
//! under `petri_*` keys.
pub mod admission;
pub mod check;
pub mod checkpoint;
pub mod engine;
pub mod hooks;
pub mod http_store;
pub mod interviewer;
pub mod petri;
pub mod platform_records;
pub mod recovery;
pub mod run_store;
pub mod runtime;
#[cfg(feature = "test-support")]
pub mod test_support;
pub mod workspace;
pub use http_store::HttpRunStore;
pub use run_store::SqliteRunStore;

View file

@ -0,0 +1,188 @@
//! Fabro's platform records as a Petri run's adapters reach them.
//!
//! The records themselves are `fabro_store::platform_records`: one table
//! beside Petri's records, one typed enum of kinds. What differs is where
//! the adapter runs. In the server process (the in-process test path, and
//! startup recovery) the table is reached directly, through
//! [`SqlitePlatformRecords`]; in a run's worker process it is reached over
//! the server's API with the worker's token, through
//! [`HttpPlatformRecords`], as the run's Petri records are. Both answer
//! the one [`PlatformRecords`] interface the hooks and recovery use.
use std::fmt;
use std::sync::Arc;
use async_trait::async_trait;
use fabro_api::types::{PetriPlatformRecord, PetriPlatformRecordAppendRequest};
use fabro_client::Client;
use fabro_store::{
PlatformRecord, PlatformRecordKind, PlatformRecordStore, RunSummaryStore, StagePosition,
StoredPlatformRecord,
};
use fabro_types::RunId;
use serde_json::Value;
/// Why a platform record could not be stored or read.
#[derive(Debug, thiserror::Error)]
pub enum PlatformRecordError {
#[error("the platform record store failed")]
Store(#[source] fabro_store::Error),
#[error("the platform record request to the server failed")]
Api(#[source] anyhow::Error),
#[error("the platform record does not encode as JSON")]
Encode(#[source] serde_json::Error),
}
/// The platform records of a run, wherever the adapter runs.
#[async_trait]
pub trait PlatformRecords: Send + Sync {
/// Store a record at the run's next seq, tied to a Petri stage when it
/// belongs to one.
async fn append(
&self,
run_id: &RunId,
record: &PlatformRecord,
position: Option<StagePosition>,
) -> Result<StoredPlatformRecord, PlatformRecordError>;
/// The run's records of one kind, in seq order.
async fn read_kind(
&self,
run_id: &RunId,
kind: PlatformRecordKind,
) -> Result<Vec<StoredPlatformRecord>, PlatformRecordError>;
}
/// The table in the server's database, with the projector's wake-up after
/// each append.
#[derive(Clone)]
pub struct SqlitePlatformRecords {
store: PlatformRecordStore,
summaries: Arc<RunSummaryStore>,
}
impl fmt::Debug for SqlitePlatformRecords {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("SqlitePlatformRecords")
.finish_non_exhaustive()
}
}
impl SqlitePlatformRecords {
#[must_use]
pub fn new(summaries: Arc<RunSummaryStore>) -> Self {
Self {
store: summaries.platform_records(),
summaries,
}
}
}
#[async_trait]
impl PlatformRecords for SqlitePlatformRecords {
async fn append(
&self,
run_id: &RunId,
record: &PlatformRecord,
position: Option<StagePosition>,
) -> Result<StoredPlatformRecord, PlatformRecordError> {
let stored = self
.store
.append(run_id, record, position)
.await
.map_err(PlatformRecordError::Store)?;
self.summaries.notify_platform_record(*run_id);
Ok(stored)
}
async fn read_kind(
&self,
run_id: &RunId,
kind: PlatformRecordKind,
) -> Result<Vec<StoredPlatformRecord>, PlatformRecordError> {
self.store
.read_kind(run_id, kind)
.await
.map_err(PlatformRecordError::Store)
}
}
/// The table as a run's worker reaches it: the server's
/// `/api/v1/runs/{id}/petri/platform-records` endpoints with the worker's
/// token.
pub struct HttpPlatformRecords {
client: Client,
}
impl fmt::Debug for HttpPlatformRecords {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("HttpPlatformRecords")
.field("server", &self.client.base_url())
.finish_non_exhaustive()
}
}
impl HttpPlatformRecords {
#[must_use]
pub fn new(client: Client) -> Self {
Self { client }
}
}
#[async_trait]
impl PlatformRecords for HttpPlatformRecords {
async fn append(
&self,
run_id: &RunId,
record: &PlatformRecord,
position: Option<StagePosition>,
) -> Result<StoredPlatformRecord, PlatformRecordError> {
let Value::Object(record) =
serde_json::to_value(record).map_err(PlatformRecordError::Encode)?
else {
return Err(PlatformRecordError::Api(anyhow::anyhow!(
"a platform record encodes as a JSON object"
)));
};
let body = PetriPlatformRecordAppendRequest {
record,
execution: position.map(|position| position.execution),
firing: position.map(|position| position.firing),
};
let stored = self
.client
.append_petri_platform_record(run_id, body)
.await
.map_err(PlatformRecordError::Api)?;
decode(stored)
}
async fn read_kind(
&self,
run_id: &RunId,
kind: PlatformRecordKind,
) -> Result<Vec<StoredPlatformRecord>, PlatformRecordError> {
let records = self
.client
.list_petri_platform_records(run_id, Some(&kind.to_string()))
.await
.map_err(PlatformRecordError::Api)?;
records.into_iter().map(decode).collect()
}
}
/// A wire record back into the store's shape.
fn decode(wire: PetriPlatformRecord) -> Result<StoredPlatformRecord, PlatformRecordError> {
let record: PlatformRecord =
serde_json::from_value(Value::Object(wire.record)).map_err(PlatformRecordError::Encode)?;
let position = match (wire.execution, wire.firing) {
(Some(execution), Some(firing)) => Some(StagePosition { execution, firing }),
_ => None,
};
Ok(StoredPlatformRecord {
seq: wire.seq,
recorded_at: wire.recorded_at,
record,
position,
})
}

View file

@ -0,0 +1,438 @@
//! Resume on restart: the recovery protocol whose rule is that the
//! workspace a resumed stage sees matches Petri's durable execution state
//! (the integration plan's F3.5).
//!
//! For a Petri run the server finds in flight at startup, once the previous
//! worker's lease is released, [`recover`] reads the durable execution
//! state through `inspect_run` and decides:
//!
//! - a run with a `checkpoint_failed` finish anywhere is reported failed and
//! not resumed: a failed checkpoint cancelled it, and nothing of it is
//! reconciled;
//! - otherwise, for every live execution, the last durable finish names the
//! snapshot its workspace must sit on: the checkpoint record's commit, or,
//! when the record was lost to the crash, the commit found by its key in the
//! workspace's snapshot repository or history, which is then recorded again;
//! - a workspace that survives is verified to sit on that commit, unchanged, or
//! reset to it; a workspace that is gone is restored from the run's snapshot
//! repository into a fresh directory;
//! - a durable finish with no snapshot fails the run with a named error rather
//! than resume it on stale files.
//!
//! Every child invocation's scope has its own snapshots, keyed by
//! execution; a nested invocation that inherits its caller's sandbox shares
//! the caller's workspace, and the workspace is brought to the newest of
//! the live executions' snapshots on it.
//!
//! A run whose workspaces are not on this host (Docker, Daytona) is resumed
//! on its retained sandbox as it was left: the snapshot side of the
//! protocol reaches only host workspaces.
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 petri_execution::host::{self, HostError};
use petri_execution::inspect::{self, ExecutionInspection, InspectError};
use petri_execution::{Access, InvocationId, RunKey, RunStore};
use petri_store::StoreError;
use tracing::{info, warn};
use crate::checkpoint::{CHECKPOINT_FAILED_CLASS, CheckpointError, CheckpointKey, RunWorkspaces};
use crate::platform_records::{PlatformRecordError, PlatformRecords};
use crate::workspace::{WorkspaceLookup, WorkspaceLookupError};
/// What recovery needs: the run, where its workspaces are, its records.
pub struct RecoveryRequest {
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,
}
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.
#[must_use]
pub fn for_run(
run_id: RunId,
run_dir: PathBuf,
store: Arc<dyn RunStore>,
records: Arc<dyn PlatformRecords>,
settings: &RunNamespace,
) -> Self {
Self {
run_id,
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,
}
}
}
/// What was done to one workspace.
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum WorkspaceAction {
/// It sat on the snapshot, unchanged.
Verified,
/// It was brought back to the snapshot.
Reset,
/// It was gone and was recreated from the snapshot repository.
Restored,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct RecoveredWorkspace {
pub workspace: String,
pub sha: String,
pub action: WorkspaceAction,
}
/// What the server does with the run next.
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum Recovery {
/// The store never held the run: it starts from its admitted graphs.
Start,
/// The run continues from its records, its live workspaces on their
/// snapshots.
Resume { workspaces: Vec<RecoveredWorkspace> },
/// The run cannot continue and is reported failed.
Failed { reason: String },
}
/// Why recovery could not decide: the records could not be read, or a
/// workspace could not be brought to its snapshot.
#[derive(Debug, thiserror::Error)]
pub enum RecoveryError {
#[error("the run's record could not be opened")]
Open(#[source] StoreError),
#[error("the run's coordinator state could not be read")]
State(#[source] HostError),
#[error("the run's record could not be inspected")]
Inspect(#[source] InspectError),
#[error("the run's workspaces could not be named")]
Lookup(#[source] WorkspaceLookupError),
#[error("the run's checkpoint records could not be read or written")]
Records(#[source] PlatformRecordError),
#[error("the workspace `{workspace}` could not be brought to its snapshot")]
Workspace {
workspace: String,
#[source]
source: CheckpointError,
},
}
/// One live execution's last durable finish and the snapshot it names.
struct Target {
execution: u64,
key: CheckpointKey,
}
/// Decide how the run continues, and bring its workspaces to their
/// snapshots.
pub async fn recover(request: RecoveryRequest) -> Result<Recovery, RecoveryError> {
let key = RunKey::new(request.run_id.to_string());
let logs = match request.store.open(&key, Access::Read).await {
Ok(logs) => logs,
Err(StoreError::NotFound { .. }) => return Ok(Recovery::Start),
Err(error) => return Err(RecoveryError::Open(error)),
};
// A record with no root invocation (the worker died between creating
// the run and declaring it) has nothing to reconcile; the worker's
// resume reports it as such.
let state = host::stored_state(&*logs)
.await
.map_err(RecoveryError::State)?;
if !state.invocations.contains_key(&InvocationId::ROOT) {
return Ok(Recovery::Resume {
workspaces: Vec::new(),
});
}
let inspection = inspect::inspect_run(&*logs)
.await
.map_err(RecoveryError::Inspect)?;
drop(logs);
if let Some(failed) = checkpoint_failure(&inspection.executions) {
return Ok(Recovery::Failed { reason: failed });
}
if !request.host_workspaces {
warn!(
run_id = %request.run_id,
"the run's workspaces are not on this host; resuming on the retained sandbox as it was left"
);
return Ok(Recovery::Resume {
workspaces: Vec::new(),
});
}
let workspaces = RunWorkspaces::new(
request.run_dir.clone(),
request.run_id.to_string(),
request.author.clone(),
&request.checkpoint,
);
let lookup = WorkspaceLookup::new(Arc::clone(&request.store), key);
let recorded = recorded_checkpoints(&*request.records, &request.run_id).await?;
// The snapshot each live execution's workspace must sit on.
let mut candidates: BTreeMap<String, Vec<(Target, String)>> = BTreeMap::new();
for execution in inspection
.executions
.iter()
.filter(|execution| execution.status == "running")
{
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(
&*request.records,
&request.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(Recovery::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 recovered = Vec::new();
for (workspace, targets) in candidates {
let sha = newest(&workspaces, &workspace, &targets).await?;
let action = bring_to(&workspaces, &workspace, &sha, &targets).await?;
info!(
run_id = %request.run_id,
workspace,
sha,
action = ?action,
"workspace brought to its durable snapshot"
);
recovered.push(RecoveredWorkspace {
workspace,
sha,
action,
});
}
Ok(Recovery::Resume {
workspaces: recovered,
})
}
/// The reason a run with a failed checkpoint is reported failed, when it
/// has one.
fn checkpoint_failure(executions: &[ExecutionInspection]) -> Option<String> {
executions.iter().find_map(|execution| {
execution
.engine
.as_ref()?
.attempts
.iter()
.find_map(|attempt| {
let failure = attempt.failure.as_ref()?;
(failure.class.as_str() == CHECKPOINT_FAILED_CLASS).then(|| {
format!(
"the checkpoint of {} (execution {} firing {} attempt {}) failed: {}",
attempt.node.as_deref().unwrap_or("a stage"),
execution.execution,
attempt.firing,
attempt.attempt,
failure.message
)
})
})
})
}
/// The last `StepFinished` of an execution's log.
fn last_finish(execution: &ExecutionInspection) -> Option<Target> {
let attempt = execution.engine.as_ref()?.attempts.last()?;
Some(Target {
execution: execution.execution.raw(),
key: CheckpointKey {
execution: execution.execution.raw(),
firing: attempt.firing,
attempt: attempt.attempt,
},
})
}
/// The run's checkpoint records by key: the workspace they name and the
/// commit.
async fn recorded_checkpoints(
records: &dyn PlatformRecords,
run_id: &RunId,
) -> Result<BTreeMap<CheckpointKey, (Option<String>, String)>, RecoveryError> {
let stored = records
.read_kind(run_id, PlatformRecordKind::Checkpoint)
.await
.map_err(RecoveryError::Records)?;
let mut recorded = BTreeMap::new();
for record in stored {
let PlatformRecord::Checkpoint(checkpoint) = record.record else {
continue;
};
let key = checkpoint
.operation
.as_ref()
.and_then(CheckpointKey::from_operation);
if let (Some(key), Some(sha)) = (key, checkpoint.git_commit_sha) {
recorded.insert(key, (checkpoint.workspace, sha));
}
}
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 workspace onto `sha`.
async fn bring_to(
workspaces: &RunWorkspaces,
workspace: &str,
sha: &str,
targets: &[(Target, String)],
) -> Result<WorkspaceAction, RecoveryError> {
let failed = |source| RecoveryError::Workspace {
workspace: workspace.to_string(),
source,
};
if workspaces.workspace_exists(workspace).await {
if workspaces.matches(workspace, sha).await.map_err(failed)? {
return Ok(WorkspaceAction::Verified);
}
workspaces.reset(workspace, sha).await.map_err(failed)?;
return Ok(WorkspaceAction::Reset);
}
let key = targets
.iter()
.find(|(_, candidate)| candidate == sha)
.map_or(targets[0].0.key, |(target, _)| target.key);
workspaces
.restore(workspace, key, sha)
.await
.map_err(failed)?;
Ok(WorkspaceAction::Restored)
}

View file

@ -1,5 +1,71 @@
//! Petri's test kit, for Fabro crates that check a store implementation
//! against Petri's contract from their own tests. Compiled only with the
//! against Petri's contract from their own tests, and an in-memory platform
//! record store for tests of the hooks and recovery. Compiled only with the
//! `test-support` feature, which a dev-dependency turns on.
use std::collections::HashMap;
use std::sync::{Mutex, MutexGuard, PoisonError};
use async_trait::async_trait;
use fabro_store::platform_records::now_ms;
use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition, StoredPlatformRecord};
use fabro_types::RunId;
pub use petri_testkit::run_store;
use crate::platform_records::{PlatformRecordError, PlatformRecords};
/// Platform records kept in memory, per run, in seq order.
#[derive(Debug, Default)]
pub struct MemoryPlatformRecords {
runs: Mutex<HashMap<RunId, Vec<StoredPlatformRecord>>>,
}
impl MemoryPlatformRecords {
#[must_use]
pub fn new() -> Self {
Self::default()
}
/// 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()
}
}
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(
&self,
run_id: &RunId,
record: &PlatformRecord,
position: Option<StagePosition>,
) -> Result<StoredPlatformRecord, PlatformRecordError> {
let mut runs = lock(&self.runs);
let records = runs.entry(*run_id).or_default();
let stored = StoredPlatformRecord {
seq: records.len() as u64 + 1,
recorded_at: now_ms(),
record: record.clone(),
position,
};
records.push(stored.clone());
Ok(stored)
}
async fn read_kind(
&self,
run_id: &RunId,
kind: PlatformRecordKind,
) -> Result<Vec<StoredPlatformRecord>, PlatformRecordError> {
Ok(self
.records(run_id)
.into_iter()
.filter(|record| record.record.kind() == kind)
.collect())
}
}

View file

@ -0,0 +1,116 @@
//! Which workspace a Petri scope runs in, from the run's own records.
//!
//! Petri names an isolated scope's workspace after its invocation and scope
//! (`invocation-<n>-scope-<m>`), and a nested invocation that inherits its
//! caller's sandbox shares the caller's workspace through the lease the
//! coordinator recorded. The hooks and recovery both need the workspace id
//! behind a scope, and both read it the same way here: the direct name
//! when its workspace exists on this host, else the lease the invocation's
//! declaration names, resolved through the resource log.
use std::sync::Arc;
use petri_execution::host::{self, HostError};
use petri_execution::{
Access, InvocationId, ResourceError, ResourceStore, RunKey, RunLogs, RunStore, SandboxBinding,
};
use petri_runtime::executor::WorkspaceId;
use petri_runtime::ir::ScopeId;
use petri_store::StoreError;
/// Why a workspace could not be named from the run's records.
#[derive(Debug, thiserror::Error)]
pub enum WorkspaceLookupError {
#[error("the run's record could not be opened")]
Open(#[source] StoreError),
#[error("the run's coordinator state could not be read")]
State(#[source] HostError),
#[error("the run's resource log could not be read")]
Resources(#[source] ResourceError),
#[error("invocation {invocation} is not in the run's record")]
UnknownInvocation { invocation: InvocationId },
}
/// The workspace id of an isolated scope: what the coordinator allocates
/// for `scope` in `invocation`.
#[must_use]
pub fn isolated_workspace(invocation: InvocationId, scope: ScopeId) -> String {
WorkspaceId::scoped(Some(&invocation.workspace_prefix()), scope)
.as_str()
.to_owned()
}
/// The workspaces of a run, read from its records through a handle that
/// holds no lease.
pub struct WorkspaceLookup {
store: Arc<dyn RunStore>,
key: RunKey,
}
impl WorkspaceLookup {
#[must_use]
pub fn new(store: Arc<dyn RunStore>, key: RunKey) -> Self {
Self { store, key }
}
async fn logs(&self) -> Result<Arc<dyn RunLogs>, WorkspaceLookupError> {
self.store
.open(&self.key, Access::Read)
.await
.map_err(WorkspaceLookupError::Open)
}
/// The workspace an invocation inherited from its caller, or `None`
/// when the invocation owns its sandboxes.
pub async fn inherited(
&self,
invocation: InvocationId,
) -> Result<Option<String>, WorkspaceLookupError> {
let logs = self.logs().await?;
let state = host::stored_state(&*logs)
.await
.map_err(WorkspaceLookupError::State)?;
let declared = state
.invocations
.get(&invocation)
.ok_or(WorkspaceLookupError::UnknownInvocation { invocation })?;
match declared.declaration.sandbox {
SandboxBinding::Isolated => Ok(None),
SandboxBinding::Inherited { lease } => {
let resources = ResourceStore::load(&logs)
.await
.map_err(WorkspaceLookupError::Resources)?;
let record = resources
.resolve(lease)
.map_err(WorkspaceLookupError::Resources)?;
Ok(Some(record.workspace.as_str().to_owned()))
}
}
}
/// Every workspace an invocation runs in: its own leases' workspaces,
/// or the one it inherited.
pub async fn of_invocation(
&self,
invocation: InvocationId,
) -> Result<Vec<String>, WorkspaceLookupError> {
if let Some(inherited) = self.inherited(invocation).await? {
return Ok(vec![inherited]);
}
let logs = self.logs().await?;
let resources = ResourceStore::load(&logs)
.await
.map_err(WorkspaceLookupError::Resources)?;
let mut workspaces: Vec<String> = resources
.records()
.filter(|record| {
record.allocation.invocation == invocation
&& record.state != petri_execution::LeaseState::Deleted
})
.map(|record| record.workspace.as_str().to_owned())
.collect();
workspaces.sort();
workspaces.dedup();
Ok(workspaces)
}
}

View file

@ -0,0 +1,468 @@
//! Fabro's hooks on a Petri run, in process: command-only bundles run
//! through `engine::run` on the host sandbox with the memory store and
//! in-memory platform records, and the checkpoint commit, its record, the
//! failure route, the fatal checkpoint, and the run-end hooks are checked
//! against the workspace's Git history and the run's records.
//!
//! Every run acquires its scope through the sandbox-driver host plugin, so
//! the tests skip when that executable is not found, unless
//! `FABRO_REQUIRE_SANDBOX_PLUGINS` is set. The crash cases of the recovery
//! protocol need a worker to kill and live in the CLI's scenario suite.
#![expect(
clippy::disallowed_methods,
reason = "the tests locate the plugin executable through the process environment and read the workspace's history with git"
)]
#![expect(clippy::print_stderr, reason = "a skipped test says why on its stderr")]
use std::collections::BTreeMap;
use std::env;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use fabro_checkpoint::author::GitAuthor;
use fabro_petri::admission::AdmittedGraphs;
use fabro_petri::check::{self, Bundle, CheckRequest, Launch};
use fabro_petri::checkpoint::{CHECKPOINT_FAILED_CLASS, CheckpointKey, RunWorkspaces};
use fabro_petri::engine::{self, Execution, RunRequest, RunStatus};
use fabro_petri::hooks::HooksSpec;
use fabro_petri::platform_records::PlatformRecords;
use fabro_petri::recovery::{self, Recovery, RecoveryRequest};
use fabro_petri::runtime::RuntimeSpec;
use fabro_petri::test_support::MemoryPlatformRecords;
use fabro_store::{PlatformRecord, PlatformRecordKind};
use fabro_types::settings::run::RunCheckpointSettings;
use fabro_types::{RunId, SandboxProviderKind};
use petri_execution::inspect::{self, RunInspection};
use petri_store::{Access, MemoryRunStore, RunKey, RunStore as _};
use tokio::fs;
use tokio::process::Command;
use tokio_util::sync::CancellationToken;
const HOST_PLUGIN: &str = "sandbox-driver-host";
const HOST_PLUGIN_OVERRIDE: &str = "PETRI_SANDBOX_HOST_PLUGIN";
const REQUIRE_ENV: &str = "FABRO_REQUIRE_SANDBOX_PLUGINS";
/// The host plugin as Petri's lookup finds it: the override variable, else
/// the executable on `PATH`. `None`, after saying so, when the test should
/// skip; a panic when the environment forbids a skip.
fn host_plugin() -> Option<PathBuf> {
let found = env::var_os(HOST_PLUGIN_OVERRIDE)
.map(PathBuf::from)
.or_else(|| {
env::split_paths(&env::var_os("PATH")?)
.map(|dir| dir.join(HOST_PLUGIN))
.find(|candidate| candidate.is_file())
});
if found.is_none() {
assert!(
env::var_os(REQUIRE_ENV).is_none(),
"{REQUIRE_ENV} is set, but {HOST_PLUGIN} is not on PATH and {HOST_PLUGIN_OVERRIDE} is unset"
);
eprintln!("skipping: {HOST_PLUGIN} is not on PATH and {HOST_PLUGIN_OVERRIDE} is unset");
}
found
}
/// A command-only bundle: the stage lines go between `start` and `exit`,
/// the edge lines after them.
fn workflow(stages: &str, edges: &str) -> String {
format!(
"digraph Hooks {{\n graph [goal=\"Check the hooks\", default_max_retries=0]\n start \
[shape=Mdiamond]\n exit [shape=Msquare]\n{stages}\n{edges}\n}}\n"
)
}
const SETTINGS: &str = "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\n";
/// The bundle admitted the way the create handler admits it.
fn admit(workflow: &str, settings: &str) -> AdmittedGraphs {
let request = CheckRequest {
bundle: Bundle {
files: BTreeMap::from([
("workflow.fabro".to_string(), workflow.to_string()),
("workflow.toml".to_string(), settings.to_string()),
]),
entrypoint: "workflow.fabro".to_string(),
project_toml: None,
},
inputs: BTreeMap::new(),
launch: Launch::default(),
runtime: RuntimeSpec::default(),
};
let admitted = check::check(&request).expect("the bundle is admitted");
AdmittedGraphs {
graph: admitted.graph,
children: admitted.children,
}
}
/// One run's pieces: the store, its platform records, where it ran.
struct Harness {
run_id: RunId,
run_dir: PathBuf,
store: Arc<MemoryRunStore>,
records: Arc<MemoryPlatformRecords>,
_root: tempfile::TempDir,
}
impl Harness {
fn new() -> Self {
let root = tempfile::tempdir().expect("a temp dir");
Self {
run_id: RunId::new(),
run_dir: root.path().join("run"),
store: Arc::new(MemoryRunStore::new()),
records: Arc::new(MemoryPlatformRecords::new()),
_root: root,
}
}
fn hooks(&self) -> HooksSpec {
HooksSpec {
records: Arc::clone(&self.records) as Arc<dyn PlatformRecords>,
author: GitAuthor::default(),
checkpoint: RunCheckpointSettings::default(),
host_workspaces: true,
test_gates: None,
}
}
/// Run the bundle to its end through the engine module, as the worker
/// does, and report what the record says.
async fn run(&self, workflow: &str, settings: &str) -> engine::RunOutcome {
let request = RunRequest {
run_id: self.run_id.to_string(),
run_dir: self.run_dir.clone(),
execution: Execution::Start(admit(workflow, settings)),
store: Arc::clone(&self.store) as Arc<dyn petri_store::RunStore>,
runtime: RuntimeSpec::default(),
provider: SandboxProviderKind::LOCAL,
cancel: CancellationToken::new(),
hooks: Some(self.hooks()),
};
engine::run(request).await.expect("the run executes")
}
async fn inspection(&self) -> RunInspection {
let logs = self
.store
.open(&RunKey::new(self.run_id.to_string()), Access::Read)
.await
.expect("the run opens for reading");
inspect::inspect_run(&*logs)
.await
.expect("the stored run inspects")
}
fn workspaces(&self) -> RunWorkspaces {
RunWorkspaces::new(
self.run_dir.clone(),
self.run_id.to_string(),
GitAuthor::default(),
&RunCheckpointSettings::default(),
)
}
/// The one workspace the run's root scope used.
async fn workspace(&self) -> String {
let mut entries = fs::read_dir(self.run_dir.join("scopes"))
.await
.expect("the scopes directory exists");
let mut names = Vec::new();
while let Some(entry) = entries.next_entry().await.expect("an entry reads") {
names.push(entry.file_name().to_string_lossy().into_owned());
}
assert_eq!(names.len(), 1, "one workspace: {names:?}");
names.remove(0)
}
fn workspace_path(&self, workspace: &str) -> PathBuf {
self.workspaces().workspace_path(workspace)
}
/// The checkpoint records, in seq order, as `(key, sha)`.
fn checkpoints(&self) -> Vec<(CheckpointKey, String)> {
self.records
.records(&self.run_id)
.into_iter()
.filter_map(|stored| match stored.record {
PlatformRecord::Checkpoint(record) => Some((
CheckpointKey::from_operation(record.operation.as_ref()?)?,
record.git_commit_sha?,
)),
_ => None,
})
.collect()
}
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,
})
.await
.expect("recovery decides")
}
}
/// `git` in a workspace, its stdout.
async fn git(path: &Path, args: &[&str]) -> String {
let output = Command::new("git")
.args(args)
.current_dir(path)
.output()
.await
.expect("git runs");
assert!(
output.status.success(),
"git {args:?} failed: {}",
String::from_utf8_lossy(&output.stderr)
);
String::from_utf8_lossy(&output.stdout).trim().to_string()
}
/// The commits on the run branch, oldest first, as `(sha, subject, key)`.
async fn commits(path: &Path) -> Vec<(String, String, Option<CheckpointKey>)> {
let log = git(path, &["log", "--reverse", "--format=%H%x00%s%x00%B%x1e"]).await;
log.split('\u{1e}')
.filter(|entry| !entry.trim().is_empty())
.map(|entry| {
let mut parts = entry.trim_start().splitn(3, '\0');
let sha = parts.next().unwrap_or_default().to_string();
let subject = parts.next().unwrap_or_default().to_string();
let body = parts.next().unwrap_or_default();
(sha, subject, CheckpointKey::from_message(body))
})
.collect()
}
/// The stages that ran, by node name, with their final status.
fn stages(inspection: &RunInspection) -> Vec<(String, String)> {
inspection
.executions
.iter()
.filter_map(|execution| execution.engine.as_ref())
.flat_map(|engine| engine.history.iter())
.map(|record| (record.node.to_string(), record.status.to_string()))
.collect()
}
/// Every finished stage is committed on the run branch with the identity
/// trailers, its platform record names the commit, and the run-end hooks
/// reached Petri's local service through Fabro's wrapper.
#[tokio::test]
async fn every_finish_is_committed_and_recorded() {
if host_plugin().is_none() {
return;
}
let harness = Harness::new();
let workflow = workflow(
" write [shape=parallelogram, script=\"echo one > out.txt\"]\n check \
[shape=parallelogram, script=\"test \\\"$(cat out.txt)\\\" = one\"]",
" start -> write -> check -> exit",
);
let settings = format!(
"{SETTINGS}\n[[run.hooks]]\nevent = \"run_complete\"\nscript = \"echo run_complete >> \
run-end.log\"\n\n[[run.hooks]]\nevent = \"sandbox_cleanup\"\nscript = \"echo \
sandbox_cleanup >> run-end.log\"\n"
);
let outcome = harness.run(&workflow, &settings).await;
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
assert!(outcome.complete, "{:?}", outcome.incomplete);
let workspace = harness.workspace().await;
let path = harness.workspace_path(&workspace);
assert_eq!(
git(&path, &["rev-parse", "--abbrev-ref", "HEAD"]).await,
format!("fabro/run/{}", harness.run_id)
);
let commits = commits(&path).await;
let subjects: Vec<&str> = commits
.iter()
.map(|(_, subject, _)| subject.as_str())
.collect();
let run_id = harness.run_id.to_string();
assert_eq!(subjects, vec![
format!("fabro({run_id}): start (success)"),
format!("fabro({run_id}): write (success)"),
format!("fabro({run_id}): check (success)"),
format!("fabro({run_id}): exit (success)"),
]);
assert!(
commits.iter().all(|(_, _, key)| key.is_some()),
"every commit carries its key: {commits:?}"
);
let checkpoints = harness.checkpoints();
assert_eq!(checkpoints.len(), 4, "{checkpoints:?}");
let by_sha: Vec<&String> = checkpoints.iter().map(|(_, sha)| sha).collect();
let committed: Vec<&String> = commits.iter().map(|(sha, _, _)| sha).collect();
assert_eq!(by_sha, committed, "each record names its stage's commit");
for ((key, _), (_, _, trailer)) in checkpoints.iter().zip(&commits) {
assert_eq!(Some(*key), *trailer);
}
assert!(
checkpoints.iter().all(|(key, _)| key.execution == 0),
"{checkpoints:?}"
);
// The snapshot repository holds every checkpoint.
let published = harness
.workspaces()
.published(&workspace)
.await
.expect("the snapshots list");
assert_eq!(published.len(), 4);
// `run_complete` and `sandbox_cleanup` ran through the forwarded
// service, with the sandbox in place.
let run_end = fs::read_to_string(path.join("run-end.log"))
.await
.expect("the run-end hooks wrote their log");
assert_eq!(run_end, "run_complete\nsandbox_cleanup\n");
}
/// A stage that fails on its own terms is committed like a successful one,
/// and its failure route runs on the committed files.
#[tokio::test]
async fn a_failed_stage_is_committed_and_its_route_sees_the_files() {
if host_plugin().is_none() {
return;
}
let harness = Harness::new();
let workflow = workflow(
" work [shape=parallelogram, script=\"echo partial > out.txt; exit 1\"]\n fix \
[shape=parallelogram, script=\"test \\\"$(cat out.txt)\\\" = partial && echo fixed >> \
out.txt\"]",
" start -> work -> exit\n work -> fix [condition=\"outcome=failed\"]\n fix -> exit",
);
let outcome = harness.run(&workflow, SETTINGS).await;
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
let inspection = harness.inspection().await;
assert_eq!(stages(&inspection), vec![
("start".to_string(), "success".to_string()),
("work".to_string(), "failure".to_string()),
("fix".to_string(), "success".to_string()),
("exit".to_string(), "success".to_string()),
]);
let workspace = harness.workspace().await;
let path = harness.workspace_path(&workspace);
let commits = commits(&path).await;
let run_id = harness.run_id.to_string();
assert_eq!(commits[1].1, format!("fabro({run_id}): work (failure)"));
assert_eq!(
git(&path, &["show", &format!("{}:out.txt", commits[1].0)]).await,
"partial",
"the failed stage's files are in its snapshot"
);
assert_eq!(
git(&path, &["show", &format!("{}:out.txt", commits[2].0)]).await,
"partial\nfixed",
"the route ran on the committed files"
);
assert_eq!(harness.checkpoints().len(), 4);
}
/// A checkpoint commit that fails is fatal: the stage's outcome is recorded
/// as `checkpoint_failed`, no route is taken, the run ends failed with the
/// checkpoint's error, and a restart reports it failed without resuming.
#[tokio::test]
async fn a_failed_checkpoint_ends_the_run_with_no_route() {
if host_plugin().is_none() {
return;
}
let harness = Harness::new();
let workflow = workflow(
" wreck [shape=parallelogram, script=\"rm -rf .git && echo garbage > .git && echo wrecked \
> out.txt\"]\n next [shape=parallelogram, script=\"echo next > next.txt\"]\n fix \
[shape=parallelogram, script=\"echo fix > fix.txt\"]",
" start -> wreck -> next -> exit\n wreck -> fix [condition=\"outcome=failed\"]\n fix -> \
exit",
);
let outcome = harness.run(&workflow, SETTINGS).await;
assert_eq!(outcome.status, RunStatus::Failed, "{outcome:?}");
let failure = outcome
.failure
.clone()
.expect("the run failed with a reason");
assert!(
failure.contains("checkpoint commit of `wreck` failed"),
"{failure}"
);
let inspection = harness.inspection().await;
let attempts: Vec<_> = inspection
.executions
.iter()
.filter_map(|execution| execution.engine.as_ref())
.flat_map(|engine| engine.attempts.iter())
.collect();
let wreck = attempts
.iter()
.find(|attempt| attempt.node.as_deref() == Some("wreck"))
.expect("the wrecked stage finished");
assert_eq!(wreck.status, "failure");
assert_eq!(
wreck.failure.as_ref().map(|failure| failure.class.as_str()),
Some(CHECKPOINT_FAILED_CLASS),
"{wreck:?}"
);
assert!(
attempts
.iter()
.all(|attempt| !matches!(attempt.node.as_deref(), Some("next" | "fix"))),
"no route ran: {:?}",
stages(&inspection)
);
let workspace = harness.workspace().await;
let path = harness.workspace_path(&workspace);
assert!(!fs::try_exists(path.join("next.txt")).await.expect("exists"));
assert!(!fs::try_exists(path.join("fix.txt")).await.expect("exists"));
// The wrecked stage has no record: its commit never landed.
let checkpoints = harness.checkpoints();
assert_eq!(checkpoints.len(), 1, "{checkpoints:?}");
// A restart finds the failed checkpoint and reports the run failed.
let recovery = harness.recover().await;
assert!(
matches!(&recovery, Recovery::Failed { reason } if reason.contains("checkpoint of wreck")),
"{recovery:?}"
);
}
/// The two ends of recovery that need no crash: a run the store never held
/// starts over, and a run that finished has nothing to bring back.
#[tokio::test]
async fn recovery_starts_an_unknown_run_and_resumes_a_finished_one() {
if host_plugin().is_none() {
return;
}
let harness = Harness::new();
assert_eq!(harness.recover().await, Recovery::Start);
let workflow = workflow(
" write [shape=parallelogram, script=\"echo one > out.txt\"]",
" start -> write -> exit",
);
let outcome = harness.run(&workflow, SETTINGS).await;
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
assert_eq!(harness.recover().await, Recovery::Resume {
workspaces: Vec::new(),
});
assert_eq!(
harness
.records
.read_kind(&harness.run_id, PlatformRecordKind::Checkpoint)
.await
.expect("the records read")
.len(),
3
);
}

View file

@ -376,6 +376,13 @@ pub struct GitIdentityRecord {
pub struct CheckpointRecord {
pub execution: u64,
pub firing: u64,
/// The attempt whose files the commit holds; absent on a record written
/// before the field existed.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub attempt: Option<u32>,
/// The Petri workspace id the commit was made in.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub workspace: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub git_commit_sha: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
@ -766,7 +773,7 @@ fn pull_request_created_record(props: &PullRequestCreatedProps) -> PullRequestCr
fn run_paired_record(props: &RunPairStartedProps) -> RunPairedRecord {
RunPairedRecord {
pair_id: props.pair_id.clone(),
pair_id: props.pair_id,
target: props.target.clone(),
}
}
@ -784,6 +791,7 @@ fn interview_answered_record(
#[cfg(test)]
mod tests {
use fabro_types::test_support::test_run_spec;
use fabro_types::{FailureReason, RunStatus, fixtures};
use serde_json::json;
@ -803,7 +811,7 @@ mod tests {
fn sample(kind: PlatformRecordKind) -> PlatformRecord {
match kind {
PlatformRecordKind::RunCreated => PlatformRecord::RunCreated(RunCreatedRecord {
spec: fabro_types::test_support::test_run_spec(),
spec: test_run_spec(),
title: Some("A run".to_string()),
parent_id: None,
retried_from: None,
@ -855,6 +863,8 @@ mod tests {
PlatformRecordKind::Checkpoint => PlatformRecord::Checkpoint(CheckpointRecord {
execution: 0,
firing: 3,
attempt: Some(1),
workspace: Some("invocation-0-scope-0".to_string()),
git_commit_sha: Some("def".to_string()),
diff_summary: Some(DiffSummary {
files_changed: 1,
@ -957,7 +967,7 @@ mod tests {
assert_eq!(json(&stored), json(&[first, second.clone()]));
assert_eq!(
json(&store.read_after(&run, 1).await.expect("the tail reads")),
json(&[second.clone()])
json(std::slice::from_ref(&second))
);
assert_eq!(
json(

View file

@ -253,7 +253,8 @@ impl RunSummaryStore {
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(hook);
}
pub(crate) fn notify_platform_record(&self, run_id: RunId) {
/// Wake the run's projector: a platform record was committed for the run.
pub fn notify_platform_record(&self, run_id: RunId) {
let hook = self
.platform_hook
.read()

View file

@ -2039,6 +2039,45 @@ impl Client {
Ok(response.into_inner().hash)
}
/// Store one of Fabro's platform records for the run, tied to a Petri
/// stage when it belongs to one, and get it back as stored.
pub async fn append_petri_platform_record(
&self,
run_id: &RunId,
body: types::PetriPlatformRecordAppendRequest,
) -> Result<types::PetriPlatformRecord> {
let response = self
.send_api(|client| async move {
client
.append_petri_platform_record()
.id(run_id.to_string())
.body(body.clone())
.send()
.await
})
.await?;
Ok(response.into_inner())
}
/// The run's platform records in `seq` order, of one kind when `kind`
/// names it.
pub async fn list_petri_platform_records(
&self,
run_id: &RunId,
kind: Option<&str>,
) -> Result<Vec<types::PetriPlatformRecord>> {
let response = self
.send_api(|client| async move {
let mut request = client.list_petri_platform_records().id(run_id.to_string());
if let Some(kind) = kind {
request = request.kind(kind);
}
request.send().await
})
.await?;
Ok(response.into_inner().records)
}
/// The blob with this digest, or `None` when the store holds no such
/// blob. A run the store does not hold is an error.
pub async fn read_petri_blob(

View file

@ -38,6 +38,10 @@ impl EnvVars {
pub const FABRO_TEST_IN_MEMORY_STORE: &'static str = "FABRO_TEST_IN_MEMORY_STORE";
pub const FABRO_TEST_DISABLE_SPA_ASSETS: &'static str = "FABRO_TEST_DISABLE_SPA_ASSETS";
pub const FABRO_TEST_MODE: &'static str = "FABRO_TEST_MODE";
/// A directory of hold and release files a test uses to pause a Petri
/// run's checkpoint at a named point (`fabro_petri::hooks`); unset
/// outside tests.
pub const FABRO_TEST_CHECKPOINT_GATES: &'static str = "FABRO_TEST_CHECKPOINT_GATES";
pub const FABRO_VERBOSE: &'static str = "FABRO_VERBOSE";
pub const FABRO_WEB_URL: &'static str = "FABRO_WEB_URL";
pub const FABRO_WORKER_TOKEN: &'static str = "FABRO_WORKER_TOKEN";
@ -222,6 +226,7 @@ mod tests {
EnvVars::FABRO_TEST_IN_MEMORY_STORE,
EnvVars::FABRO_TEST_DISABLE_SPA_ASSETS,
EnvVars::FABRO_TEST_MODE,
EnvVars::FABRO_TEST_CHECKPOINT_GATES,
EnvVars::FABRO_VERBOSE,
EnvVars::FABRO_WEB_URL,
EnvVars::FABRO_WORKER_TOKEN,