Merge branch 'petri-integration-adapters' into petri-integration

This commit is contained in:
Bryan Helmkamp 2026-09-17 22:40:30 -04:00
commit 53b16b591d
No known key found for this signature in database
25 changed files with 2869 additions and 99 deletions

View file

@ -15,6 +15,12 @@ leak-timeout = "500ms"
filter = "package(fabro-workflow)"
slow-timeout = { period = "2s", terminate-after = 3 }
# fabro-petri's adapter tests run whole workflows on the host sandbox
# through the sandbox-driver plugin, and one of them calls the twin.
[[profile.default.overrides]]
filter = "package(fabro-petri)"
slow-timeout = { period = "5s", terminate-after = 4 }
# Real descendant regressions include bounded reaping and process probes.
# Leave room for their own watchdogs to run fail-safe fixture cleanup.
[[profile.default.overrides]]
@ -59,3 +65,7 @@ leak-timeout = "2s"
filter = "package(fabro-workflow)"
slow-timeout = { period = "30s", terminate-after = 4 }
[[profile.ci.overrides]]
filter = "package(fabro-petri)"
slow-timeout = { period = "30s", terminate-after = 4 }

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

@ -251,17 +251,22 @@ async fn wait_for_http_ready(base_url: &str, child: &mut Child) {
/// A workspace holding a command-only bundle whose `workflow.toml` names
/// Petri, with the given stage script.
fn write_petri_workspace(context: &fabro_test::TestContext, script: &str) -> PathBuf {
let workspace = context.temp_dir.join("petri-workspace");
std::fs::create_dir_all(&workspace).expect("the workspace creates");
std::fs::write(
workspace.join("workflow.fabro"),
format!(
write_petri_workflow(
context,
&format!(
"digraph Command {{\n graph [goal=\"Run one command\", default_max_retries=0]\n start \
[shape=Mdiamond]\n exit [shape=Msquare]\n say [shape=parallelogram, \
script=\"{script}\", max_retries=0]\n start -> say -> exit\n}}\n"
),
)
.expect("the workflow writes");
}
/// A workspace holding the given workflow with a `workflow.toml` that names
/// Petri.
fn write_petri_workflow(context: &fabro_test::TestContext, dot: &str) -> PathBuf {
let workspace = context.temp_dir.join("petri-workspace");
std::fs::create_dir_all(&workspace).expect("the workspace creates");
std::fs::write(workspace.join("workflow.fabro"), dot).expect("the workflow writes");
std::fs::write(
workspace.join("workflow.toml"),
"_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\nengine = \"petri\"\n\n[run]\ngoal \
@ -271,12 +276,22 @@ fn write_petri_workspace(context: &fabro_test::TestContext, script: &str) -> Pat
workspace
}
/// `fabro run --detach` against the server: the run is created and started,
/// and its id comes back.
/// `fabro run --detach --auto-approve` against the server: the run is
/// created and started, and its id comes back.
fn run_detached(
context: &fabro_test::TestContext,
server: &RunningServer,
workspace: &Path,
) -> String {
run_detached_with(context, server, workspace, &["--auto-approve"])
}
/// `fabro run --detach` against the server with extra arguments.
fn run_detached_with(
context: &fabro_test::TestContext,
server: &RunningServer,
workspace: &Path,
extra: &[&str],
) -> String {
let target = server.target();
seed_dev_token_auth(
@ -287,15 +302,9 @@ fn run_detached(
let output = context
.run_cmd()
.current_dir(workspace)
.args([
"--server",
&target,
"--detach",
"--auto-approve",
"--environment",
"local",
"workflow.toml",
])
.args(["--server", &target, "--detach"])
.args(extra)
.args(["--environment", "local", "workflow.toml"])
.output()
.expect("the detached run executes");
assert!(
@ -566,3 +575,227 @@ async fn run_status_offline(server: &RunningServer) -> Option<String> {
.ok()
.map(|response| response.status().to_string())
}
/// The run's pending questions, as the API lists them.
async fn questions(server: &RunningServer, run_id: &str) -> Vec<serde_json::Value> {
run_json(server, &format!("runs/{run_id}/questions")).await["data"]
.as_array()
.cloned()
.expect("the questions list is an array")
}
/// Wait until `count` questions are pending at once.
async fn wait_for_questions(
server: &RunningServer,
run_id: &str,
count: usize,
) -> Vec<serde_json::Value> {
let deadline = Instant::now() + RUN_TIMEOUT;
loop {
let pending = questions(server, run_id).await;
if pending.len() >= count {
return pending;
}
assert!(
Instant::now() < deadline,
"run {run_id} did not ask {count} question(s); pending: {pending:?}"
);
tokio::time::sleep(POLL).await;
}
}
/// Answer a question through the API, as the web app and the CLI do.
async fn answer(server: &RunningServer, run_id: &str, question_id: &str, body: serde_json::Value) {
let response = fabro_test::test_http_client()
.post(format!(
"{}/api/v1/runs/{run_id}/questions/{question_id}/answer",
server.api_base_url
))
.bearer_auth(TEST_DEV_TOKEN)
.json(&body)
.send()
.await
.expect("the answer sends");
let status = response.status();
let body = response.text().await.unwrap_or_default();
assert_eq!(
status,
fabro_http::StatusCode::NO_CONTENT,
"POST /api/v1/runs/{run_id}/questions/{question_id}/answer: {body}"
);
}
/// A yes/no gate whose branches each leave a marker file.
fn gate_dot(markers: &Path, gate_attrs: &str) -> String {
format!(
"digraph Gate {{\n graph [goal=\"Ask before running\"]\n start [shape=Mdiamond]\n \
exit [shape=Msquare]\n gate [shape=hexagon, label=\"Go?\", \
question_type=\"yes_no\"{gate_attrs}]\n yes [shape=parallelogram, script=\"touch \
{dir}/yes\"]\n no [shape=parallelogram, script=\"touch {dir}/no\"]\n start -> gate\n \
gate -> yes [label=\"[Y] Yes\"]\n gate -> no [label=\"[N] No\"]\n yes -> exit\n no \
-> exit\n}}\n",
dir = markers.display()
)
}
/// Two gates as the branches of one parallel node; the join's results are
/// written out, so each gate's answer is read from its branch result.
fn two_gates_dot(markers: &Path) -> String {
format!(
"digraph Gates {{\n graph [goal=\"Ask twice at once\"]\n start [shape=Mdiamond]\n \
exit [shape=Msquare]\n fan [shape=component]\n a [shape=hexagon, label=\"A?\", \
question_type=\"yes_no\"]\n b [shape=hexagon, label=\"B?\", \
question_type=\"yes_no\"]\n join [shape=tripleoctagon]\n report \
[shape=parallelogram, script=\"cat > {dir}/results.json\", \
stdin_source=\"context.parallel.results\"]\n start -> fan\n fan -> a\n fan -> b\n \
a -> join [label=\"[Y] Yes\"]\n a -> join [label=\"[N] No\"]\n b -> join [label=\"[Y] \
Yes\"]\n b -> join [label=\"[N] No\"]\n join -> report -> exit\n}}\n",
dir = markers.display()
)
}
/// A human gate in the worker asks through the server: the question is
/// listed by the questions API with the gate's stage and options, the
/// answer reaches the worker over its control channel and routes the gate,
/// and the run's stream records the interview.
#[tokio::test(flavor = "multi_thread")]
async fn a_human_gate_in_the_worker_is_answered_through_the_api() {
if host_plugin().is_none() {
return;
}
let context = test_context!();
let server = RunningServer::start().await;
let markers = context.temp_dir.join("markers");
std::fs::create_dir_all(&markers).expect("the marker dir creates");
let workspace = write_petri_workflow(&context, &gate_dot(&markers, ""));
let run_id = run_detached_with(&context, &server, &workspace, &[]);
let pending = wait_for_questions(&server, &run_id, 1).await;
let question = &pending[0];
assert_eq!(question["stage"], "gate", "{question}");
assert_eq!(question["question_type"], "yes_no", "{question}");
let question_id = question["id"].as_str().expect("an id").to_string();
answer(
&server,
&run_id,
&question_id,
serde_json::json!({ "kind": "no" }),
)
.await;
let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await;
assert_eq!(
status,
"succeeded",
"server stderr:\n{}",
server.stderr_text()
);
assert!(
markers.join("no").exists() && !markers.join("yes").exists(),
"the no branch ran"
);
let names = run_events(&server, &run_id).await;
let names = event_names(&names);
assert!(
names.contains(&"interview.started") && names.contains(&"interview.completed"),
"{names:?}"
);
assert!(questions(&server, &run_id).await.is_empty());
let store = server.petri_store().await;
let outcome = engine::outcome_of(&store, &run_id)
.await
.expect("the run's Petri record inspects");
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
server.shutdown();
}
/// Two branches of a parallel node ask at once; each answer, given through
/// the API in the other order, binds to its own branch.
#[tokio::test(flavor = "multi_thread")]
async fn two_parallel_gates_in_the_worker_each_bind_their_own_answer() {
if host_plugin().is_none() {
return;
}
let context = test_context!();
let server = RunningServer::start().await;
let markers = context.temp_dir.join("markers");
std::fs::create_dir_all(&markers).expect("the marker dir creates");
let workspace = write_petri_workflow(&context, &two_gates_dot(&markers));
let run_id = run_detached_with(&context, &server, &workspace, &[]);
let pending = wait_for_questions(&server, &run_id, 2).await;
let id_of = |stage: &str| {
pending
.iter()
.find(|question| question["stage"] == stage)
.and_then(|question| question["id"].as_str())
.unwrap_or_else(|| panic!("`{stage}` is pending: {pending:?}"))
.to_string()
};
let (a, b) = (id_of("a"), id_of("b"));
assert_ne!(a, b);
for question in &pending {
assert_eq!(question["question_type"], "yes_no", "{question}");
}
// A yes/no question takes `yes` or `no`, as the API validates it.
answer(&server, &run_id, &b, serde_json::json!({ "kind": "yes" })).await;
answer(&server, &run_id, &a, serde_json::json!({ "kind": "no" })).await;
let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await;
assert_eq!(
status,
"succeeded",
"server stderr:\n{}",
server.stderr_text()
);
let results: serde_json::Value = serde_json::from_str(
&std::fs::read_to_string(markers.join("results.json")).expect("the join wrote its results"),
)
.expect("the results parse");
let results = results.as_array().expect("a list of branch results");
assert_eq!(results.len(), 2, "{results:?}");
assert_eq!(results[0]["id"], "a");
assert_eq!(results[0]["context_updates"]["human.gate.selected"], "N");
assert_eq!(results[1]["id"], "b");
assert_eq!(results[1]["context_updates"]["human.gate.selected"], "Y");
server.shutdown();
}
/// A gate nobody answers expires on its own deadline: the run takes the
/// gate's default, the stream records the timeout, and nothing stays
/// pending.
#[tokio::test(flavor = "multi_thread")]
async fn an_unanswered_gate_in_the_worker_expires_with_its_default() {
if host_plugin().is_none() {
return;
}
let context = test_context!();
let server = RunningServer::start().await;
let markers = context.temp_dir.join("markers");
std::fs::create_dir_all(&markers).expect("the marker dir creates");
let workspace = write_petri_workflow(
&context,
&gate_dot(&markers, ", timeout=\"2s\", human.default_choice=\"no\""),
);
let run_id = run_detached_with(&context, &server, &workspace, &[]);
let pending = wait_for_questions(&server, &run_id, 1).await;
assert_eq!(pending[0]["timeout_seconds"], 2.0, "{}", pending[0]);
let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await;
assert_eq!(
status,
"succeeded",
"server stderr:\n{}",
server.stderr_text()
);
assert!(
markers.join("no").exists() && !markers.join("yes").exists(),
"the default ran"
);
let events = run_events(&server, &run_id).await;
let names = event_names(&events);
assert!(names.contains(&"interview.timeout"), "{names:?}");
assert!(questions(&server, &run_id).await.is_empty());
server.shutdown();
}

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

@ -390,3 +390,145 @@ async fn an_unknown_model_is_refused_at_create_with_attractor_model_unknown() {
"expected the admission diagnostic in the detail, got {body}"
);
}
/// A yes/no gate whose branches each leave a marker file.
fn gate_dot(markers: &std::path::Path) -> String {
format!(
r#"digraph Gate {{
graph [goal="Ask before running"]
start [shape=Mdiamond]
exit [shape=Msquare]
gate [shape=hexagon, label="Go?", question_type="yes_no"]
yes [shape=parallelogram, script="touch {dir}/yes"]
no [shape=parallelogram, script="touch {dir}/no"]
start -> gate
gate -> yes [label="[Y] Yes"]
gate -> no [label="[N] No"]
yes -> exit
no -> exit
}}"#,
dir = markers.display()
)
}
/// The run's first pending question, once one is listed.
async fn wait_for_question(app: &axum::Router, run_id: &str) -> serde_json::Value {
for _ in 0..600 {
let req = Request::builder()
.method("GET")
.uri(api(&format!("/runs/{run_id}/questions")))
.body(Body::empty())
.expect("questions request should build");
let response = app
.clone()
.oneshot(req)
.await
.expect("questions request routes");
let body = response_json(
response,
StatusCode::OK,
format!("GET /api/v1/runs/{run_id}/questions"),
)
.await;
if let Some(question) = body["data"].as_array().and_then(|items| items.first()) {
return question.clone();
}
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
panic!("run {run_id} never asked a question");
}
/// A human gate in a Petri run asks through the questions API and is
/// answered through it: the question is listed with the gate's stage and
/// options, the answer routes the gate, and the run's stream records the
/// interview as a legacy stage's would.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_human_gate_is_answered_through_the_questions_api() {
if host_plugin().is_none() {
return;
}
let workspace = tempfile::tempdir().expect("workspace tempdir");
let markers = tempfile::tempdir().expect("marker tempdir");
let settings = settings_from_toml(
"_version = 1\n\n[run.environment]\nid = \"local\"\n\n[server.execution]\nengine = \
\"petri\"\n",
);
let state = test_app_state_with_options(settings, 5);
let app = test_app_with_scheduler(Arc::clone(&state));
let dot = gate_dot(markers.path());
let version_id = register_version(&app, &[
("workflow.fabro", &dot),
("workflow.toml", PLAIN_SETTINGS),
])
.await;
let run_id =
create_and_start_run_from_intent(&app, intent(&version_id, workspace.path())).await;
let question = wait_for_question(&app, &run_id).await;
assert_eq!(question["stage"], "gate", "{question}");
assert_eq!(question["text"], "Go?", "{question}");
assert_eq!(question["question_type"], "yes_no", "{question}");
let keys: Vec<&str> = question["options"]
.as_array()
.expect("options")
.iter()
.filter_map(|option| option["key"].as_str())
.collect();
assert_eq!(keys, vec!["Y", "N"], "{question}");
let question_id = question["id"].as_str().expect("an id").to_string();
assert!(question_id.starts_with("gate."), "{question_id}");
let req = Request::builder()
.method("POST")
.uri(api(&format!(
"/runs/{run_id}/questions/{question_id}/answer"
)))
.header("content-type", "application/json")
.body(Body::from(r#"{"kind":"no"}"#))
.expect("answer request should build");
let response = app
.clone()
.oneshot(req)
.await
.expect("answer request routes");
crate::helpers::response_status(
response,
StatusCode::NO_CONTENT,
format!("POST /api/v1/runs/{run_id}/questions/{question_id}/answer"),
)
.await;
let status = wait_for_run_status(&app, &run_id, &["succeeded", "failed"]).await;
let run = run_json(&app, &run_id).await;
assert_eq!(status, "succeeded", "run: {run}");
assert!(
markers.path().join("no").exists() && !markers.path().join("yes").exists(),
"the no branch ran"
);
let outcome = petri_outcome(&state, &run_id).await;
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
let state_body = {
let req = Request::builder()
.method("GET")
.uri(api(&format!("/runs/{run_id}/state")))
.body(Body::empty())
.expect("state request should build");
response_json(
app.clone()
.oneshot(req)
.await
.expect("state request routes"),
StatusCode::OK,
format!("GET /api/v1/runs/{run_id}/state"),
)
.await
};
assert!(
state_body["pending_interviews"]
.as_object()
.is_some_and(serde_json::Map::is_empty),
"the answered question is no longer pending: {}",
state_body["pending_interviews"]
);
}

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
@ -48,6 +51,7 @@ 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
tempfile = "3"
tokio = { workspace = true, features = ["macros", "rt-multi-thread"] }

View file

@ -36,8 +36,27 @@ Every adapter the integration plan describes lands here.
through `inspect_run` and mapped to the conclusion Fabro's read side
records. The run's worker process runs it over `HttpRunStore`; the server
runs it in its own process only under its test override, over
`SqliteRunStore`. `interviewer::Unattended` fails any question until the
interview adapter lands.
`SqliteRunStore`. The caller supplies the interviewer, and the secret
provider and blob table when it has them.
- `interview`: Petri's `Interviewer` over Fabro's questions API and the
worker's control channel. A human gate's question is posted as the
`interview.started` event a legacy `human` stage emits (through the
worker's run event sink, or the run's database in the server process), so
`GET /runs/{id}/questions`, the web app and Slack list it; the answer
posted to `/questions/{qid}/answer` reaches the worker's control
interviewer over the control bus (or the in-process one directly) under
the same id, and is mapped onto Petri's answer. The question id is
derived from Petri's identity (node, execution, firing, occurrence, ask).
An expired or cancelled question is completed as `interview.timeout` or
`interview.interrupted`; an auto-approved run answers itself. The module
docs mark the hook points the read side takes over.
- `secrets`: Petri's `SecretProvider` over the vault's token entries, so a
`{{ secrets.NAME }}` reference resolves at spawn into a command's
environment and is masked in every record; a sensitive answer registers
as a dynamic secret.
- `blobs`: Petri's `OutputStore` over Fabro's `blobs` table, through the
server's `BlobStore` or the worker's client, so a large stage value
leaves the records for the table under `blob://sha256/<hex>`.
- `HttpRunStore`: the same store as a run's worker process reaches it, over
the server's `/api/v1/runs/{id}/petri/*` endpoints with the worker's token.
The server answers from its `SqliteRunStore`, so the lease and the
@ -48,8 +67,8 @@ Every adapter the integration plan describes lands here.
- `petri`: the Petri store vocabulary re-exported for the server, which
answers the worker endpoints from a `SqliteRunStore` without naming a Petri
package in its own manifest.
- The platform adapters the plan adds after it: hooks, interviews over
Fabro's API, secrets, output storage, run tools, the event projection.
- The platform adapters the plan adds after it: hooks, run tools, the
event projection.
A run goes to Petri when its workflow version's `workflow.toml` names
`engine = "petri"` in `[workflow]`, or when the server's
@ -83,6 +102,21 @@ Integration tests live under `tests/`:
(`petri_testkit::run_store::conformance`) against `SqliteRunStore`, plus the
operator release, lease exclusivity, a crash between appends, and blob
interoperation with Fabro's `BlobStore`.
- `interview.rs` runs human gates through the engine assembly with the
interview adapter over a control interviewer: a gate answered under the
posted id, two parallel gates each bound to their own answer, an expiry
with the gate's default, an auto-approved run, and a cancelled run.
- `secrets.rs` resolves a `{{ secrets.NAME }}` reference from a vault into
a command's environment over `SqliteRunStore` and checks the value is in
no `petri_records` row while the masked output is.
- `blobs.rs` offloads a command's large output to the `blobs` table and
reads it back by the `blob://sha256/<hex>` reference a record carries.
- `model.rs` runs the `hello` bundle against the OpenAI twin with a model
client over a vault that holds the key, and checks the skills step
searched the configured Fabro home.
Those four need the host plugin like `runs.rs` does, and `model.rs` also
starts the twin.
The conformance suite over `HttpRunStore` needs a server to talk to, so it
lives with the server's integration tests
@ -99,13 +133,17 @@ The server's end-to-end coverage is `lib/apps/fabro-server/tests/it/scenario/pet
the `hello` bundle on the OpenAI twin and a command-only bundle run to
completion through the create handler and the scheduler, in the server
process under its test override, under the version flag and under the
server setting, and Petri's diagnostics refuse a run at create. The
server's `petri_runs` unit tests cover the lease ending at worker exit and
the restart reconcile that relaunches a worker in resume mode.
server setting, a human gate is answered through the questions API, and
Petri's diagnostics refuse a run at create. The server's `petri_runs` unit
tests cover the lease ending at worker exit and the restart reconcile that
relaunches a worker in resume mode.
The worker path is covered with the real binary in
`lib/apps/fabro-cli/tests/it/scenario/petri.rs`: a command-only Petri run
executes in the worker a foreground server launched, its records reach
`petri_records` over the HTTP store and its lease ends with the worker; and
a run whose server and worker are both killed mid-stage resumes in a new
worker after the server restarts, with one `run.completed`.
worker after the server restarts, with one `run.completed`; a human gate in
the worker is answered through the questions API over the control channel;
two parallel gates each bind their own answer; and an unanswered gate
expires with its default.

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"));
}
}

View file

@ -0,0 +1,111 @@
//! A large stage value leaves the run's records for Fabro's blob table
//! under `blob://sha256/<hex>`, and comes back from the same table.
//!
//! The run takes its host scope through the sandbox-driver host plugin, so
//! the test skips, and says why, when the executable is not found, unless
//! `FABRO_REQUIRE_SANDBOX_PLUGINS` is set.
mod support;
use std::sync::Arc;
use fabro_petri::SqliteRunStore;
use fabro_petri::blobs::Blobs;
use fabro_petri::check::Launch;
use fabro_petri::engine::{self, RunStatus};
use fabro_petri::runtime::RuntimeSpec;
use fabro_store::{BlobStore, test_support};
use fabro_types::BlobHash;
use petri_attractor_steps::blobs::{BLOB_REF_PREFIX, OFFLOAD_THRESHOLD, parse_blob_ref};
use support::{SETTINGS, Silent, admit, all_records, host_plugin, no_questions, run_request};
/// One line of the command's output.
const LINE: &str = "xxxxxxxx";
/// The command prints `lines` lines, more than the offload threshold in
/// all.
fn workflow(lines: usize) -> String {
format!(
r#"digraph Big {{
graph [goal="Print a lot"]
start [shape=Mdiamond]
exit [shape=Msquare]
say [shape=parallelogram, script="yes {LINE} | head -n {lines}"]
start -> say -> exit
}}"#
)
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_large_output_round_trips_through_the_blob_table() {
if host_plugin().is_none() {
return;
}
let root = tempfile::tempdir().expect("a temp dir");
let pool = test_support::in_memory_pool_with(&[
fabro_db::BLOBS_MIGRATION_SQL,
fabro_db::PETRI_RECORDS_MIGRATION_SQL,
]);
let store = Arc::new(SqliteRunStore::new(pool.clone()));
let blobs = Arc::new(BlobStore::new(pool.clone()));
let lines = OFFLOAD_THRESHOLD / (LINE.len() + 1) + 512;
let expected = format!("{LINE}\n").repeat(lines);
let workflow = workflow(lines);
let runtime = RuntimeSpec::default();
let graphs = admit(
&[("workflow.fabro", &workflow), ("workflow.toml", SETTINGS)],
Launch::default(),
&runtime,
);
let mut request = run_request(
"big",
&root.path().join("run"),
graphs,
store.clone(),
runtime,
no_questions(Arc::new(Silent)),
);
request.blobs = Some(blobs.clone());
let outcome = engine::run(request).await.expect("the run ends");
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
let records = all_records(store.as_ref(), "big").await;
let rendered: Vec<String> = records.iter().map(ToString::to_string).collect();
let inline = serde_json::to_string(&expected).expect("encodes");
let inline = inline.trim_matches('"');
assert!(
rendered.iter().all(|record| !record.contains(inline)),
"the output stayed inline in a record"
);
let reference = rendered
.iter()
.find_map(|record| {
let start = record.find(BLOB_REF_PREFIX)?;
let tail = &record[start..];
let end = tail.find(['"', '#']).unwrap_or(tail.len());
Some(tail[..end].to_string())
})
.expect("a record carries the reference");
let digest = parse_blob_ref(&reference).expect("a well-formed reference");
let hash: BlobHash = digest.parse().expect("a blob hash");
let bytes = Blobs::read(blobs.as_ref(), &hash)
.await
.expect("the table reads")
.expect("the blob is in the table");
assert_eq!(
String::from_utf8(bytes.to_vec()).expect("text"),
expected,
"the blob is the output byte for byte"
);
assert_eq!(BlobHash::new(&bytes), hash, "content-addressed");
let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM blobs")
.fetch_one(&pool)
.await
.expect("the blob table counts");
assert!(count >= 1, "the blob is a row of Fabro's table");
// The run directory's own store was not used: nothing under it holds
// the digest.
let local = root.path().join("run").join("blobs").join(digest);
assert!(!local.exists(), "the local store was bypassed");
}

View file

@ -0,0 +1,470 @@
//! Petri's human gates through Fabro's interview adapter: a question is
//! posted as Fabro's `interview.started`, the answer submitted to the
//! control interviewer under the posted id reaches the gate, two parallel
//! gates each get their own answer, an expired question is completed as a
//! timeout with the gate's default, an auto-approved run answers itself,
//! and a cancelled run interrupts its question.
//!
//! Every run takes its host scope through the sandbox-driver host plugin,
//! so the tests skip, and say why, when the executable is not found,
//! unless `FABRO_REQUIRE_SANDBOX_PLUGINS` is set.
mod support;
use std::path::Path;
use std::sync::{Arc, Mutex};
use fabro_interview::{Answer as LegacyAnswer, ControlInterviewer};
use fabro_petri::check::Launch;
use fabro_petri::engine::{self, RunStatus};
use fabro_petri::interview::{
Approval, AskedQuestion, FabroInterviewer, QuestionNotice, QuestionSink, engine_submission,
};
use fabro_petri::runtime::RuntimeSpec;
use fabro_types::{Principal, QuestionType, SystemActorKind};
use petri_execution::{Delivery, InterviewReceipt, RECEIPT_FILE, ReplyRecord};
use petri_store::MemoryRunStore;
use support::{SETTINGS, admit, all_records, host_plugin, run_request, wait_until};
use tokio::fs;
/// A board of every notice the adapter posted.
#[derive(Default)]
struct Board {
notices: Mutex<Vec<QuestionNotice>>,
}
impl Board {
fn notices(&self) -> Vec<QuestionNotice> {
self.notices.lock().expect("not poisoned").clone()
}
/// Whether a notice other than `Asked` names `question_id`.
fn ended(&self, question_id: &str) -> bool {
self.notices().iter().any(|notice| match notice {
QuestionNotice::Asked(_) => false,
QuestionNotice::Answered {
question_id: id, ..
}
| QuestionNotice::Expired {
question_id: id, ..
}
| QuestionNotice::Interrupted {
question_id: id, ..
} => id == question_id,
})
}
/// The notices once the end of `question_id` is posted, which lands on
/// a task of its own.
async fn wait_ended(&self, question_id: &str) -> Vec<QuestionNotice> {
wait_until(&format!("`{question_id}` to end"), || {
self.ended(question_id)
})
.await;
self.notices()
}
fn asked(&self, stage: &str) -> Option<AskedQuestion> {
self.notices().into_iter().find_map(|notice| match notice {
QuestionNotice::Asked(asked) if asked.stage == stage => Some(asked),
_ => None,
})
}
/// The question `stage` asked, once it is posted.
async fn wait_asked(&self, stage: &str) -> AskedQuestion {
wait_until(&format!("`{stage}` to ask"), || self.asked(stage).is_some()).await;
self.asked(stage).expect("asked")
}
}
#[async_trait::async_trait]
impl QuestionSink for Board {
async fn post(&self, notice: QuestionNotice) -> anyhow::Result<()> {
self.notices.lock().expect("not poisoned").push(notice);
Ok(())
}
}
/// One yes/no gate whose branches leave a marker file each.
fn one_gate(markers: &Path, gate_attrs: &str) -> String {
format!(
r#"digraph G {{
start [shape=Mdiamond]
exit [shape=Msquare]
gate [shape=hexagon, label="Go?", question_type="yes_no"{gate_attrs}]
yes [shape=parallelogram, script="touch {dir}/yes"]
no [shape=parallelogram, script="touch {dir}/no"]
start -> gate
gate -> yes [label="[Y] Yes"]
gate -> no [label="[N] No"]
yes -> exit
no -> exit
}}"#,
dir = markers.display()
)
}
/// Two gates as the branches of one parallel node; the join's results are
/// written out, so each gate's answer is read from its branch result.
fn two_gates(markers: &Path) -> String {
format!(
r#"digraph G {{
start [shape=Mdiamond]
exit [shape=Msquare]
fan [shape=component]
a [shape=hexagon, label="A?", question_type="yes_no"]
b [shape=hexagon, label="B?", question_type="yes_no"]
join [shape=tripleoctagon]
report [shape=parallelogram, script="cat > {dir}/results.json", stdin_source="context.parallel.results"]
start -> fan
fan -> a
fan -> b
a -> join [label="[Y] Yes"]
a -> join [label="[N] No"]
b -> join [label="[Y] Yes"]
b -> join [label="[N] No"]
join -> report -> exit
}}"#,
dir = markers.display()
)
}
struct Gate {
_root: tempfile::TempDir,
markers: std::path::PathBuf,
run_dir: std::path::PathBuf,
store: Arc<MemoryRunStore>,
control: Arc<ControlInterviewer>,
board: Arc<Board>,
}
impl Gate {
fn new() -> Self {
let root = tempfile::tempdir().expect("a temp dir");
let markers = root.path().join("markers");
std::fs::create_dir_all(&markers).expect("the marker dir creates");
Self {
run_dir: root.path().join("run"),
markers,
_root: root,
store: Arc::new(MemoryRunStore::new()),
control: Arc::new(ControlInterviewer::new()),
board: Arc::new(Board::default()),
}
}
fn interviewer(&self, approval: Approval) -> FabroInterviewer {
FabroInterviewer::new(Arc::clone(&self.control), self.board.clone(), approval)
}
fn marker(&self, name: &str) -> bool {
self.markers.join(name).exists()
}
async fn receipt(&self) -> InterviewReceipt {
let text = fs::read_to_string(self.run_dir.join(RECEIPT_FILE))
.await
.expect("the receipt was written");
serde_json::from_str(&text).expect("the receipt parses")
}
}
/// The question is posted with Fabro's type, options and stage; the answer
/// submitted under the posted id, as the API delivers it, routes the gate.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_gate_answered_under_the_posted_id_routes_on_the_answer() {
if host_plugin().is_none() {
return;
}
let gate = Gate::new();
let workflow = one_gate(&gate.markers, "");
let runtime = RuntimeSpec::default();
let graphs = admit(
&[("workflow.fabro", &workflow), ("workflow.toml", SETTINGS)],
Launch::default(),
&runtime,
);
let request = run_request(
"gate",
&gate.run_dir,
graphs,
gate.store.clone(),
runtime,
gate.interviewer(Approval::Prompt),
);
let answer = {
let board = gate.board.clone();
let control = gate.control.clone();
tokio::spawn(async move {
let asked = board.wait_asked("gate").await;
control
.submit(&asked.question_id, engine_submission(LegacyAnswer::no()))
.await
.expect("the answer is accepted");
asked
})
};
let outcome = engine::run(request).await.expect("the run ends");
let asked = answer.await.expect("the answer task ends");
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
assert!(
gate.marker("no") && !gate.marker("yes"),
"the no branch ran"
);
assert_eq!(asked.stage, "gate");
assert_eq!(asked.text, "Go?");
assert_eq!(asked.question_type, QuestionType::YesNo);
assert_eq!(
asked
.options
.iter()
.map(|option| (option.key.as_str(), option.label.as_str()))
.collect::<Vec<_>>(),
vec![("Y", "[Y] Yes"), ("N", "[N] No")]
);
assert!(
asked.question_id.starts_with("gate.x0.f"),
"{}",
asked.question_id
);
assert_eq!(asked.identity.node, "gate");
assert_eq!(asked.identity.invocation_path, "/");
let notices = gate.board.notices();
assert!(
matches!(
&notices[1],
QuestionNotice::Answered { question_id, answer, actor: Principal::System { system_kind: SystemActorKind::Engine }, .. }
if *question_id == asked.question_id && answer == "N"
),
"{notices:?}"
);
let receipt = gate.receipt().await;
assert!(receipt.is_clean(), "{:?}", receipt.errors);
assert_eq!(receipt.questions.len(), 1);
assert_eq!(receipt.questions[0].delivery, Delivery::Delivered);
assert_eq!(receipt.questions[0].reply, ReplyRecord::Answered {
choice: Some("N".to_string()),
choices: Vec::new(),
text: None,
});
}
/// Two branches ask at once; each answer, submitted under its own id in
/// the other order, lands on its own branch.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn two_parallel_gates_each_bind_their_own_answer() {
if host_plugin().is_none() {
return;
}
let gate = Gate::new();
let workflow = two_gates(&gate.markers);
let runtime = RuntimeSpec::default();
let graphs = admit(
&[("workflow.fabro", &workflow), ("workflow.toml", SETTINGS)],
Launch::default(),
&runtime,
);
let request = run_request(
"gates",
&gate.run_dir,
graphs,
gate.store.clone(),
runtime,
gate.interviewer(Approval::Prompt),
);
let answers = {
let board = gate.board.clone();
let control = gate.control.clone();
tokio::spawn(async move {
// Both are pending before either is answered, and `b` first.
let a = board.wait_asked("a").await;
let b = board.wait_asked("b").await;
assert_ne!(a.question_id, b.question_id);
control
.submit(&b.question_id, engine_submission(LegacyAnswer::yes()))
.await
.expect("b's answer is accepted");
control
.submit(&a.question_id, engine_submission(LegacyAnswer::no()))
.await
.expect("a's answer is accepted");
(a, b)
})
};
let outcome = engine::run(request).await.expect("the run ends");
let (a, b) = answers.await.expect("the answer task ends");
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
let results: serde_json::Value = serde_json::from_str(
&fs::read_to_string(gate.markers.join("results.json"))
.await
.expect("the join wrote its results"),
)
.expect("the results parse");
let results = results.as_array().expect("a list of branch results");
assert_eq!(results.len(), 2, "{results:?}");
assert_eq!(results[0]["id"], "a");
assert_eq!(results[0]["context_updates"]["human.gate.selected"], "N");
assert_eq!(results[1]["id"], "b");
assert_eq!(results[1]["context_updates"]["human.gate.selected"], "Y");
assert!(
a.identity.invocation_path.starts_with("/branch:"),
"{}",
a.identity.invocation_path
);
assert_ne!(a.identity.invocation_path, b.identity.invocation_path);
let receipt = gate.receipt().await;
assert!(receipt.is_clean(), "{:?}", receipt.errors);
assert_eq!(receipt.questions.len(), 2);
}
/// The gate's deadline passes with no answer: the adapter completes the
/// question as a timeout, the receipt says the gate took its default, and
/// the default's branch runs.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn an_unanswered_question_expires_with_the_gates_default() {
if host_plugin().is_none() {
return;
}
let gate = Gate::new();
let workflow = one_gate(
&gate.markers,
r#", timeout="300ms", human.default_choice="no""#,
);
let runtime = RuntimeSpec::default();
let graphs = admit(
&[("workflow.fabro", &workflow), ("workflow.toml", SETTINGS)],
Launch::default(),
&runtime,
);
let request = run_request(
"expiry",
&gate.run_dir,
graphs,
gate.store.clone(),
runtime,
gate.interviewer(Approval::Prompt),
);
let outcome = engine::run(request).await.expect("the run ends");
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
assert!(gate.marker("no") && !gate.marker("yes"), "the default ran");
let asked = gate.board.asked("gate").expect("asked");
let notices = gate.board.wait_ended(&asked.question_id).await;
assert_eq!(asked.timeout_seconds, Some(0.3));
assert!(
matches!(
&notices[1],
QuestionNotice::Expired { question_id, stage, .. }
if *question_id == asked.question_id && stage == "gate"
),
"{notices:?}"
);
let receipt = gate.receipt().await;
assert!(receipt.is_clean(), "{:?}", receipt.errors);
assert_eq!(receipt.questions[0].reply, ReplyRecord::TimedOut {
default: Some("N".to_string()),
});
assert_eq!(receipt.questions[0].delivery, Delivery::Expired);
}
/// An auto-approved run answers its gate at once, attributed to the
/// engine, and still posts the question and its answer.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn an_auto_approved_run_answers_yes_at_once() {
if host_plugin().is_none() {
return;
}
let gate = Gate::new();
let workflow = one_gate(&gate.markers, "");
let runtime = RuntimeSpec::default();
let graphs = admit(
&[("workflow.fabro", &workflow), ("workflow.toml", SETTINGS)],
Launch::default(),
&runtime,
);
let request = run_request(
"auto",
&gate.run_dir,
graphs,
gate.store.clone(),
runtime,
gate.interviewer(Approval::Auto),
);
let outcome = engine::run(request).await.expect("the run ends");
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
assert!(
gate.marker("yes") && !gate.marker("no"),
"the yes branch ran"
);
let notices = gate.board.notices();
assert_eq!(notices.len(), 2, "{notices:?}");
assert!(
matches!(
&notices[1],
QuestionNotice::Answered { answer, actor: Principal::System { system_kind: SystemActorKind::Engine }, .. }
if answer == "Y"
),
"{notices:?}"
);
}
/// A run cancelled while its gate waits: the adapter returns promptly, the
/// question is interrupted, the gate fails closed and the run is cancelled.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_cancelled_run_interrupts_its_pending_question() {
if host_plugin().is_none() {
return;
}
let gate = Gate::new();
let workflow = one_gate(&gate.markers, "");
let runtime = RuntimeSpec::default();
let graphs = admit(
&[("workflow.fabro", &workflow), ("workflow.toml", SETTINGS)],
Launch::default(),
&runtime,
);
let request = run_request(
"cancel",
&gate.run_dir,
graphs,
gate.store.clone(),
runtime,
gate.interviewer(Approval::Prompt),
);
let cancel = request.cancel.clone();
let canceller = {
let board = gate.board.clone();
tokio::spawn(async move {
board.wait_asked("gate").await;
cancel.cancel();
})
};
let outcome = engine::run(request).await.expect("the run ends");
canceller.await.expect("the cancel task ends");
assert_eq!(outcome.status, RunStatus::Cancelled, "{outcome:?}");
assert!(!gate.marker("yes") && !gate.marker("no"), "no branch ran");
let asked = gate.board.asked("gate").expect("asked");
let notices = gate.board.wait_ended(&asked.question_id).await;
assert!(
matches!(
&notices[1],
QuestionNotice::Interrupted { reason, .. } if reason == "cancelled"
),
"{notices:?}"
);
let receipt = gate.receipt().await;
assert_eq!(receipt.questions[0].reply, ReplyRecord::Cancelled);
// Nothing the adapter posted names the answer a person never gave.
let records = all_records(gate.store.as_ref(), "cancel").await;
assert!(!records.is_empty());
}

