diff --git a/Cargo.lock b/Cargo.lock index 670f22e8b..82600c564 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -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", diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 25ba828cf..305efece0 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -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 diff --git a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs index 2ec71756c..cf9aa20f0 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -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 { + 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. diff --git a/lib/apps/fabro-server/src/server/handler/petri.rs b/lib/apps/fabro-server/src/server/handler/petri.rs index cc6824147..0c88e160f 100644 --- a/lib/apps/fabro-server/src/server/handler/petri.rs +++ b/lib/apps/fabro-server/src/server/handler/petri.rs @@ -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> { 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, +} + async fn open_run( RequireWorkerRunScoped(id): RequireWorkerRunScoped, State(state): State>, @@ -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>, + Query(query): Query, +) -> Response { + let store = state.stores.run_summaries.platform_records(); + let records = match query.kind.as_deref() { + Some(kind) => match kind.parse::() { + 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::, _>>() + { + 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>, + Json(request): Json, +) -> 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 { + 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 `." diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 931a8a459..da9ad3091 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -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, 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, 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, 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, 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, diff --git a/lib/apps/fabro-server/src/spawn_env.rs b/lib/apps/fabro-server/src/spawn_env.rs index 346053c93..351939d7c 100644 --- a/lib/apps/fabro-server/src/spawn_env.rs +++ b/lib/apps/fabro-server/src/spawn_env.rs @@ -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]; diff --git a/lib/components/fabro-petri/Cargo.toml b/lib/components/fabro-petri/Cargo.toml index 1933209e3..82967dca2 100644 --- a/lib/components/fabro-petri/Cargo.toml +++ b/lib/components/fabro-petri/Cargo.toml @@ -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"] } diff --git a/lib/components/fabro-petri/src/checkpoint.rs b/lib/components/fabro-petri/src/checkpoint.rs new file mode 100644 index 000000000..b3cc30dcf --- /dev/null +++ b/lib/components/fabro-petri/src/checkpoint.rs @@ -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//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/.git`, under an immutable ref +//! per checkpoint (`refs/checkpoints///`). 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 { + 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 { + 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 { + 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, + 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//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 { + 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, 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, 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 { + 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, 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 { + 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, 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, CheckpointError> { + self.git_status(path, "rev-parse", &["rev-parse", "-q", "--verify", "HEAD"]) + .await + } + + async fn is_clean(&self, path: &Path) -> Result { + 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>( + &self, + cwd: &Path, + action: &str, + args: &[S], + ) -> Result { + 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>( + &self, + cwd: &Path, + action: &str, + args: &[S], + ) -> Result, 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>( + &self, + cwd: &Path, + action: &str, + args: &[S], + ) -> Result { + 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"); + } +} diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index 0287ec699..1b3f28db5 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -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, } /// The recorded status of a finished run. @@ -140,17 +148,37 @@ pub async fn run(request: RunRequest) -> Result { 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); + } 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 { 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. diff --git a/lib/components/fabro-petri/src/hooks.rs b/lib/components/fabro-petri/src/hooks.rs new file mode 100644 index 000000000..2d3acf3a6 --- /dev/null +++ b/lib/components/fabro-petri/src/hooks.rs @@ -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, + 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, +} + +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, 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) -> Self { + self.test_gates = gates; + self + } +} + +/// Fabro's `ExecutionHooks`, around the hooks the runtime installed. +pub struct FabroHooks { + inner: Arc, + run_id: RunId, + records: Arc, + workspaces: RunWorkspaces, + lookup: WorkspaceLookup, + host_workspaces: bool, + test_gates: Option, + handle: OnceLock, + /// The workspace and commit of every checkpoint this process made. + committed: Mutex>, + /// Which checkpoints have their platform record, loaded from the store + /// once and kept up to date with every append. + recorded: Mutex>, + recorded_loaded: OnceCell<()>, + /// Inherited workspaces resolved through the run's records. + inherited: Mutex>>, + /// The checkpoint failure that ended the run, when one did. + failure: Mutex>, + 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, + run_id: RunId, + run_key: RunKey, + run_dir: PathBuf, + store: Arc, + ) -> 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 { + 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 { + 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, 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(mutex: &Mutex) -> 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 { + 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 { + self.inner.after_record(context, recorded).await + } + + async fn transition( + &self, + context: &HookContext, + transition: Transition, + ) -> Result { + 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 { + 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 { + 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 + } +} diff --git a/lib/components/fabro-petri/src/lib.rs b/lib/components/fabro-petri/src/lib.rs index 443032ac8..0ab08b96d 100644 --- a/lib/components/fabro-petri/src/lib.rs +++ b/lib/components/fabro-petri/src/lib.rs @@ -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; diff --git a/lib/components/fabro-petri/src/platform_records.rs b/lib/components/fabro-petri/src/platform_records.rs new file mode 100644 index 000000000..5e71fadaf --- /dev/null +++ b/lib/components/fabro-petri/src/platform_records.rs @@ -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, + ) -> Result; + + /// The run's records of one kind, in seq order. + async fn read_kind( + &self, + run_id: &RunId, + kind: PlatformRecordKind, + ) -> Result, 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, +} + +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) -> Self { + Self { + store: summaries.platform_records(), + summaries, + } + } +} + +#[async_trait] +impl PlatformRecords for SqlitePlatformRecords { + async fn append( + &self, + run_id: &RunId, + record: &PlatformRecord, + position: Option, + ) -> Result { + 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, 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, + ) -> Result { + 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, 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 { + 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, + }) +} diff --git a/lib/components/fabro-petri/src/recovery.rs b/lib/components/fabro-petri/src/recovery.rs new file mode 100644 index 000000000..b80f4b26c --- /dev/null +++ b/lib/components/fabro-petri/src/recovery.rs @@ -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, + pub records: Arc, + 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, + records: Arc, + 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 }, + /// 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 { + 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> = 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 { + 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 { + 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, 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 { + 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 { + 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) +} diff --git a/lib/components/fabro-petri/src/test_support.rs b/lib/components/fabro-petri/src/test_support.rs index 8f677c1d0..873fbd035 100644 --- a/lib/components/fabro-petri/src/test_support.rs +++ b/lib/components/fabro-petri/src/test_support.rs @@ -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>>, +} + +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 { + lock(&self.runs).get(run_id).cloned().unwrap_or_default() + } +} + +fn lock(mutex: &Mutex) -> 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, + ) -> Result { + 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, PlatformRecordError> { + Ok(self + .records(run_id) + .into_iter() + .filter(|record| record.record.kind() == kind) + .collect()) + } +} diff --git a/lib/components/fabro-petri/src/workspace.rs b/lib/components/fabro-petri/src/workspace.rs new file mode 100644 index 000000000..9081a022e --- /dev/null +++ b/lib/components/fabro-petri/src/workspace.rs @@ -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--scope-`), 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, + key: RunKey, +} + +impl WorkspaceLookup { + #[must_use] + pub fn new(store: Arc, key: RunKey) -> Self { + Self { store, key } + } + + async fn logs(&self) -> Result, 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, 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, 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 = 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) + } +} diff --git a/lib/components/fabro-petri/tests/hooks.rs b/lib/components/fabro-petri/tests/hooks.rs new file mode 100644 index 000000000..6493ca9c6 --- /dev/null +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -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 { + 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, + records: Arc, + _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, + 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, + 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, + records: Arc::clone(&self.records) as Arc, + 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)> { + 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 + ); +} diff --git a/lib/components/fabro-store/src/platform_records.rs b/lib/components/fabro-store/src/platform_records.rs index d2ae7625f..88f42ab5d 100644 --- a/lib/components/fabro-store/src/platform_records.rs +++ b/lib/components/fabro-store/src/platform_records.rs @@ -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, + /// The Petri workspace id the commit was made in. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub workspace: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub git_commit_sha: Option, #[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( diff --git a/lib/components/fabro-store/src/run_summary_store.rs b/lib/components/fabro-store/src/run_summary_store.rs index 8b51a41fd..cbe428d5d 100644 --- a/lib/components/fabro-store/src/run_summary_store.rs +++ b/lib/components/fabro-store/src/run_summary_store.rs @@ -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() diff --git a/lib/foundation/fabro-client/src/client.rs b/lib/foundation/fabro-client/src/client.rs index 03f16e479..c4bd9a22f 100644 --- a/lib/foundation/fabro-client/src/client.rs +++ b/lib/foundation/fabro-client/src/client.rs @@ -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 { + 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> { + 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( diff --git a/lib/foundation/fabro-static/src/env_vars.rs b/lib/foundation/fabro-static/src/env_vars.rs index 11178bef6..a14c8ce7e 100644 --- a/lib/foundation/fabro-static/src/env_vars.rs +++ b/lib/foundation/fabro-static/src/env_vars.rs @@ -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,