From 3aadc21739d2678a610d670e0b7c07d153b99b4d Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 17 Sep 2026 22:31:35 -0400 Subject: [PATCH] Test the interview, secret, blob and home adapters through the engine Integration tests in `fabro-petri` run workflows through `engine::run` on the host sandbox: a gate answered under the posted question id, two parallel gates each bound to their own answer, an expired question completed as a timeout with the gate's default, an auto-approved run, a cancelled run; a secret resolved from a vault into a command and masked in every `petri_records` row; a command's large output round-tripped through the `blobs` table under `blob://sha256/`; and the `hello` bundle on the OpenAI twin with a model client over a vault that holds the key, whose skills step searched the configured home. Co-Authored-By: Claude Fable 5.1 --- lib/components/fabro-petri/README.md | 54 +- lib/components/fabro-petri/tests/blobs.rs | 111 +++++ lib/components/fabro-petri/tests/interview.rs | 470 ++++++++++++++++++ lib/components/fabro-petri/tests/model.rs | 113 +++++ lib/components/fabro-petri/tests/secrets.rs | 138 +++++ .../fabro-petri/tests/support/mod.rs | 178 +++++++ 6 files changed, 1056 insertions(+), 8 deletions(-) create mode 100644 lib/components/fabro-petri/tests/blobs.rs create mode 100644 lib/components/fabro-petri/tests/interview.rs create mode 100644 lib/components/fabro-petri/tests/model.rs create mode 100644 lib/components/fabro-petri/tests/secrets.rs create mode 100644 lib/components/fabro-petri/tests/support/mod.rs diff --git a/lib/components/fabro-petri/README.md b/lib/components/fabro-petri/README.md index 9b8134b11..6ed4fa1ca 100644 --- a/lib/components/fabro-petri/README.md +++ b/lib/components/fabro-petri/README.md @@ -34,8 +34,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/`. - `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 @@ -46,8 +65,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 @@ -78,6 +97,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/` 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 @@ -94,13 +128,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. diff --git a/lib/components/fabro-petri/tests/blobs.rs b/lib/components/fabro-petri/tests/blobs.rs new file mode 100644 index 000000000..bad4600db --- /dev/null +++ b/lib/components/fabro-petri/tests/blobs.rs @@ -0,0 +1,111 @@ +//! A large stage value leaves the run's records for Fabro's blob table +//! under `blob://sha256/`, 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 = 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"); +} diff --git a/lib/components/fabro-petri/tests/interview.rs b/lib/components/fabro-petri/tests/interview.rs new file mode 100644 index 000000000..7415a9e19 --- /dev/null +++ b/lib/components/fabro-petri/tests/interview.rs @@ -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>, +} + +impl Board { + fn notices(&self) -> Vec { + 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 { + wait_until(&format!("`{question_id}` to end"), || { + self.ended(question_id) + }) + .await; + self.notices() + } + + fn asked(&self, stage: &str) -> Option { + 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, + control: Arc, + board: Arc, +} + +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![("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!( + ¬ices[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!( + ¬ices[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!( + ¬ices[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!( + ¬ices[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()); +} diff --git a/lib/components/fabro-petri/tests/model.rs b/lib/components/fabro-petri/tests/model.rs new file mode 100644 index 000000000..34db3b65b --- /dev/null +++ b/lib/components/fabro-petri/tests/model.rs @@ -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}" + ); +} diff --git a/lib/components/fabro-petri/tests/secrets.rs b/lib/components/fabro-petri/tests/secrets.rs new file mode 100644 index 000000000..c799613f9 --- /dev/null +++ b/lib/components/fabro-petri/tests/secrets.rs @@ -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 = 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:?}" + ); +} diff --git a/lib/components/fabro-petri/tests/support/mod.rs b/lib/components/fabro-petri/tests/support/mod.rs new file mode 100644 index 000000000..206ba09da --- /dev/null +++ b/lib/components/fabro-petri/tests/support/mod.rs @@ -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 { + 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, + 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) -> 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 { + 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; + } +}