View file

@ -0,0 +1,113 @@
//! A model call from a Petri run authenticates through Fabro's vault, and
//! the skills step reads the Fabro home the runtime was given.
//!
//! The `hello` bundle's agent stage calls the OpenAI twin through a model
//! client built over a vault that holds the key; the twin requires a
//! bearer token and logs requests under it, so a request logged under the
//! vault's key proves the key came from the vault. The run takes its host
//! scope through the sandbox-driver host plugin, so the test skips, and
//! says why, when the executable is not found, unless
//! `FABRO_REQUIRE_SANDBOX_PLUGINS` is set.
mod support;
use std::collections::HashMap;
use std::sync::Arc;
use fabro_auth::VaultCredentialSource;
use fabro_llm::test_support::test_catalog_with_provider_base_url;
use fabro_petri::check::Launch;
use fabro_petri::engine::{self, RunStatus};
use fabro_petri::runtime::{self, RuntimeSpec};
use fabro_test::{TwinScenario, TwinScenarios, twin_openai};
use fabro_types::SecretType;
use fabro_vault::Vault;
use lithos_llm::catalog::ProviderId;
use petri_store::MemoryRunStore;
use support::{Silent, all_records, hello_bundle, host_plugin, no_questions, run_request};
use tokio::fs;
use tokio::sync::RwLock as AsyncRwLock;
const OPENAI_MODEL: &str = "gpt-5.4";
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_model_call_authenticates_through_the_vault_and_skills_read_the_home() {
if host_plugin().is_none() {
return;
}
let twin = twin_openai().await;
let namespace = format!("{}::{}", module_path!(), line!());
TwinScenarios::new(&namespace)
.scenario(TwinScenario::responses(OPENAI_MODEL).text("A haiku, added."))
.load(twin)
.await;
let root = tempfile::tempdir().expect("a temp dir");
let home = root.path().join("fabro-home");
std::fs::create_dir_all(home.join("skills")).expect("the skills dir creates");
// The vault holds the key; nothing in the environment does.
let mut vault = Vault::from_entries(HashMap::new());
vault
.set("OPENAI_API_KEY", &namespace, SecretType::Token, None)
.expect("a detached vault takes an entry");
let credentials = Arc::new(VaultCredentialSource::vault_only(Arc::new(
AsyncRwLock::new(vault),
)));
let catalog = test_catalog_with_provider_base_url("openai", &twin.base_url);
let client = runtime::model_client(catalog, credentials, None, &[ProviderId::new("openai")])
.expect("the model client builds")
.expect("openai is eligible");
let runtime = RuntimeSpec {
model_client: Some(client),
fabro_home: Some(home.clone()),
..RuntimeSpec::default()
};
let workflow = fs::read_to_string(hello_bundle().join("workflow.fabro"))
.await
.expect("the hello workflow is checked in");
let settings = fs::read_to_string(hello_bundle().join("workflow.toml"))
.await
.expect("the hello settings are checked in");
let graphs = support::admit(
&[("workflow.fabro", &workflow), ("workflow.toml", &settings)],
Launch {
model: Some(OPENAI_MODEL.to_string()),
..Launch::default()
},
&runtime,
);
let store = Arc::new(MemoryRunStore::new());
let request = run_request(
"hello",
&root.path().join("run"),
graphs,
store.clone(),
runtime,
no_questions(Arc::new(Silent)),
);
let outcome = engine::run(request).await.expect("the run ends");
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
let logs = twin.request_logs(&namespace).await;
let requests = logs["requests"]
.as_array()
.expect("twin request logs are an array");
assert!(
requests
.iter()
.any(|request| request["model"] == OPENAI_MODEL),
"the stage should have called the twin with the vault's key, got {logs}"
);
let records = all_records(store.as_ref(), "hello").await;
let resolved = records
.iter()
.find(|record| record.to_string().contains("\"attractor.skills\""))
.unwrap_or_else(|| panic!("the skills step recorded what it searched: {records:?}"));
let configured = home.join("skills").display().to_string();
assert!(
resolved.to_string().contains(&configured),
"the configured home is searched: {resolved}"
);
}

