mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-01 02:04:24 +00:00
Merge branch 'petri-integration-tools' into petri-integration
# Conflicts: # Cargo.lock # lib/apps/fabro-cli/src/commands/run/petri_worker.rs # lib/apps/fabro-cli/tests/it/scenario/petri.rs # lib/apps/fabro-server/src/server/petri_runs.rs # lib/components/fabro-petri/Cargo.toml # lib/components/fabro-petri/src/lib.rs
This commit is contained in:
commit
5d93f20fb0
13 changed files with 1043 additions and 28 deletions
3
Cargo.lock
generated
3
Cargo.lock
generated
|
|
@ -2903,11 +2903,14 @@ dependencies = [
|
|||
"fabro-petri",
|
||||
"fabro-store",
|
||||
"fabro-test",
|
||||
"fabro-tool",
|
||||
"fabro-types",
|
||||
"fabro-util",
|
||||
"fabro-vault",
|
||||
"fabro-workflow",
|
||||
"httpmock",
|
||||
"lithos-llm",
|
||||
"pebble-coding-agent",
|
||||
"petri-attractor-steps",
|
||||
"petri-execution",
|
||||
"petri-frontend-attractor",
|
||||
|
|
|
|||
|
|
@ -125,6 +125,7 @@ fabro-mcp = { path = "../../components/fabro-mcp", features = ["test-support"] }
|
|||
fabro-build-support = { path = "../../foundation/build-support" }
|
||||
fabro-sandbox = { path = "../../components/fabro-sandbox", features = ["test-support"] }
|
||||
fabro-server = { path = "../fabro-server", features = ["test-support"] }
|
||||
fabro-petri = { path = "../../components/fabro-petri", features = ["test-support"] }
|
||||
fabro-workflow = { path = "../../components/fabro-workflow", features = ["test-support"] }
|
||||
fabro-types = { path = "../../foundation/fabro-types", features = ["clap", "test-support"] }
|
||||
insta = { workspace = true, features = ["filters"] }
|
||||
|
|
|
|||
|
|
@ -39,6 +39,11 @@
|
|||
//! 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.
|
||||
//! Fabro's run tools go to every agent session of the run when the run's
|
||||
//! settings enable them (`[run.agent] fabro_tools`) and the worker token
|
||||
//! carries the `agent:run_tools` scope the server issues for such a run,
|
||||
//! the same gate the legacy worker applies; they bind to the worker's
|
||||
//! client and the run id, as the legacy worker binds them.
|
||||
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::Arc;
|
||||
|
|
@ -67,6 +72,7 @@ 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 fabro_workflow::services::FabroRunToolServices;
|
||||
use tokio::sync::RwLock as AsyncRwLock;
|
||||
use tokio_util::sync::CancellationToken;
|
||||
use tracing::{info, warn};
|
||||
|
|
@ -149,7 +155,14 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
|
|||
|
||||
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 run_tools = run_tool_services(&worker);
|
||||
let runtime = runtime_spec(
|
||||
&vault,
|
||||
&worker.run_state,
|
||||
worker.fabro_home.clone(),
|
||||
run_tools,
|
||||
)
|
||||
.await?;
|
||||
let execution = match worker.mode {
|
||||
RunWorkerMode::Start => {
|
||||
let client = worker.client.clone_for_reuse();
|
||||
|
|
@ -281,14 +294,42 @@ fn test_checkpoint_gates() -> Option<PathBuf> {
|
|||
std::env::var_os(EnvVars::FABRO_TEST_CHECKPOINT_GATES).map(PathBuf::from)
|
||||
}
|
||||
|
||||
/// Fabro's run tools for the run's agent sessions, when the run's settings
|
||||
/// enable them and the worker token carries the scope; `None` otherwise.
|
||||
/// The server issues the scope from the same setting, so the two agree
|
||||
/// unless the token was issued for another run.
|
||||
fn run_tool_services(worker: &PetriWorker<'_>) -> Option<FabroRunToolServices> {
|
||||
let enabled = worker.run_state.spec.settings.run.agent.fabro_tools;
|
||||
let scoped = runner::fabro_run_tools_enabled_from_worker_token(worker.worker_token);
|
||||
if !enabled || !scoped {
|
||||
info!(
|
||||
run_id = %worker.run_id,
|
||||
enabled,
|
||||
scoped,
|
||||
"Fabro's run tools are not registered on this Petri run"
|
||||
);
|
||||
return None;
|
||||
}
|
||||
let services = runner::build_fabro_run_tool_services(
|
||||
worker.worker_token,
|
||||
worker.client.clone_for_reuse(),
|
||||
worker.run_id,
|
||||
);
|
||||
if services.is_some() {
|
||||
info!(run_id = %worker.run_id, "Fabro's run tools are registered on this Petri run");
|
||||
}
|
||||
services
|
||||
}
|
||||
|
||||
/// 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, the run's mode, and the Fabro
|
||||
/// home the server named.
|
||||
/// the providers whose credentials resolve, the run's mode, the Fabro
|
||||
/// home the server named, and the run tools when the run has them.
|
||||
async fn runtime_spec(
|
||||
vault: &Arc<AsyncRwLock<Vault>>,
|
||||
run_state: &RunProjection,
|
||||
fabro_home: Option<PathBuf>,
|
||||
run_tools: Option<FabroRunToolServices>,
|
||||
) -> Result<RuntimeSpec> {
|
||||
let catalog =
|
||||
command_context::load_cli_catalog().context("failed to build worker LLM catalog")?;
|
||||
|
|
@ -310,5 +351,6 @@ async fn runtime_spec(
|
|||
model_client,
|
||||
dry_run: run_state.spec.settings.run.execution.mode == RunMode::DryRun,
|
||||
fabro_home,
|
||||
run_tools,
|
||||
})
|
||||
}
|
||||
|
|
|
|||
|
|
@ -214,7 +214,7 @@ struct WorkerTokenScopeClaim {
|
|||
scope: String,
|
||||
}
|
||||
|
||||
fn fabro_run_tools_enabled_from_worker_token(worker_token: &str) -> bool {
|
||||
pub(super) fn fabro_run_tools_enabled_from_worker_token(worker_token: &str) -> bool {
|
||||
// Local tool registration only. The server validates the token signature and
|
||||
// scopes.
|
||||
insecure_decode::<WorkerTokenScopeClaim>(worker_token)
|
||||
|
|
@ -234,7 +234,7 @@ fn worker_scope_has_run_tools(scope_claim: &str) -> bool {
|
|||
has_run_worker && has_agent_run_tools
|
||||
}
|
||||
|
||||
fn build_fabro_run_tool_services(
|
||||
pub(super) fn build_fabro_run_tool_services(
|
||||
worker_token: &str,
|
||||
client: fabro_client::Client,
|
||||
current_run_id: RunId,
|
||||
|
|
|
|||
|
|
@ -9,6 +9,7 @@ mod auth;
|
|||
mod exec;
|
||||
mod lifecycle;
|
||||
mod petri;
|
||||
mod petri_tools;
|
||||
mod server_lifecycle;
|
||||
mod smoke;
|
||||
|
||||
|
|
|
|||
|
|
@ -10,6 +10,9 @@
|
|||
//! executable is not found, unless `FABRO_REQUIRE_SANDBOX_PLUGINS` is set.
|
||||
//! The plugin's path override crosses into the server and its workers the
|
||||
//! way `PATH` does.
|
||||
//!
|
||||
//! The harness here (the server, the detached run, the status and event
|
||||
//! reads) is shared with the run-tools scenarios in `petri_tools.rs`.
|
||||
|
||||
#![expect(
|
||||
clippy::disallowed_methods,
|
||||
|
|
@ -39,6 +42,7 @@ use fabro_test::{
|
|||
apply_test_isolation, expect_reqwest_json, fabro_snapshot, isolated_storage_dir, test_context,
|
||||
};
|
||||
use fabro_types::RunId;
|
||||
use fabro_vault::{SecretType, Vault};
|
||||
|
||||
use crate::cmd::support::created_run_id;
|
||||
use crate::support::{TEST_DEV_TOKEN, TEST_SESSION_SECRET, seed_dev_token_auth};
|
||||
|
|
@ -51,7 +55,7 @@ const POLL: Duration = Duration::from_millis(50);
|
|||
/// 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.
|
||||
fn host_plugin() -> Option<PathBuf> {
|
||||
pub(super) fn host_plugin() -> Option<PathBuf> {
|
||||
let found = env::var_os(EnvVars::PETRI_SANDBOX_HOST_PLUGIN)
|
||||
.map(PathBuf::from)
|
||||
.or_else(|| {
|
||||
|
|
@ -75,20 +79,28 @@ fn host_plugin() -> Option<PathBuf> {
|
|||
|
||||
/// A foreground server on its own disk storage, dev-token auth, started
|
||||
/// from the compiled `fabro` binary. Dropping it kills the process.
|
||||
struct RunningServer {
|
||||
child: Option<Child>,
|
||||
home_root: tempfile::TempDir,
|
||||
_storage_root: tempfile::TempDir,
|
||||
storage_dir: PathBuf,
|
||||
config_path: PathBuf,
|
||||
port: u16,
|
||||
api_base_url: String,
|
||||
pub(super) struct RunningServer {
|
||||
child: Option<Child>,
|
||||
home_root: tempfile::TempDir,
|
||||
_storage_root: tempfile::TempDir,
|
||||
pub(super) storage_dir: PathBuf,
|
||||
config_path: PathBuf,
|
||||
port: u16,
|
||||
pub(super) api_base_url: String,
|
||||
/// The checkpoint gate directory the server forwards to its workers.
|
||||
gates_dir: PathBuf,
|
||||
gates_dir: PathBuf,
|
||||
}
|
||||
|
||||
impl RunningServer {
|
||||
async fn start() -> Self {
|
||||
pub(super) async fn start() -> Self {
|
||||
Self::start_with("", &[]).await
|
||||
}
|
||||
|
||||
/// Start with `settings` appended to the server's settings file (the
|
||||
/// workers read the same file through `FABRO_CONFIG`) and `secrets`
|
||||
/// in the vault before the first launch, so the server and its workers
|
||||
/// see them from the start.
|
||||
pub(super) async fn start_with(settings: &str, secrets: &[(&str, &str)]) -> Self {
|
||||
let home_root = tempfile::tempdir_in("/tmp").expect("home tempdir");
|
||||
let storage_root = isolated_storage_dir();
|
||||
let storage_dir = storage_root.path().join("storage");
|
||||
|
|
@ -96,9 +108,18 @@ impl RunningServer {
|
|||
let config_path = home_root.path().join("settings.toml");
|
||||
std::fs::write(
|
||||
&config_path,
|
||||
"_version = 1\n\n[server.auth]\nmethods = [\"dev-token\"]\n",
|
||||
format!("_version = 1\n\n[server.auth]\nmethods = [\"dev-token\"]\n{settings}"),
|
||||
)
|
||||
.expect("the server settings write");
|
||||
if !secrets.is_empty() {
|
||||
let mut vault = Vault::load(Storage::new(&storage_dir).secrets_path())
|
||||
.expect("the server vault loads");
|
||||
for (name, value) in secrets {
|
||||
vault
|
||||
.set(name, value, SecretType::Token, None)
|
||||
.expect("the secret stores in the server vault");
|
||||
}
|
||||
}
|
||||
let runtime_directory = Storage::new(&storage_dir).runtime_directory();
|
||||
envfile::merge_env_file(&runtime_directory.env_path(), [
|
||||
("SESSION_SECRET", TEST_SESSION_SECRET),
|
||||
|
|
@ -165,13 +186,13 @@ impl RunningServer {
|
|||
Stdio::from(file)
|
||||
}
|
||||
|
||||
fn stderr_text(&self) -> String {
|
||||
pub(super) fn stderr_text(&self) -> String {
|
||||
std::fs::read_to_string(self.storage_dir.with_file_name("server.stderr.log"))
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
/// The `--server` target a CLI command reaches this server at.
|
||||
fn target(&self) -> String {
|
||||
pub(super) fn target(&self) -> String {
|
||||
format!("{}/api/v1", self.api_base_url)
|
||||
}
|
||||
|
||||
|
|
@ -183,7 +204,7 @@ impl RunningServer {
|
|||
let _ = child.wait();
|
||||
}
|
||||
|
||||
fn shutdown(mut self) {
|
||||
pub(super) fn shutdown(mut self) {
|
||||
let mut stop = Command::new(env!("CARGO_BIN_EXE_fabro"));
|
||||
apply_test_isolation(&mut stop, self.home_root.path());
|
||||
stop.args(["server", "stop"])
|
||||
|
|
@ -211,7 +232,7 @@ impl RunningServer {
|
|||
|
||||
/// Petri's store over the server's database, read beside the server:
|
||||
/// what `petri inspect` would see.
|
||||
async fn petri_store(&self) -> SqliteRunStore {
|
||||
pub(super) async fn petri_store(&self) -> SqliteRunStore {
|
||||
let database = fabro_db::Database::connect(Storage::new(&self.storage_dir).sqlite_path())
|
||||
.await
|
||||
.expect("the server database opens");
|
||||
|
|
@ -405,8 +426,9 @@ fn run_detached(
|
|||
run_detached_with(context, server, workspace, &["--auto-approve"])
|
||||
}
|
||||
|
||||
/// `fabro run --detach` against the server with extra arguments.
|
||||
fn run_detached_with(
|
||||
/// `fabro run --detach` against the server with extra arguments, such as
|
||||
/// `--auto-approve` or the model to run the workflow's agents on.
|
||||
pub(super) fn run_detached_with(
|
||||
context: &fabro_test::TestContext,
|
||||
server: &RunningServer,
|
||||
workspace: &Path,
|
||||
|
|
@ -435,7 +457,7 @@ fn run_detached_with(
|
|||
created_run_id(&output)
|
||||
}
|
||||
|
||||
async fn run_json(server: &RunningServer, path: &str) -> serde_json::Value {
|
||||
pub(super) async fn run_json(server: &RunningServer, path: &str) -> serde_json::Value {
|
||||
let response = fabro_test::test_http_client()
|
||||
.get(format!("{}/api/v1/{path}", server.api_base_url))
|
||||
.bearer_auth(TEST_DEV_TOKEN)
|
||||
|
|
@ -457,7 +479,11 @@ async fn run_status(server: &RunningServer, run_id: &str) -> String {
|
|||
.to_string()
|
||||
}
|
||||
|
||||
async fn wait_for_status(server: &RunningServer, run_id: &str, expected: &[&str]) -> String {
|
||||
pub(super) async fn wait_for_status(
|
||||
server: &RunningServer,
|
||||
run_id: &str,
|
||||
expected: &[&str],
|
||||
) -> String {
|
||||
let deadline = Instant::now() + RUN_TIMEOUT;
|
||||
loop {
|
||||
let status = run_status(server, run_id).await;
|
||||
|
|
|
|||
453
lib/apps/fabro-cli/tests/it/scenario/petri_tools.rs
Normal file
453
lib/apps/fabro-cli/tests/it/scenario/petri_tools.rs
Normal file
|
|
@ -0,0 +1,453 @@
|
|||
//! Fabro's run tools inside a Petri run (integration plan item F3.4): a
|
||||
//! workflow that enables `[run.agent] fabro_tools` runs on Petri in the
|
||||
//! worker the server launched, and the agent stage's model, the twin, calls
|
||||
//! the run tools the worker registered through Petri's host tool
|
||||
//! capability. The agent creates a child run from inside the Petri run; a
|
||||
//! `[[run.hooks]]` hook blocks a run tool; a sub-agent calls an inherited
|
||||
//! run tool. Each call is read back from Petri's record of the run, under
|
||||
//! the stage it served.
|
||||
//!
|
||||
//! The harness is `petri.rs`'s: a foreground server on disk storage with
|
||||
//! the `openai` provider repointed at the twin, its key in the vault, and
|
||||
//! the run started with `fabro run --detach`. The runs take their host
|
||||
//! scope through the sandbox-driver host plugin, so the tests skip, and say
|
||||
//! why, when it is not found.
|
||||
|
||||
#![expect(
|
||||
clippy::disallowed_methods,
|
||||
reason = "these scenarios stage workspaces with sync std::fs and start a real server subprocess"
|
||||
)]
|
||||
#![expect(
|
||||
clippy::print_stderr,
|
||||
reason = "a scenario says where it is, and why it skipped"
|
||||
)]
|
||||
|
||||
use std::collections::BTreeMap;
|
||||
use std::path::PathBuf;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use fabro_petri::engine::{self, RunStatus};
|
||||
use fabro_petri::host_tools::recorded::{self, ExecutionId, InvocationId, ToolCall};
|
||||
use fabro_static::EnvVars;
|
||||
use fabro_test::{TwinScenario, TwinScenarios, TwinToolCall, test_context, twin_openai};
|
||||
use fabro_types::{WorkflowPath, WorkflowVersion};
|
||||
use serde_json::{Value, json};
|
||||
|
||||
use super::petri::{RunningServer, host_plugin, run_detached_with, run_json, wait_for_status};
|
||||
use crate::support::TEST_DEV_TOKEN;
|
||||
|
||||
const MODEL: &str = "gpt-5.4";
|
||||
/// How every scenario starts its run: approved up front, on the twin's
|
||||
/// model.
|
||||
const RUN_ARGS: &[&str] = &["--auto-approve", "--provider", "openai", "--model", MODEL];
|
||||
const POLL: Duration = Duration::from_millis(50);
|
||||
const RUN_TIMEOUT: Duration = Duration::from_mins(1);
|
||||
|
||||
/// What the stage asks of its agent; every request of the stage's own
|
||||
/// session carries it, and no request of a sub-agent does.
|
||||
const PROMPT: &str = "Work with the Fabro run tools as instructed.";
|
||||
/// The task the stage hands a sub-agent; every request of the child's
|
||||
/// session carries it.
|
||||
const TASK: &str = "Helper: search the Fabro runs and report what you find.";
|
||||
|
||||
/// The child workflow the agent starts: one command stage, on Petri.
|
||||
const CHILD_DOT: &str = r#"digraph Child {
|
||||
graph [goal="Run one command", default_max_retries=0]
|
||||
start [shape=Mdiamond]
|
||||
exit [shape=Msquare]
|
||||
say [shape=parallelogram, script="echo hello from the child", max_retries=0]
|
||||
start -> say -> exit
|
||||
}"#;
|
||||
const CHILD_SETTINGS: &str =
|
||||
"_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\nengine = \"petri\"\n";
|
||||
|
||||
/// A `[[run.hooks]]` entry that blocks every `fabro_run_search` call.
|
||||
const BLOCKING_HOOK: &str = r#"
|
||||
[[run.hooks]]
|
||||
name = "no-run-search"
|
||||
event = "pre_tool_use"
|
||||
script = '''if grep -q 'fabro_run_search' "$FABRO_HOOK_CONTEXT"; then echo '{"decision":"block","reason":"run tools are not allowed here"}'; exit 2; fi'''
|
||||
"#;
|
||||
|
||||
/// A server whose `openai` provider is the twin, keyed by `namespace`.
|
||||
async fn server_on_twin(twin_base_url: &str, namespace: &str) -> RunningServer {
|
||||
RunningServer::start_with(
|
||||
&format!("\n[llm.providers.openai]\nbase_url = \"{twin_base_url}\"\n"),
|
||||
&[(EnvVars::OPENAI_API_KEY, namespace)],
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// A workspace holding a one-stage agent workflow on Petri with the run
|
||||
/// tools enabled, and `extra_settings` appended to its `workflow.toml`.
|
||||
fn write_agent_workspace(context: &fabro_test::TestContext, extra_settings: &str) -> PathBuf {
|
||||
let workspace = context.temp_dir.join("tools-workspace");
|
||||
std::fs::create_dir_all(&workspace).expect("the workspace creates");
|
||||
std::fs::write(
|
||||
workspace.join("workflow.fabro"),
|
||||
format!(
|
||||
"digraph Tools {{\n graph [goal=\"Use the run tools\", default_max_retries=0]\n \
|
||||
start [shape=Mdiamond]\n exit [shape=Msquare]\n work [shape=box, \
|
||||
prompt=\"{PROMPT}\", max_retries=0]\n start -> work -> exit\n}}\n"
|
||||
),
|
||||
)
|
||||
.expect("the workflow writes");
|
||||
std::fs::write(
|
||||
workspace.join("workflow.toml"),
|
||||
format!(
|
||||
"_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\nengine = \"petri\"\n\n[run]\n\
|
||||
goal = \"Use the run tools\"\n\n[run.agent]\nfabro_tools = true\n{extra_settings}"
|
||||
),
|
||||
)
|
||||
.expect("the settings write");
|
||||
workspace
|
||||
}
|
||||
|
||||
/// Register the child workflow as a version through the server's API; its
|
||||
/// id, for the agent's `fabro_run_create` call.
|
||||
async fn register_child_version(server: &RunningServer) -> String {
|
||||
let path = |name: &str| WorkflowPath::new(name).expect("the fixture path is valid");
|
||||
let version = WorkflowVersion::new(
|
||||
path("workflow.fabro"),
|
||||
BTreeMap::from([
|
||||
(path("workflow.fabro"), CHILD_DOT.to_string()),
|
||||
(path("workflow.toml"), CHILD_SETTINGS.to_string()),
|
||||
]),
|
||||
BTreeMap::new(),
|
||||
)
|
||||
.expect("the child version is valid");
|
||||
let response = fabro_test::test_http_client()
|
||||
.post(format!("{}/api/v1/workflow-versions", server.api_base_url))
|
||||
.bearer_auth(TEST_DEV_TOKEN)
|
||||
.json(&version)
|
||||
.send()
|
||||
.await
|
||||
.expect("the registration sends");
|
||||
let body = fabro_test::expect_reqwest_json(
|
||||
response,
|
||||
fabro_http::StatusCode::CREATED,
|
||||
"POST /api/v1/workflow-versions",
|
||||
)
|
||||
.await;
|
||||
body["workflow_version_id"]
|
||||
.as_str()
|
||||
.expect("the registration names the version")
|
||||
.to_string()
|
||||
}
|
||||
|
||||
fn run_stage_scenario() -> TwinScenario {
|
||||
TwinScenario::responses(MODEL).input_contains(PROMPT)
|
||||
}
|
||||
|
||||
fn child_scenario() -> TwinScenario {
|
||||
TwinScenario::responses(MODEL).input_contains(TASK)
|
||||
}
|
||||
|
||||
/// The run's Petri outcome: succeeded and whole, or the test says why not.
|
||||
async fn assert_petri_succeeded(server: &RunningServer, run_id: &str) {
|
||||
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:?}");
|
||||
assert!(outcome.complete, "{:?}", outcome.incomplete);
|
||||
}
|
||||
|
||||
/// Every completed call of `tool` in the run's Petri record.
|
||||
async fn recorded_calls(server: &RunningServer, run_id: &str, tool: &str) -> Vec<ToolCall> {
|
||||
let store = server.petri_store().await;
|
||||
recorded::tool_calls(&store, run_id, tool)
|
||||
.await
|
||||
.expect("the run's Petri record replays")
|
||||
}
|
||||
|
||||
/// The twin's request log for `namespace`: the input text of each request
|
||||
/// the stage's agent or its sub-agents made, in order. The server's own
|
||||
/// request for a run title goes to the same twin and is left out.
|
||||
async fn request_inputs(twin: &fabro_test::TwinOpenAi, namespace: &str) -> Vec<String> {
|
||||
let logs = twin.request_logs(namespace).await;
|
||||
logs["requests"]
|
||||
.as_array()
|
||||
.expect("the twin request log is an array")
|
||||
.iter()
|
||||
.map(|request| {
|
||||
request["input_text"]
|
||||
.as_str()
|
||||
.unwrap_or_default()
|
||||
.to_string()
|
||||
})
|
||||
.filter(|input| input.contains(PROMPT) || input.contains(TASK))
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// Approve `run_id` as a user, through the server's API.
|
||||
async fn approve_run(server: &RunningServer, run_id: &str) {
|
||||
let response = fabro_test::test_http_client()
|
||||
.post(format!(
|
||||
"{}/api/v1/runs/{run_id}/approve",
|
||||
server.api_base_url
|
||||
))
|
||||
.bearer_auth(TEST_DEV_TOKEN)
|
||||
.send()
|
||||
.await
|
||||
.expect("the approval sends");
|
||||
fabro_test::expect_reqwest_json(
|
||||
response,
|
||||
fabro_http::StatusCode::OK,
|
||||
format!("POST /api/v1/runs/{run_id}/approve"),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
/// The runs whose parent is `parent_id`.
|
||||
async fn children_of(server: &RunningServer, parent_id: &str) -> Vec<Value> {
|
||||
run_json(server, &format!("runs?parent_id={parent_id}")).await["data"]
|
||||
.as_array()
|
||||
.cloned()
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
async fn wait_for_children(server: &RunningServer, parent_id: &str) -> Vec<Value> {
|
||||
let deadline = Instant::now() + RUN_TIMEOUT;
|
||||
loop {
|
||||
let children = children_of(server, parent_id).await;
|
||||
if !children.is_empty() {
|
||||
return children;
|
||||
}
|
||||
assert!(
|
||||
Instant::now() < deadline,
|
||||
"no child run of {parent_id} appeared"
|
||||
);
|
||||
tokio::time::sleep(POLL).await;
|
||||
}
|
||||
}
|
||||
|
||||
/// The agent calls `fabro_run_create` from inside the Petri run: the child
|
||||
/// run is created under the Petri run as its parent and runs to its end,
|
||||
/// the model reads the tool's answer, the run succeeds, and the call is in
|
||||
/// Petri's record under the stage.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn an_agent_starts_a_child_run_with_a_run_tool_inside_a_petri_run() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let context = test_context!();
|
||||
let twin = twin_openai().await;
|
||||
let namespace = format!("{}::{}", module_path!(), line!());
|
||||
let server = server_on_twin(&twin.base_url, &namespace).await;
|
||||
let child_version = register_child_version(&server).await;
|
||||
// The child runs in the local environment, which serves a folder
|
||||
// target and not a `none` one.
|
||||
let child_workspace = context.temp_dir.join("child-workspace");
|
||||
std::fs::create_dir_all(&child_workspace).expect("the child workspace creates");
|
||||
TwinScenarios::new(namespace.clone())
|
||||
.scenario(run_stage_scenario().tool_call(TwinToolCall::new(
|
||||
"fabro_run_create",
|
||||
json!({
|
||||
"runs": [{
|
||||
"workflow_version_id": child_version,
|
||||
"target": {"kind": "folder", "path": child_workspace},
|
||||
"environment_id": "local",
|
||||
"args": {"auto_approve": true},
|
||||
}],
|
||||
}),
|
||||
)))
|
||||
.scenario(run_stage_scenario().text("The child run is on its way."))
|
||||
.load(twin)
|
||||
.await;
|
||||
let workspace = write_agent_workspace(&context, "");
|
||||
let run_id = run_detached_with(&context, &server, &workspace, RUN_ARGS);
|
||||
|
||||
eprintln!("run {run_id} started");
|
||||
let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await;
|
||||
eprintln!("run {run_id} is {status}");
|
||||
let run = run_json(&server, &format!("runs/{run_id}")).await;
|
||||
assert_eq!(
|
||||
status,
|
||||
"succeeded",
|
||||
"run: {run}\nserver stderr:\n{}",
|
||||
server.stderr_text()
|
||||
);
|
||||
assert_petri_succeeded(&server, &run_id).await;
|
||||
|
||||
let calls = recorded_calls(&server, &run_id, "fabro_run_create").await;
|
||||
assert_eq!(calls.len(), 1, "one call in the record: {calls:?}");
|
||||
let call = &calls[0];
|
||||
assert_eq!(call.node, "work", "recorded under the stage");
|
||||
assert_eq!(call.invocation, Some(InvocationId::ROOT));
|
||||
assert_eq!(call.execution, Some(ExecutionId::new(0)));
|
||||
assert!(
|
||||
call.parent_session.is_none(),
|
||||
"the stage's own session called"
|
||||
);
|
||||
assert_eq!(call.payload["is_error"], false, "{:?}", call.payload);
|
||||
|
||||
let children = wait_for_children(&server, &run_id).await;
|
||||
assert_eq!(children.len(), 1, "one child run: {children:?}");
|
||||
let child_id = children[0]["id"]
|
||||
.as_str()
|
||||
.expect("the child run has an id")
|
||||
.to_string();
|
||||
eprintln!("child run {child_id} found");
|
||||
assert_eq!(children[0]["parent_id"], run_id, "{:?}", children[0]);
|
||||
// A run a worker creates waits for a person's approval, as it does
|
||||
// when the legacy worker's agent creates one; the test is that person.
|
||||
assert_eq!(
|
||||
children[0]["lifecycle"]["status"]["reason"], "approval_required",
|
||||
"{:?}",
|
||||
children[0]["lifecycle"]
|
||||
);
|
||||
approve_run(&server, &child_id).await;
|
||||
let child_status = wait_for_status(&server, &child_id, &["succeeded", "failed"]).await;
|
||||
eprintln!("child run {child_id} is {child_status}");
|
||||
assert_eq!(
|
||||
child_status,
|
||||
"succeeded",
|
||||
"child: {}",
|
||||
run_json(&server, &format!("runs/{child_id}")).await
|
||||
);
|
||||
|
||||
let inputs = request_inputs(twin, &namespace).await;
|
||||
assert_eq!(inputs.len(), 2, "{inputs:?}");
|
||||
assert!(
|
||||
inputs[1].contains(&child_id),
|
||||
"the model read the tool's answer naming the child run: {}",
|
||||
inputs[1]
|
||||
);
|
||||
server.shutdown();
|
||||
}
|
||||
|
||||
/// A `pre_tool_use` hook from `[[run.hooks]]` blocks a run tool as it
|
||||
/// blocks Pebble's: the tool never runs, the model reads the reason, the
|
||||
/// run goes on, and Petri's record holds the hook's report and the denied
|
||||
/// call.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn a_run_hook_blocks_a_run_tool_inside_a_petri_run() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let context = test_context!();
|
||||
let twin = twin_openai().await;
|
||||
let namespace = format!("{}::{}", module_path!(), line!());
|
||||
let server = server_on_twin(&twin.base_url, &namespace).await;
|
||||
TwinScenarios::new(namespace.clone())
|
||||
.scenario(
|
||||
run_stage_scenario()
|
||||
.tool_call(TwinToolCall::new("fabro_run_search", json!({ "first": 5 }))),
|
||||
)
|
||||
.scenario(run_stage_scenario().text("The search was refused."))
|
||||
.load(twin)
|
||||
.await;
|
||||
let workspace = write_agent_workspace(&context, BLOCKING_HOOK);
|
||||
let run_id = run_detached_with(&context, &server, &workspace, RUN_ARGS);
|
||||
|
||||
let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await;
|
||||
let run = run_json(&server, &format!("runs/{run_id}")).await;
|
||||
assert_eq!(
|
||||
status,
|
||||
"succeeded",
|
||||
"run: {run}\nserver stderr:\n{}",
|
||||
server.stderr_text()
|
||||
);
|
||||
assert_petri_succeeded(&server, &run_id).await;
|
||||
|
||||
let calls = recorded_calls(&server, &run_id, "fabro_run_search").await;
|
||||
assert_eq!(calls.len(), 1, "{calls:?}");
|
||||
assert_eq!(calls[0].node, "work");
|
||||
assert_eq!(calls[0].payload["is_error"], true, "{:?}", calls[0].payload);
|
||||
assert_eq!(
|
||||
calls[0].payload["error_kind"], "denied",
|
||||
"{:?}",
|
||||
calls[0].payload
|
||||
);
|
||||
assert!(
|
||||
children_of(&server, &run_id).await.is_empty(),
|
||||
"the blocked tool created nothing"
|
||||
);
|
||||
|
||||
let store = server.petri_store().await;
|
||||
let reports = recorded::hook_reports(&store, &run_id, "pre_tool_use")
|
||||
.await
|
||||
.expect("the record replays");
|
||||
assert_eq!(reports.len(), 1, "one pre_tool_use report: {reports:?}");
|
||||
assert_eq!(reports[0]["node"], "work");
|
||||
let report = serde_json::to_string(&reports[0]["report"]).expect("the report serializes");
|
||||
assert!(
|
||||
report.contains("run tools are not allowed here"),
|
||||
"the report carries the block: {report}"
|
||||
);
|
||||
|
||||
let inputs = request_inputs(twin, &namespace).await;
|
||||
assert_eq!(inputs.len(), 2, "{inputs:?}");
|
||||
assert!(
|
||||
inputs[1].contains("run tools are not allowed here"),
|
||||
"the model saw the block reason: {}",
|
||||
inputs[1]
|
||||
);
|
||||
server.shutdown();
|
||||
}
|
||||
|
||||
/// A sub-agent the stage spawns inherits the run tools: the child's call
|
||||
/// runs under the stage, and Petri records it under the stage naming the
|
||||
/// parent session.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn a_sub_agent_calls_an_inherited_run_tool_inside_a_petri_run() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let context = test_context!();
|
||||
let twin = twin_openai().await;
|
||||
let namespace = format!("{}::{}", module_path!(), line!());
|
||||
let server = server_on_twin(&twin.base_url, &namespace).await;
|
||||
// The queue answers each request with its first unspent match: the
|
||||
// stage's requests carry `PROMPT` and the child's carry `TASK`, so the
|
||||
// stage spawns, then waits, then finishes, while the child searches
|
||||
// and reports, whichever order the two sessions ask in.
|
||||
TwinScenarios::new(namespace.clone())
|
||||
.scenario(
|
||||
run_stage_scenario()
|
||||
.tool_call(TwinToolCall::new("spawn_agent", json!({ "task": TASK }))),
|
||||
)
|
||||
.scenario(
|
||||
child_scenario()
|
||||
.tool_call(TwinToolCall::new("fabro_run_search", json!({ "first": 5 }))),
|
||||
)
|
||||
.scenario(child_scenario().text("Found the runs."))
|
||||
.scenario(run_stage_scenario().tool_call(TwinToolCall::new("wait", json!({}))))
|
||||
.scenario(run_stage_scenario().text("The helper searched."))
|
||||
.load(twin)
|
||||
.await;
|
||||
let workspace = write_agent_workspace(&context, "");
|
||||
let run_id = run_detached_with(&context, &server, &workspace, RUN_ARGS);
|
||||
|
||||
let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await;
|
||||
let run = run_json(&server, &format!("runs/{run_id}")).await;
|
||||
assert_eq!(
|
||||
status,
|
||||
"succeeded",
|
||||
"run: {run}\nserver stderr:\n{}",
|
||||
server.stderr_text()
|
||||
);
|
||||
assert_petri_succeeded(&server, &run_id).await;
|
||||
|
||||
let calls = recorded_calls(&server, &run_id, "fabro_run_search").await;
|
||||
assert_eq!(calls.len(), 1, "{calls:?}");
|
||||
let call = &calls[0];
|
||||
assert_eq!(call.node, "work", "recorded under the parent stage");
|
||||
assert_eq!(call.invocation, Some(InvocationId::ROOT));
|
||||
assert!(
|
||||
call.parent_session.is_some(),
|
||||
"the child's call names its parent session: {call:?}"
|
||||
);
|
||||
assert_eq!(call.payload["is_error"], false, "{:?}", call.payload);
|
||||
|
||||
let inputs = request_inputs(twin, &namespace).await;
|
||||
let child_inputs: Vec<&String> = inputs.iter().filter(|input| input.contains(TASK)).collect();
|
||||
assert_eq!(child_inputs.len(), 2, "the child asked twice: {inputs:?}");
|
||||
assert!(
|
||||
child_inputs[1].contains(&run_id),
|
||||
"the child read the search answer naming this run: {}",
|
||||
child_inputs[1]
|
||||
);
|
||||
server.shutdown();
|
||||
}
|
||||
|
|
@ -103,6 +103,9 @@ pub(crate) fn runtime_spec(
|
|||
model_client,
|
||||
dry_run,
|
||||
fabro_home: Some(Home::from_env().root().to_path_buf()),
|
||||
// The in-process test path has no worker client to bind the run
|
||||
// tools to; like the legacy in-process path, it runs without them.
|
||||
run_tools: None,
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -30,6 +30,7 @@ fabro-vault = { path = "../../foundation/fabro-vault" }
|
|||
fabro-workflow = { path = "../fabro-workflow" }
|
||||
fabro-checkpoint = { path = "../fabro-checkpoint" }
|
||||
fabro-util = { path = "../../foundation/fabro-util" }
|
||||
pebble-coding-agent.workspace = true
|
||||
petri_runtime.workspace = true
|
||||
petri_execution.workspace = true
|
||||
petri_store.workspace = true
|
||||
|
|
@ -52,6 +53,9 @@ tracing.workspace = true
|
|||
|
||||
[dev-dependencies]
|
||||
fabro-petri = { path = ".", features = ["test-support"] }
|
||||
fabro-tool = { path = "../fabro-tool" }
|
||||
httpmock = "0.8"
|
||||
pebble-coding-agent = { workspace = true, features = ["test-util"] }
|
||||
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"] }
|
||||
|
|
|
|||
167
lib/components/fabro-petri/src/host_tools.rs
Normal file
167
lib/components/fabro-petri/src/host_tools.rs
Normal file
|
|
@ -0,0 +1,167 @@
|
|||
//! Fabro's run tools inside a Petri run: the adapter from Petri's host tool
|
||||
//! capability to `register_fabro_run_tools` (integration plan item F3.4).
|
||||
//!
|
||||
//! Petri's native agent step asks the [`HostTools`] capability for the
|
||||
//! host's tools once per agent session, with a [`HostToolContext`] naming
|
||||
//! the stage the session serves: the run key, the invocation and execution,
|
||||
//! the node, the firing and the attempt. This module answers with the same
|
||||
//! tools the legacy worker registers on a stage's Pebble builder,
|
||||
//! `fabro_run_create`, `fabro_run_get` and the rest, built by
|
||||
//! `register_fabro_run_tools` over the same [`FabroRunToolServices`]: the
|
||||
//! worker's authenticated client and the run id every child run is parented
|
||||
//! to. `fabro exec` and Ask Fabro sessions keep registering the tools on
|
||||
//! their builders directly; this adapter is only for a run Petri executes.
|
||||
//!
|
||||
//! From there Petri treats the tools as any other: the model sees their
|
||||
//! definitions beside Pebble's, every call passes through the run's tool
|
||||
//! hooks (a `pre_tool_use` hook from `[[run.hooks]]` can block one), Pebble
|
||||
//! reports the call on its event stream, and Petri records it under the
|
||||
//! stage. A sub-agent inherits them through Pebble's own rule, since the
|
||||
//! registration marks every run tool `allow_in_subagents`.
|
||||
//!
|
||||
//! # Identity
|
||||
//!
|
||||
//! The run tools need one identity: the Fabro run id, which is Petri's run
|
||||
//! key for the run (`RunRequest::run_id`) and `FabroRunToolServices::
|
||||
//! current_run_id`. It is the parent link of every child run a stage
|
||||
//! creates. No run tool records a stage on the effects it creates, so
|
||||
//! nothing here derives Fabro's old `StageId` (`node@visit`); the stage a
|
||||
//! call came from is Petri's own record of the call, under the stage key
|
||||
//! `(run, execution, firing)`, and this adapter logs that key with the node
|
||||
//! and attempt when it builds a session's tools.
|
||||
//!
|
||||
//! The builder refuses a context whose run key is not the run the services
|
||||
//! were built for: the tools would parent child runs to the wrong run. That
|
||||
//! cannot happen in the worker, which builds both from one run id, so it is
|
||||
//! logged as an error and the session gets no run tools rather than the
|
||||
//! wrong ones.
|
||||
|
||||
use fabro_workflow::handler::llm::register_fabro_run_tools;
|
||||
use fabro_workflow::services::FabroRunToolServices;
|
||||
use pebble_coding_agent::tools::RegisteredTool;
|
||||
use petri_attractor_steps::host_tools::{HostToolContext, HostTools};
|
||||
use tracing::{debug, error};
|
||||
|
||||
/// The `HostTools` capability that gives every native agent session of the
|
||||
/// run Fabro's run tools, bound to `services`. Register it on the runtime
|
||||
/// the run executes with; `RuntimeSpec::run_tools` does.
|
||||
#[must_use]
|
||||
pub fn capability(services: FabroRunToolServices) -> HostTools {
|
||||
HostTools::new().with(move |context| tools_for_stage(&services, context))
|
||||
}
|
||||
|
||||
/// The run tools for the session `context` names: what
|
||||
/// `register_fabro_run_tools` builds for the legacy worker, or nothing when
|
||||
/// the context's run is not the one `services` serves.
|
||||
#[must_use]
|
||||
pub fn tools_for_stage(
|
||||
services: &FabroRunToolServices,
|
||||
context: &HostToolContext,
|
||||
) -> Vec<RegisteredTool> {
|
||||
let run_id = services.current_run_id.to_string();
|
||||
if context.run.as_str() != run_id {
|
||||
error!(
|
||||
run = %context.run,
|
||||
services_run_id = %run_id,
|
||||
node = %context.node,
|
||||
"the Petri run key is not the run the Fabro run tools serve; the session gets no run tools"
|
||||
);
|
||||
return Vec::new();
|
||||
}
|
||||
debug!(
|
||||
run = %context.run,
|
||||
invocation = %context.invocation,
|
||||
execution = %context.execution,
|
||||
firing = %context.firing,
|
||||
node = %context.node,
|
||||
attempt = ?context.attempt,
|
||||
"registering Fabro's run tools on a Petri agent session"
|
||||
);
|
||||
register_fabro_run_tools(services)
|
||||
}
|
||||
|
||||
/// What a test reads back from a Petri run's record about the run tools,
|
||||
/// without depending on the Petri packages itself.
|
||||
#[cfg(feature = "test-support")]
|
||||
pub mod recorded {
|
||||
use petri_attractor_steps::hooks::REPORT_EVENT;
|
||||
use petri_execution::events::{RunEvent, replay_run};
|
||||
use petri_execution::{Access, RunKey, RunStore};
|
||||
pub use petri_execution::{ExecutionId, InvocationId};
|
||||
use serde_json::Value;
|
||||
|
||||
/// One completed tool call as Petri recorded it: the stage it was
|
||||
/// recorded under and Pebble's completion payload.
|
||||
#[derive(Clone, Debug)]
|
||||
pub struct ToolCall {
|
||||
/// The node's instance name.
|
||||
pub node: String,
|
||||
pub invocation: Option<InvocationId>,
|
||||
pub execution: Option<ExecutionId>,
|
||||
/// The parent session of a sub-agent's call; `None` for a call of
|
||||
/// the stage's own session.
|
||||
pub parent_session: Option<String>,
|
||||
/// Pebble's `ToolCallCompleted` payload (`tool_name`, `is_error`,
|
||||
/// `error_kind`, the output).
|
||||
pub payload: Value,
|
||||
}
|
||||
|
||||
/// Every event of the run, replayed from its record.
|
||||
async fn events(store: &dyn RunStore, run_id: &str) -> anyhow::Result<Vec<RunEvent>> {
|
||||
let logs = store
|
||||
.open(&RunKey::new(run_id), Access::Read)
|
||||
.await
|
||||
.map_err(anyhow::Error::new)?;
|
||||
replay_run(&*logs).await.map_err(anyhow::Error::new)
|
||||
}
|
||||
|
||||
/// Every completed call of `tool` in the run's record, in record order.
|
||||
pub async fn tool_calls(
|
||||
store: &dyn RunStore,
|
||||
run_id: &str,
|
||||
tool: &str,
|
||||
) -> anyhow::Result<Vec<ToolCall>> {
|
||||
Ok(events(store, run_id)
|
||||
.await?
|
||||
.iter()
|
||||
.filter_map(|event| {
|
||||
let custom = event.custom()?;
|
||||
if custom["kind"] != "pebble" {
|
||||
return None;
|
||||
}
|
||||
let envelope = custom.get("event")?;
|
||||
let payload = envelope["event"].get("ToolCallCompleted")?;
|
||||
if payload["tool_name"] != tool {
|
||||
return None;
|
||||
}
|
||||
Some(ToolCall {
|
||||
node: event
|
||||
.subject
|
||||
.as_ref()
|
||||
.map(|subject| subject.node.name.to_string())
|
||||
.unwrap_or_default(),
|
||||
invocation: event.context.invocation,
|
||||
execution: event.context.execution,
|
||||
parent_session: envelope["parent_session_id"].as_str().map(str::to_owned),
|
||||
payload: payload.clone(),
|
||||
})
|
||||
})
|
||||
.collect())
|
||||
}
|
||||
|
||||
/// Every hook report for `event` (`pre_tool_use`, say) in the run's
|
||||
/// record, as Petri's hook service recorded it.
|
||||
pub async fn hook_reports(
|
||||
store: &dyn RunStore,
|
||||
run_id: &str,
|
||||
event: &str,
|
||||
) -> anyhow::Result<Vec<Value>> {
|
||||
Ok(events(store, run_id)
|
||||
.await?
|
||||
.iter()
|
||||
.filter_map(RunEvent::custom)
|
||||
.filter(|value| value["kind"] == REPORT_EVENT && value["event"] == event)
|
||||
.cloned()
|
||||
.collect())
|
||||
}
|
||||
}
|
||||
|
|
@ -43,7 +43,8 @@
|
|||
//! to a worker;
|
||||
//! - [`platform_records`]: Fabro's platform records as the adapters reach them,
|
||||
//! in the server's database or over its API from a worker;
|
||||
//! - the platform adapter still to come: the run tools.
|
||||
//! - [`host_tools`]: Fabro's run tools on every native agent session of a run,
|
||||
//! through Petri's `HostTools` capability.
|
||||
//!
|
||||
//! The Petri packages are pinned by revision in the workspace `Cargo.toml`
|
||||
//! under `petri_*` keys.
|
||||
|
|
@ -54,6 +55,7 @@ pub mod check;
|
|||
pub mod checkpoint;
|
||||
pub mod engine;
|
||||
pub mod hooks;
|
||||
pub mod host_tools;
|
||||
pub mod http_store;
|
||||
pub mod interview;
|
||||
pub mod petri;
|
||||
|
|
|
|||
|
|
@ -5,13 +5,16 @@
|
|||
//! carrying the server's settings layer, the Attractor step kinds (the real
|
||||
//! ones, or the simulated registry for a dry run), the model client as the
|
||||
//! `PebbleClient` capability so Petri's admission pass pins every LLM node's
|
||||
//! route, and the Fabro home for the skills step. Nothing here knows about a
|
||||
//! run: the store and the run options are added by the caller.
|
||||
//! route, the Fabro home for the skills step, and, at execution, Fabro's run
|
||||
//! tools as the `HostTools` capability when the run enables them. Nothing
|
||||
//! here knows about a run's record: the store and the run options are added
|
||||
//! by the caller.
|
||||
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
|
||||
use fabro_http::HttpClient;
|
||||
use fabro_workflow::services::FabroRunToolServices;
|
||||
use lithos_llm::Client;
|
||||
use lithos_llm::catalog::{Catalog, ProviderId};
|
||||
use lithos_llm::client::ClientBuildError;
|
||||
|
|
@ -22,6 +25,8 @@ use petri_frontend_fabro::Fabro;
|
|||
use petri_runtime::Runtime;
|
||||
use tracing::debug;
|
||||
|
||||
use crate::host_tools;
|
||||
|
||||
/// What every Petri runtime Fabro builds is configured with.
|
||||
#[derive(Clone, Default)]
|
||||
pub struct RuntimeSpec {
|
||||
|
|
@ -39,6 +44,11 @@ pub struct RuntimeSpec {
|
|||
/// The Fabro home the skills step reads; `None` leaves it to Petri's
|
||||
/// own lookup (`FABRO_HOME`, else `$HOME/.fabro`).
|
||||
pub fabro_home: Option<PathBuf>,
|
||||
/// Fabro's run tools for every native agent session of the run, when
|
||||
/// the run enables them (`[run.agent] fabro_tools` and the worker
|
||||
/// token's `agent:run_tools` scope); `None` gives the sessions Pebble's
|
||||
/// tools alone. See [`crate::host_tools`].
|
||||
pub run_tools: Option<FabroRunToolServices>,
|
||||
}
|
||||
|
||||
impl RuntimeSpec {
|
||||
|
|
@ -60,6 +70,9 @@ impl RuntimeSpec {
|
|||
if let Some(home) = home {
|
||||
runtime = runtime.capability(home);
|
||||
}
|
||||
if let Some(services) = &self.run_tools {
|
||||
runtime = runtime.capability(host_tools::capability(services.clone()));
|
||||
}
|
||||
if for_execution && self.dry_run {
|
||||
petri_attractor_steps::register_stubs(runtime)
|
||||
} else {
|
||||
|
|
|
|||
300
lib/components/fabro-petri/tests/host_tools.rs
Normal file
300
lib/components/fabro-petri/tests/host_tools.rs
Normal file
|
|
@ -0,0 +1,300 @@
|
|||
//! Fabro's run tools on a Petri run from this crate (integration plan item
|
||||
//! F3.4): `RuntimeSpec::run_tools` installs the adapter as Petri's host
|
||||
//! tool capability, a workflow with one agent stage runs on the real step
|
||||
//! registry against a scripted model, and the stage's session gets the
|
||||
//! tools the legacy worker registers, bound to the run: the model is
|
||||
//! advertised every run tool, its `fabro_run_create` call reaches Fabro's
|
||||
//! API with the Petri run as the child's parent, the API's answer comes
|
||||
//! back to the model, and the call is in the run's record under the stage.
|
||||
//!
|
||||
//! Every run takes its scope through the sandbox-driver host plugin, so
|
||||
//! the tests skip when that executable is not found, unless
|
||||
//! `FABRO_REQUIRE_SANDBOX_PLUGINS` is set.
|
||||
|
||||
#![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")]
|
||||
|
||||
use std::env;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
use fabro_petri::host_tools::recorded::{self, ExecutionId, InvocationId};
|
||||
use fabro_petri::runtime::RuntimeSpec;
|
||||
use fabro_tool::fabro_client::ClientBackend;
|
||||
use fabro_types::{BlobHash, RunId, WorkflowVersionId};
|
||||
use fabro_workflow::handler::llm::register_fabro_run_tools;
|
||||
use fabro_workflow::services::FabroRunToolServices;
|
||||
use httpmock::{Method, MockServer};
|
||||
use lithos_llm::types::Request;
|
||||
use pebble_coding_agent::test_support::{
|
||||
ScriptedCall, ScriptedProvider, scripted_client, text_response, tool_call_response,
|
||||
};
|
||||
use petri_execution::host::{self, HostRun};
|
||||
use petri_runtime::executor::Retention;
|
||||
use petri_runtime::frontend::CompileInputs;
|
||||
use petri_runtime::ir::RunStatus;
|
||||
use petri_runtime::{RunOptions, Runtime};
|
||||
use petri_store::{MemoryRunStore, RunKey, RunStore};
|
||||
use serde_json::json;
|
||||
use tokio::fs;
|
||||
|
||||
const HOST_PLUGIN: &str = "sandbox-driver-host";
|
||||
const HOST_PLUGIN_OVERRIDE: &str = "PETRI_SANDBOX_HOST_PLUGIN";
|
||||
const REQUIRE_ENV: &str = "FABRO_REQUIRE_SANDBOX_PLUGINS";
|
||||
|
||||
/// One agent stage on the native backend, pinned to the scripted model.
|
||||
const AGENT_WORKFLOW: &str = r#"digraph Agent {
|
||||
graph [goal="Start a child run", backend="api", default_max_retries=0]
|
||||
start [shape=Mdiamond]
|
||||
exit [shape=Msquare]
|
||||
work [shape=box, prompt="Start the child run", model="test/model", max_retries=0]
|
||||
start -> work -> exit
|
||||
}"#;
|
||||
|
||||
const AGENT_SETTINGS: &str = "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\n";
|
||||
|
||||
/// 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.
|
||||
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
|
||||
}
|
||||
|
||||
/// Write the agent bundle into `<root>/.fabro/workflows/agent`; the
|
||||
/// workflow file.
|
||||
async fn install_bundle(root: &Path) -> PathBuf {
|
||||
let bundle = root.join(".fabro").join("workflows").join("agent");
|
||||
fs::create_dir_all(&bundle)
|
||||
.await
|
||||
.expect("the bundle directory is creatable");
|
||||
fs::write(bundle.join("workflow.fabro"), AGENT_WORKFLOW)
|
||||
.await
|
||||
.expect("the workflow is writable");
|
||||
fs::write(bundle.join("workflow.toml"), AGENT_SETTINGS)
|
||||
.await
|
||||
.expect("the settings are writable");
|
||||
bundle.join("workflow.fabro")
|
||||
}
|
||||
|
||||
/// The run tools' services over `server`, as the worker binds them: the
|
||||
/// client backend and the run the tools serve.
|
||||
fn services(server: &MockServer, run_id: RunId) -> FabroRunToolServices {
|
||||
let client = fabro_client::Client::new_no_proxy(&server.url("")).expect("the client builds");
|
||||
FabroRunToolServices {
|
||||
backend: Arc::new(ClientBackend::new(Arc::new(client))),
|
||||
current_run_id: run_id,
|
||||
}
|
||||
}
|
||||
|
||||
/// The scripted model: one `fabro_run_create` call, then a closing line.
|
||||
fn scripted_model(version_id: &str) -> (lithos_llm::Client, Arc<ScriptedProvider>) {
|
||||
scripted_client(vec![
|
||||
ScriptedCall::response(tool_call_response(
|
||||
"fabro_run_create",
|
||||
"create",
|
||||
json!({
|
||||
"runs": [{
|
||||
"workflow_version_id": version_id,
|
||||
"target": {"kind": "none"},
|
||||
"args": {"auto_approve": false},
|
||||
}],
|
||||
}),
|
||||
)),
|
||||
ScriptedCall::response(text_response("Asked for the child run.")),
|
||||
])
|
||||
}
|
||||
|
||||
/// The runtime the worker would build for the run: the scripted model as
|
||||
/// the model client and `run_tools` as the run tools.
|
||||
fn runtime(
|
||||
run_dir: &Path,
|
||||
run_id: &str,
|
||||
model_client: lithos_llm::Client,
|
||||
run_tools: Option<FabroRunToolServices>,
|
||||
store: &Arc<MemoryRunStore>,
|
||||
) -> Runtime {
|
||||
let mut options = RunOptions::new(run_dir);
|
||||
options.grace = Duration::from_secs(2);
|
||||
options.retention = Retention::Never;
|
||||
options.echo = false;
|
||||
options.run_key = Some(RunKey::new(run_id));
|
||||
RuntimeSpec {
|
||||
model_client: Some(model_client),
|
||||
run_tools,
|
||||
..RuntimeSpec::default()
|
||||
}
|
||||
.runtime(true)
|
||||
.store(Arc::clone(store) as Arc<dyn RunStore>)
|
||||
.options(options)
|
||||
}
|
||||
|
||||
/// Lower the bundle and run it to its end.
|
||||
async fn run(rt: &Runtime, workflow: &Path) {
|
||||
let lowered = rt
|
||||
.check(workflow, None, None, &CompileInputs::new())
|
||||
.expect("the workflow file loads");
|
||||
let graph = lowered
|
||||
.graph
|
||||
.unwrap_or_else(|| panic!("the workflow lowers: {:?}", lowered.diagnostics));
|
||||
let report = host::run_configured(rt, HostRun::new(graph), |_, _| {})
|
||||
.await
|
||||
.expect("the run completes");
|
||||
assert_eq!(
|
||||
report.status,
|
||||
RunStatus::Success,
|
||||
"errors: {:?}; history: {:#?}",
|
||||
report.state.errors(),
|
||||
report.state.history()
|
||||
);
|
||||
}
|
||||
|
||||
/// The tools a request advertised, as `(name, description)`.
|
||||
fn advertised(request: &Request) -> Vec<(String, String)> {
|
||||
request
|
||||
.tools()
|
||||
.iter()
|
||||
.map(|tool| (tool.name.clone(), tool.description.clone()))
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// The stage's session is given every run tool the legacy worker
|
||||
/// registers, by the same names and descriptions; its `fabro_run_create`
|
||||
/// call reaches Fabro's API with the Petri run as the parent; the API's
|
||||
/// answer reaches the model; the call is in the record under the stage.
|
||||
#[tokio::test]
|
||||
async fn a_petri_stage_calls_a_run_tool_bound_to_the_run() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let root = tempfile::tempdir().expect("a temp dir");
|
||||
let workflow = install_bundle(root.path()).await;
|
||||
let run_id = RunId::new();
|
||||
let version_id: WorkflowVersionId = BlobHash::new(b"child workflow").into();
|
||||
let server = MockServer::start_async().await;
|
||||
// Admission rejection proves the tool reached the canonical API with
|
||||
// the Petri run as the child's parent, and nothing was created.
|
||||
let create = server
|
||||
.mock_async(|when, then| {
|
||||
when.method(Method::POST)
|
||||
.path("/api/v1/runs")
|
||||
.json_body(json!({
|
||||
"workflow_version_id": version_id,
|
||||
"target": {"kind": "none"},
|
||||
"parent_id": run_id,
|
||||
"args": {"auto_approve": false},
|
||||
}));
|
||||
then.status(422).body("native admission rejection");
|
||||
})
|
||||
.await;
|
||||
let services = services(&server, run_id);
|
||||
let legacy: Vec<(String, String)> = register_fabro_run_tools(&services)
|
||||
.iter()
|
||||
.map(|tool| {
|
||||
(
|
||||
tool.definition().name.clone(),
|
||||
tool.definition().description.clone(),
|
||||
)
|
||||
})
|
||||
.collect();
|
||||
let (client, provider) = scripted_model(&version_id.to_string());
|
||||
let store = Arc::new(MemoryRunStore::new());
|
||||
let rt = runtime(
|
||||
&root.path().join("run"),
|
||||
&run_id.to_string(),
|
||||
client,
|
||||
Some(services),
|
||||
&store,
|
||||
);
|
||||
|
||||
run(&rt, &workflow).await;
|
||||
|
||||
let requests = provider.requests();
|
||||
assert_eq!(requests.len(), 2, "one tool call, one closing turn");
|
||||
let tools = advertised(&requests[0]);
|
||||
assert!(!legacy.is_empty());
|
||||
for tool in &legacy {
|
||||
assert!(
|
||||
tools.contains(tool),
|
||||
"the model was advertised {tool:?} as the legacy worker registers it: {tools:?}"
|
||||
);
|
||||
}
|
||||
assert!(
|
||||
tools.iter().any(|(name, _)| name == "shell"),
|
||||
"Pebble's own tools stay: {tools:?}"
|
||||
);
|
||||
let answer = serde_json::to_string(&requests[1]).expect("the request serializes");
|
||||
assert!(
|
||||
answer.contains("native admission rejection"),
|
||||
"the model read the API's answer: {answer}"
|
||||
);
|
||||
create.assert_calls_async(1).await;
|
||||
|
||||
let calls = recorded::tool_calls(store.as_ref(), &run_id.to_string(), "fabro_run_create")
|
||||
.await
|
||||
.expect("the record replays");
|
||||
assert_eq!(calls.len(), 1, "{calls:?}");
|
||||
assert_eq!(calls[0].node, "work", "recorded under the stage");
|
||||
assert_eq!(calls[0].invocation, Some(InvocationId::ROOT));
|
||||
assert_eq!(calls[0].execution, Some(ExecutionId::new(0)));
|
||||
assert!(calls[0].parent_session.is_none());
|
||||
assert_eq!(calls[0].payload["is_error"], true, "{:?}", calls[0].payload);
|
||||
}
|
||||
|
||||
/// Services bound to another run give the stage no run tools: the model
|
||||
/// is not advertised them, and its call is refused as an unknown tool
|
||||
/// rather than parenting a child run to the wrong run.
|
||||
#[tokio::test]
|
||||
async fn services_for_another_run_give_the_stage_no_run_tools() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let root = tempfile::tempdir().expect("a temp dir");
|
||||
let workflow = install_bundle(root.path()).await;
|
||||
let run_id = RunId::new();
|
||||
let server = MockServer::start_async().await;
|
||||
let create = server
|
||||
.mock_async(|when, then| {
|
||||
when.method(Method::POST).path("/api/v1/runs");
|
||||
then.status(500);
|
||||
})
|
||||
.await;
|
||||
let services = services(&server, RunId::new());
|
||||
let (client, provider) = scripted_model("0000");
|
||||
let store = Arc::new(MemoryRunStore::new());
|
||||
let rt = runtime(
|
||||
&root.path().join("run"),
|
||||
&run_id.to_string(),
|
||||
client,
|
||||
Some(services),
|
||||
&store,
|
||||
);
|
||||
|
||||
run(&rt, &workflow).await;
|
||||
|
||||
let requests = provider.requests();
|
||||
assert_eq!(requests.len(), 2);
|
||||
let tools = advertised(&requests[0]);
|
||||
assert!(
|
||||
!tools.iter().any(|(name, _)| name.starts_with("fabro_")),
|
||||
"no run tool was advertised: {tools:?}"
|
||||
);
|
||||
create.assert_calls_async(0).await;
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue