diff --git a/Cargo.lock b/Cargo.lock index 670f22e8b..fb759b79a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2896,9 +2896,13 @@ dependencies = [ "fabro-client", "fabro-db", "fabro-http", + "fabro-interview", "fabro-llm", "fabro-store", + "fabro-test", "fabro-types", + "fabro-vault", + "fabro-workflow", "lithos-llm", "petri-attractor-steps", "petri-execution", diff --git a/lib/apps/fabro-cli/src/args.rs b/lib/apps/fabro-cli/src/args.rs index 09d7a01dc..42c0b7af6 100644 --- a/lib/apps/fabro-cli/src/args.rs +++ b/lib/apps/fabro-cli/src/args.rs @@ -1089,6 +1089,11 @@ pub(crate) struct RunWorkerArgs { /// Worker mode #[arg(long, value_enum)] pub(crate) mode: RunWorkerMode, + + /// The Fabro home the server runs under, for the skills a Petri run's + /// agents read + #[arg(long, hide = true)] + pub(crate) fabro_home: Option, } #[derive(Args, Debug, Clone, Default)] diff --git a/lib/apps/fabro-cli/src/commands/run/mod.rs b/lib/apps/fabro-cli/src/commands/run/mod.rs index 47ae4ad9f..f59559070 100644 --- a/lib/apps/fabro-cli/src/commands/run/mod.rs +++ b/lib/apps/fabro-cli/src/commands/run/mod.rs @@ -101,6 +101,7 @@ pub(crate) async fn dispatch( run_dir, run_id, mode, + fabro_home, }) => { let worker_token = worker_token .filter(|token| !token.trim().is_empty()) @@ -109,8 +110,16 @@ pub(crate) async fn dispatch( })?; let run_span = tracing::info_span!("run", id = %run_id); Box::pin( - runner::execute(run_id, server, storage_dir, run_dir, mode, &worker_token) - .instrument(run_span), + runner::execute( + run_id, + server, + storage_dir, + run_dir, + mode, + fabro_home, + &worker_token, + ) + .instrument(run_span), ) .await } 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..b805a041d 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -17,18 +17,24 @@ //! (`run.starting`, `run.running`, then `run.completed` or `run.failed`) //! through the client, as the legacy worker does. //! -//! Of the server's controls, cancel is wired: the control channel's cancel -//! and `SIGTERM`/`SIGINT` fire one token, which cancels Petri's root -//! invocation politely. Pause, unpause and steer are received and ignored -//! with a warning until their Petri adapters land. A control channel that -//! 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. +//! Of the server's controls, cancel and answers are wired: the control +//! channel's cancel and `SIGTERM`/`SIGINT` fire one token, which cancels +//! Petri's root invocation politely, and an `interview.answer` message +//! reaches the control interviewer the run's questions wait on +//! (`fabro_petri::interview`), so a human gate answered through the API +//! continues. Pause, unpause and steer are received and ignored with a +//! warning until their Petri adapters land. A control channel that 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. //! //! 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 //! catalog and vault snapshot for the providers whose credentials resolve, -//! the same eligible set the legacy worker's LLM backend uses. +//! the same eligible set the legacy worker's LLM backend uses. The same +//! vault snapshot is the run's secret provider, the run's blobs go to the +//! server's blob table through the worker's client, and the Fabro home the +//! server named on the command line is the home the skills step reads. use std::path::{Path, PathBuf}; use std::sync::Arc; @@ -39,17 +45,22 @@ use fabro_auth::VaultCredentialSource; use fabro_client::{Client, ServerTarget}; use fabro_interview::ControlInterviewer; use fabro_llm::credentials::{CredentialProvider, readiness}; +use fabro_petri::blobs::ClientBlobs; use fabro_petri::engine::{self, Conclusion, Execution, RunRequest}; +use fabro_petri::interview::{Approval, EventSinkQuestions, FabroInterviewer}; use fabro_petri::petri::OwnerId; use fabro_petri::runtime::{self, RuntimeSpec}; +use fabro_petri::secrets::VaultSecrets; use fabro_petri::{HttpRunStore, admission}; use fabro_store::RunProjection; -use fabro_types::settings::run::RunMode; +use fabro_types::settings::run::{ApprovalMode, RunMode}; use fabro_types::{FailureReason, RunId, RunTiming, StageOutcome, SuccessReason}; +use fabro_vault::Vault; use fabro_workflow::Error as WorkflowError; use fabro_workflow::event::{self as workflow_event, Emitter, Event, RunEventSink}; use fabro_workflow::run_control::RunControlState; use fabro_workflow::runtime_store::RunStoreHandle; +use tokio::sync::RwLock as AsyncRwLock; use tokio_util::sync::CancellationToken; use tracing::{info, warn}; @@ -69,6 +80,9 @@ pub(super) struct PetriWorker<'a> { pub(super) storage_dir: &'a Path, pub(super) run_dir: PathBuf, pub(super) mode: RunWorkerMode, + /// The Fabro home the server named; `None` falls back to Petri's own + /// lookup of the worker's environment. + pub(super) fabro_home: Option, pub(super) worker_token: &'a str, } @@ -102,7 +116,7 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { worker.target.clone(), run_id, worker.worker_token.to_owned(), - interviewer, + Arc::clone(&interviewer), cancel_token.clone(), steering_hub, run_control, @@ -110,14 +124,25 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { control_manager.wait_for_first_connection().await?; warn!( run_id = %run_id, - "a Petri run answers cancel only: pause, unpause and steer are not wired yet and are ignored" + "a Petri run answers cancel and questions only: pause, unpause and steer are not wired \ + yet and are ignored" ); let sink = RunEventSink::map( runner::stamp_system_worker, RunEventSink::backend(worker.run_store.clone()), ); + let approval = if worker.run_state.spec.settings.run.execution.approval == ApprovalMode::Auto { + Approval::Auto + } else { + Approval::Prompt + }; + let questions = Arc::new(EventSinkQuestions::new(sink.clone(), run_id)); + let petri_interviewer = FabroInterviewer::new(interviewer, questions, approval); + let observers = vec![petri_interviewer.observer()]; - let runtime = runtime_spec(worker.storage_dir, &worker.run_state).await?; + let vault = runner::load_worker_vault(worker.storage_dir).await?; + let secrets = VaultSecrets::from_vault(&*vault.read().await); + let runtime = runtime_spec(&vault, &worker.run_state, worker.fabro_home.clone()).await?; let execution = match worker.mode { RunWorkerMode::Start => { let client = worker.client.clone_for_reuse(); @@ -156,6 +181,13 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { .provider .clone(), cancel: cancel_token.clone(), + interviewer: Arc::new(petri_interviewer), + observers, + secrets: Some(Arc::new(secrets)), + blobs: Some(Arc::new(ClientBlobs::new( + worker.client.clone_for_reuse(), + run_id, + ))), }; let run = Box::pin(engine::run(request)); tokio::pin!(run); @@ -229,12 +261,17 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { /// 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. -async fn runtime_spec(storage_dir: &Path, run_state: &RunProjection) -> Result { +/// the providers whose credentials resolve, the run's mode, and the Fabro +/// home the server named. +async fn runtime_spec( + vault: &Arc>, + run_state: &RunProjection, + fabro_home: Option, +) -> Result { let catalog = command_context::load_cli_catalog().context("failed to build worker LLM catalog")?; - let vault = runner::load_worker_vault(storage_dir).await?; - let credentials: Arc = Arc::new(VaultCredentialSource::new(vault)); + let credentials: Arc = + Arc::new(VaultCredentialSource::new(Arc::clone(vault))); let ready = readiness(catalog.enabled_providers(), credentials.as_ref()).await; for (provider, issue) in &ready.issues { warn!(provider = %provider, error = %issue, "model provider credentials unusable"); @@ -250,6 +287,6 @@ async fn runtime_spec(storage_dir: &Path, run_state: &RunProjection) -> Result, worker_token: &str, ) -> Result<()> { let _ = fabro_proc::title_init(); @@ -96,6 +97,7 @@ pub(crate) async fn execute( storage_dir: &storage_dir, run_dir, mode, + fabro_home, worker_token, })) .await; diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 21a57cf24..720a767d0 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -3788,6 +3788,7 @@ fn worker_launch_spec( fabro_log, active_config_path: state.active_config_path().to_path_buf(), github_app_private_key, + fabro_home: fabro_config::Home::from_env().root().to_path_buf(), }) } diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 931a8a459..05edebd26 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -16,8 +16,11 @@ //! (`crate::petri_runs`). Under the test override that replaces the handler //! registry, [`execute`] runs the same engine in the server process over the //! run store in the server's database, so the scenario tests need no -//! worker binary. No stage or agent event is projected either way, which is -//! the read-side item that follows. +//! worker binary; its questions go to an in-process control interviewer +//! the answer endpoint reaches directly, its secrets come from a snapshot +//! of the server's vault, and its blobs go to the server's blob store. No +//! stage or agent event is projected either way, which is 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. @@ -26,14 +29,17 @@ use std::collections::{BTreeMap, HashSet}; use std::sync::Arc; use std::time::Instant; -use fabro_config::{SettingsLayer, Storage}; +use fabro_config::{Home, SettingsLayer, Storage}; +use fabro_interview::ControlInterviewer; 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::interview::{Approval, DatabaseQuestions, FabroInterviewer}; use fabro_petri::petri::StoreError; use fabro_petri::runtime::{self, RuntimeSpec}; +use fabro_petri::secrets::VaultSecrets; use fabro_petri::{SqliteRunStore, admission}; -use fabro_types::settings::run::RunMode; +use fabro_types::settings::run::{ApprovalMode, RunMode}; use fabro_types::{ Engine, PetriAdmission, RunId, RunRunnableSource, RunTarget, RunTiming, ServerSettings, StageOutcome, @@ -41,13 +47,14 @@ use fabro_types::{ use fabro_util::error as error_util; use fabro_validate::{Diagnostic as FabroDiagnostic, Severity}; use fabro_workflow::Error as WorkflowError; +use fabro_workflow::event::Emitter; use fabro_workflow::run_status::{FailureReason, RunStatus, SuccessReason}; use lithos_llm::catalog::ProviderId; use tokio::task; use tokio_util::sync::CancellationToken; use tracing::{error, info, warn}; -use super::{AppState, RunExecutionMode, clear_live_run_state, workflow_event}; +use super::{AppState, RunAnswerTransport, RunExecutionMode, clear_live_run_state, workflow_event}; use crate::petri_runs::PetriRuns; use crate::run_compiler::{PreparedRun, RunCompilerError}; @@ -89,7 +96,7 @@ pub(crate) fn runtime_spec( settings_toml, model_client, dry_run, - fabro_home: None, + fabro_home: Some(Home::from_env().root().to_path_buf()), } } @@ -309,6 +316,22 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { } RunExecutionMode::Resume => Execution::Resume, }; + // The run's secrets: a snapshot of the server's vault, as a worker + // takes one at launch. + let vault = match state.stores.vault.snapshot().await { + Ok(snapshot) => snapshot.into_vault(), + Err(err) => { + let message = error_util::collect_chain(&err).join(": "); + fail_before_execution( + &state, + &run_store, + run_id, + &format!("the vault could not be read for the run: {message}"), + ) + .await; + return; + } + }; let started = Instant::now(); for event in [ workflow_event::Event::RunStarting, @@ -327,14 +350,32 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { return; } } + // The answer endpoint reaches this interviewer directly, as it does + // for a legacy run in this process. + let interviewer = Arc::new(ControlInterviewer::new()); + let steering_hub = Arc::new(fabro_workflow::SteeringHub::new(Arc::new(Emitter::new( + run_id, + )))); { let mut runs = state.runs.lock().expect("runs lock poisoned"); if let Some(managed_run) = runs.get_mut(&run_id) { if managed_run.status == RunStatus::Starting { managed_run.status = RunStatus::Running; + managed_run.answer_transport = Some(RunAnswerTransport::InProcess { + interviewer: Arc::clone(&interviewer), + steering_hub, + }); } } } + let approval = if run_state.spec.settings.run.execution.approval == ApprovalMode::Auto { + Approval::Auto + } else { + Approval::Prompt + }; + let questions = Arc::new(DatabaseQuestions::new(run_store.clone(), run_id)); + let petri_interviewer = FabroInterviewer::new(interviewer, questions, approval); + let observers = vec![petri_interviewer.observer()]; let (_, eligible) = state.resolve_llm_client_with_ready_ids().await; let dry_run = run_state.spec.settings.run.execution.mode == RunMode::DryRun; let request = RunRequest { @@ -345,6 +386,10 @@ 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, + interviewer: Arc::new(petri_interviewer), + observers, + secrets: Some(Arc::new(VaultSecrets::from_vault(&vault))), + blobs: Some(state.store_ref().blobs()), }; let result = Box::pin(engine::run(request)).await; let timing = RunTiming { diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index be9f73077..75333bfe1 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -2334,6 +2334,8 @@ fn worker_command_sets_worker_args() { run_id.to_string(), "--mode".to_string(), "resume".to_string(), + "--fabro-home".to_string(), + fabro_config::Home::from_env().root().display().to_string(), ]); } diff --git a/lib/apps/fabro-server/src/worker_runtime.rs b/lib/apps/fabro-server/src/worker_runtime.rs index f2779f055..7b4fa39b6 100644 --- a/lib/apps/fabro-server/src/worker_runtime.rs +++ b/lib/apps/fabro-server/src/worker_runtime.rs @@ -48,6 +48,9 @@ pub(crate) struct WorkerLaunchSpec { pub(crate) fabro_log: Option, pub(crate) active_config_path: PathBuf, pub(crate) github_app_private_key: Option, + /// The Fabro home the server resolved, so a Petri run's skills step + /// reads the same home whatever the worker's environment says. + pub(crate) fabro_home: PathBuf, } pub(crate) struct StartedWorker { @@ -89,6 +92,8 @@ impl LocalWorkerRuntime { .arg(spec.run_id.to_string()) .arg("--mode") .arg(spec.mode) + .arg("--fabro-home") + .arg(&spec.fabro_home) .stdin(Stdio::null()) .stdout(worker_stdout) .stderr(Stdio::piped()); diff --git a/lib/components/fabro-petri/Cargo.toml b/lib/components/fabro-petri/Cargo.toml index 1933209e3..b73620de8 100644 --- a/lib/components/fabro-petri/Cargo.toml +++ b/lib/components/fabro-petri/Cargo.toml @@ -23,8 +23,11 @@ fabro-api = { path = "../../foundation/fabro-api" } fabro-client = { path = "../../foundation/fabro-client" } fabro-db = { path = "../../foundation/fabro-db" } fabro-http.workspace = true +fabro-interview = { path = "../fabro-interview" } fabro-store = { path = "../fabro-store" } fabro-types = { path = "../../foundation/fabro-types" } +fabro-vault = { path = "../../foundation/fabro-vault" } +fabro-workflow = { path = "../fabro-workflow" } petri_runtime.workspace = true petri_execution.workspace = true petri_store.workspace = true @@ -49,5 +52,6 @@ tracing.workspace = true 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"] } +fabro-test.workspace = true petri_testkit.workspace = true tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } diff --git a/lib/components/fabro-petri/src/blobs.rs b/lib/components/fabro-petri/src/blobs.rs new file mode 100644 index 000000000..bee52c0ee --- /dev/null +++ b/lib/components/fabro-petri/src/blobs.rs @@ -0,0 +1,180 @@ +//! Petri's `OutputStore` over Fabro's blob table. +//! +//! A stage value above Petri's offload threshold leaves the run context for +//! the run's blob store and is replaced by the reference +//! `blob://sha256/`, Fabro's spelling; a later step hydrates it back +//! through the same store. Petri's default store is a directory under the +//! run directory. Installed instead is [`RunBlobs`], which carries every +//! blob to Fabro's `blobs` table: in the server process through its +//! [`fabro_store::BlobStore`], and in a run's worker process through the +//! worker's client ([`ClientBlobs`]), whose run blob endpoints the server +//! answers from the same table. Either way the digest is the same SHA-256 +//! hex Fabro's [`BlobHash`] renders, so a reference a Petri record carries +//! names a row Fabro's own readers can fetch. + +use std::sync::Arc; + +use bytes::Bytes; +use fabro_client::Client; +use fabro_types::{BlobHash, RunId}; +use petri_attractor_steps::blobs::{BlobError, BlobStore, OutputStore}; + +/// Fabro's content-addressed blob table, as a run reaches it. +#[async_trait::async_trait] +pub trait Blobs: Send + Sync { + /// Store `bytes` and return the hash that names them. + async fn write(&self, bytes: &[u8]) -> anyhow::Result; + + /// The bytes behind a hash, or `None` when the table has none. + async fn read(&self, hash: &BlobHash) -> anyhow::Result>; +} + +#[async_trait::async_trait] +impl Blobs for fabro_store::BlobStore { + async fn write(&self, bytes: &[u8]) -> anyhow::Result { + Self::write(self, bytes).await.map_err(anyhow::Error::new) + } + + async fn read(&self, hash: &BlobHash) -> anyhow::Result> { + Self::read(self, hash).await.map_err(anyhow::Error::new) + } +} + +/// The blob table as a run's worker reaches it: the run's blob endpoints, +/// with the worker's token. +pub struct ClientBlobs { + client: Client, + run_id: RunId, +} + +impl ClientBlobs { + #[must_use] + pub fn new(client: Client, run_id: RunId) -> Self { + Self { client, run_id } + } +} + +#[async_trait::async_trait] +impl Blobs for ClientBlobs { + async fn write(&self, bytes: &[u8]) -> anyhow::Result { + self.client.write_run_blob(&self.run_id, bytes).await + } + + async fn read(&self, hash: &BlobHash) -> anyhow::Result> { + self.client.read_run_blob(&self.run_id, hash).await + } +} + +/// Petri's blob store over Fabro's blob table. +pub struct RunBlobs { + blobs: Arc, +} + +impl RunBlobs { + #[must_use] + pub fn new(blobs: Arc) -> Self { + Self { blobs } + } + + /// The capability a runtime installs so every offloaded value goes to + /// the table. + #[must_use] + pub fn output_store(blobs: Arc) -> OutputStore { + OutputStore(Arc::new(Self::new(blobs))) + } +} + +/// The store's refusal, with the cause chain on one line: Petri's error +/// carries text, not a source. +fn refused(digest: &str, error: &anyhow::Error) -> BlobError { + BlobError::Store { + digest: digest.to_string(), + message: format!("{error:#}"), + } +} + +#[async_trait::async_trait] +impl BlobStore for RunBlobs { + async fn put(&self, bytes: &[u8]) -> Result { + let expected = BlobHash::new(bytes); + let hash = self + .blobs + .write(bytes) + .await + .map_err(|error| refused(&expected.to_string(), &error))?; + Ok(hash.to_string()) + } + + async fn get(&self, digest: &str) -> Result>, BlobError> { + let Ok(hash) = digest.parse::() else { + // Not a digest the table can hold, so nothing is behind it. + return Ok(None); + }; + let bytes = self + .blobs + .read(&hash) + .await + .map_err(|error| refused(digest, &error))?; + Ok(bytes.map(|bytes| bytes.to_vec())) + } +} + +#[cfg(test)] +mod tests { + use std::sync::Mutex; + + use petri_attractor_steps::blobs::{blob_ref, hydrate, offload_above}; + use petri_runtime::ir::Value; + + use super::*; + + /// A table in memory. + #[derive(Default)] + struct MemoryBlobs { + rows: Mutex)>>, + } + + #[async_trait::async_trait] + impl Blobs for MemoryBlobs { + async fn write(&self, bytes: &[u8]) -> anyhow::Result { + let hash = BlobHash::new(bytes); + self.rows + .lock() + .expect("not poisoned") + .push((hash, bytes.to_vec())); + Ok(hash) + } + + async fn read(&self, hash: &BlobHash) -> anyhow::Result> { + Ok(self + .rows + .lock() + .expect("not poisoned") + .iter() + .find(|(stored, _)| stored == hash) + .map(|(_, bytes)| Bytes::copy_from_slice(bytes))) + } + } + + #[tokio::test] + async fn a_value_round_trips_through_the_table_under_fabros_reference() { + let table = Arc::new(MemoryBlobs::default()); + let store = RunBlobs::new(table.clone()); + let mut value = Value::String("x".repeat(10)); + let reference = offload_above(&mut value, &store, 0) + .await + .expect("offloaded"); + let hex = BlobHash::new(b"xxxxxxxxxx").to_string(); + assert_eq!(reference, blob_ref(&hex)); + assert_eq!(value, Value::String(reference)); + assert_eq!(hydrate(value, &store).await, Value::String("x".repeat(10))); + assert_eq!(table.rows.lock().expect("not poisoned").len(), 1); + } + + #[tokio::test] + async fn an_unknown_digest_and_a_malformed_one_are_absent() { + let store = RunBlobs::new(Arc::new(MemoryBlobs::default())); + assert_eq!(store.get(&"a".repeat(64)).await.expect("reads"), None); + assert_eq!(store.get("not-a-digest").await.expect("reads"), None); + } +} diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index 0287ec699..8c630d4eb 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -15,12 +15,15 @@ //! `inspect_run` over a read handle of the same store, so what the caller //! 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 -//! rides the caller's token: when it fires, the root invocation is cancelled -//! politely and Petri records why. +//! What the caller supplies beyond the runtime: the interviewer its +//! questions go to ([`interview`](crate::interview) in the worker and the +//! server), the secret provider over the vault ([`secrets`](crate::secrets)) +//! and the blob table ([`blobs`](crate::blobs)) when it has them. What the +//! standalone runner's defaults give the run: Petri's local hook service +//! for `[[run.hooks]]`, no `ExecutionHooks` of Fabro's own, no host tools, +//! and `Retention::Always` for every workspace, Fabro's default. +//! 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 @@ -39,17 +42,19 @@ use fabro_types::{FailureReason, 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, + Access, CancelReason, ExecutionObserver, InterviewDispatcher, Interviewer, InvocationId, + RECEIPT_FILE, RunKey, RunStore, }; -use petri_runtime::executor::Retention; +use petri_runtime::executor::{Retention, SecretProvider}; use petri_runtime::{RunOptions, SandboxBackend}; use tokio::fs; use tokio_util::sync::CancellationToken; use tracing::{debug, info, warn}; use crate::admission::AdmittedGraphs; -use crate::interviewer::Unattended; +use crate::blobs::{Blobs, RunBlobs}; use crate::runtime::RuntimeSpec; +use crate::secrets::SharedSecrets; /// How the run is entered: fresh, from the admitted graphs, or continued /// from its records. @@ -66,18 +71,30 @@ pub enum Execution { pub struct RunRequest { /// The Fabro run id, which becomes Petri's run key: the run's identity /// in the store and the label on every sandbox of the run. - pub run_id: String, + pub run_id: String, /// Where the run's workspaces, step output and blobs live. - pub run_dir: PathBuf, - pub execution: Execution, + pub run_dir: PathBuf, + pub execution: Execution, /// The run's durable record: the worker's HTTP store, or the server's /// SQLite store under the test override. - pub store: Arc, - pub runtime: RuntimeSpec, + pub store: Arc, + pub runtime: RuntimeSpec, /// The sandbox provider Fabro resolved for the run's environment. - pub provider: SandboxProviderKind, + pub provider: SandboxProviderKind, /// Fires to cancel the run. - pub cancel: CancellationToken, + pub cancel: CancellationToken, + /// Where the run's questions go. + pub interviewer: Arc, + /// The caller's observers of every record, registered ahead of the + /// interview dispatcher: the interviewer's own expiry observer among + /// them. + pub observers: Vec>, + /// Where `{{ secrets.NAME }}` references resolve from; `None` leaves + /// every secret unknown. + pub secrets: Option>, + /// Where offloaded stage values go; `None` keeps Petri's local store + /// under the run directory. + pub blobs: Option>, } /// The recorded status of a finished run. @@ -140,13 +157,19 @@ 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); + if let Some(secrets) = request.secrets { + runtime = runtime.secrets(SharedSecrets(secrets)); + } + if let Some(blobs) = request.blobs { + runtime = runtime.capability(RunBlobs::output_store(blobs)); + } - let dispatcher = InterviewDispatcher::new(Arc::new(Unattended)); + let dispatcher = InterviewDispatcher::new(request.interviewer); let cancel = request.cancel.clone(); let mut cancel_task = None; let with_handle = |handle: petri_execution::CoordinatorHandle, secrets| { @@ -157,12 +180,15 @@ pub async fn run(request: RunRequest) -> Result { handle.cancel_root_for(CancelReason::Control); })); }; + let mut observers = request.observers; + observers.push(Arc::new(dispatcher.clone())); let result = match request.execution { Execution::Start(graphs) => { info!(run_id = %request.run_id, backend = %backend, "Starting Petri run"); - let host_run = HostRun::new(graphs.graph) - .with_children(graphs.children) - .observe(Arc::new(dispatcher.clone())); + let mut host_run = HostRun::new(graphs.graph).with_children(graphs.children); + for observer in observers { + host_run = host_run.observe(observer); + } Box::pin(host::run_configured(&runtime, host_run, with_handle)).await } Execution::Resume => { @@ -171,7 +197,7 @@ pub async fn run(request: RunRequest) -> Result { Box::pin(host::resume_configured( &runtime, Vec::new(), - vec![Arc::new(dispatcher.clone())], + observers, with_handle, )) .await diff --git a/lib/components/fabro-petri/src/interview.rs b/lib/components/fabro-petri/src/interview.rs new file mode 100644 index 000000000..fb62019d6 --- /dev/null +++ b/lib/components/fabro-petri/src/interview.rs @@ -0,0 +1,862 @@ +//! Petri's `Interviewer` over Fabro's questions API and the worker's +//! control channel. +//! +//! A human gate in a Petri run asks through Petri's interview boundary: the +//! dispatcher hands this adapter one [`InterviewRequest`] per question, on +//! its own task, with the question's identity (invocation path, execution, +//! firing, attempt, node, occurrence, ask). The adapter surfaces the +//! question to Fabro the way a legacy `human` stage does, waits for the +//! answer the way the legacy worker does, and hands Petri the reply. +//! +//! # How a question reaches a person +//! +//! The legacy stage emits `interview.started` on the run's event stream; +//! the read side keeps it in the projection's `pending_interviews`, keyed +//! by question id, and that is what `GET /runs/{id}/questions`, the web +//! app's interview dock and the Slack integration read pending questions +//! from. This adapter posts the same event through a [`QuestionSink`]: the +//! worker's [`EventSinkQuestions`] appends it over the run event sink the +//! worker already carries lifecycle events on, and the server's in-process +//! path appends it through [`DatabaseQuestions`]. The question id is +//! Fabro's key for the question and is derived from Petri's identity +//! ([`question_id`]); the node name is the event's `stage`, and the Fabro +//! question type, options, freeform flag, deadline and review target are +//! mapped from Petri's [`Question`]. +//! +//! # How the answer comes back +//! +//! `POST /runs/{id}/questions/{qid}/answer` validates the answer against +//! the pending record and delivers it to the run: over the worker control +//! bus as an `interview.answer` message, which the worker's control +//! manager applies to its [`ControlInterviewer`] by question id, or +//! straight to that interviewer for a run in the server process. The +//! adapter waits on that interviewer under the same id, so an answer +//! submitted before the wait began is buffered and one submitted after it +//! is delivered. The legacy answer shape is mapped onto Petri's +//! [`Answer`]: `yes` and `no` name the gate's affirmative and negative +//! choices by key, a selection names its key, a multi-selection its keys, +//! free text is text. A cancelled or interrupted answer ends the interview +//! without one: Petri's gate fails closed on it. +//! +//! # Expiry and cancellation +//! +//! The gate owns its answer deadline (Fabro's default when the node names +//! none) and reports the expiry itself; the dispatcher then fires the +//! adapter's cancel token, as it does when the firing ends without an +//! answer or the run is cancelled. The adapter returns promptly with +//! [`InterviewReply::Cancelled`] and posts `interview.timeout` when the +//! gate reported the expiry, else `interview.interrupted`, so the pending +//! question clears from Fabro's view. The expiry report is seen by the +//! adapter's own observer ([`FabroInterviewer::observer`]), which the run +//! registers ahead of the dispatcher so the report is noted before the +//! token fires. The dispatcher races the reply against the same token and +//! may drop the reply future the moment the token fires, so the notice is +//! posted from a guard that runs whether the future completes or is +//! dropped, on a task of its own. The dispatcher's own record of the +//! outcome (`TimedOut` with the default taken, `Cancelled`, `Late`) is the +//! authoritative one and reaches the receipt. +//! +//! # Auto-approval +//! +//! A run whose `[run.execution] approval` is `auto` answers every question +//! at once as the legacy runner's auto-approve interviewer does (`yes`, +//! the first option, or `auto-approved` text), attributed to the engine. +//! The question is still posted and completed, so the run's stream shows +//! what was decided. +//! +//! # Hook points for the read side +//! +//! The events posted here are the interim bridge to Fabro's read side. +//! Once the projection over Petri's records derives pending questions from +//! the `question` and `question_expired` records and the delivered answer, +//! the sink can become a no-op: the [`QuestionSink`] is the one seam to +//! replace. Two Fabro facts a Petri record does not carry are marked in +//! [`FabroInterviewer::reply`]: who answered (`AnswerSubmission::actor`, +//! carried on `interview.completed` for now) and the Fabro question id +//! that Petri's identity was mapped to. Both belong in a platform record +//! keyed on the same identity when that record kind exists. + +use std::collections::HashSet; +use std::sync::{Arc, Mutex, PoisonError}; +use std::time::{Duration, Instant}; + +use fabro_interview::{ + Answer as LegacyAnswer, AnswerSubmission, AnswerValue, AutoApproveInterviewer, + ControlInterviewer, Interviewer as LegacyInterviewer, Question as LegacyQuestion, +}; +use fabro_store::RunDatabase; +use fabro_types::{ + InterviewOption, Principal, QuestionType, ReviewTarget, ReviewTargetKind, RunId, + SystemActorKind, +}; +use fabro_workflow::event::{self as workflow_event, Event, RunEventSink}; +use petri_execution::{ + CoordinatorRecord, ExecutionId, ExecutionObserver, InterviewError, InterviewReply, + InterviewRequest, Interviewer, +}; +use petri_runtime::engine::{EngineState, Event as EngineEvent, EventRecord}; +use petri_runtime::steps::{Answer, Question, QuestionExpired, QuestionOption}; +use tokio::runtime::Handle; +use tokio_util::sync::CancellationToken; +use tracing::{debug, warn}; + +/// Whether a run answers its own questions. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum Approval { + /// A person answers, through the API. + Prompt, + /// The engine answers at once, as `--auto-approve` does. + Auto, +} + +/// Petri's identity for one question, as the read side keys it. +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct QuestionIdentity { + pub invocation_path: String, + pub execution: u64, + pub firing: u64, + pub attempt: u32, + pub node: String, + pub occurrence: u32, + pub ask: u32, +} + +impl QuestionIdentity { + fn of(request: &InterviewRequest) -> Self { + Self { + invocation_path: request.invocation_path.clone(), + execution: request.execution.raw(), + firing: request.firing.raw(), + attempt: request.attempt.raw(), + node: request.node.to_string(), + occurrence: request.occurrence, + ask: request.ask, + } + } +} + +/// A question as Fabro shows it: the fields of `interview.started`. +#[derive(Clone, Debug, PartialEq)] +pub struct AskedQuestion { + pub question_id: String, + pub identity: QuestionIdentity, + pub text: String, + pub stage: String, + pub question_type: QuestionType, + pub options: Vec, + pub allow_freeform: bool, + pub timeout_seconds: Option, + pub review_target: Option, +} + +/// What the adapter tells Fabro about a question, in the order it happens. +#[derive(Clone, Debug, PartialEq)] +pub enum QuestionNotice { + Asked(AskedQuestion), + Answered { + question_id: String, + text: String, + /// The answer as Fabro records it: the word, the key, the keys, or + /// the text; a sensitive answer is masked. + answer: String, + actor: Principal, + duration_ms: u64, + }, + Expired { + question_id: String, + text: String, + stage: String, + duration_ms: u64, + }, + Interrupted { + question_id: String, + text: String, + stage: String, + reason: String, + duration_ms: u64, + }, +} + +impl QuestionNotice { + /// The run event the legacy `human` stage emits for the same fact. + #[must_use] + pub fn into_event(self) -> Event { + match self { + Self::Asked(asked) => Event::InterviewStarted { + question_id: asked.question_id, + question: asked.text, + stage: asked.stage, + question_type: asked.question_type.to_string(), + options: asked.options, + allow_freeform: asked.allow_freeform, + timeout_seconds: asked.timeout_seconds, + context_display: None, + review_target: asked.review_target, + }, + Self::Answered { + question_id, + text, + answer, + actor, + duration_ms, + } => Event::InterviewCompleted { + actor: Some(actor), + question_id, + question: text, + answer, + duration_ms, + }, + Self::Expired { + question_id, + text, + stage, + duration_ms, + } => Event::InterviewTimeout { + actor: None, + question_id, + question: text, + stage, + duration_ms, + }, + Self::Interrupted { + question_id, + text, + stage, + reason, + duration_ms, + } => Event::InterviewInterrupted { + actor: None, + question_id, + question: text, + stage, + reason, + duration_ms, + }, + } + } + + fn question_id(&self) -> &str { + match self { + Self::Asked(asked) => &asked.question_id, + Self::Answered { question_id, .. } + | Self::Expired { question_id, .. } + | Self::Interrupted { question_id, .. } => question_id, + } + } +} + +/// Where the adapter posts what happens to a question: the run's event +/// stream, whichever way the process reaches it. +#[async_trait::async_trait] +pub trait QuestionSink: Send + Sync { + async fn post(&self, notice: QuestionNotice) -> anyhow::Result<()>; +} + +/// The worker's sink: the run event sink its lifecycle events go through. +pub struct EventSinkQuestions { + sink: RunEventSink, + run_id: RunId, +} + +impl EventSinkQuestions { + #[must_use] + pub fn new(sink: RunEventSink, run_id: RunId) -> Self { + Self { sink, run_id } + } +} + +#[async_trait::async_trait] +impl QuestionSink for EventSinkQuestions { + async fn post(&self, notice: QuestionNotice) -> anyhow::Result<()> { + workflow_event::append_event_to_sink(&self.sink, &self.run_id, ¬ice.into_event()) + .await + .map_err(anyhow::Error::new) + } +} + +/// The server's sink for a run in its own process: the run's database. +pub struct DatabaseQuestions { + store: RunDatabase, + run_id: RunId, +} + +impl DatabaseQuestions { + #[must_use] + pub fn new(store: RunDatabase, run_id: RunId) -> Self { + Self { store, run_id } + } +} + +#[async_trait::async_trait] +impl QuestionSink for DatabaseQuestions { + async fn post(&self, notice: QuestionNotice) -> anyhow::Result<()> { + workflow_event::append_event(&self.store, &self.run_id, ¬ice.into_event()).await + } +} + +/// The questions whose expiry the gate reported, by execution and Petri +/// question id: an observer the run registers ahead of the dispatcher. +#[derive(Default)] +pub struct Expiries { + expired: Mutex>, +} + +impl Expiries { + fn contains(&self, execution: ExecutionId, question: &str) -> bool { + self.expired + .lock() + .unwrap_or_else(PoisonError::into_inner) + .contains(&(execution, question.to_string())) + } +} + +impl ExecutionObserver for Expiries { + fn on_engine_record( + &self, + execution: ExecutionId, + record: &EventRecord, + _recorded_at: u64, + _state: &EngineState, + ) { + let EngineEvent::StepProgressRecorded { ev, .. } = &record.event else { + return; + }; + if let Some(expired) = QuestionExpired::from_event(ev) { + self.expired + .lock() + .unwrap_or_else(PoisonError::into_inner) + .insert((execution, expired.question)); + } + } + + fn on_lifecycle(&self, _record: &CoordinatorRecord) {} +} + +/// The interviewer a Fabro run installs. +pub struct FabroInterviewer { + answers: Arc, + sink: Arc, + approval: Approval, + expiries: Arc, +} + +impl FabroInterviewer { + /// Over the control interviewer the run's answers are delivered to, + /// and the sink its questions are posted through. + #[must_use] + pub fn new( + answers: Arc, + sink: Arc, + approval: Approval, + ) -> Self { + Self { + answers, + sink, + approval, + expiries: Arc::new(Expiries::default()), + } + } + + /// The observer that sees a gate report a question's expiry. A run + /// registers it ahead of the interview dispatcher, so the adapter + /// tells an expiry from an interruption when the dispatcher ends its + /// wait. + #[must_use] + pub fn observer(&self) -> Arc { + self.expiries.clone() + } + + /// Post a notice; a failure after the question was asked is logged, + /// since the answer, not the notice, is what the run depends on. + async fn post(&self, notice: QuestionNotice) { + let question_id = notice.question_id().to_string(); + if let Err(error) = self.sink.post(notice).await { + warn!( + question_id = %question_id, + error = format!("{error:#}"), + "a question notice could not be posted" + ); + } + } +} + +#[async_trait::async_trait] +impl Interviewer for FabroInterviewer { + async fn reply(&self, request: InterviewRequest, cancel: CancellationToken) -> InterviewReply { + let asked = asked_question(&request); + let question_id = asked.question_id.clone(); + let text = asked.text.clone(); + let stage = asked.stage.clone(); + let legacy = legacy_question(&asked); + // HOOK POINT (read side): the mapping from Petri's identity + // (`asked.identity`) to Fabro's question id is a platform fact + // worth a record keyed on that identity; today it lives only in + // the `interview.started` event posted here. + if let Err(error) = self.sink.post(QuestionNotice::Asked(asked)).await { + return InterviewReply::Failed(InterviewError::with_source( + format!("question `{question_id}` could not be published to Fabro"), + AnyhowError(error), + )); + } + let mut outstanding = Outstanding { + sink: Arc::clone(&self.sink), + expiries: Arc::clone(&self.expiries), + execution: request.execution, + question: request.question.id.clone(), + question_id: question_id.clone(), + text: text.clone(), + stage: stage.clone(), + started: Instant::now(), + open: true, + }; + let submission = match self.approval { + Approval::Auto => Some(AutoApproveInterviewer::engine().ask(legacy).await), + Approval::Prompt => tokio::select! { + submission = self.answers.ask(legacy) => Some(submission), + () = cancel.cancelled() => None, + }, + }; + let duration_ms = millis(outstanding.started.elapsed()); + let Some(submission) = submission else { + // The dispatcher ended the wait: the gate expired the question, + // the firing finished, the run was cancelled, or the run ended. + outstanding.close_unanswered("cancelled"); + return InterviewReply::Cancelled; + }; + // HOOK POINT (read side): `submission.actor` is who answered, a + // Fabro fact Petri's answer record does not carry; it rides on + // `interview.completed` until a platform record holds it. + let Some(answer) = petri_answer(&submission.answer, &request.question) else { + outstanding.close_unanswered(&reason_of(&submission.answer.value)); + return InterviewReply::Cancelled; + }; + debug!(question_id = %question_id, actor = ?submission.actor, "question answered"); + outstanding.open = false; + self.post(QuestionNotice::Answered { + question_id, + text, + answer: describe(&answer, &request.question), + actor: submission.actor, + duration_ms, + }) + .await; + InterviewReply::Answered(answer) + } +} + +/// A question the adapter is waiting on. When the wait ends without an +/// answer, whether the adapter saw the cancel or the dispatcher dropped +/// the reply future first, the end of the question is posted from here +/// on its own task: `interview.timeout` when the gate reported the +/// expiry, else `interview.interrupted`. +struct Outstanding { + sink: Arc, + expiries: Arc, + execution: ExecutionId, + /// Petri's question id, as the expiry report names it. + question: String, + question_id: String, + text: String, + stage: String, + started: Instant, + open: bool, +} + +impl Outstanding { + /// End the question without an answer, for `reason` unless the gate + /// reported the expiry. + fn close_unanswered(&mut self, reason: &str) { + if !self.open { + return; + } + self.open = false; + let duration_ms = millis(self.started.elapsed()); + let expired = self.expiries.contains(self.execution, &self.question); + let notice = if expired { + QuestionNotice::Expired { + question_id: self.question_id.clone(), + text: self.text.clone(), + stage: self.stage.clone(), + duration_ms, + } + } else { + QuestionNotice::Interrupted { + question_id: self.question_id.clone(), + text: self.text.clone(), + stage: self.stage.clone(), + reason: reason.to_string(), + duration_ms, + } + }; + let sink = Arc::clone(&self.sink); + let question_id = self.question_id.clone(); + let post = async move { + if let Err(error) = sink.post(notice).await { + warn!( + question_id = %question_id, + error = format!("{error:#}"), + "the end of a question could not be posted" + ); + } + }; + if let Ok(handle) = Handle::try_current() { + handle.spawn(post); + } else { + warn!( + question_id = %self.question_id, + "no runtime to post the end of a question from" + ); + } + } +} + +impl Drop for Outstanding { + fn drop(&mut self) { + self.close_unanswered("cancelled"); + } +} + +/// An `anyhow` error as a source for Petri's interview error. +#[derive(Debug)] +struct AnyhowError(anyhow::Error); + +impl std::fmt::Display for AnyhowError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(f, "{:#}", self.0) + } +} + +impl std::error::Error for AnyhowError {} + +/// Fabro's id for a question, from Petri's identity: the node, then the +/// execution and firing (unique in the run), the occurrence and the ask +/// (a re-asked question is a new one). Only URL-safe characters, so the +/// id travels in the answer endpoint's path as it is. +#[must_use] +pub fn question_id(identity: &QuestionIdentity) -> String { + let node: String = identity + .node + .chars() + .map(|c| { + if c.is_ascii_alphanumeric() || c == '_' || c == '-' { + c + } else { + '_' + } + }) + .collect(); + format!( + "{node}.x{}.f{}.q{}.a{}", + identity.execution, identity.firing, identity.occurrence, identity.ask + ) +} + +/// The question as Fabro shows it. +fn asked_question(request: &InterviewRequest) -> AskedQuestion { + let identity = QuestionIdentity::of(request); + let question = &request.question; + AskedQuestion { + question_id: question_id(&identity), + identity, + text: question.text.clone(), + stage: request.node.to_string(), + question_type: question_type(question), + options: question + .options + .iter() + .map(|option| InterviewOption { + key: option.key.clone(), + label: option.label.clone(), + description: None, + preview: None, + }) + .collect(), + allow_freeform: question.freeform, + timeout_seconds: question + .timeout_ms + .map(|ms| Duration::from_millis(ms).as_secs_f64()), + review_target: question.reference.as_ref().and_then(|reference| { + let kind = match reference.kind.as_deref() { + None | Some("document") => ReviewTargetKind::Document, + Some(other) => { + warn!( + kind = other, + "review target kind is not one Fabro shows; showing a document" + ); + ReviewTargetKind::Document + } + }; + ReviewTarget::new(&reference.label, &reference.url, kind) + .inspect_err(|error| { + warn!(error = %error, "review target could not be shown"); + }) + .ok() + }), + } +} + +/// Fabro's question type: the one the gate names, else what the shape +/// implies. +fn question_type(question: &Question) -> QuestionType { + question + .kind + .as_deref() + .and_then(|kind| kind.parse().ok()) + .unwrap_or(if question.options.is_empty() { + QuestionType::Freeform + } else { + QuestionType::MultipleChoice + }) +} + +/// The legacy question the control interviewer waits under: only the id +/// matters to it; the rest is what the auto-approve interviewer decides on. +fn legacy_question(asked: &AskedQuestion) -> LegacyQuestion { + let mut question = LegacyQuestion::new(asked.text.clone(), asked.question_type); + question.id.clone_from(&asked.question_id); + question.options.clone_from(&asked.options); + question.allow_freeform = asked.allow_freeform; + question.timeout_seconds = asked.timeout_seconds; + question.stage.clone_from(&asked.stage); + question.review_target.clone_from(&asked.review_target); + question +} + +/// Petri's answer for a legacy one, or `None` when the person or the +/// engine ended the interview without one. +fn petri_answer(answer: &LegacyAnswer, question: &Question) -> Option { + match &answer.value { + AnswerValue::Yes => Some(Answer::choice(&affirmative_key(question))), + AnswerValue::No => Some(Answer::choice(&negative_key(question))), + AnswerValue::Selected(key) => Some(Answer::choice(key)), + AnswerValue::MultiSelected(keys) => Some(Answer::choices(keys.iter().cloned())), + AnswerValue::Text(text) => Some(Answer::text(text.clone())), + AnswerValue::Cancelled + | AnswerValue::Interrupted + | AnswerValue::Skipped + | AnswerValue::Timeout => None, + } +} + +/// Whether a choice is the affirmative one of a yes/no gate, as the gate +/// itself matches a `yes` answer: key `y` or `yes`, or label `yes`. +fn is_affirmative(option: &QuestionOption) -> bool { + option.key.eq_ignore_ascii_case("y") + || option.key.eq_ignore_ascii_case("yes") + || strip_accelerator(&option.label).eq_ignore_ascii_case("yes") +} + +fn is_negative(option: &QuestionOption) -> bool { + option.key.eq_ignore_ascii_case("n") + || option.key.eq_ignore_ascii_case("no") + || strip_accelerator(&option.label).eq_ignore_ascii_case("no") +} + +/// The key a `yes` answer names: the affirmative choice, else the word +/// itself for the gate to match. +fn affirmative_key(question: &Question) -> String { + question + .options + .iter() + .find(|option| is_affirmative(option)) + .map_or_else(|| "yes".to_string(), |option| option.key.clone()) +} + +/// The key a `no` answer names: the negative choice, else the first choice +/// that is not affirmative, else the word itself. +fn negative_key(question: &Question) -> String { + question + .options + .iter() + .find(|option| is_negative(option)) + .or_else(|| { + question + .options + .iter() + .find(|option| !is_affirmative(option)) + }) + .map_or_else(|| "no".to_string(), |option| option.key.clone()) +} + +/// A label without its `[K] ` accelerator prefix. +fn strip_accelerator(label: &str) -> &str { + let trimmed = label.trim(); + match trimmed + .strip_prefix('[') + .and_then(|rest| rest.split_once(']')) + { + Some((_, rest)) => rest.trim(), + None => trimmed, + } +} + +/// The answer as `interview.completed` records it. A sensitive text +/// answer is never written out: the dispatcher registers it as a secret. +fn describe(answer: &Answer, question: &Question) -> String { + if !answer.choices.is_empty() { + return answer.choices.join(", "); + } + if let Some(choice) = &answer.choice { + return choice.clone(); + } + match &answer.text { + Some(_) if question.sensitive => "***".to_string(), + Some(serde_json::Value::String(text)) => text.clone(), + Some(other) => other.to_string(), + None => String::new(), + } +} + +fn reason_of(value: &AnswerValue) -> String { + match value { + AnswerValue::Cancelled => "cancelled", + AnswerValue::Interrupted => "interrupted", + AnswerValue::Skipped => "skipped", + AnswerValue::Timeout => "timeout", + _ => "unanswered", + } + .to_string() +} + +fn millis(elapsed: Duration) -> u64 { + u64::try_from(elapsed.as_millis()).unwrap_or(u64::MAX) +} + +/// An engine actor, for callers that answer on the run's behalf. +#[must_use] +pub fn engine_actor() -> Principal { + Principal::System { + system_kind: SystemActorKind::Engine, + } +} + +/// A submission on the run's behalf. +#[must_use] +pub fn engine_submission(answer: LegacyAnswer) -> AnswerSubmission { + AnswerSubmission::new(answer, engine_actor()) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn yes_no() -> Question { + let mut question = Question::new("gate#3", "Go?"); + question.options = vec![ + QuestionOption { + key: "Y".into(), + label: "[Y] Yes".into(), + }, + QuestionOption { + key: "N".into(), + label: "[N] No".into(), + }, + ]; + question.kind = Some("yes_no".into()); + question + } + + #[test] + fn a_question_id_is_url_safe_and_names_the_identity() { + let identity = QuestionIdentity { + invocation_path: "/branch:fan@2:0:a".into(), + execution: 2, + firing: 3, + attempt: 1, + node: "approve plan".into(), + occurrence: 1, + ask: 2, + }; + assert_eq!(question_id(&identity), "approve_plan.x2.f3.q1.a2"); + } + + #[test] + fn yes_and_no_name_the_gates_choices_by_key() { + let question = yes_no(); + assert_eq!( + petri_answer(&LegacyAnswer::yes(), &question), + Some(Answer::choice("Y")) + ); + assert_eq!( + petri_answer(&LegacyAnswer::no(), &question), + Some(Answer::choice("N")) + ); + let mut approve = Question::new("q", "Ship?"); + approve.options = vec![ + QuestionOption { + key: "A".into(), + label: "Approve".into(), + }, + QuestionOption { + key: "R".into(), + label: "Reject".into(), + }, + ]; + assert_eq!( + petri_answer(&LegacyAnswer::yes(), &approve), + Some(Answer::choice("yes")), + "no affirmative choice: the word reaches the gate to match" + ); + assert_eq!( + petri_answer(&LegacyAnswer::no(), &approve), + Some(Answer::choice("A")), + "the first choice that is not affirmative" + ); + } + + #[test] + fn selections_text_and_refusals_map_to_petris_shapes() { + let question = yes_no(); + assert_eq!( + petri_answer( + &LegacyAnswer { + value: AnswerValue::Selected("N".into()), + selected_option: None, + text: None, + }, + &question + ), + Some(Answer::choice("N")) + ); + assert_eq!( + petri_answer( + &LegacyAnswer::multi_selected(vec!["A".into(), "B".into()]), + &question + ), + Some(Answer::choices(["A", "B"])) + ); + assert_eq!( + petri_answer(&LegacyAnswer::text("ship it"), &question), + Some(Answer::text("ship it")) + ); + for ended in [ + LegacyAnswer::cancelled(), + LegacyAnswer::interrupted(), + LegacyAnswer::skipped(), + LegacyAnswer::timeout(), + ] { + assert_eq!(petri_answer(&ended, &question), None); + } + } + + #[test] + fn the_question_type_is_the_gates_else_the_shapes() { + assert_eq!(question_type(&yes_no()), QuestionType::YesNo); + let mut choice = yes_no(); + choice.kind = None; + assert_eq!(question_type(&choice), QuestionType::MultipleChoice); + let mut free = Question::new("q", "Name?"); + free.freeform = true; + assert_eq!(question_type(&free), QuestionType::Freeform); + } + + #[test] + fn a_sensitive_text_answer_is_described_masked() { + let mut question = Question::new("q", "Token?"); + question.sensitive = true; + assert_eq!(describe(&Answer::text("hunter2"), &question), "***"); + question.sensitive = false; + assert_eq!(describe(&Answer::text("hunter2"), &question), "hunter2"); + assert_eq!(describe(&Answer::choices(["A", "B"]), &question), "A, B"); + } +} diff --git a/lib/components/fabro-petri/src/interviewer.rs b/lib/components/fabro-petri/src/interviewer.rs deleted file mode 100644 index 594843ba0..000000000 --- a/lib/components/fabro-petri/src/interviewer.rs +++ /dev/null @@ -1,24 +0,0 @@ -//! The interviewer of a run nobody is watching. -//! -//! Until the questions adapter over Fabro's API lands (F3.2), a Petri run in -//! the server has no way to reach a person. A human gate that asks anyway -//! gets a failure that says so, the gate fails closed, and the reason -//! reaches the interview receipt, instead of a question that waits forever. - -use petri_execution::{InterviewError, InterviewReply, InterviewRequest, Interviewer}; -use tokio_util::sync::CancellationToken; - -/// Fails every question with a clear error. -#[derive(Clone, Copy, Debug, Default)] -pub struct Unattended; - -#[async_trait::async_trait] -impl Interviewer for Unattended { - async fn reply(&self, request: InterviewRequest, _cancel: CancellationToken) -> InterviewReply { - InterviewReply::Failed(InterviewError::new(format!( - "node `{}` asked a question, but a Petri run has no interviewer yet: questions reach \ - nobody until the interview adapter lands", - request.node - ))) - } -} diff --git a/lib/components/fabro-petri/src/lib.rs b/lib/components/fabro-petri/src/lib.rs index 443032ac8..93de7a1ce 100644 --- a/lib/components/fabro-petri/src/lib.rs +++ b/lib/components/fabro-petri/src/lib.rs @@ -19,24 +19,33 @@ //! - [`engine`]: a run executed by Petri, started or resumed, in the run's //! worker process over the HTTP store (or in the server process under its //! test override), with the outcome read from its record; -//! - [`interviewer`]: the interviewer of a run nobody is watching; +//! - [`interview`]: Petri's interviewer over Fabro's questions API and the +//! worker's control channel, so a human gate's question reaches the same +//! places a legacy stage's does and its answer comes back the same way; +//! - [`secrets`]: Petri's secret provider over Fabro's vault, so a `{{ +//! secrets.NAME }}` reference resolves from the vault at spawn and is masked +//! in every record; +//! - [`blobs`]: Petri's output store over Fabro's blob table, so a large stage +//! value lives in `blobs` under `blob://sha256/`; //! - [`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. +//! - the platform adapters still to come: hooks, 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 blobs; pub mod check; pub mod engine; pub mod http_store; -pub mod interviewer; +pub mod interview; pub mod petri; pub mod run_store; pub mod runtime; +pub mod secrets; #[cfg(feature = "test-support")] pub mod test_support; diff --git a/lib/components/fabro-petri/src/secrets.rs b/lib/components/fabro-petri/src/secrets.rs new file mode 100644 index 000000000..1f52328ea --- /dev/null +++ b/lib/components/fabro-petri/src/secrets.rs @@ -0,0 +1,170 @@ +//! Petri's `SecretProvider` over Fabro's vault. +//! +//! A run's commands reach a secret as a `{"$secret": "NAME"}` reference +//! that Petri resolves at spawn, straight into the child's environment; the +//! value never enters a record. Resolving a secret registers it with the +//! run's masker, and every record and log line is masked before it is +//! appended, so a value that was resolved cannot appear in `petri_records`. +//! This provider is what makes the vault the place those names resolve +//! from, in the worker (over the vault snapshot the worker loads from the +//! server storage) and in the server process under its test override. +//! +//! Only `Token` entries resolve, as the legacy runner resolves +//! `{{ secrets.NAME }}` (`fabro_auth::vault_get_token`): an OAuth record +//! or a file-shaped secret is not a value a command's environment should +//! carry, so such a name is unknown here. +//! +//! [`SecretProvider::register`] is served: a human gate's sensitive answer +//! is registered under `answer:` before its reference is +//! delivered, and lives as long as the provider, which is the run. + +use std::sync::Arc; + +use fabro_types::SecretType; +use fabro_vault::Vault; +use petri_runtime::executor::{MapSecrets, Masker, Secret, SecretError, SecretProvider}; + +/// The vault's token entries, as Petri's secret provider for one run. +pub struct VaultSecrets { + inner: MapSecrets, +} + +impl VaultSecrets { + /// A provider over the vault's `Token` entries as they are now: the + /// worker holds a snapshot, so a later change to the vault is not seen + /// by a running run, as with the legacy runner. + #[must_use] + pub fn from_vault(vault: &Vault) -> Self { + let pairs = vault + .entries() + .iter() + .filter(|(_, entry)| entry.secret_type == SecretType::Token) + .map(|(name, entry)| (name.as_str(), entry.value.as_str())) + .collect::>(); + Self::from_pairs(&pairs) + } + + /// A provider over the given names and values. + #[must_use] + pub fn from_pairs(pairs: &[(&str, &str)]) -> Self { + Self { + inner: MapSecrets::from_pairs(pairs), + } + } +} + +impl SecretProvider for VaultSecrets { + fn resolve(&self, name: &str) -> Result { + self.inner.resolve(name) + } + + fn register(&self, name: &str, value: &str) -> Result<(), SecretError> { + self.inner.register(name, value) + } + + fn masker(&self) -> Masker { + self.inner.masker() + } +} + +/// A shared provider, installed on a runtime that takes its provider by +/// value: the run's engine assembly holds the provider as a trait object +/// so a caller can hand in any implementation. +pub struct SharedSecrets(pub Arc); + +impl SecretProvider for SharedSecrets { + fn resolve(&self, name: &str) -> Result { + self.0.resolve(name) + } + + fn register(&self, name: &str, value: &str) -> Result<(), SecretError> { + self.0.register(name, value) + } + + fn masker(&self) -> Masker { + self.0.masker() + } +} + +#[cfg(test)] +mod tests { + use std::collections::HashMap; + + use super::*; + + fn vault() -> Vault { + let mut vault = Vault::from_entries(HashMap::new()); + vault + .set("TOKEN", "hunter2-hunter2", SecretType::Token, None) + .expect("a detached vault takes an entry"); + vault + .set( + "OAUTH", + r#"{"access_token":"oauth-secret-value"}"#, + SecretType::Oauth, + None, + ) + .expect("a detached vault takes an entry"); + vault + } + + #[test] + fn a_token_entry_resolves_and_is_masked_afterwards() { + let secrets = VaultSecrets::from_vault(&vault()); + let masker = secrets.masker(); + assert!( + !masker.contains_secret("hunter2-hunter2"), + "nothing resolved yet" + ); + let secret = secrets.resolve("TOKEN").expect("the token resolves"); + assert_eq!(secret.expose(), "hunter2-hunter2"); + assert_eq!(masker.mask("got hunter2-hunter2"), "got ***"); + } + + #[test] + fn a_non_token_entry_and_an_unknown_name_are_unknown() { + let secrets = VaultSecrets::from_vault(&vault()); + assert!(matches!( + secrets.resolve("OAUTH"), + Err(SecretError::Unknown(name)) if name == "OAUTH" + )); + assert!(matches!( + secrets.resolve("MISSING"), + Err(SecretError::Unknown(name)) if name == "MISSING" + )); + } + + #[test] + fn a_dynamic_secret_registers_once_and_masks_at_once() { + let secrets = VaultSecrets::from_vault(&vault()); + secrets + .register("answer:gate#3", "sensitive-answer") + .expect("a new name registers"); + assert_eq!( + secrets.masker().mask("said sensitive-answer"), + "said ***", + "registration feeds the masker before any resolution" + ); + assert_eq!( + secrets + .resolve("answer:gate#3") + .expect("registered") + .expose(), + "sensitive-answer" + ); + assert!(matches!( + secrets.register("TOKEN", "shadow"), + Err(SecretError::Duplicate(_)) + )); + } + + #[test] + fn a_shared_provider_delegates() { + let shared = SharedSecrets(Arc::new(VaultSecrets::from_vault(&vault()))); + assert_eq!( + shared.resolve("TOKEN").expect("resolves").expose(), + "hunter2-hunter2" + ); + assert!(shared.masker().contains_secret("hunter2-hunter2")); + } +}