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/<hex>`; 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 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-17 22:31:35 -04:00
parent 223e10ea20
commit 3aadc21739
No known key found for this signature in database
6 changed files with 1056 additions and 8 deletions

View file

@ -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/<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
@ -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/<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
@ -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.

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