Install Petri's interview, secret, blob and home adapters in a Fabro run

A Petri run in the worker, and in the server under its test override, now
gets Fabro's platform adapters instead of the standalone defaults:

- `fabro_petri::interview`: Petri's `Interviewer` over the questions API
  and the worker's control channel. A human gate's question is posted as
  the `interview.started` event a legacy stage emits, keyed by an id
  derived from Petri's identity (node, execution, firing, occurrence,
  ask), so the API, the web app and Slack list it; the answer posted to
  the questions endpoint reaches the control interviewer the adapter waits
  on and is mapped onto Petri's answer. An expiry the gate reports is
  completed as `interview.timeout`, a cancel as `interview.interrupted`,
  and an auto-approved run answers itself. The hook points the read side
  takes over are marked.
- `fabro_petri::secrets`: Petri's `SecretProvider` over the vault's token
  entries, so `{{ secrets.NAME }}` resolves at spawn and is masked in every
  record; a sensitive answer registers as a dynamic secret.
- `fabro_petri::blobs`: Petri's `OutputStore` over Fabro's `blobs` table,
  through the server's blob store or the worker's client.
- The Fabro home the server resolved travels to the worker as
  `--fabro-home`, so the skills step reads it whatever the worker's
  environment says.

`engine::RunRequest` takes the interviewer, its observers, the secret
provider and the blob table from the caller; `interviewer::Unattended` is
gone.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-17 22:31:35 -04:00
parent 16354186fa
commit 223e10ea20
No known key found for this signature in database
16 changed files with 1411 additions and 74 deletions

4
Cargo.lock generated
View file

@ -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",

View file

@ -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<PathBuf>,
}
#[derive(Args, Debug, Clone, Default)]

View file

@ -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
}

View file

@ -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<PathBuf>,
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<RuntimeSpec> {
/// the providers whose credentials resolve, the run's mode, and the Fabro
/// home the server named.
async fn runtime_spec(
vault: &Arc<AsyncRwLock<Vault>>,
run_state: &RunProjection,
fabro_home: Option<PathBuf>,
) -> Result<RuntimeSpec> {
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<dyn CredentialProvider> = Arc::new(VaultCredentialSource::new(vault));
let credentials: Arc<dyn CredentialProvider> =
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<R
settings_toml: None,
model_client,
dry_run: run_state.spec.settings.run.execution.mode == RunMode::DryRun,
fabro_home: None,
fabro_home,
})
}

View file

@ -74,6 +74,7 @@ pub(crate) async fn execute(
storage_dir: PathBuf,
run_dir: PathBuf,
mode: RunWorkerMode,
fabro_home: Option<PathBuf>,
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;

View file

@ -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(),
})
}

View file

@ -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<AppState>, 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<AppState>, 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<AppState>, 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 {

View file

@ -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(),
]);
}

View file

@ -48,6 +48,9 @@ pub(crate) struct WorkerLaunchSpec {
pub(crate) fabro_log: Option<String>,
pub(crate) active_config_path: PathBuf,
pub(crate) github_app_private_key: Option<String>,
/// 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());

View file

@ -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"] }

View file

@ -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/<hex>`, 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<BlobHash>;
/// The bytes behind a hash, or `None` when the table has none.
async fn read(&self, hash: &BlobHash) -> anyhow::Result<Option<Bytes>>;
}
#[async_trait::async_trait]
impl Blobs for fabro_store::BlobStore {
async fn write(&self, bytes: &[u8]) -> anyhow::Result<BlobHash> {
Self::write(self, bytes).await.map_err(anyhow::Error::new)
}
async fn read(&self, hash: &BlobHash) -> anyhow::Result<Option<Bytes>> {
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<BlobHash> {
self.client.write_run_blob(&self.run_id, bytes).await
}
async fn read(&self, hash: &BlobHash) -> anyhow::Result<Option<Bytes>> {
self.client.read_run_blob(&self.run_id, hash).await
}
}
/// Petri's blob store over Fabro's blob table.
pub struct RunBlobs {
blobs: Arc<dyn Blobs>,
}
impl RunBlobs {
#[must_use]
pub fn new(blobs: Arc<dyn Blobs>) -> Self {
Self { blobs }
}
/// The capability a runtime installs so every offloaded value goes to
/// the table.
#[must_use]
pub fn output_store(blobs: Arc<dyn Blobs>) -> 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<String, BlobError> {
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<Option<Vec<u8>>, BlobError> {
let Ok(hash) = digest.parse::<BlobHash>() 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<Vec<(BlobHash, Vec<u8>)>>,
}
#[async_trait::async_trait]
impl Blobs for MemoryBlobs {
async fn write(&self, bytes: &[u8]) -> anyhow::Result<BlobHash> {
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<Option<Bytes>> {
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);
}
}

View file

@ -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<dyn RunStore>,
pub runtime: RuntimeSpec,
pub store: Arc<dyn RunStore>,
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<dyn Interviewer>,
/// The caller's observers of every record, registered ahead of the
/// interview dispatcher: the interviewer's own expiry observer among
/// them.
pub observers: Vec<Arc<dyn ExecutionObserver>>,
/// Where `{{ secrets.NAME }}` references resolve from; `None` leaves
/// every secret unknown.
pub secrets: Option<Arc<dyn SecretProvider>>,
/// Where offloaded stage values go; `None` keeps Petri's local store
/// under the run directory.
pub blobs: Option<Arc<dyn Blobs>>,
}
/// The recorded status of a finished run.
@ -140,13 +157,19 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
options.run_key = Some(key.clone());
options.retention = Retention::Always;
options.sandbox.backend = backend;
let runtime = request
let mut runtime = request
.runtime
.runtime(true)
.store(Arc::clone(&request.store))
.options(options);
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<RunOutcome, RunError> {
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<RunOutcome, RunError> {
Box::pin(host::resume_configured(
&runtime,
Vec::new(),
vec![Arc::new(dispatcher.clone())],
observers,
with_handle,
))
.await

View file

@ -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<InterviewOption>,
pub allow_freeform: bool,
pub timeout_seconds: Option<f64>,
pub review_target: Option<ReviewTarget>,
}
/// 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, &notice.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, &notice.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<HashSet<(ExecutionId, String)>>,
}
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<ControlInterviewer>,
sink: Arc<dyn QuestionSink>,
approval: Approval,
expiries: Arc<Expiries>,
}
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<ControlInterviewer>,
sink: Arc<dyn QuestionSink>,
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<dyn ExecutionObserver> {
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<dyn QuestionSink>,
expiries: Arc<Expiries>,
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<Answer> {
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");
}
}

View file

@ -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
)))
}
}

View file

@ -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/<hex>`;
//! - [`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;

View file

@ -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:<question id>` 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::<Vec<_>>();
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<Secret, SecretError> {
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<dyn SecretProvider>);
impl SecretProvider for SharedSecrets {
fn resolve(&self, name: &str) -> Result<Secret, SecretError> {
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"));
}
}