View file

@ -0,0 +1,138 @@
//! A `{{ secrets.NAME }}` reference resolves from the vault into a
//! command's environment, and the value never reaches `petri_records`:
//! Petri masks every record before it is appended.
//!
//! The run takes its host scope through the sandbox-driver host plugin, so
//! the test skips, and says why, when the executable is not found, unless
//! `FABRO_REQUIRE_SANDBOX_PLUGINS` is set.
mod support;
use std::collections::HashMap;
use std::sync::Arc;
use fabro_petri::SqliteRunStore;
use fabro_petri::check::Launch;
use fabro_petri::engine::{self, RunStatus};
use fabro_petri::runtime::RuntimeSpec;
use fabro_petri::secrets::VaultSecrets;
use fabro_store::test_support;
use fabro_types::SecretType;
use fabro_vault::Vault;
use support::{Silent, admit, host_plugin, no_questions, run_request};
const TOKEN: &str = "hunter2-hunter2-hunter2";
/// A command that checks the secret reached its environment and then
/// prints it, so the value would land in a log line if nothing masked it.
const WORKFLOW: &str = r#"digraph Secret {
graph [goal="Use a secret"]
start [shape=Mdiamond]
exit [shape=Msquare]
say [shape=parallelogram, script="test \"$TOKEN\" = hunter2-hunter2-hunter2 && echo \"token is $TOKEN\""]
start -> say -> exit
}"#;
const SETTINGS: &str = r#"_version = 1
[workflow]
graph = "workflow.fabro"
[run.environment]
id = "local"
[environments.local]
provider = "local"
[environments.local.env]
TOKEN = "{{ secrets.TOKEN }}"
"#;
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_secret_reaches_the_command_and_is_masked_in_every_record() {
if host_plugin().is_none() {
return;
}
let root = tempfile::tempdir().expect("a temp dir");
let pool = test_support::in_memory_pool_with(&[
fabro_db::BLOBS_MIGRATION_SQL,
fabro_db::PETRI_RECORDS_MIGRATION_SQL,
]);
let store = Arc::new(SqliteRunStore::new(pool.clone()));
let mut vault = Vault::from_entries(HashMap::new());
vault
.set("TOKEN", TOKEN, SecretType::Token, None)
.expect("a detached vault takes an entry");
let runtime = RuntimeSpec::default();
let graphs = admit(
&[("workflow.fabro", WORKFLOW), ("workflow.toml", SETTINGS)],
Launch::default(),
&runtime,
);
let mut request = run_request(
"secret",
&root.path().join("run"),
graphs,
store.clone(),
runtime,
no_questions(Arc::new(Silent)),
);
request.secrets = Some(Arc::new(VaultSecrets::from_vault(&vault)));
let outcome = engine::run(request).await.expect("the run ends");
assert_eq!(
outcome.status,
RunStatus::Success,
"the command saw the secret: {outcome:?}"
);
let records: Vec<String> = sqlx::query_scalar("SELECT record_json FROM petri_records")
.fetch_all(&pool)
.await
.expect("the records read");
assert!(!records.is_empty());
assert!(
records.iter().all(|record| !record.contains(TOKEN)),
"the secret's value is in a record"
);
assert!(
records.iter().any(|record| record.contains("token is ***")),
"the command's output was masked, not dropped"
);
}
/// Without a provider the reference resolves to nothing and the command
/// fails on the missing secret, as the standalone runner's does; the run
/// ends the way Fabro's failure policy for a command ends it.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_secret_nobody_provides_fails_the_command() {
if host_plugin().is_none() {
return;
}
let root = tempfile::tempdir().expect("a temp dir");
let store = Arc::new(petri_store::MemoryRunStore::new());
let runtime = RuntimeSpec::default();
let graphs = admit(
&[("workflow.fabro", WORKFLOW), ("workflow.toml", SETTINGS)],
Launch::default(),
&runtime,
);
let request = run_request(
"unprovided",
&root.path().join("run"),
graphs,
store,
runtime,
no_questions(Arc::new(Silent)),
);
let outcome = engine::run(request).await.expect("the run ends");
assert!(
outcome
.failure
.as_deref()
.is_some_and(|failure| failure.contains("no secret named `TOKEN`")),
"{outcome:?}"
);
}

View file

@ -0,0 +1,178 @@
//! What the adapter tests share: the host plugin lookup, a bundle admitted
//! through `check`, a run request over the engine assembly, and the run's
//! records read back from its store.
#![allow(
dead_code,
reason = "each test file uses the part of the support it needs"
)]
use std::collections::BTreeMap;
use std::env;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::{Duration, Instant};
use fabro_petri::admission::AdmittedGraphs;
use fabro_petri::check::{self, Bundle, CheckRequest, Launch};
use fabro_petri::engine::{Execution, RunRequest};
use fabro_petri::interview::{Approval, FabroInterviewer, QuestionNotice, QuestionSink};
use fabro_petri::runtime::RuntimeSpec;
use fabro_types::SandboxProviderKind;
use petri_execution::inspect;
use petri_store::{Access, LogId, RunKey, RunStore};
use tokio::time::sleep;
use tokio_util::sync::CancellationToken;
const HOST_PLUGIN: &str = "sandbox-driver-host";
const HOST_PLUGIN_OVERRIDE: &str = "PETRI_SANDBOX_HOST_PLUGIN";
const REQUIRE_ENV: &str = "FABRO_REQUIRE_SANDBOX_PLUGINS";
pub(crate) const POLL: Duration = Duration::from_millis(10);
pub(crate) const PATIENCE: Duration = Duration::from_secs(30);
/// The host plugin as Petri's lookup finds it: the override variable, else
/// the executable on `PATH`. `None`, after saying so, when the test should
/// skip; a panic when the environment forbids a skip.
#[expect(
clippy::disallowed_methods,
reason = "the tests locate the plugin executable through the process environment"
)]
#[expect(clippy::print_stderr, reason = "a skipped test says why on its stderr")]
pub(crate) fn host_plugin() -> Option<PathBuf> {
let found = env::var_os(HOST_PLUGIN_OVERRIDE)
.map(PathBuf::from)
.or_else(|| {
env::split_paths(&env::var_os("PATH")?)
.map(|dir| dir.join(HOST_PLUGIN))
.find(|candidate| candidate.is_file())
});
if found.is_none() {
assert!(
env::var_os(REQUIRE_ENV).is_none(),
"{REQUIRE_ENV} is set, but {HOST_PLUGIN} is not on PATH and {HOST_PLUGIN_OVERRIDE} is unset"
);
eprintln!("skipping: {HOST_PLUGIN} is not on PATH and {HOST_PLUGIN_OVERRIDE} is unset");
}
found
}
/// The `.fabro/workflows/hello` bundle checked into this repository.
pub(crate) fn hello_bundle() -> PathBuf {
Path::new(env!("CARGO_MANIFEST_DIR")).join("../../../.fabro/workflows/hello")
}
pub(crate) const SETTINGS: &str = "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\n";
pub(crate) fn bundle(files: &[(&str, &str)]) -> Bundle {
Bundle {
files: files
.iter()
.map(|(path, text)| ((*path).to_string(), (*text).to_string()))
.collect(),
entrypoint: "workflow.fabro".to_string(),
project_toml: None,
}
}
/// Admit a bundle as the create handler does, with the given launch.
pub(crate) fn admit(
files: &[(&str, &str)],
launch: Launch,
runtime: &RuntimeSpec,
) -> AdmittedGraphs {
let request = CheckRequest {
bundle: bundle(files),
inputs: BTreeMap::new(),
launch,
runtime: runtime.clone(),
};
let admitted = check::check(&request)
.unwrap_or_else(|error| panic!("the workflow is admitted: {error:?}"));
AdmittedGraphs {
graph: admitted.graph,
children: admitted.children,
}
}
/// A run request over the engine assembly, on the host sandbox, with a
/// fresh cancel token and nothing installed beyond the interviewer.
pub(crate) fn run_request(
run_id: &str,
run_dir: &Path,
graphs: AdmittedGraphs,
store: Arc<dyn RunStore>,
runtime: RuntimeSpec,
interviewer: FabroInterviewer,
) -> RunRequest {
RunRequest {
run_id: run_id.to_string(),
run_dir: run_dir.to_path_buf(),
execution: Execution::Start(graphs),
store,
runtime,
provider: SandboxProviderKind::LOCAL,
cancel: CancellationToken::new(),
observers: vec![interviewer.observer()],
interviewer: Arc::new(interviewer),
secrets: None,
blobs: None,
}
}
/// An interviewer whose answers nobody delivers, for runs that ask nothing.
pub(crate) fn no_questions(sink: Arc<dyn QuestionSink>) -> FabroInterviewer {
FabroInterviewer::new(
Arc::new(fabro_interview::ControlInterviewer::new()),
sink,
Approval::Prompt,
)
}
/// A sink that drops every notice.
pub(crate) struct Silent;
#[async_trait::async_trait]
impl QuestionSink for Silent {
async fn post(&self, _notice: QuestionNotice) -> anyhow::Result<()> {
Ok(())
}
}
/// Every record of every log of a stored run, as JSON, in log order.
pub(crate) async fn all_records(store: &dyn RunStore, run_id: &str) -> Vec<serde_json::Value> {
let logs = store
.open(&RunKey::new(run_id), Access::Read)
.await
.expect("the run opens for reading");
let inspection = inspect::inspect_run(&*logs)
.await
.expect("the stored run inspects");
let mut ids = vec![LogId::Coordinator, LogId::Resources];
ids.extend(
inspection
.executions
.iter()
.map(|execution| LogId::Execution(execution.execution)),
);
let mut records = Vec::new();
for id in ids {
records.extend(
logs.read(&id)
.await
.expect("the log reads")
.into_iter()
.map(|record| record.record),
);
}
records
}
/// Wait until `condition` holds, polling, or fail after [`PATIENCE`].
pub(crate) async fn wait_until(what: &str, mut condition: impl FnMut() -> bool) {
let deadline = Instant::now() + PATIENCE;
while !condition() {
assert!(Instant::now() < deadline, "timed out waiting for {what}");
sleep(POLL).await;
}
}