feat(workflow): enforce strict api/acp backends (#307)

## Summary

This PR makes agent execution a strict two-backend contract: API-backed
stages use Fabro-owned model/provider auth, while ACP-backed stages
launch a user-supplied stdio process that owns its own auth and tools.
That removes the legacy CLI backend and prevents ACP execution from
accidentally resolving or forwarding provider credentials.

## Changes

- Replaces the old `api`/`cli`/`acp` backend model with `AgentBackend {
api, acp }`, with `backend=\"cli\"` rejected and migrated toward
explicit ACP process configuration.
- Splits ACP process configuration into `acp.command` for shell command
strings and `acp.config` for JSON stdio configs, while rejecting legacy
`acp_command`.
- Restricts ACP to `agent` nodes and rejects API-only attributes such as
`model`, `provider`, `reasoning_effort`, `max_tokens`, and `speed` on
ACP nodes.
- Deletes the workflow CLI runtime, CLI credential resolver surface, CLI
live smoke tests, and `agent.cli.*` event handling.
- Updates ACP events and projections to report process identity
(`command`, optional `config_name`) rather than provider/model metadata.
- Updates import/stylesheet propagation, CLI workflow smoke coverage,
server steering tests, and web model extraction for the new
event/backend contract.

## Validation

- `cargo check -p fabro-auth -p fabro-acp -p fabro-workflow -p fabro-cli
--all-targets`
- `cargo nextest run -p fabro-auth -p fabro-acp -p fabro-validate -p
fabro-store -p fabro-workflow --lib`
- `cargo nextest run -p fabro-acp`
- `cargo nextest run -p fabro-cli --test it
workflow::acp::acp_backend_workflow`
- `cargo nextest run -p fabro-workflow --test it
codergen_without_backend_simulated`
- `cargo nextest run -p fabro-workflow --test it
import_e2e_through_engine`
- `cargo nextest run -p fabro-workflow --test it stylesheet_application`
- `cargo nextest run -p fabro-server
steer_with_active_acp_stage_returns_non_steerable_conflict`
- `cargo nextest run -p fabro-server
active_acp_stage_marker_clears_on_terminal_paths`
- `cargo nextest run -p fabro-types
agent_backend_accepts_only_api_and_acp`
- `cd apps/fabro-web && bun test app/routes/run-stages.test.ts`
- `cd apps/fabro-web && bun run typecheck`
- `cargo +nightly-2026-04-14 fmt --check --all`
- `cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D
warnings`

---

[![Compound
Engineering](https://img.shields.io/badge/Compound_Engineering-6366f1)](https://github.com/EveryInc/compound-engineering-plugin)
🤖 Generated with GPT-5 via [Codex](https://openai.com/codex)

---------

Co-authored-by: Peter Bell <4843+PeterBell@users.noreply.github.com>
This commit is contained in:
Bryan Helmkamp 2026-05-18 10:20:56 -07:00 • committed by GitHub
parent c6547f126a
commit 29b7cc0de0
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
48 changed files with 1124 additions and 4256 deletions

7
Cargo.lock generated
View file

@ -1584,20 +1584,17 @@ version = "0.237.0-nightly.0"
dependencies = [
"agent-client-protocol",
"agent-client-protocol-tokio",
"bytes",
"fabro-model",
"fabro-sandbox",
"fabro-types",
"fabro-util",
"futures",
"serde",
"serde_json",
"shlex",
"tempfile",
"thiserror 2.0.18",
"tokio",
"tokio-util",
"tracing",
"uuid",
]
[[package]]
@ -1678,7 +1675,6 @@ dependencies = [
"httpmock",
"serde",
"serde_json",
"shlex",
"tempfile",
"thiserror 2.0.18",
"tokio",
@ -2508,6 +2504,7 @@ dependencies = [
name = "fabro-validate"
version = "0.237.0-nightly.0"
dependencies = [
"fabro-acp",
"fabro-graphviz",
"fabro-model",
"fabro-types",

View file

@ -368,13 +368,12 @@ describe("eventsToActivity", () => {
properties: { model: "claude-opus-4-5" },
}),
envelope(2, {
event: "agent.cli.started",
event: "agent.session.activated",
stage_id: "agent@1",
node_id: "agent",
properties: {
provider: "anthropic",
model: "claude-sonnet-4-6",
command: "claude",
},
}),
];

View file

@ -490,7 +490,6 @@ export function buildThreadDnaItems(
const STAGE_MODEL_EVENT_NAMES = new Set([
"stage.prompt",
"agent.session.activated",
"agent.cli.started",
]);
export function extractStageModel(

View file

@ -6,7 +6,7 @@ description: "Delegate subtasks to child agent sessions"
An agent can spawn **sub-agents** to delegate work to independent child sessions. Each sub-agent gets its own LLM session and tool access, runs concurrently with the parent, and returns its result when finished.
<Note>
Sub-agents are only available with the [API backend](/core-concepts/agents#api-backend-default) (the default). Agents using the [CLI backend](/core-concepts/agents#cli-backend) or [ACP backend](/core-concepts/agents#acp-backend) cannot spawn Fabro sub-agents.
Sub-agents are only available with the [API backend](/core-concepts/agents#api-backend-default) (the default). Agents using the [ACP backend](/core-concepts/agents#acp-backend) cannot spawn Fabro sub-agents.
</Note>
## Tools

View file

@ -6,7 +6,7 @@ description: "Built-in tools for file I/O, shell commands, search, and web acces
Every agent in Fabro has access to a set of built-in tools for interacting with the codebase and environment. Tools execute inside the agent's [sandbox](/execution/environments) — whether that's the local machine, a Docker container, or a Daytona VM — so the same tool calls work identically regardless of provider.
<Note>
The tools described on this page apply to the **API backend** (the default). When using the [CLI backend](/core-concepts/agents#cli-backend) or [ACP backend](/core-concepts/agents#acp-backend), the external agent process provides its own tools — Fabro's built-in tools are not used.
The tools described on this page apply to the **API backend** (the default). When using the [ACP backend](/core-concepts/agents#acp-backend), the external agent process provides its own tools — Fabro's built-in tools are not used.
</Note>
## Core tools

View file

@ -18,7 +18,7 @@ This loop continues until the model stops calling tools, indicating it considers
## Backends
Every agent and prompt node uses a **backend** that determines how Fabro interacts with the LLM. There are three options:
Every agent node uses a **backend** that determines how Fabro interacts with the LLM or external agent process. Prompt nodes always use the API backend. Agent nodes support two backend values: `api` and `acp`.
### API backend (default)
@ -29,66 +29,48 @@ Fabro manages the agent loop directly — it calls the LLM provider's API, execu
- Provider failover
- All [built-in tools](/agents/tools) and [MCP](/agents/mcp) integrations
### CLI backend
Fabro delegates execution to a legacy external coding assistant CLI. The CLI tool manages its own tool loop internally — Fabro sends the prompt, waits for the CLI to finish, and tracks file changes via `git diff` before and after execution.
The CLI is selected automatically based on the node's provider:
| Provider | CLI tool |
|---|---|
| Anthropic | `claude` |
| OpenAI | `codex` |
| Gemini | `gemini` |
Fabro does not install these CLIs at runtime. Install the selected CLI in the sandbox image or run setup steps before the workflow reaches a `backend="cli"` node.
Set the CLI backend on a node with `backend="cli"` or via a [model stylesheet](/workflows/stylesheets):
```dot
implement [label="Implement", backend="cli"]
```
```
// Stylesheet
* { backend: cli; }
```
### ACP backend
Fabro can also run Agent Client Protocol (ACP) stdio agents with `backend="acp"`. ACP agents run inside the active Fabro sandbox, so local and Docker runs keep the same workspace isolation, secret forwarding, cancellation, and file-change tracking behavior as other agent stages.
Fabro can run Agent Client Protocol (ACP) stdio agents with `backend="acp"`. ACP agents run inside the active Fabro sandbox, so local and Docker runs keep the same workspace isolation, cancellation, and file-change tracking behavior as other agent stages.
Set ACP on a node with `backend="acp"` and an explicit `acp_command`:
ACP stages do not use Fabro model/provider credentials. The ACP process owns its auth, tools, and model behavior. Configure the process with exactly one of:
- `acp.command`: a shell command string
- `acp.config`: a JSON stdio ACP config
```dot
implement [label="Implement", backend="acp", acp_command="python3 tools/fake_acp_agent.py"]
implement [label="Implement", backend="acp", acp.command="python3 tools/fake_acp_agent.py"]
```
Fabro does not install ACP agents, Node.js, npm, or `npx` at runtime. The command must already be available in the sandbox image, repository, or setup steps. You can use `npx ...@latest` as an explicit `acp_command` if that is the behavior you want, but Fabro will treat it like any other user-supplied command.
```dot
implement [
label="Implement"
backend="acp"
acp.config="{\"type\":\"stdio\",\"name\":\"agent\",\"command\":\"python3\",\"args\":[\"tools/fake_acp_agent.py\"]}"
]
```
ACP v1 does not have a portable model-selection request. Fabro records the selected provider and model in events and run projections, but model-specific ACP behavior must be encoded in the chosen command for now. ACP is supported with local and Docker sandboxes; Daytona does not expose bidirectional stdio yet, so ACP nodes fail there with an explicit unsupported-provider error.
Fabro does not install ACP agents, Node.js, npm, or `npx` at runtime. Commands must already be available in the sandbox image, repository, or setup steps. You can use `npx ...@latest` as an explicit `acp.command` if that is the behavior you want, but Fabro will treat it like any other user-supplied command.
The legacy `acp_command` attribute is rejected; use `acp.command` for shell commands or `acp.config` for JSON stdio configs. ACP is supported with local and Docker sandboxes; Daytona does not expose bidirectional stdio yet, so ACP nodes fail there with an explicit unsupported-provider error.
### Comparison
| Capability | API backend | CLI backend | ACP backend |
|---|---|---|---|
| Tools | Fabro built-in tools + MCP | CLI's own tool set | ACP agent's own tool set |
| Session caching | Supported (`fidelity` + `thread_id`) | Not supported | Agent-dependent |
| Sub-agents | Supported | Not supported | Not supported through Fabro tools |
| Provider failover | Supported | Not supported | Not supported |
| File tracking | Tool call events | `git diff` before/after | `git diff` before/after |
### When to use the CLI backend
- **CLI-specific tools** — leverage tool implementations built into a specific CLI (e.g. Claude Code's computer use, Codex's sandboxed execution)
- **CLI-only models** — use models that are only available through a CLI tool, not via API
- **Existing workflows** — integrate a CLI tool you already depend on without rewriting its configuration
| Capability | API backend | ACP backend |
|---|---|---|
| Tools | Fabro built-in tools + MCP | ACP agent's own tool set |
| Session caching | Supported (`fidelity` + `thread_id`) | Agent-dependent |
| Sub-agents | Supported | Not supported through Fabro tools |
| Provider failover | Supported | Not supported |
| Model/provider credentials | Resolved by Fabro | Not used by Fabro |
| File tracking | Tool call events | `git diff` before/after |
### When to use the ACP backend
- **Protocol adapters** — run ACP-compatible coding agents through a stable stdio protocol
- **Sandbox parity** — keep agent process execution inside Fabro's local or Docker sandbox
- **Custom agents** — use `acp_command` for a checked-in or preinstalled ACP adapter
- **Command-owned auth** — let the ACP command manage its own provider login, model selection, and credentials
- **Custom agents** — use `acp.command` or `acp.config` for a checked-in or preinstalled ACP adapter
## Tools

View file

@ -206,8 +206,9 @@ Start nodes can also be identified by ID (`start` or `Start`). Exit nodes can be
| `model` | String | Explicit model ID (overrides stylesheet) |
| `provider` | String | Explicit provider name (overrides stylesheet). Auto-inferred from the model catalog when omitted. |
| `project_memory` | Boolean | When `true` (default), prompt nodes discover and include project docs (`AGENTS.md`, `CLAUDE.md`, etc.) as a system prompt. Set to `false` to disable. |
| `backend` | String | Agent execution backend: `api` (default), `cli`, or `acp`. `api` runs Fabro's tool loop through provider APIs; `cli` delegates to the legacy provider CLI; `acp` runs an Agent Client Protocol stdio agent inside the active sandbox. See [Agents — Backends](/core-concepts/agents#backends). |
| `acp_command` | String | Required for nodes with `backend="acp"`. The value must be a stdio ACP command available in the sandbox. Fabro records model selection but does not send it through stable ACP v1. |
| `backend` | String | Agent execution backend: `api` (default) or `acp`. `api` runs Fabro's tool loop through provider APIs; `acp` runs an Agent Client Protocol stdio agent inside the active sandbox. Prompt nodes are API-only. See [Agents — Backends](/core-concepts/agents#backends). |
| `acp.command` | String | Shell command for nodes with `backend="acp"`. Mutually exclusive with `acp.config`. The value is always parsed as a command string, not JSON. |
| `acp.config` | String | JSON stdio ACP config for nodes with `backend="acp"`. Mutually exclusive with `acp.command`. |
### Command nodes

View file

@ -73,7 +73,8 @@ A small set of attributes on the placeholder node propagate as **defaults** to e
| `reasoning_effort` |
| `speed` |
| `backend` |
| `acp_command` |
| `acp.command` |
| `acp.config` |
| `fidelity` |
| `max_retries` |
| `thread_id` |

View file

@ -7,6 +7,16 @@ license.workspace = true
description = "Agent Client Protocol backend support for Fabro"
[features]
default = ["runtime"]
runtime = [
"dep:fabro-sandbox",
"dep:fabro-types",
"dep:fabro-util",
"dep:futures",
"dep:tokio",
"dep:tokio-util",
"dep:tracing",
]
test-support = []
[lib]
@ -18,19 +28,16 @@ workspace = true
[dependencies]
agent-client-protocol.workspace = true
agent-client-protocol-tokio.workspace = true
fabro-model = { path = "../fabro-model" }
fabro-sandbox = { path = "../fabro-sandbox" }
fabro-types = { path = "../fabro-types" }
fabro-util = { path = "../fabro-util" }
bytes.workspace = true
serde.workspace = true
fabro-sandbox = { path = "../fabro-sandbox", optional = true }
fabro-types = { path = "../fabro-types", optional = true }
fabro-util = { path = "../fabro-util", optional = true }
serde_json.workspace = true
shlex = "1"
thiserror.workspace = true
tokio.workspace = true
tokio-util = { workspace = true, features = ["compat", "io"] }
futures.workspace = true
uuid.workspace = true
tracing.workspace = true
tokio = { workspace = true, optional = true }
tokio-util = { workspace = true, features = ["compat", "io"], optional = true }
futures = { workspace = true, optional = true }
tracing = { workspace = true, optional = true }
[dev-dependencies]
fabro-sandbox = { path = "../fabro-sandbox", features = ["test-support"] }

View file

@ -1,19 +1,102 @@
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::str::FromStr;
use agent_client_protocol::schema::McpServer;
use agent_client_protocol::schema::{McpServer, McpServerStdio};
use agent_client_protocol_tokio::AcpAgent;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AcpCommand {
display: String,
pub struct AcpProcessSpec {
name: Option<String>,
program: PathBuf,
args: Vec<String>,
env: HashMap<String, String>,
}
impl AcpCommand {
impl AcpProcessSpec {
pub fn from_attrs(
legacy_command: Option<&str>,
command: Option<&str>,
config: Option<&str>,
) -> Result<Self, AcpCommandError> {
if legacy_command.is_some() {
return Err(AcpCommandError::LegacyCommandAttribute);
}
match (command, config) {
(Some(command), None) => Self::from_command_attr(command),
(None, Some(config)) => Self::from_config_attr(config),
(None, None) | (Some(_), Some(_)) => Err(AcpCommandError::MissingOverride),
}
}
pub fn from_command_attr(raw: &str) -> Result<Self, AcpCommandError> {
let trimmed = raw.trim();
if trimmed.is_empty() {
return Err(AcpCommandError::EmptyOverride);
}
let parts = shlex::split(trimmed).ok_or(AcpCommandError::InvalidCommandString)?;
let agent =
AcpAgent::from_args(parts).map_err(|_| AcpCommandError::InvalidCommandString)?;
let mut spec = Self::from_server(agent.into_server())?;
spec.name = None;
Ok(spec)
}
pub fn from_config_attr(raw: &str) -> Result<Self, AcpCommandError> {
let trimmed = raw.trim();
if trimmed.is_empty() {
return Err(AcpCommandError::EmptyOverride);
}
let server = parse_config_server(trimmed)?;
Self::from_server(server)
}
fn from_server(server: McpServer) -> Result<Self, AcpCommandError> {
match server {
McpServer::Stdio(stdio) => Self::from_stdio_config(stdio),
_ => Err(AcpCommandError::UnsupportedTransport),
}
}
fn from_stdio_config(stdio: McpServerStdio) -> Result<Self, AcpCommandError> {
if stdio.command.as_os_str().is_empty() {
return Err(AcpCommandError::InvalidConfigShape("missing command"));
}
let env = stdio
.env
.into_iter()
.map(|env| (env.name, env.value))
.collect();
Ok(Self::from_stdio_parts(
Some(stdio.name),
stdio.command,
stdio.args,
env,
))
}
fn from_stdio_parts(
name: Option<String>,
program: PathBuf,
args: Vec<String>,
env: HashMap<String, String>,
) -> Self {
Self {
name,
program,
args,
env,
}
}
#[must_use]
pub fn name(&self) -> Option<&str> {
self.name.as_deref()
}
#[must_use]
pub fn program(&self) -> &Path {
&self.program
@ -29,101 +112,70 @@ impl AcpCommand {
&self.env
}
#[must_use]
pub fn display(&self) -> &str {
&self.display
}
#[must_use]
pub fn to_shell_command(&self) -> String {
render_command(&self.program, &self.args)
}
}
impl std::fmt::Display for AcpCommand {
impl std::fmt::Display for AcpProcessSpec {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.display)
f.write_str(&self.to_shell_command())
}
}
#[derive(Debug, thiserror::Error)]
pub enum AcpCommandError {
#[error("acp_command must not be empty")]
#[error("acp_command is no longer supported; use acp.command or acp.config")]
LegacyCommandAttribute,
#[error("ACP process attribute must not be empty")]
EmptyOverride,
#[error(
"acp_command is required for backend=\"acp\" because Fabro does not install ACP agents"
)]
#[error("backend=\"acp\" requires exactly one of acp.command or acp.config")]
MissingOverride,
#[error("only stdio ACP commands are supported")]
UnsupportedTransport,
#[error("failed to parse acp_command")]
Parse(#[source] agent_client_protocol::Error),
}
impl From<agent_client_protocol::Error> for AcpCommandError {
fn from(error: agent_client_protocol::Error) -> Self {
Self::Parse(error)
}
}
pub fn resolve_acp_command(override_command: Option<&str>) -> Result<AcpCommand, AcpCommandError> {
if let Some(raw) = override_command {
let trimmed = raw.trim();
if trimmed.is_empty() {
return Err(AcpCommandError::EmptyOverride);
}
return parse_acp_command(trimmed);
}
Err(AcpCommandError::MissingOverride)
}
fn parse_acp_command(raw: &str) -> Result<AcpCommand, AcpCommandError> {
reject_non_stdio_json_transport(raw)?;
let agent = AcpAgent::from_str(raw)?;
let McpServer::Stdio(stdio) = agent.into_server() else {
return Err(AcpCommandError::UnsupportedTransport);
};
let program = stdio.command;
let args = stdio.args;
let display = render_command(&program, &args);
Ok(AcpCommand {
display,
program,
args,
env: stdio
.env
.into_iter()
.map(|env| (env.name, env.value))
.collect(),
})
#[error("failed to parse acp.command as a shell command")]
InvalidCommandString,
#[error("failed to parse acp.config as JSON")]
InvalidConfigJson(#[source] serde_json::Error),
#[error("invalid acp.config shape: {0}")]
InvalidConfigShape(&'static str),
}
fn render_command(program: &Path, args: &[String]) -> String {
std::iter::once(program.to_string_lossy().into_owned())
.chain(args.iter().cloned())
.map(|part| fabro_sandbox::shell_quote(&part))
.map(|part| shell_quote(&part))
.collect::<Vec<_>>()
.join(" ")
}
fn reject_non_stdio_json_transport(raw: &str) -> Result<(), AcpCommandError> {
let trimmed = raw.trim_start();
if !trimmed.starts_with('{') {
return Ok(());
}
let Ok(value) = serde_json::from_str::<serde_json::Value>(trimmed) else {
return Ok(());
};
fn parse_config_server(raw: &str) -> Result<McpServer, AcpCommandError> {
let mut value: serde_json::Value =
serde_json::from_str(raw).map_err(AcpCommandError::InvalidConfigJson)?;
match value.get("type").and_then(serde_json::Value::as_str) {
Some("stdio") | None => Ok(()),
Some(_) => Err(AcpCommandError::UnsupportedTransport),
Some("stdio") | None => {}
Some(_) => return Err(AcpCommandError::UnsupportedTransport),
}
if let Some(object) = value.as_object_mut() {
object
.entry("args".to_string())
.or_insert_with(|| serde_json::Value::Array(Vec::new()));
object
.entry("env".to_string())
.or_insert_with(|| serde_json::Value::Array(Vec::new()));
}
serde_json::from_value(value).map_err(AcpCommandError::InvalidConfigJson)
}
fn shell_quote(s: &str) -> String {
shlex::try_quote(s).map_or_else(
|_| format!("'{}'", s.replace('\'', "'\\''")),
|quoted| quoted.to_string(),
)
}
#[cfg(test)]
@ -133,41 +185,53 @@ mod tests {
use super::*;
#[test]
fn missing_acp_command_is_rejected() {
let err = resolve_acp_command(None).unwrap_err();
assert!(
err.to_string()
.contains("acp_command is required for backend=\"acp\"")
);
}
#[test]
fn explicit_acp_command_overrides_provider_default() {
let command = resolve_acp_command(Some("python fake_agent.py")).unwrap();
fn command_attr_parses_shell_command() {
let command = AcpProcessSpec::from_command_attr("python fake_agent.py").unwrap();
assert_eq!(command.to_string(), "python fake_agent.py");
assert_eq!(command.name(), None);
assert_eq!(command.program(), Path::new("python"));
assert_eq!(command.args(), &["fake_agent.py".to_string()]);
}
#[test]
fn blank_acp_command_is_rejected() {
let err = resolve_acp_command(Some(" ")).unwrap_err();
assert!(err.to_string().contains("acp_command must not be empty"));
fn command_attr_parses_leading_env_assignments() {
let command = AcpProcessSpec::from_command_attr(
"RUST_LOG=debug TOKEN='secret value' python fake_agent.py",
)
.unwrap();
assert_eq!(command.to_string(), "python fake_agent.py");
assert_eq!(command.program(), Path::new("python"));
assert_eq!(
command.env().get("RUST_LOG").map(String::as_str),
Some("debug")
);
assert_eq!(
command.env().get("TOKEN").map(String::as_str),
Some("secret value")
);
}
#[test]
fn json_stdio_acp_command_is_supported() {
fn blank_acp_process_attr_is_rejected() {
let err = AcpProcessSpec::from_command_attr(" ").unwrap_err();
assert!(err.to_string().contains("must not be empty"));
}
#[test]
fn json_stdio_acp_config_is_supported() {
let raw = r#"{"type":"stdio","name":"fake","command":"python","args":["fake agent.py"],"env":[{"name":"MODE","value":"test"}]}"#;
let command = resolve_acp_command(Some(raw)).unwrap();
let command = AcpProcessSpec::from_config_attr(raw).unwrap();
assert_eq!(command.name(), Some("fake"));
assert_eq!(command.program(), Path::new("python"));
assert_eq!(command.args(), &["fake agent.py".to_string()]);
assert_eq!(command.env().get("MODE").map(String::as_str), Some("test"));
}
#[test]
fn json_stdio_acp_command_display_omits_env_contents() {
fn json_stdio_acp_config_display_omits_env_contents() {
let raw = r#"{"type":"stdio","name":"fake","command":"agent","args":["--flag","two words"],"env":[{"name":"OPENAI_API_KEY","value":"secret-key"}]}"#;
let command = resolve_acp_command(Some(raw)).unwrap();
let command = AcpProcessSpec::from_config_attr(raw).unwrap();
assert_eq!(
command.env().get("OPENAI_API_KEY").map(String::as_str),
@ -179,12 +243,40 @@ mod tests {
}
#[test]
fn non_stdio_acp_command_is_rejected() {
fn non_stdio_acp_config_is_rejected() {
let raw = r#"{"type":"http","name":"remote","url":"https://example.test/acp"}"#;
let err = resolve_acp_command(Some(raw)).unwrap_err();
let err = AcpProcessSpec::from_config_attr(raw).unwrap_err();
assert!(
err.to_string()
.contains("only stdio ACP commands are supported")
);
}
#[test]
fn command_attr_is_always_shell_command_even_when_json_shaped() {
let command = AcpProcessSpec::from_command_attr(r#"{"type":"stdio"}"#).unwrap();
assert_ne!(command.program(), Path::new("stdio"));
assert!(command.args().is_empty());
}
#[test]
fn config_attr_requires_json_stdio_config() {
let command = AcpProcessSpec::from_config_attr(
r#"{"type":"stdio","name":"fake","command":"python3","args":["agent.py"]}"#,
)
.unwrap();
assert_eq!(command.name(), Some("fake"));
assert_eq!(command.program(), Path::new("python3"));
assert_eq!(command.args(), &["agent.py".to_string()]);
assert!(AcpProcessSpec::from_config_attr("python3 agent.py").is_err());
assert!(
AcpProcessSpec::from_config_attr(
r#"{"type":"http","name":"remote","url":"https://example.test/acp"}"#
)
.is_err()
);
}
}

View file

@ -1,12 +1,18 @@
pub mod command;
#[cfg(feature = "runtime")]
pub mod error;
#[cfg(feature = "runtime")]
pub mod session;
#[cfg(any(test, feature = "test-support"))]
pub mod test_support;
#[cfg(feature = "runtime")]
mod transport;
pub use command::{AcpCommand, AcpCommandError, resolve_acp_command};
pub use command::{AcpCommandError, AcpProcessSpec};
#[cfg(feature = "runtime")]
pub use error::{AcpError, AcpProcessExit};
#[cfg(feature = "runtime")]
pub use session::{AcpRunRequest, AcpRunResult, render_stop_reason, run_acp_turn};

View file

@ -14,12 +14,12 @@ use fabro_util::time::elapsed_ms;
use tokio::time::{sleep, timeout};
use tokio_util::sync::CancellationToken;
use crate::command::AcpCommand;
use crate::command::AcpProcessSpec;
use crate::error::AcpError;
use crate::transport::{SandboxAcpTransport, TransportState};
pub struct AcpRunRequest {
pub command: AcpCommand,
pub command: AcpProcessSpec,
pub prompt: String,
pub cwd: String,
pub timeout_ms: Option<u64>,

View file

@ -35,6 +35,19 @@ if os.environ.get("ACP_PID_RECORD"):
with open(os.environ["ACP_PID_RECORD"], "w", encoding="utf-8") as record:
record.write(str(os.getpid()))
if os.environ.get("ACP_ENV_RECORD"):
keys = [
key.strip()
for key in os.environ.get(
"ACP_ENV_RECORD_KEYS",
"ANTHROPIC_API_KEY,OPENAI_API_KEY,GEMINI_API_KEY",
).split(",")
if key.strip()
]
snapshot = {key: os.environ[key] for key in keys if key in os.environ}
with open(os.environ["ACP_ENV_RECORD"], "w", encoding="utf-8") as record:
record.write(json.dumps(snapshot, sort_keys=True))
def handle_sigterm(signum, frame):
if os.environ.get("ACP_LINGER_TERMINATED"):
with open(os.environ["ACP_LINGER_TERMINATED"], "w", encoding="utf-8") as record:

View file

@ -20,7 +20,7 @@ use tokio::sync::Mutex as TokioMutex;
use tokio::time::timeout;
use tokio_util::compat::{TokioAsyncReadCompatExt, TokioAsyncWriteCompatExt};
use crate::command::AcpCommand;
use crate::command::AcpProcessSpec;
use crate::error::AcpProcessExit;
const CLEAN_EXIT_PROTOCOL_GRACE: Duration = Duration::from_millis(500);
@ -89,7 +89,7 @@ impl TransportState {
}
pub(crate) struct SandboxAcpTransport {
command: AcpCommand,
command: AcpProcessSpec,
cwd: String,
env: HashMap<String, String>,
sandbox: Arc<dyn Sandbox>,
@ -98,7 +98,7 @@ pub(crate) struct SandboxAcpTransport {
impl SandboxAcpTransport {
pub(crate) fn new(
command: AcpCommand,
command: AcpProcessSpec,
cwd: String,
env: HashMap<String, String>,
sandbox: Arc<dyn Sandbox>,

View file

@ -4,7 +4,7 @@ use std::sync::Arc;
use std::time::Duration;
use agent_client_protocol::schema::StopReason;
use fabro_acp::{AcpError, AcpRunRequest, AcpRunResult, resolve_acp_command, run_acp_turn};
use fabro_acp::{AcpError, AcpProcessSpec, AcpRunRequest, AcpRunResult, run_acp_turn};
use fabro_sandbox::test_support::{MockSandbox, MockStdioProcess};
use fabro_sandbox::{LocalSandbox, Sandbox, shell_quote};
use fabro_util::error::collect_chain;
@ -31,7 +31,7 @@ use test_support::fake_acp_agent_script;
async fn stdio_spawn_failure_returns_sandbox_error() {
const SANDBOX_FAILURE: &str = "ACP backend requires bidirectional stdio; the Daytona sandbox provider does not support it yet";
let command = resolve_acp_command(Some("fake-acp-agent")).expect("resolve ACP command");
let command = AcpProcessSpec::from_command_attr("fake-acp-agent").expect("parse ACP command");
let mut sandbox = MockSandbox::linux();
sandbox.stdio_process_error = Some(SANDBOX_FAILURE.to_string());
let sandbox: Arc<dyn Sandbox> = Arc::new(sandbox);
@ -67,7 +67,7 @@ async fn clean_stdio_exit_after_final_response_completes_turn() {
let sandbox = MockSandbox::linux();
sandbox.set_stdio_process(mock_acp_stdio_process("end_turn"));
let sandbox: Arc<dyn Sandbox> = Arc::new(sandbox);
let command = resolve_acp_command(Some("mock-acp-agent")).expect("resolve ACP command");
let command = AcpProcessSpec::from_command_attr("mock-acp-agent").expect("parse ACP command");
let result = run_acp_turn(AcpRunRequest {
command,
@ -96,7 +96,7 @@ async fn session_lifecycle_initializes_sends_prompt_and_aggregates_text() {
.expect("write fake ACP agent");
let raw_command = format!("python3 {}", shell_quote(&script_path.to_string_lossy()));
let command = resolve_acp_command(Some(&raw_command)).expect("resolve ACP command");
let command = AcpProcessSpec::from_command_attr(&raw_command).expect("parse ACP command");
let sandbox: Arc<dyn Sandbox> = Arc::new(LocalSandbox::new(tempdir.path().to_path_buf()));
let result = run_acp_turn(AcpRunRequest {
@ -461,7 +461,7 @@ async fn run_fake_agent_with_activity(
.await
.expect("write fake ACP agent");
let raw_command = format!("python3 {}", shell_quote(&script_path.to_string_lossy()));
let command = resolve_acp_command(Some(&raw_command)).expect("resolve ACP command");
let command = AcpProcessSpec::from_command_attr(&raw_command).expect("parse ACP command");
let sandbox: Arc<dyn Sandbox> = Arc::new(LocalSandbox::new(tempdir.to_path_buf()));
run_acp_turn(AcpRunRequest {

View file

@ -23,7 +23,6 @@ fabro-types = { path = "../fabro-types" }
fabro-vault = { path = "../fabro-vault" }
serde.workspace = true
serde_json.workspace = true
shlex = "1"
thiserror.workspace = true
tokio.workspace = true

View file

@ -16,8 +16,8 @@ pub use credential_source::{CredentialSource, ResolvedCredentials};
pub use env_source::EnvCredentialSource;
pub use refresh::refresh_oauth_credential;
pub use resolve::{
ApiCredential, CliAgentKind, CliCredential, CredentialResolver, CredentialUsage, EnvLookup,
ResolveError, ResolvedCredential, auth_issue_message, build_api_key_header,
ApiCredential, CredentialResolver, CredentialUsage, EnvLookup, ResolveError,
ResolvedCredential, auth_issue_message, build_api_key_header,
configured_providers_from_process_env,
};
pub use strategy::{

View file

@ -5,7 +5,6 @@ use fabro_model::catalog::CatalogProvider;
use fabro_model::{ApiKeyHeaderPolicy, Catalog, CredentialRef, HeaderValueRef, ProviderId};
use fabro_static::EnvVars;
use fabro_vault::{SecretType, Vault};
use shlex::try_quote;
use tokio::sync::RwLock as AsyncRwLock;
use tokio::task::spawn_blocking;
@ -17,17 +16,9 @@ use crate::vault_ext::{VaultLookupError, vault_get_oauth, vault_get_token, vault
pub type EnvLookup = Arc<dyn Fn(&str) -> Option<String> + Send + Sync>;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CliAgentKind {
Claude,
Codex,
Gemini,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CredentialUsage {
ApiRequest,
CliAgent(CliAgentKind),
}
#[derive(Debug, Clone, PartialEq, Eq)]
@ -125,16 +116,9 @@ fn auth_header_for_catalog_provider(
Ok(build_api_key_header(auth.header.clone(), key))
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CliCredential {
pub env_vars: HashMap<String, String>,
pub login_command: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ResolvedCredential {
Api(ApiCredential),
Cli(CliCredential),
}
#[derive(Debug, thiserror::Error)]
@ -212,14 +196,14 @@ impl CredentialResolver {
pub async fn resolve(
&self,
provider: impl Into<ProviderId>,
usage: CredentialUsage,
_usage: CredentialUsage,
catalog: &Catalog,
) -> Result<ResolvedCredential, ResolveError> {
let provider_id = provider.into();
let Some(catalog_provider) = catalog.provider(&provider_id) else {
return Err(ResolveError::NotConfigured(provider_id));
};
if usage == CredentialUsage::ApiRequest && catalog_provider.auth.is_none() {
if catalog_provider.auth.is_none() {
let vault = self.vault.read().await;
return self
.api_credential_from_provider_auth(&vault, catalog_provider, catalog)
@ -274,14 +258,8 @@ impl CredentialResolver {
};
let vault = self.vault.read().await;
match usage {
CredentialUsage::ApiRequest => self
.to_api_credential(&vault, &provider_id, &secret, catalog)
.map(ResolvedCredential::Api),
CredentialUsage::CliAgent(kind) => Ok(ResolvedCredential::Cli(
Self::to_cli_credential(&provider_id, &secret, kind, catalog),
)),
}
self.to_api_credential(&vault, &provider_id, &secret, catalog)
.map(ResolvedCredential::Api)
}
#[must_use]
@ -468,49 +446,6 @@ impl CredentialResolver {
project_id: None,
})
}
fn to_cli_credential(
provider_id: &ProviderId,
secret: &ResolvedSecret,
kind: CliAgentKind,
catalog: &Catalog,
) -> CliCredential {
let mut env_vars = HashMap::new();
let is_openai = provider_id == &ProviderId::openai();
let login_command = match (is_openai, secret, kind) {
(true, ResolvedSecret::ApiKey(key), CliAgentKind::Codex) => {
env_vars.insert(EnvVars::OPENAI_API_KEY.to_string(), key.clone());
Some(codex_login_command(key))
}
(true, ResolvedSecret::OAuth { credential, .. }, CliAgentKind::Codex) => {
env_vars.insert(
EnvVars::OPENAI_API_KEY.to_string(),
credential.tokens.access_token.clone(),
);
if let Some(account_id) = &credential.account_id {
env_vars.insert(EnvVars::CHATGPT_ACCOUNT_ID.to_string(), account_id.clone());
}
Some(codex_login_command(&credential.tokens.access_token))
}
(_, ResolvedSecret::ApiKey(key), _) => {
if let Some(name) = primary_api_key_env_var(provider_id, catalog) {
env_vars.insert(name.to_string(), key.clone());
}
None
}
(_, ResolvedSecret::OAuth { credential, .. }, _) => {
if let Some(name) = primary_api_key_env_var(provider_id, catalog) {
env_vars.insert(name.to_string(), credential.tokens.access_token.clone());
}
None
}
};
CliCredential {
env_vars,
login_command,
}
}
}
fn vault_lookup_error(provider: &ProviderId, name: &str, err: VaultLookupError) -> ResolveError {
@ -545,32 +480,8 @@ pub async fn configured_providers_from_process_env(
}
}
}
fn primary_api_key_env_var<'a>(provider: &ProviderId, catalog: &'a Catalog) -> Option<&'a str> {
catalog
.provider(provider)?
.auth
.as_ref()?
.credentials
.iter()
.find_map(|credential_ref| match credential_ref {
CredentialRef::Env(name) => Some(name.as_str()),
CredentialRef::Vault(_) => None,
})
}
fn codex_login_command(api_key: &str) -> String {
let quoted =
try_quote(api_key).map_or_else(|_| api_key.to_string(), std::borrow::Cow::into_owned);
format!(
"export PATH=\"$HOME/.local/bin:$PATH\" && printf '%s\\n' {quoted} | codex login --with-api-key"
)
}
#[cfg(test)]
mod tests {
#[cfg(unix)]
use std::os::unix::fs::PermissionsExt;
use chrono::{Duration, Utc};
use fabro_model::catalog::LlmCatalogSettings;
use httpmock::Method::POST;
@ -628,9 +539,7 @@ mod tests {
.await
.unwrap();
let ResolvedCredential::Api(api) = resolved else {
panic!("expected api credential");
};
let ResolvedCredential::Api(api) = resolved;
assert_eq!(
api.auth_header,
Some(ApiKeyHeader::Bearer("env-key".to_string()))
@ -658,9 +567,7 @@ mod tests {
.await
.unwrap();
let ResolvedCredential::Api(api) = resolved else {
panic!("expected api credential");
};
let ResolvedCredential::Api(api) = resolved;
assert_eq!(
api.auth_header,
Some(ApiKeyHeader::Bearer("expired-access".to_string()))
@ -702,17 +609,15 @@ mod tests {
let resolver = test_resolver(vault, Arc::new(|_| None));
let catalog = default_catalog();
let ResolvedCredential::Api(api) = resolver
let resolved = resolver
.resolve(
ProviderId::anthropic(),
CredentialUsage::ApiRequest,
&catalog,
)
.await
.unwrap()
else {
panic!("expected api credential");
};
.unwrap();
let ResolvedCredential::Api(api) = resolved;
assert_eq!(
api.auth_header,
@ -764,9 +669,7 @@ reasoning = false
.await
.unwrap();
let ResolvedCredential::Api(api) = resolved else {
panic!("expected api credential");
};
let ResolvedCredential::Api(api) = resolved;
assert_eq!(
api.auth_header,
Some(ApiKeyHeader::Bearer("compat-key".to_string()))
@ -777,141 +680,6 @@ reasoning = false
);
}
#[tokio::test]
async fn openai_codex_cli_credential_includes_login_command_and_account_id() {
let dir = tempfile::tempdir().unwrap();
let mut vault = Vault::load(dir.path().join("secrets.json")).unwrap();
vault_set_oauth(
&mut vault,
crate::OPENAI_CODEX_VAULT_SECRET_NAME,
&oauth_credential(
"https://auth.openai.com/oauth/token".to_string(),
Utc::now() + Duration::hours(1),
),
)
.unwrap();
let resolver = test_resolver(vault, Arc::new(|_| None));
let catalog = default_catalog();
let ResolvedCredential::Cli(cli) = resolver
.resolve(
ProviderId::openai(),
CredentialUsage::CliAgent(CliAgentKind::Codex),
&catalog,
)
.await
.unwrap()
else {
panic!("expected cli credential");
};
assert_eq!(
cli.env_vars.get("OPENAI_API_KEY").map(String::as_str),
Some("expired-access")
);
assert_eq!(
cli.env_vars.get("CHATGPT_ACCOUNT_ID").map(String::as_str),
Some("acct_123")
);
assert!(
cli.login_command
.as_deref()
.is_some_and(|command| command.contains("codex login --with-api-key"))
);
}
#[tokio::test]
async fn openai_api_key_cli_fallback_has_no_account_id() {
let dir = tempfile::tempdir().unwrap();
let mut vault = Vault::load(dir.path().join("secrets.json")).unwrap();
vault_set_token(&mut vault, "OPENAI_API_KEY", "openai-key").unwrap();
let resolver = test_resolver(vault, Arc::new(|_| None));
let catalog = default_catalog();
let ResolvedCredential::Cli(cli) = resolver
.resolve(
ProviderId::openai(),
CredentialUsage::CliAgent(CliAgentKind::Codex),
&catalog,
)
.await
.unwrap()
else {
panic!("expected cli credential");
};
assert_eq!(
cli.env_vars.get("OPENAI_API_KEY").map(String::as_str),
Some("openai-key")
);
assert!(!cli.env_vars.contains_key("CHATGPT_ACCOUNT_ID"));
assert!(cli.login_command.is_some());
}
#[cfg(unix)]
#[tokio::test]
#[expect(
clippy::disallowed_methods,
reason = "integration-style test: writes and reads a fake codex script via sync std::fs to \
verify the login_command string passes stdin correctly"
)]
async fn openai_api_key_cli_login_command_executes_codex_from_local_bin() {
let dir = tempfile::tempdir().unwrap();
let local_bin = dir.path().join(".local/bin");
std::fs::create_dir_all(&local_bin).unwrap();
let codex_path = local_bin.join("codex");
std::fs::write(
&codex_path,
"#!/bin/sh\nprintf '%s\\n' \"$@\" > \"$HOME/codex-args.txt\"\ncat > \"$HOME/codex-stdin.txt\"\n",
)
.unwrap();
let mut permissions = std::fs::metadata(&codex_path).unwrap().permissions();
permissions.set_mode(0o755);
std::fs::set_permissions(&codex_path, permissions).unwrap();
let mut vault = Vault::load(dir.path().join("secrets.json")).unwrap();
vault_set_token(&mut vault, "OPENAI_API_KEY", "openai-key").unwrap();
let resolver = test_resolver(vault, Arc::new(|_| None));
let catalog = default_catalog();
let ResolvedCredential::Cli(cli) = resolver
.resolve(
ProviderId::openai(),
CredentialUsage::CliAgent(CliAgentKind::Codex),
&catalog,
)
.await
.unwrap()
else {
panic!("expected cli credential");
};
#[allow(
clippy::disallowed_methods,
reason = "This test shells through /bin/sh to verify the configured login command."
)]
let status = std::process::Command::new("/bin/sh")
.arg("-lc")
.arg(cli.login_command.unwrap())
.env("HOME", dir.path())
.env("PATH", "/usr/bin:/bin")
.status()
.unwrap();
assert!(status.success());
assert_eq!(
std::fs::read_to_string(dir.path().join("codex-args.txt")).unwrap(),
"login\n--with-api-key\n"
);
assert_eq!(
std::fs::read_to_string(dir.path().join("codex-stdin.txt"))
.unwrap()
.trim_end(),
"openai-key"
);
}
#[tokio::test]
async fn with_env_lookup_overrides_vault_settings() {
let dir = tempfile::tempdir().unwrap();
@ -935,13 +703,11 @@ reasoning = false
);
let catalog = default_catalog();
let ResolvedCredential::Api(api) = resolver
let resolved = resolver
.resolve(ProviderId::openai(), CredentialUsage::ApiRequest, &catalog)
.await
.unwrap()
else {
panic!("expected api credential");
};
.unwrap();
let ResolvedCredential::Api(api) = resolved;
assert_eq!(api.org_id.as_deref(), Some("env-org"));
}
@ -1002,9 +768,7 @@ reasoning = false
.await
.unwrap();
let ResolvedCredential::Api(api) = resolved else {
panic!("expected api credential");
};
let ResolvedCredential::Api(api) = resolved;
assert_eq!(api.provider, ProviderId::new("acme"));
assert_eq!(
api.auth_header,
@ -1068,22 +832,17 @@ reasoning = false
let resolver = CredentialResolver::with_env_lookup(Arc::clone(&vault), Arc::new(|_| None));
let catalog = default_catalog();
let ResolvedCredential::Cli(cli) = resolver
.resolve(
ProviderId::openai(),
CredentialUsage::CliAgent(CliAgentKind::Codex),
&catalog,
)
let resolved = resolver
.resolve(ProviderId::openai(), CredentialUsage::ApiRequest, &catalog)
.await
.unwrap()
else {
panic!("expected cli credential");
};
.unwrap();
let ResolvedCredential::Api(api) = resolved;
assert_eq!(
cli.env_vars.get("OPENAI_API_KEY").map(String::as_str),
Some("new-access")
api.auth_header,
Some(ApiKeyHeader::Bearer("new-access".to_string()))
);
assert!(api.codex_mode);
let stored = {
let vault = vault.read().await;
@ -1116,11 +875,7 @@ reasoning = false
let catalog = default_catalog();
let err = resolver
.resolve(
ProviderId::openai(),
CredentialUsage::CliAgent(CliAgentKind::Codex),
&catalog,
)
.resolve(ProviderId::openai(), CredentialUsage::ApiRequest, &catalog)
.await
.unwrap_err();

View file

@ -52,7 +52,6 @@ impl CredentialSource for VaultCredentialSource {
.await
{
Ok(ResolvedCredential::Api(credential)) => credentials.push(credential),
Ok(ResolvedCredential::Cli(_)) => {}
Err(ResolveError::NotConfigured(_)) if provider.auth.is_some() => {}
Err(err) => auth_issues.push((provider.id.clone(), err)),
}

View file

@ -19,9 +19,8 @@ fn acp_backend_workflow() {
"[server.auth]\nmethods = [\"dev-token\"]\n",
);
context.isolated_server();
seed_openai_vault(&context.storage_dir);
let fake_agent = write_fake_acp_agent(&context);
let acp_command = fake_acp_command_attr(&fake_agent);
let acp_config = fake_acp_config_attr(&fake_agent);
let workflow = context.temp_dir.join("acp_backend.fabro");
context.write_temp(
"acp_backend.fabro",
@ -29,7 +28,7 @@ fn acp_backend_workflow() {
r#"digraph ACP {{
graph [goal="Exercise ACP backend"]
start [shape=Mdiamond]
work [type="agent", backend="acp", provider="openai", model="fake-acp", prompt="write hello.txt", acp_command={acp_command}]
work [type="agent", backend="acp", prompt="write hello.txt", acp.config={acp_config}]
exit [shape=Msquare]
start -> work
work -> exit
@ -78,34 +77,37 @@ fn acp_backend_workflow() {
assert!(
stages.values().any(|stage| {
stage["provider_used"]["mode"] == "acp"
&& stage["provider_used"]["provider"] == "openai"
&& stage["provider_used"]["config_name"] == "fake"
&& stage["provider_used"].get("provider").is_none()
}),
"run projection should include ACP provider metadata: {stages:?}"
"run projection should include ACP process metadata without provider: {stages:?}"
);
}
#[test]
fn acp_prompt_workflow_uses_acp_backend() {
fn acp_backend_does_not_inject_registered_provider_credentials() {
let mut context = test_context!();
context.write_home(
".fabro/settings.toml",
"[server.auth]\nmethods = [\"dev-token\"]\n",
);
context.isolated_server();
seed_openai_vault(&context.storage_dir);
seed_anthropic_vault(&context.storage_dir);
let fake_agent = write_fake_acp_agent(&context);
let acp_command = fake_acp_command_attr(&fake_agent);
let workflow = context.temp_dir.join("acp_prompt_backend.fabro");
let env_record = context.temp_dir.join("acp-env.json");
let acp_config = fake_acp_config_attr_recording_env(&fake_agent, &env_record);
let workflow = context.temp_dir.join("acp_provider_env.fabro");
context.write_temp(
"acp_prompt_backend.fabro",
"acp_provider_env.fabro",
format!(
r#"digraph ACP {{
graph [goal="Exercise ACP prompt backend"]
graph [goal="Exercise ACP backend with stored provider credentials"]
start [shape=Mdiamond]
prompt [type="prompt", backend="acp", provider="openai", model="fake-acp", project_memory=false, prompt="write hello.txt", acp_command={acp_command}]
work [type="agent", backend="acp", prompt="write hello.txt", acp.config={acp_config}]
exit [shape=Msquare]
start -> prompt
prompt -> exit
start -> work
work -> exit
}}"#
),
);
@ -113,6 +115,9 @@ fn acp_prompt_workflow_uses_acp_backend() {
context
.run_cmd()
.env_remove("ANTHROPIC_API_KEY")
.env_remove("OPENAI_API_KEY")
.env_remove("GEMINI_API_KEY")
.args(["--auto-approve", "--sandbox", "local"])
.arg(&workflow)
.assert()
@ -122,45 +127,12 @@ fn acp_prompt_workflow_uses_acp_backend() {
let conclusion = read_conclusion(&run_dir);
assert_eq!(conclusion["status"].as_str(), Some("succeeded"));
let events = run_events(&run_dir);
assert!(has_event(&run_dir, "agent.acp.started"));
assert!(has_event(&run_dir, "agent.acp.completed"));
assert!(
!has_event(&run_dir, "agent.session.activated"),
"ACP prompt should not activate an API-mode agent session"
);
let completed = events
.iter()
.find_map(|event| match &event.event.body {
EventBody::StageCompleted(props)
if event.event.node_id.as_deref() == Some("prompt") =>
{
Some(props)
}
_ => None,
})
.expect("prompt stage should complete");
assert_eq!(completed.response.as_deref(), Some("hello from acp"));
let state = serde_json::to_value(run_state(&run_dir)).expect("run state should serialize");
let stages = state["stages"]
.as_object()
.expect("run state should contain stages");
assert!(
stages.values().any(|stage| {
stage["provider_used"]["mode"] == "acp"
&& stage["provider_used"]["provider"] == "openai"
}),
"run projection should include ACP provider metadata: {stages:?}"
);
}
fn seed_openai_vault(storage_dir: &std::path::Path) {
let mut vault =
Vault::load(Storage::new(storage_dir).secrets_path()).expect("test vault should load");
vault
.set("OPENAI_API_KEY", "test-openai-key", SecretType::Token, None)
.expect("OpenAI credential should store in test vault");
let recorded_env: serde_json::Value = serde_json::from_str(
&std::fs::read_to_string(&env_record)
.expect("fake ACP agent should record its environment"),
)
.expect("fake ACP environment record should be valid JSON");
assert_eq!(recorded_env, serde_json::json!({}));
}
fn write_fake_acp_agent(context: &fabro_test::TestContext) -> std::path::PathBuf {
@ -168,16 +140,52 @@ fn write_fake_acp_agent(context: &fabro_test::TestContext) -> std::path::PathBuf
context.temp_dir.join("fake_acp_agent.py")
}
fn fake_acp_command_attr(script_path: &std::path::Path) -> String {
let command = serde_json::json!({
fn fake_acp_config_attr(script_path: &std::path::Path) -> String {
fake_acp_config_attr_with_env(script_path, Vec::new())
}
fn fake_acp_config_attr_recording_env(
script_path: &std::path::Path,
env_record: &std::path::Path,
) -> String {
fake_acp_config_attr_with_env(script_path, vec![
serde_json::json!({"name": "ACP_ENV_RECORD", "value": env_record.to_string_lossy()}),
serde_json::json!({
"name": "ACP_ENV_RECORD_KEYS",
"value": "ANTHROPIC_API_KEY,OPENAI_API_KEY,GEMINI_API_KEY",
}),
])
}
fn fake_acp_config_attr_with_env(
script_path: &std::path::Path,
extra_env: Vec<serde_json::Value>,
) -> String {
let mut env = vec![serde_json::json!({"name": "ACP_MODE", "value": "write_file"})];
env.extend(extra_env);
let config = serde_json::json!({
"type": "stdio",
"name": "fake",
"command": "python3",
"args": [script_path.to_string_lossy()],
"env": [{"name": "ACP_MODE", "value": "write_file"}],
"env": env,
})
.to_string();
format!("{command:?}")
format!("{config:?}")
}
fn seed_anthropic_vault(storage_dir: &std::path::Path) {
let mut vault =
Vault::load(Storage::new(storage_dir).secrets_path()).expect("test vault should load");
vault
.set(
"ANTHROPIC_API_KEY",
"vault-anthropic-key",
SecretType::Token,
None,
)
.expect("Anthropic credential should store in test vault");
}
fn init_git_repo(dir: &std::path::Path) {

View file

@ -12,7 +12,6 @@ mod dry_run_examples;
mod full_stack;
mod hooks;
mod human_gate;
mod real_cli;
use std::path::{Path, PathBuf};
use std::time::Duration;

View file

@ -1,71 +0,0 @@
use std::sync::Arc;
use fabro_graphviz::graph::{AttrValue, Node};
use fabro_model::ProviderId;
use fabro_workflow::context::Context;
use fabro_workflow::event::Emitter;
use fabro_workflow::handler::agent::{CodergenBackend, CodergenResult, CodergenRunRequest};
use fabro_workflow::handler::llm::cli::AgentCliBackend;
/// Run a real CLI tool via LocalSandbox and verify the full flow.
async fn run_real_cli_test(provider: ProviderId, model: &str) {
let workspace = tempfile::tempdir().expect("real CLI test workspace should create");
let env: Arc<dyn fabro_agent::Sandbox> = Arc::new(fabro_agent::LocalSandbox::new(
workspace.path().to_path_buf(),
));
let backend = AgentCliBackend::new_from_env(model.to_string(), provider.clone());
let mut node = Node::new("real_cli_test");
node.attrs.insert(
"prompt".to_string(),
AttrValue::String("What is 2+2? Reply with just the number.".to_string()),
);
let context = Context::new();
let emitter = Arc::new(Emitter::default());
let result = backend
.run(CodergenRunRequest {
node: &node,
prompt: "What is 2+2? Reply with just the number.",
context: &context,
thread_id: None,
emitter: &emitter,
sandbox: &env,
tool_hooks: None,
cancel_token: tokio_util::sync::CancellationToken::new(),
})
.await
.unwrap_or_else(|_| panic!("CLI backend ({provider}/{model}) should succeed"));
match result {
CodergenResult::Text { text, usage, .. } => {
assert!(
text.contains('4'),
"{provider}/{model}: expected response to contain '4', got: {text}"
);
let usage = usage.unwrap_or_else(|| panic!("{provider}/{model}: should have usage"));
let tokens = usage.tokens();
assert!(
tokens.input_tokens > 0,
"{provider}/{model}: input_tokens should be > 0, got {}",
tokens.input_tokens
);
}
CodergenResult::Full(_) => panic!("expected Text result from {provider}/{model}"),
}
}
#[fabro_macros::e2e_test(live("ANTHROPIC_API_KEY"))]
async fn real_cli_claude() {
run_real_cli_test(ProviderId::anthropic(), "haiku").await;
}
#[fabro_macros::e2e_test(live("OPENAI_API_KEY"))]
async fn real_cli_codex() {
run_real_cli_test(ProviderId::openai(), "").await;
}
#[fabro_macros::e2e_test(live("GEMINI_API_KEY"))]
async fn real_cli_gemini() {
run_real_cli_test(ProviderId::gemini(), "gemini-2.5-flash").await;
}

View file

@ -2483,18 +2483,17 @@ fn update_live_run_from_event(state: &AppState, run_id: RunId, event: &RunEvent)
}
}
}
// Track non-steerable agent stages. CLI/ACP started/completed are
// Track non-steerable agent stages. ACP started/completed are
// coarser and sometimes fail to emit terminal events on error paths;
// stage.completed/stage.failed below are the backstops.
EventBody::AgentCliStarted(_) | EventBody::AgentAcpStarted(_) => {
EventBody::AgentAcpStarted(_) => {
if let Some(stage_id) = event.stage_id.as_ref() {
managed_run
.active_non_steerable_agent_stages
.insert(stage_id.clone());
}
}
EventBody::AgentCliCompleted(_)
| EventBody::AgentAcpCompleted(_)
EventBody::AgentAcpCompleted(_)
| EventBody::AgentAcpCancelled(_)
| EventBody::AgentAcpTimedOut(_) => {
if let Some(stage_id) = &event.stage_id {
@ -2504,7 +2503,7 @@ fn update_live_run_from_event(state: &AppState, run_id: RunId, event: &RunEvent)
}
}
// Stage lifecycle backstop: cover both completion and failure
// paths so a failing CLI stage doesn't strand its entry.
// paths so a failing ACP stage doesn't strand its entry.
EventBody::StageCompleted(_) | EventBody::StageFailed(_) => {
if let Some(stage_id) = &event.stage_id {
managed_run.active_api_stages.remove(stage_id);

View file

@ -8551,12 +8551,10 @@ async fn steer_with_active_acp_stage_returns_non_steerable_conflict() {
);
let started = acp_event_for_stage(&run_id, &workflow_event::Event::AgentAcpStarted {
node_id: "agent".to_string(),
visit: 1,
mode: "acp".to_string(),
provider: "openai".to_string(),
model: "fake-acp".to_string(),
command: "python fake_agent.py".to_string(),
node_id: "agent".to_string(),
visit: 1,
command: "python fake_agent.py".to_string(),
config_name: None,
});
update_live_run_from_event(&state, run_id, &started);
@ -8640,12 +8638,10 @@ async fn active_acp_stage_marker_clears_on_terminal_paths() {
Some(RunAnswerTransport::Subprocess { control_tx }),
);
let started = acp_event_for_stage(&run_id, &workflow_event::Event::AgentAcpStarted {
node_id: "agent".to_string(),
visit: 1,
mode: "acp".to_string(),
provider: "openai".to_string(),
model: "fake-acp".to_string(),
command: "python fake_agent.py".to_string(),
node_id: "agent".to_string(),
visit: 1,
command: "python fake_agent.py".to_string(),
config_name: None,
});
update_live_run_from_event(&state, run_id, &started);
let terminal = acp_event_for_stage(&run_id, &terminal_event);

View file

@ -3,9 +3,8 @@ use std::str::FromStr;
use chrono::{DateTime, Utc};
use fabro_types::run_event::{
AgentAcpStartedProps, AgentCliStartedProps, AgentSessionActivatedProps,
CheckpointCompletedProps, RunCompletedProps, RunFailedProps, StageCompletedProps,
StagePromptProps,
AgentAcpStartedProps, AgentSessionActivatedProps, CheckpointCompletedProps, RunCompletedProps,
RunFailedProps, StageCompletedProps, StagePromptProps,
};
use fabro_types::settings::run::RunSandboxSettings;
use fabro_types::{
@ -376,13 +375,6 @@ impl RunProjectionReducer for RunProjection {
};
stage.provider_used = Some(provider_used_from_agent_session_activated(props));
}
EventBody::AgentCliStarted(props) => {
let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq)
else {
return Ok(());
};
stage.provider_used = Some(provider_used_from_agent_cli_started(props));
}
EventBody::AgentAcpStarted(props) => {
let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq)
else {
@ -412,42 +404,6 @@ impl RunProjectionReducer for RunProjection {
stage.termination = Some(props.termination);
stage.script_timing = Some(script_timing);
}
EventBody::AgentCliCompleted(props) => {
let Some(stage) = stage_at_current_visit(self, stored, event.seq) else {
return Ok(());
};
apply_agent_terminal(
"agent.cli",
stage,
props,
merge_agent_cli_output(&props.stdout, &props.stderr),
CommandTermination::Exited,
)?;
}
EventBody::AgentCliCancelled(props) => {
let Some(stage) = stage_at_current_visit(self, stored, event.seq) else {
return Ok(());
};
apply_agent_terminal(
"agent.cli",
stage,
props,
merge_agent_cli_output(&props.stdout, &props.stderr),
CommandTermination::Cancelled,
)?;
}
EventBody::AgentCliTimedOut(props) => {
let Some(stage) = stage_at_current_visit(self, stored, event.seq) else {
return Ok(());
};
apply_agent_terminal(
"agent.cli",
stage,
props,
merge_agent_cli_output(&props.stdout, &props.stderr),
CommandTermination::TimedOut,
)?;
}
EventBody::AgentAcpCompleted(props) => {
let Some(stage) = stage_at_current_visit(self, stored, event.seq) else {
return Ok(());
@ -456,7 +412,7 @@ impl RunProjectionReducer for RunProjection {
"agent.acp",
stage,
props,
merge_agent_cli_output(&props.stdout, &props.stderr),
merge_agent_process_output(&props.stdout, &props.stderr),
CommandTermination::Exited,
)?;
}
@ -468,7 +424,7 @@ impl RunProjectionReducer for RunProjection {
"agent.acp",
stage,
props,
merge_agent_cli_output(&props.stdout, &props.stderr),
merge_agent_process_output(&props.stdout, &props.stderr),
CommandTermination::Cancelled,
)?;
}
@ -480,7 +436,7 @@ impl RunProjectionReducer for RunProjection {
"agent.acp",
stage,
props,
merge_agent_cli_output(&props.stdout, &props.stderr),
merge_agent_process_output(&props.stdout, &props.stderr),
CommandTermination::TimedOut,
)?;
}
@ -893,25 +849,13 @@ fn provider_used_from_agent_session_activated(props: &AgentSessionActivatedProps
Value::Object(provider_used)
}
fn provider_used_from_agent_cli_started(props: &AgentCliStartedProps) -> Value {
provider_used_from_agent_process_started("cli", &props.provider, &props.model, &props.command)
}
fn provider_used_from_agent_acp_started(props: &AgentAcpStartedProps) -> Value {
provider_used_from_agent_process_started("acp", &props.provider, &props.model, &props.command)
}
fn provider_used_from_agent_process_started(
mode: &str,
provider: &str,
model: &str,
command: &str,
) -> Value {
let mut provider_used = serde_json::Map::new();
provider_used.insert("mode".to_string(), Value::String(mode.to_string()));
provider_used.insert("provider".to_string(), Value::String(provider.to_string()));
provider_used.insert("model".to_string(), Value::String(model.to_string()));
provider_used.insert("command".to_string(), Value::String(command.to_string()));
provider_used.insert("mode".to_string(), Value::String("acp".to_string()));
provider_used.insert("command".to_string(), Value::String(props.command.clone()));
if let Some(config_name) = props.config_name.clone() {
provider_used.insert("config_name".to_string(), Value::String(config_name));
}
Value::Object(provider_used)
}
@ -931,7 +875,7 @@ fn apply_agent_terminal(
Ok(())
}
fn merge_agent_cli_output(stdout: &str, stderr: &str) -> String {
fn merge_agent_process_output(stdout: &str, stderr: &str) -> String {
match (stdout.is_empty(), stderr.is_empty()) {
(true, true) => String::new(),
(false, true) => stdout.to_string(),
@ -948,8 +892,7 @@ mod tests {
use fabro_types::run_event::run::RunFailedProps;
use fabro_types::run_event::{
AgentAcpCancelledProps, AgentAcpCompletedProps, AgentAcpStartedProps,
AgentAcpTimedOutProps, AgentCliCancelledProps, AgentCliCompletedProps,
AgentCliTimedOutProps, AgentMessageProps, AgentSessionActivatedProps,
AgentAcpTimedOutProps, AgentMessageProps, AgentSessionActivatedProps,
AgentSessionEndedProps, AgentSessionStartedProps, CheckpointCompletedProps,
InterviewCompletedProps, InterviewOption, InterviewStartedProps, RunControlEffectProps,
StageCompletedProps, StageFailedProps, StagePromptProps, StageRetryingProps,
@ -1367,11 +1310,9 @@ mod tests {
.apply_event(&test_stage_event(
4,
EventBody::AgentAcpStarted(AgentAcpStartedProps {
visit: 1,
mode: "acp".to_string(),
provider: "openai".to_string(),
model: "fake-acp".to_string(),
command: "python fake_agent.py".to_string(),
visit: 1,
command: "python fake_agent.py".to_string(),
config_name: Some("fake".to_string()),
}),
stage_id.clone(),
))
@ -1382,9 +1323,8 @@ mod tests {
stage.provider_used.as_ref().unwrap(),
&json!({
"mode": "acp",
"provider": "openai",
"model": "fake-acp",
"command": "python fake_agent.py"
"command": "python fake_agent.py",
"config_name": "fake"
})
);
}
@ -1460,88 +1400,6 @@ mod tests {
assert_eq!(stage.termination, Some(CommandTermination::TimedOut));
}
#[test]
fn agent_cli_completed_updates_stage_output_projection() {
let mut state = initialized_projection();
let stage_id = StageId::new("code", 1);
start_stage(&mut state, &stage_id);
state
.apply_event(&test_stage_event(
4,
EventBody::AgentCliCompleted(AgentCliCompletedProps {
stdout: "done".to_string(),
stderr: "warn".to_string(),
exit_code: 0,
duration_ms: 42,
}),
stage_id.clone(),
))
.unwrap();
let stage = state.stage(&stage_id).unwrap();
assert_eq!(stage.output.as_deref(), Some("done\nwarn"));
assert_eq!(stage.termination, Some(CommandTermination::Exited));
assert_eq!(
stage.script_timing.as_ref().unwrap()["duration_ms"],
serde_json::json!(42)
);
}
#[test]
fn agent_cli_cancelled_updates_stage_output_projection() {
let mut state = initialized_projection();
let stage_id = StageId::new("code", 1);
start_stage(&mut state, &stage_id);
state
.apply_event(&test_stage_event(
4,
EventBody::AgentCliCancelled(AgentCliCancelledProps {
stdout: "partial".to_string(),
stderr: "cancelled".to_string(),
duration_ms: 7,
}),
stage_id.clone(),
))
.unwrap();
let stage = state.stage(&stage_id).unwrap();
assert_eq!(stage.output.as_deref(), Some("partial\ncancelled"));
assert_eq!(stage.termination, Some(CommandTermination::Cancelled));
assert_eq!(
stage.script_timing.as_ref().unwrap()["duration_ms"],
serde_json::json!(7)
);
}
#[test]
fn agent_cli_timed_out_updates_stage_output_projection() {
let mut state = initialized_projection();
let stage_id = StageId::new("code", 1);
start_stage(&mut state, &stage_id);
state
.apply_event(&test_stage_event(
4,
EventBody::AgentCliTimedOut(AgentCliTimedOutProps {
stdout: "partial".to_string(),
stderr: "timeout".to_string(),
duration_ms: 600,
}),
stage_id.clone(),
))
.unwrap();
let stage = state.stage(&stage_id).unwrap();
assert_eq!(stage.output.as_deref(), Some("partial\ntimeout"));
assert_eq!(stage.termination, Some(CommandTermination::TimedOut));
assert_eq!(
stage.script_timing.as_ref().unwrap()["duration_ms"],
serde_json::json!(600)
);
}
#[test]
fn stage_completed_event_captures_duration_and_usage_per_visit() {
let mut state = initialized_projection();

View file

@ -3,7 +3,7 @@ use std::time::Duration;
use serde::{Deserialize, Serialize};
use crate::LlmBackend;
use crate::AgentBackend;
/// Typed attribute values for nodes, edges, and graph-level attributes.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
@ -258,15 +258,25 @@ impl Node {
}
#[must_use]
pub fn llm_backend(&self) -> Option<Result<LlmBackend, strum::ParseError>> {
pub fn agent_backend(&self) -> Option<Result<AgentBackend, strum::ParseError>> {
self.backend().map(str::parse)
}
#[must_use]
pub fn acp_command(&self) -> Option<&str> {
pub fn legacy_acp_command_attr(&self) -> Option<&str> {
self.str_attr("acp_command")
}
#[must_use]
pub fn acp_command_attr(&self) -> Option<&str> {
self.str_attr("acp.command")
}
#[must_use]
pub fn acp_config_attr(&self) -> Option<&str> {
self.str_attr("acp.config")
}
#[must_use]
pub fn selection(&self) -> &str {
self.str_attr("selection").unwrap_or("deterministic")

View file

@ -61,7 +61,7 @@ pub use graph::{
shape_to_handler_type,
};
pub use interview::{InterviewQuestionRecord, QuestionType};
pub use llm_backend::LlmBackend;
pub use llm_backend::AgentBackend;
pub use manifest_path::{ManifestPath, ManifestPathParseError};
pub use outcome::{
FailureCategory, FailureDetail, NodeResult, Outcome, OutcomeMeta, StageOutcome, StageState,

View file

@ -18,15 +18,27 @@ use strum::{Display, EnumString, IntoStaticStr, VariantArray, VariantNames};
)]
#[serde(rename_all = "snake_case")]
#[strum(serialize_all = "snake_case")]
pub enum LlmBackend {
pub enum AgentBackend {
Api,
Cli,
Acp,
}
impl LlmBackend {
impl AgentBackend {
#[must_use]
pub fn expected_values() -> String {
<Self as VariantNames>::VARIANTS.join(", ")
}
}
#[cfg(test)]
mod tests {
use super::AgentBackend;
#[test]
fn agent_backend_accepts_only_api_and_acp() {
assert_eq!("api".parse::<AgentBackend>().unwrap(), AgentBackend::Api);
assert_eq!("acp".parse::<AgentBackend>().unwrap(), AgentBackend::Acp);
assert!("cli".parse::<AgentBackend>().is_err());
assert_eq!(AgentBackend::expected_values(), "api, acp");
}
}

View file

@ -218,44 +218,12 @@ pub struct CommandCompletedProps {
pub live_streaming: bool,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentCliStartedProps {
pub visit: u32,
pub mode: String,
pub provider: String,
pub model: String,
pub command: String,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentCliCompletedProps {
pub stdout: String,
pub stderr: String,
pub exit_code: i32,
pub duration_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentCliCancelledProps {
pub stdout: String,
pub stderr: String,
pub duration_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentCliTimedOutProps {
pub stdout: String,
pub stderr: String,
pub duration_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentAcpStartedProps {
pub visit: u32,
pub mode: String,
pub provider: String,
pub model: String,
pub command: String,
pub visit: u32,
pub command: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub config_name: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]

View file

@ -286,14 +286,6 @@ pub enum EventBody {
CommandStarted(CommandStartedProps),
#[serde(rename = "command.completed")]
CommandCompleted(CommandCompletedProps),
#[serde(rename = "agent.cli.started")]
AgentCliStarted(AgentCliStartedProps),
#[serde(rename = "agent.cli.completed")]
AgentCliCompleted(AgentCliCompletedProps),
#[serde(rename = "agent.cli.cancelled")]
AgentCliCancelled(AgentCliCancelledProps),
#[serde(rename = "agent.cli.timed_out")]
AgentCliTimedOut(AgentCliTimedOutProps),
#[serde(rename = "agent.acp.started")]
AgentAcpStarted(AgentAcpStartedProps),
#[serde(rename = "agent.acp.completed")]
@ -498,10 +490,6 @@ impl EventBody {
Self::CliEnsureFailed(_) => "cli.ensure.failed",
Self::CommandStarted(_) => "command.started",
Self::CommandCompleted(_) => "command.completed",
Self::AgentCliStarted(_) => "agent.cli.started",
Self::AgentCliCompleted(_) => "agent.cli.completed",
Self::AgentCliCancelled(_) => "agent.cli.cancelled",
Self::AgentCliTimedOut(_) => "agent.cli.timed_out",
Self::AgentAcpStarted(_) => "agent.acp.started",
Self::AgentAcpCompleted(_) => "agent.acp.completed",
Self::AgentAcpCancelled(_) => "agent.acp.cancelled",
@ -653,10 +641,6 @@ fn is_known_event_name(event: &str) -> bool {
| "cli.ensure.failed"
| "command.started"
| "command.completed"
| "agent.cli.started"
| "agent.cli.completed"
| "agent.cli.cancelled"
| "agent.cli.timed_out"
| "agent.acp.started"
| "agent.acp.completed"
| "agent.acp.cancelled"

View file

@ -13,6 +13,7 @@ doctest = false
workspace = true
[dependencies]
fabro-acp = { path = "../fabro-acp", default-features = false }
fabro-graphviz = { path = "../fabro-graphviz" }
fabro-model = { path = "../fabro-model" }
fabro-types = { path = "../fabro-types" }

View file

@ -1,5 +1,6 @@
use fabro_acp::{AcpCommandError, AcpProcessSpec};
use fabro_graphviz::graph::{Graph, Node};
use fabro_types::LlmBackend;
use fabro_types::AgentBackend;
use crate::{Diagnostic, LintRule, Severity};
@ -18,38 +19,16 @@ impl LintRule for Rule {
let mut diagnostics = Vec::new();
for node in graph.nodes.values() {
if let Some(backend) = node.backend() {
match node.llm_backend() {
match node.agent_backend() {
Some(Err(_)) => {
let expected = LlmBackend::expected_values();
diagnostics.push(Diagnostic {
rule: self.name().to_string(),
severity: Severity::Error,
message: format!(
"unsupported LLM backend \"{backend}\"; expected one of: {expected}"
),
node_id: Some(node.id.clone()),
edge: None,
fix: Some(format!("Use one of: {expected}")),
..Diagnostic::default()
});
diagnostics.push(unsupported_backend_diagnostic(
self.name(),
node,
backend,
));
}
Some(Ok(LlmBackend::Acp)) if acp_command_missing(node) => {
diagnostics.push(Diagnostic {
rule: self.name().to_string(),
severity: Severity::Error,
message: "backend=\"acp\" requires acp_command because Fabro does \
not install ACP agents"
.to_string(),
node_id: Some(node.id.clone()),
edge: None,
fix: Some(
"Set acp_command to a stdio ACP command available in the sandbox"
.to_string(),
),
..Diagnostic::default()
});
Some(Ok(AgentBackend::Acp)) => {
diagnostics.extend(validate_acp_node(self.name(), node));
}
Some(Ok(_)) | None => {}
}
@ -59,11 +38,150 @@ impl LintRule for Rule {
}
}
fn acp_command_missing(node: &Node) -> bool {
match node.acp_command() {
Some(command) => command.trim().is_empty(),
None => true,
fn unsupported_backend_diagnostic(rule: &str, node: &Node, backend: &str) -> Diagnostic {
if backend == "cli" {
return Diagnostic {
rule: rule.to_string(),
severity: Severity::Error,
message: "backend=\"cli\" is no longer supported; external agents must be launched \
through backend=\"acp\" with acp.command or acp.config"
.to_string(),
node_id: Some(node.id.clone()),
edge: None,
fix: Some(
"Use backend=\"api\" for Fabro-owned provider execution, or backend=\"acp\" with \
acp.command/acp.config for a user-supplied ACP process"
.to_string(),
),
..Diagnostic::default()
};
}
let expected = AgentBackend::expected_values();
Diagnostic {
rule: rule.to_string(),
severity: Severity::Error,
message: format!("unsupported agent backend \"{backend}\"; expected one of: {expected}"),
node_id: Some(node.id.clone()),
edge: None,
fix: Some(format!("Use one of: {expected}")),
..Diagnostic::default()
}
}
fn validate_acp_node(rule: &str, node: &Node) -> Vec<Diagnostic> {
let mut diagnostics = Vec::new();
if node.handler_type() != Some("agent") {
diagnostics.push(Diagnostic {
rule: rule.to_string(),
severity: Severity::Error,
message: "backend=\"acp\" is only valid on agent nodes; prompt nodes are API-only"
.to_string(),
node_id: Some(node.id.clone()),
edge: None,
fix: Some("Use backend=\"api\" on prompt nodes".to_string()),
..Diagnostic::default()
});
}
if let Err(error) = AcpProcessSpec::from_attrs(
node.legacy_acp_command_attr(),
node.acp_command_attr(),
node.acp_config_attr(),
) {
diagnostics.push(acp_process_diagnostic(rule, node, &error));
}
let api_only_attrs = api_only_attrs_present(node);
if !api_only_attrs.is_empty() {
diagnostics.push(Diagnostic {
rule: rule.to_string(),
severity: Severity::Error,
message: format!(
"backend=\"acp\" does not support API-only attributes: {}",
api_only_attrs.join(", ")
),
node_id: Some(node.id.clone()),
edge: None,
fix: Some("Remove API model/provider/control attributes from ACP nodes".to_string()),
..Diagnostic::default()
});
}
diagnostics
}
fn acp_process_diagnostic(rule: &str, node: &Node, error: &AcpCommandError) -> Diagnostic {
match error {
AcpCommandError::LegacyCommandAttribute => Diagnostic {
rule: rule.to_string(),
severity: Severity::Error,
message: "acp_command is no longer supported; use acp.command for shell commands or \
acp.config for JSON stdio ACP configs"
.to_string(),
node_id: Some(node.id.clone()),
edge: None,
fix: Some("Rename acp_command to acp.command".to_string()),
..Diagnostic::default()
},
AcpCommandError::EmptyOverride
| AcpCommandError::MissingOverride
| AcpCommandError::InvalidCommandString => Diagnostic {
rule: rule.to_string(),
severity: Severity::Error,
message: render_acp_process_error(error),
node_id: Some(node.id.clone()),
edge: None,
fix: Some(
"Set acp.command to a shell command, or acp.config to a JSON stdio ACP config"
.to_string(),
),
..Diagnostic::default()
},
AcpCommandError::InvalidConfigJson(_)
| AcpCommandError::InvalidConfigShape(_)
| AcpCommandError::UnsupportedTransport => Diagnostic {
rule: rule.to_string(),
severity: Severity::Error,
message: format!(
"acp.config must be a JSON stdio ACP config: {}",
render_acp_process_error(error)
),
node_id: Some(node.id.clone()),
edge: None,
fix: Some(
"Provide a JSON config with type=\"stdio\", command, and optional args".to_string(),
),
..Diagnostic::default()
},
}
}
fn render_acp_process_error(error: &AcpCommandError) -> String {
error.to_string()
}
fn api_only_attrs_present(node: &Node) -> Vec<&'static str> {
const API_ONLY_ATTRS: &[&str] = &[
"model",
"provider",
"reasoning_effort",
"max_tokens",
"speed",
];
API_ONLY_ATTRS
.iter()
.copied()
.filter(|attr| node.attrs.contains_key(*attr))
.collect()
}
#[cfg(test)]
@ -75,8 +193,8 @@ mod tests {
use crate::{LintRule, Severity};
#[test]
fn backend_valid_accepts_absent_api_and_cli() {
for backend in [None, Some("api"), Some("cli")] {
fn backend_valid_accepts_absent_and_api() {
for backend in [None, Some("api")] {
let mut graph = minimal_graph();
let mut node = Node::new("work");
if let Some(backend) = backend {
@ -107,12 +225,12 @@ mod tests {
assert!(
diagnostics[0]
.message
.contains("unsupported LLM backend \"codex\"; expected one of: api, cli, acp")
.contains("unsupported agent backend \"codex\"; expected one of: api, acp")
);
}
#[test]
fn backend_valid_requires_acp_command_for_acp_backend() {
fn backend_valid_requires_acp_process_attr_for_acp_backend() {
let mut graph = minimal_graph();
let mut node = Node::new("work");
node.attrs
@ -122,23 +240,208 @@ mod tests {
let diagnostics = Rule.apply(&graph);
assert_eq!(diagnostics.len(), 1);
assert_eq!(diagnostics[0].severity, Severity::Error);
assert!(diagnostics[0].message.contains(
"backend=\"acp\" requires acp_command because Fabro does not install ACP agents"
));
assert!(
diagnostics[0]
.message
.contains("requires exactly one of acp.command or acp.config")
);
}
#[test]
fn backend_valid_accepts_acp_backend_with_acp_command() {
fn backend_valid_accepts_acp_backend_with_acp_command_attr() {
let mut graph = minimal_graph();
let mut node = Node::new("work");
node.attrs
.insert("backend".to_string(), AttrValue::String("acp".to_string()));
node.attrs.insert(
"acp_command".to_string(),
"acp.command".to_string(),
AttrValue::String("agent-acp".to_string()),
);
graph.nodes.insert("work".to_string(), node);
assert!(Rule.apply(&graph).is_empty());
}
#[test]
fn backend_valid_rejects_cli_backend_with_migration_guidance() {
let mut graph = minimal_graph();
let mut node = Node::new("work");
node.attrs
.insert("backend".to_string(), AttrValue::String("cli".to_string()));
graph.nodes.insert("work".to_string(), node);
let diagnostics = Rule.apply(&graph);
assert_eq!(diagnostics.len(), 1);
assert_eq!(diagnostics[0].severity, Severity::Error);
assert!(
diagnostics[0]
.message
.contains("backend=\"cli\" is no longer supported")
);
assert!(
diagnostics[0]
.fix
.as_deref()
.unwrap()
.contains("backend=\"acp\"")
);
}
#[test]
fn backend_valid_requires_exactly_one_acp_process_attr() {
let mut missing = minimal_graph();
let mut missing_node = Node::new("missing");
missing_node
.attrs
.insert("backend".to_string(), AttrValue::String("acp".to_string()));
missing.nodes.insert("missing".to_string(), missing_node);
let diagnostics = Rule.apply(&missing);
assert_eq!(diagnostics.len(), 1);
assert!(
diagnostics[0]
.message
.contains("requires exactly one of acp.command or acp.config")
);
let mut both = minimal_graph();
let mut both_node = Node::new("both");
both_node
.attrs
.insert("backend".to_string(), AttrValue::String("acp".to_string()));
both_node.attrs.insert(
"acp.command".to_string(),
AttrValue::String("python3 agent.py".to_string()),
);
both_node.attrs.insert(
"acp.config".to_string(),
AttrValue::String(
r#"{"type":"stdio","name":"agent","command":"python3","args":["agent.py"]}"#
.to_string(),
),
);
both.nodes.insert("both".to_string(), both_node);
let diagnostics = Rule.apply(&both);
assert_eq!(diagnostics.len(), 1);
assert!(
diagnostics[0]
.message
.contains("requires exactly one of acp.command or acp.config")
);
}
#[test]
fn backend_valid_rejects_legacy_acp_command_attr() {
let mut graph = minimal_graph();
let mut node = Node::new("work");
node.attrs
.insert("backend".to_string(), AttrValue::String("acp".to_string()));
node.attrs.insert(
"acp_command".to_string(),
AttrValue::String("python3 agent.py".to_string()),
);
graph.nodes.insert("work".to_string(), node);
let diagnostics = Rule.apply(&graph);
assert_eq!(diagnostics.len(), 1);
assert!(
diagnostics[0]
.message
.contains("acp_command is no longer supported")
);
}
#[test]
fn backend_valid_rejects_acp_on_prompt_nodes_and_api_only_attrs() {
let mut graph = minimal_graph();
let mut node = Node::new("prompt");
node.attrs
.insert("type".to_string(), AttrValue::String("prompt".to_string()));
node.attrs
.insert("backend".to_string(), AttrValue::String("acp".to_string()));
node.attrs.insert(
"acp.command".to_string(),
AttrValue::String("python3 agent.py".to_string()),
);
node.attrs.insert(
"model".to_string(),
AttrValue::String("gpt-5.4".to_string()),
);
node.attrs.insert(
"reasoning_effort".to_string(),
AttrValue::String("high".to_string()),
);
graph.nodes.insert("prompt".to_string(), node);
let diagnostics = Rule.apply(&graph);
assert_eq!(diagnostics.len(), 2);
assert!(diagnostics.iter().any(|diagnostic| {
diagnostic
.message
.contains("backend=\"acp\" is only valid on agent nodes")
}));
assert!(diagnostics.iter().any(|diagnostic| {
diagnostic
.message
.contains("backend=\"acp\" does not support API-only attributes")
}));
}
#[test]
fn backend_valid_rejects_invalid_acp_config_but_accepts_json_shaped_command() {
let mut command_graph = minimal_graph();
let mut command_node = Node::new("command");
command_node
.attrs
.insert("backend".to_string(), AttrValue::String("acp".to_string()));
command_node.attrs.insert(
"acp.command".to_string(),
AttrValue::String(r#"{"type":"stdio"}"#.to_string()),
);
command_graph
.nodes
.insert("command".to_string(), command_node);
assert!(Rule.apply(&command_graph).is_empty());
let mut config_graph = minimal_graph();
let mut config_node = Node::new("config");
config_node
.attrs
.insert("backend".to_string(), AttrValue::String("acp".to_string()));
config_node.attrs.insert(
"acp.config".to_string(),
AttrValue::String("python3 agent.py".to_string()),
);
config_graph.nodes.insert("config".to_string(), config_node);
let diagnostics = Rule.apply(&config_graph);
assert_eq!(diagnostics.len(), 1);
assert!(
diagnostics[0]
.message
.contains("acp.config must be a JSON stdio ACP config")
);
}
#[test]
fn backend_valid_rejects_invalid_acp_command() {
let mut graph = minimal_graph();
let mut node = Node::new("command");
node.attrs
.insert("backend".to_string(), AttrValue::String("acp".to_string()));
node.attrs.insert(
"acp.command".to_string(),
AttrValue::String("python 'unterminated".to_string()),
);
graph.nodes.insert("command".to_string(), node);
let diagnostics = Rule.apply(&graph);
assert_eq!(diagnostics.len(), 1);
assert!(
diagnostics[0]
.message
.contains("failed to parse acp.command as a shell command")
);
}
}

View file

@ -1028,32 +1028,6 @@ fn event_body_from_event(event: &Event) -> EventBody {
output_bytes: *output_bytes,
live_streaming: *live_streaming,
}),
Event::AgentCliStarted {
visit,
mode,
provider,
model,
command,
..
} => EventBody::AgentCliStarted(fabro_types::AgentCliStartedProps {
visit: *visit,
mode: mode.clone(),
provider: provider.clone(),
model: model.clone(),
command: command.clone(),
}),
Event::AgentCliCompleted {
stdout,
stderr,
exit_code,
duration_ms,
..
} => EventBody::AgentCliCompleted(fabro_types::AgentCliCompletedProps {
stdout: stdout.clone(),
stderr: stderr.clone(),
exit_code: *exit_code,
duration_ms: *duration_ms,
}),
Event::AgentSessionStarted {
provider, model, ..
} => EventBody::AgentSessionStarted(fabro_types::AgentSessionStartedProps {
@ -1096,39 +1070,15 @@ fn event_body_from_event(event: &Event) -> EventBody {
count: *count,
})
}
Event::AgentCliCancelled {
stdout,
stderr,
duration_ms,
..
} => EventBody::AgentCliCancelled(fabro_types::AgentCliCancelledProps {
stdout: stdout.clone(),
stderr: stderr.clone(),
duration_ms: *duration_ms,
}),
Event::AgentCliTimedOut {
stdout,
stderr,
duration_ms,
..
} => EventBody::AgentCliTimedOut(fabro_types::AgentCliTimedOutProps {
stdout: stdout.clone(),
stderr: stderr.clone(),
duration_ms: *duration_ms,
}),
Event::AgentAcpStarted {
visit,
mode,
provider,
model,
command,
config_name,
..
} => EventBody::AgentAcpStarted(fabro_types::AgentAcpStartedProps {
visit: *visit,
mode: mode.clone(),
provider: provider.clone(),
model: model.clone(),
command: command.clone(),
visit: *visit,
command: command.clone(),
config_name: config_name.clone(),
}),
Event::AgentAcpCompleted {
stdout,
@ -2077,48 +2027,6 @@ mod tests {
assert_eq!(message.billing.total_usd_micros, None);
}
#[test]
fn agent_cli_cancelled_maps_to_event_body_with_node_id() {
let stored = to_run_event(&fixtures::RUN_1, &Event::AgentCliCancelled {
node_id: "code".to_string(),
stdout: "out".to_string(),
stderr: "err".to_string(),
duration_ms: 42,
});
assert_eq!(stored.event_name(), "agent.cli.cancelled");
assert_eq!(stored.node_id.as_deref(), Some("code"));
match &stored.body {
EventBody::AgentCliCancelled(props) => {
assert_eq!(props.stdout, "out");
assert_eq!(props.stderr, "err");
assert_eq!(props.duration_ms, 42);
}
other => panic!("expected AgentCliCancelled, got {other:?}"),
}
}
#[test]
fn agent_cli_timed_out_maps_to_event_body_with_node_id() {
let stored = to_run_event(&fixtures::RUN_1, &Event::AgentCliTimedOut {
node_id: "code".to_string(),
stdout: "out".to_string(),
stderr: "err".to_string(),
duration_ms: 99,
});
assert_eq!(stored.event_name(), "agent.cli.timed_out");
assert_eq!(stored.node_id.as_deref(), Some("code"));
match &stored.body {
EventBody::AgentCliTimedOut(props) => {
assert_eq!(props.stdout, "out");
assert_eq!(props.stderr, "err");
assert_eq!(props.duration_ms, 99);
}
other => panic!("expected AgentCliTimedOut, got {other:?}"),
}
}
#[test]
fn agent_acp_events_map_to_event_bodies_with_stage_scope() {
let scope = StageScope {
@ -2131,12 +2039,10 @@ mod tests {
let started = to_run_event_at(
&fixtures::RUN_1,
&Event::AgentAcpStarted {
node_id: "code".to_string(),
visit: 2,
mode: "acp".to_string(),
provider: "openai".to_string(),
model: "fake-acp".to_string(),
command: "python fake_agent.py".to_string(),
node_id: "code".to_string(),
visit: 2,
command: "python fake_agent.py".to_string(),
config_name: Some("fake".to_string()),
},
Utc::now(),
Some(&scope),
@ -2149,10 +2055,8 @@ mod tests {
match &started.body {
EventBody::AgentAcpStarted(props) => {
assert_eq!(props.visit, 2);
assert_eq!(props.mode, "acp");
assert_eq!(props.provider, "openai");
assert_eq!(props.model, "fake-acp");
assert_eq!(props.command, "python fake_agent.py");
assert_eq!(props.config_name.as_deref(), Some("fake"));
}
other => panic!("expected AgentAcpStarted, got {other:?}"),
}

View file

@ -536,14 +536,6 @@ pub enum Event {
output_bytes: u64,
live_streaming: bool,
},
AgentCliStarted {
node_id: String,
visit: u32,
mode: String,
provider: String,
model: String,
command: String,
},
/// A top-level agent session object started its lifecycle.
AgentSessionStarted {
session_id: String,
@ -606,32 +598,12 @@ pub enum Event {
#[serde(default, skip_serializing_if = "Option::is_none")]
visit: Option<u32>,
},
AgentCliCompleted {
node_id: String,
stdout: String,
stderr: String,
exit_code: i32,
duration_ms: u64,
},
AgentCliCancelled {
node_id: String,
stdout: String,
stderr: String,
duration_ms: u64,
},
AgentCliTimedOut {
node_id: String,
stdout: String,
stderr: String,
duration_ms: u64,
},
AgentAcpStarted {
node_id: String,
visit: u32,
mode: String,
provider: String,
model: String,
command: String,
node_id: String,
visit: u32,
command: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
config_name: Option<String>,
},
AgentAcpCompleted {
node_id: String,
@ -1356,22 +1328,6 @@ impl Event {
"Command completed"
);
}
Self::AgentCliStarted {
node_id,
provider,
model,
..
} => {
debug!(node_id, provider, model, "Agent CLI started");
}
Self::AgentCliCompleted {
node_id,
exit_code,
duration_ms,
..
} => {
debug!(node_id, exit_code, duration_ms, "Agent CLI completed");
}
Self::AgentSessionStarted {
session_id,
provider,
@ -1412,27 +1368,13 @@ impl Event {
Self::AgentSteerDropped { reason, count, .. } => {
warn!(?reason, count, "Steer dropped");
}
Self::AgentCliCancelled {
node_id,
duration_ms,
..
} => {
debug!(node_id, duration_ms, "Agent CLI cancelled");
}
Self::AgentCliTimedOut {
node_id,
duration_ms,
..
} => {
debug!(node_id, duration_ms, "Agent CLI timed out");
}
Self::AgentAcpStarted {
node_id,
provider,
model,
command,
config_name,
..
} => {
debug!(node_id, provider, model, "Agent ACP started");
debug!(node_id, command, ?config_name, "Agent ACP started");
}
Self::AgentAcpCompleted {
node_id,

View file

@ -125,8 +125,6 @@ pub fn event_name(event: &Event) -> &'static str {
Event::Failover { .. } => "agent.failover",
Event::CommandStarted { .. } => "command.started",
Event::CommandCompleted { .. } => "command.completed",
Event::AgentCliStarted { .. } => "agent.cli.started",
Event::AgentCliCompleted { .. } => "agent.cli.completed",
Event::AgentSessionStarted { .. } => "agent.session.started",
Event::AgentSessionActivated { .. } => "agent.session.activated",
Event::AgentSessionDeactivated { .. } => "agent.session.deactivated",
@ -134,8 +132,6 @@ pub fn event_name(event: &Event) -> &'static str {
Event::AgentInterruptInjected { .. } => "agent.interrupt.injected",
Event::AgentSteerBuffered { .. } => "agent.steer.buffered",
Event::AgentSteerDropped { .. } => "agent.steer.dropped",
Event::AgentCliCancelled { .. } => "agent.cli.cancelled",
Event::AgentCliTimedOut { .. } => "agent.cli.timed_out",
Event::AgentAcpStarted { .. } => "agent.acp.started",
Event::AgentAcpCompleted { .. } => "agent.acp.completed",
Event::AgentAcpCancelled { .. } => "agent.acp.cancelled",

View file

@ -121,10 +121,6 @@ fn stored_event_fields_for_variant(event: &Event) -> StoredEventFields {
| Event::PromptCompleted { node_id, .. }
| Event::CommandStarted { node_id, .. }
| Event::CommandCompleted { node_id, .. }
| Event::AgentCliStarted { node_id, .. }
| Event::AgentCliCompleted { node_id, .. }
| Event::AgentCliCancelled { node_id, .. }
| Event::AgentCliTimedOut { node_id, .. }
| Event::AgentAcpCompleted { node_id, .. }
| Event::AgentAcpCancelled { node_id, .. }
| Event::AgentAcpTimedOut { node_id, .. } => node_stored_fields(Some(node_id.clone())),

View file

@ -4,62 +4,28 @@ use std::collections::HashMap;
use std::sync::Arc;
use async_trait::async_trait;
use fabro_acp::{
AcpCommandError, AcpError, AcpRunRequest, render_stop_reason, resolve_acp_command,
};
use fabro_acp::{AcpCommandError, AcpError, AcpProcessSpec, AcpRunRequest, render_stop_reason};
use fabro_agent::{Sandbox, StaticEnvProvider, ToolEnvProvider};
use fabro_auth::CredentialResolver;
use fabro_graphviz::graph::Node;
use fabro_model::{Catalog, ProviderId};
use fabro_util::time::elapsed_ms;
use tokio_util::sync::CancellationToken;
use super::super::agent::{CodergenBackend, CodergenResult, CodergenRunRequest, OneShotRequest};
use super::cli::AgentCli;
use super::launch_env::{AgentLaunchEnvRequest, resolve_agent_launch_env};
use super::{changed_files, routing};
use super::changed_files;
use crate::error::Error;
use crate::event::{Emitter, Event, StageScope};
use crate::event::{Emitter, Event, RunNoticeCode, RunNoticeLevel, StageScope};
pub struct AgentAcpBackend {
model: String,
provider_id: ProviderId,
tool_env: Option<Arc<dyn ToolEnvProvider>>,
tool_env: Option<Arc<dyn ToolEnvProvider>>,
github_token_refresh_managed: bool,
resolver: Option<CredentialResolver>,
catalog: Arc<Catalog>,
}
impl AgentAcpBackend {
#[must_use]
pub fn new(
model: String,
provider_id: impl Into<ProviderId>,
resolver: CredentialResolver,
) -> Self {
let provider_id = provider_id.into();
let catalog = default_catalog();
pub fn new() -> Self {
Self {
model,
provider_id,
tool_env: None,
tool_env: None,
github_token_refresh_managed: false,
resolver: Some(resolver),
catalog,
}
}
#[must_use]
pub fn new_from_env(model: String, provider_id: impl Into<ProviderId>) -> Self {
let provider_id = provider_id.into();
let catalog = default_catalog();
Self {
model,
provider_id,
tool_env: None,
github_token_refresh_managed: false,
resolver: None,
catalog,
}
}
@ -80,12 +46,6 @@ impl AgentAcpBackend {
self
}
#[must_use]
pub fn with_catalog(mut self, catalog: Arc<Catalog>) -> Self {
self.catalog = catalog;
self
}
async fn run_turn(
&self,
node: &Node,
@ -95,53 +55,29 @@ impl AgentAcpBackend {
sandbox: &Arc<dyn Sandbox>,
cancel_token: CancellationToken,
) -> Result<CodergenResult, Error> {
let files_before = changed_files::detect_changed_files(sandbox).await;
let model = node.model().unwrap_or(&self.model);
let provider = routing::resolve_node_provider_context(
self.catalog.as_ref(),
&self.provider_id,
&self.model,
node,
)?;
let provider_id = provider.provider_id;
let profile_kind = provider.profile_kind;
let command =
resolve_acp_command(node.acp_command()).map_err(acp_command_error_to_workflow)?;
let launch_env = resolve_agent_launch_env(AgentLaunchEnvRequest {
provider_id: provider_id.clone(),
cli: AgentCli::for_profile_kind(profile_kind),
catalog: self.catalog.as_ref(),
resolver: self.resolver.as_ref(),
tool_env: self.tool_env.as_ref(),
github_token_refresh_managed: self.github_token_refresh_managed,
stage_label: "ACP",
emitter,
sandbox,
cancel_token: &cancel_token,
})
.await?;
let process_spec = resolve_acp_process_spec(node)?;
let config_name = process_spec.name().map(str::to_string);
let launch_env = self.resolve_launch_env(emitter).await?;
let on_activity = {
let emitter = Arc::clone(emitter);
Arc::new(move || emitter.touch()) as Arc<dyn Fn() + Send + Sync>
};
let command_display = command.to_string();
let command_display = process_spec.to_string();
emitter.emit_scoped(
&Event::AgentAcpStarted {
node_id: node.id.clone(),
visit: stage_scope.visit,
mode: "acp".to_string(),
provider: provider_id.to_string(),
model: model.to_string(),
command: command_display,
node_id: node.id.clone(),
visit: stage_scope.visit,
command: command_display,
config_name,
},
stage_scope,
);
let files_before = changed_files::detect_changed_files(sandbox).await;
let launch_start = std::time::Instant::now();
let result = match fabro_acp::run_acp_turn(AcpRunRequest {
command,
command: process_spec,
prompt,
cwd: sandbox.working_directory().to_string(),
timeout_ms: node.timeout().map(crate::millis_u64),
@ -224,10 +160,33 @@ impl AgentAcpBackend {
last_file_touched,
})
}
async fn resolve_launch_env(
&self,
emitter: &Arc<Emitter>,
) -> Result<HashMap<String, String>, Error> {
let Some(provider) = &self.tool_env else {
return Ok(HashMap::new());
};
if self.github_token_refresh_managed {
emitter.notice(
RunNoticeLevel::Info,
RunNoticeCode::GithubTokenRefreshLimited,
"ACP agent stages receive workflow env at process launch; stages running beyond \
token expiry may need to be retried.",
);
}
provider
.resolve()
.await
.map_err(|err| Error::handler_with_anyhow("Failed to resolve ACP agent env", err))
}
}
fn default_catalog() -> Arc<Catalog> {
Arc::new(Catalog::from_builtin().expect("default catalog should build"))
impl Default for AgentAcpBackend {
fn default() -> Self {
Self::new()
}
}
#[async_trait]
@ -245,38 +204,46 @@ impl CodergenBackend for AgentAcpBackend {
.await
}
async fn one_shot(&self, request: OneShotRequest<'_>) -> Result<CodergenResult, Error> {
let prompt = match request.system_prompt.filter(|prompt| !prompt.is_empty()) {
Some(system_prompt) => format!("System:\n{system_prompt}\n\nUser:\n{}", request.prompt),
None => request.prompt.to_string(),
};
self.run_turn(
request.node,
prompt,
request.emitter,
request.stage_scope,
request.sandbox,
request.cancel_token,
)
.await
async fn one_shot(&self, _request: OneShotRequest<'_>) -> Result<CodergenResult, Error> {
Err(Error::Validation(
"backend=\"acp\" is only valid on agent nodes; prompt nodes are API-only".to_string(),
))
}
}
fn acp_command_error_to_workflow(error: AcpCommandError) -> Error {
fn acp_process_error_to_workflow(error: AcpCommandError) -> Error {
match error {
AcpCommandError::EmptyOverride => Error::handler("acp_command must not be empty"),
AcpCommandError::MissingOverride => Error::handler(
"acp_command is required for backend=\"acp\" because Fabro does not install ACP agents",
),
AcpCommandError::LegacyCommandAttribute => {
Error::handler("acp_command is no longer supported; use acp.command or acp.config")
}
AcpCommandError::EmptyOverride => Error::handler("ACP process attribute must not be empty"),
AcpCommandError::MissingOverride => {
Error::handler("backend=\"acp\" requires exactly one of acp.command or acp.config")
}
AcpCommandError::UnsupportedTransport => {
Error::handler("only stdio ACP commands are supported")
}
AcpCommandError::Parse(source) => {
Error::handler_with_source("Failed to resolve ACP command", source)
AcpCommandError::InvalidCommandString => {
Error::handler("Failed to parse acp.command as a shell command")
}
AcpCommandError::InvalidConfigJson(source) => {
Error::handler_with_source("Failed to parse acp.config as JSON", source)
}
AcpCommandError::InvalidConfigShape(message) => {
Error::handler(format!("Invalid acp.config shape: {message}"))
}
}
}
fn resolve_acp_process_spec(node: &Node) -> Result<AcpProcessSpec, Error> {
AcpProcessSpec::from_attrs(
node.legacy_acp_command_attr(),
node.acp_command_attr(),
node.acp_config_attr(),
)
.map_err(acp_process_error_to_workflow)
}
fn acp_error_to_workflow(error: AcpError) -> Error {
match error {
AcpError::Cancelled => Error::Cancelled,
@ -307,17 +274,14 @@ mod tests {
use fabro_acp::{AcpError, AcpProcessExit};
use fabro_agent::{LocalSandbox, Sandbox, shell_quote};
use fabro_graphviz::graph::{AttrValue, Node};
use fabro_model::ProviderId;
use fabro_sandbox::test_support::MockSandbox;
use fabro_types::{CommandTermination, EventBody, ExecOutputTail};
use tokio_util::sync::CancellationToken;
use super::{AgentAcpBackend, acp_error_to_workflow};
use crate::context::Context;
use crate::event::{Emitter, StageScope};
use crate::handler::agent::{
CodergenBackend, CodergenResult, CodergenRunRequest, OneShotRequest,
};
use crate::event::Emitter;
use crate::handler::agent::{CodergenBackend, CodergenResult, CodergenRunRequest};
#[tokio::test]
async fn acp_backend_run_sends_prompt_and_returns_text() {
@ -329,28 +293,20 @@ mod tests {
.unwrap();
let mut node = Node::new("work");
node.attrs.insert(
"provider".to_string(),
AttrValue::String("openai".to_string()),
);
node.attrs.insert(
"model".to_string(),
AttrValue::String("fake-acp".to_string()),
);
node.attrs
.insert("backend".to_string(), AttrValue::String("acp".to_string()));
node.attrs.insert(
"acp_command".to_string(),
"acp.command".to_string(),
AttrValue::String(format!(
"python3 {}",
shell_quote(&script_path.to_string_lossy())
)),
);
let backend =
AgentAcpBackend::new_from_env("fake-acp".to_string(), ProviderId::openai()).with_env(
HashMap::from([("ACP_MODE".to_string(), "write_file".to_string())]),
);
let backend = AgentAcpBackend::new().with_env(HashMap::from([(
"ACP_MODE".to_string(),
"write_file".to_string(),
)]));
let sandbox: Arc<dyn Sandbox> = Arc::new(LocalSandbox::new(tempdir.path().to_path_buf()));
let emitter = Arc::new(Emitter::default());
let context = Context::new();
@ -381,63 +337,93 @@ mod tests {
}
#[tokio::test]
async fn acp_backend_one_shot_combines_system_prompt_and_uses_passed_sandbox() {
async fn acp_backend_accepts_acp_command_attribute_without_model_or_provider() {
let tempdir = tempfile::tempdir().unwrap();
init_git(tempdir.path());
let script_path = tempdir.path().join("fake_acp_agent.py");
let prompt_record_path = tempdir.path().join("prompt.json");
tokio::fs::write(&script_path, fake_acp_agent_script())
.await
.unwrap();
let mut node = Node::new("prompt");
node.attrs.insert(
"provider".to_string(),
AttrValue::String("openai".to_string()),
);
let mut node = Node::new("work");
node.attrs
.insert("backend".to_string(), AttrValue::String("acp".to_string()));
node.attrs.insert(
"acp_command".to_string(),
"acp.command".to_string(),
AttrValue::String(format!(
"python3 {}",
shell_quote(&script_path.to_string_lossy())
)),
);
let backend = AgentAcpBackend::new_from_env("fake-acp".to_string(), ProviderId::openai())
.with_env(HashMap::from([
(
"ACP_PROMPT_RECORD".to_string(),
prompt_record_path.to_string_lossy().into_owned(),
),
("ACP_MODE".to_string(), "write_file".to_string()),
]));
let backend = AgentAcpBackend::new().with_env(HashMap::from([(
"ACP_MODE".to_string(),
"write_file".to_string(),
)]));
let sandbox: Arc<dyn Sandbox> = Arc::new(LocalSandbox::new(tempdir.path().to_path_buf()));
let emitter = Arc::new(Emitter::default());
let context = Context::new();
let stage_scope = StageScope::for_handler(&context, "prompt");
let result = backend
.one_shot(OneShotRequest {
node: &node,
prompt: "User prompt",
system_prompt: Some("System prompt"),
emitter: &emitter,
stage_scope: &stage_scope,
sandbox: &sandbox,
cancel_token: CancellationToken::new(),
.run(CodergenRunRequest {
node: &node,
prompt: "write hello",
context: &context,
thread_id: None,
emitter: &emitter,
sandbox: &sandbox,
tool_hooks: None,
cancel_token: CancellationToken::new(),
})
.await
.unwrap();
assert!(matches!(result, CodergenResult::Text { .. }));
let recorded = tokio::fs::read_to_string(prompt_record_path).await.unwrap();
assert!(recorded.contains("System:\\nSystem prompt\\n\\nUser:\\nUser prompt"));
assert_eq!(
tokio::fs::read_to_string(tempdir.path().join("hello.txt"))
.await
.unwrap(),
"hello from sandbox\n"
let CodergenResult::Text { text, .. } = result else {
panic!("expected text result");
};
assert_eq!(text, "hello from acp");
}
#[tokio::test]
async fn acp_backend_does_not_forward_provider_credentials() {
let mut sandbox = MockSandbox::linux();
sandbox.stdio_process_error = Some("stop before ACP handshake".to_string());
let sandbox = Arc::new(sandbox);
let sandbox_dyn: Arc<dyn Sandbox> = sandbox.clone();
let mut node = Node::new("work");
node.attrs
.insert("backend".to_string(), AttrValue::String("acp".to_string()));
node.attrs.insert(
"acp.command".to_string(),
AttrValue::String("fake-acp-agent".to_string()),
);
let backend = AgentAcpBackend::new();
let emitter = Arc::new(Emitter::default());
let context = Context::new();
let result = backend
.run(CodergenRunRequest {
node: &node,
prompt: "write hello",
context: &context,
thread_id: None,
emitter: &emitter,
sandbox: &sandbox_dyn,
tool_hooks: None,
cancel_token: CancellationToken::new(),
})
.await;
assert!(result.is_err());
let captured = sandbox
.captured_env_vars
.lock()
.expect("captured env lock poisoned")
.clone()
.unwrap_or_default();
assert!(!captured.contains_key("OPENAI_API_KEY"));
assert!(!captured.contains_key("ANTHROPIC_API_KEY"));
assert!(!captured.contains_key("GEMINI_API_KEY"));
}
#[tokio::test]
@ -450,21 +436,17 @@ mod tests {
let mut node = Node::new("work");
node.attrs.insert(
"provider".to_string(),
AttrValue::String("openai".to_string()),
);
node.attrs.insert(
"acp_command".to_string(),
"acp.command".to_string(),
AttrValue::String(format!(
"python3 {}",
shell_quote(&script_path.to_string_lossy())
)),
);
let backend =
AgentAcpBackend::new_from_env("fake-acp".to_string(), ProviderId::openai()).with_env(
HashMap::from([("ACP_STOP_REASON".to_string(), "cancelled".to_string())]),
);
let backend = AgentAcpBackend::new().with_env(HashMap::from([(
"ACP_STOP_REASON".to_string(),
"cancelled".to_string(),
)]));
let sandbox: Arc<dyn Sandbox> = Arc::new(LocalSandbox::new(tempdir.path().to_path_buf()));
let emitter = Arc::new(Emitter::default());
let context = Context::new();
@ -506,16 +488,12 @@ mod tests {
})
.to_string();
let mut node = Node::new("work");
node.attrs.insert(
"provider".to_string(),
AttrValue::String("openai".to_string()),
);
node.attrs
.insert("backend".to_string(), AttrValue::String("acp".to_string()));
node.attrs
.insert("acp_command".to_string(), AttrValue::String(raw_command));
.insert("acp.config".to_string(), AttrValue::String(raw_command));
let backend = AgentAcpBackend::new_from_env("fake-acp".to_string(), ProviderId::openai());
let backend = AgentAcpBackend::new();
let sandbox: Arc<dyn Sandbox> = Arc::new(LocalSandbox::new(tempdir.path().to_path_buf()));
let emitter = Arc::new(Emitter::default());
let events = Arc::new(Mutex::new(Vec::new()));
@ -554,20 +532,16 @@ mod tests {
}
#[tokio::test]
async fn acp_backend_requires_explicit_acp_command() {
async fn acp_backend_requires_explicit_process_attr() {
let sandbox = MockSandbox::linux();
let sandbox = Arc::new(sandbox);
let sandbox_dyn: Arc<dyn Sandbox> = sandbox.clone();
let mut node = Node::new("work");
node.attrs.insert(
"provider".to_string(),
AttrValue::String("openai".to_string()),
);
node.attrs
.insert("backend".to_string(), AttrValue::String("acp".to_string()));
let backend = AgentAcpBackend::new_from_env("fake-acp".to_string(), ProviderId::openai());
let backend = AgentAcpBackend::new();
let emitter = Arc::new(Emitter::default());
let context = Context::new();
let result = backend
@ -583,11 +557,11 @@ mod tests {
})
.await;
let Err(err) = result else {
panic!("ACP without acp_command should fail");
panic!("ACP without process attr should fail");
};
assert!(
err.to_string()
.contains("acp_command is required for backend=\"acp\"")
.contains("requires exactly one of acp.command or acp.config")
);
assert!(
sandbox
@ -595,7 +569,7 @@ mod tests {
.lock()
.expect("captured env lock poisoned")
.is_none(),
"ACP process should not launch when acp_command is missing"
"ACP process should not launch when process attr is missing"
);
}
@ -609,21 +583,17 @@ mod tests {
let sandbox_dyn: Arc<dyn Sandbox> = sandbox.clone();
let mut node = Node::new("work");
node.attrs.insert(
"provider".to_string(),
AttrValue::String("openai".to_string()),
);
node.attrs
.insert("backend".to_string(), AttrValue::String("acp".to_string()));
node.attrs.insert(
"acp_command".to_string(),
"acp.command".to_string(),
AttrValue::String("fake-acp-agent".to_string()),
);
let backend =
AgentAcpBackend::new_from_env("fake-acp".to_string(), ProviderId::openai()).with_env(
HashMap::from([("OPENAI_API_KEY".to_string(), "test-key".to_string())]),
);
let backend = AgentAcpBackend::new().with_env(HashMap::from([(
"WORKFLOW_ENV".to_string(),
"test-value".to_string(),
)]));
let emitter = Arc::new(Emitter::default());
let context = Context::new();
let result = backend

File diff suppressed because it is too large Load diff

View file

@ -1,121 +0,0 @@
use std::collections::HashMap;
use std::sync::Arc;
use fabro_agent::{Sandbox, ToolEnvProvider};
use fabro_auth::{CliAgentKind, CredentialResolver, CredentialUsage, ResolvedCredential};
use fabro_model::{Catalog, CredentialRef, ProviderId};
use tokio_util::sync::CancellationToken;
use super::cli::{AgentCli, process_env_var};
use crate::error::Error;
use crate::event::{Emitter, RunNoticeCode, RunNoticeLevel};
pub(crate) struct AgentLaunchEnvRequest<'a> {
pub provider_id: ProviderId,
pub cli: AgentCli,
pub catalog: &'a Catalog,
pub resolver: Option<&'a CredentialResolver>,
pub tool_env: Option<&'a Arc<dyn ToolEnvProvider>>,
pub github_token_refresh_managed: bool,
pub stage_label: &'static str,
pub emitter: &'a Arc<Emitter>,
pub sandbox: &'a Arc<dyn Sandbox>,
pub cancel_token: &'a CancellationToken,
}
pub(crate) async fn resolve_agent_launch_env(
request: AgentLaunchEnvRequest<'_>,
) -> Result<HashMap<String, String>, Error> {
let cli_agent = match request.cli {
AgentCli::Claude => CliAgentKind::Claude,
AgentCli::Codex => CliAgentKind::Codex,
AgentCli::Gemini => CliAgentKind::Gemini,
};
let mut launch_env = if let Some(resolver) = request.resolver {
let resolved = resolver
.resolve(
request.provider_id.clone(),
CredentialUsage::CliAgent(cli_agent),
request.catalog,
)
.await
.map_err(|err| {
Error::handler_with_source(
format!("Failed to resolve {} credential", request.stage_label),
err,
)
})?;
let ResolvedCredential::Cli(cli_credential) = resolved else {
return Err(Error::handler("Expected CLI credential".to_string()));
};
if let Some(login_cmd) = &cli_credential.login_command {
let login_result = request
.sandbox
.exec_command(
login_cmd,
30_000,
None,
None,
Some(request.cancel_token.child_token()),
)
.await
.map_err(|err| {
Error::handler_with_source(
format!("{} credential login failed", request.stage_label),
err,
)
})?;
if !login_result.is_success() {
tracing::warn!(
exit_code = login_result.display_exit_code(),
stage = request.stage_label,
"{} credential login failed: {}",
request.stage_label,
login_result.stderr
);
}
}
cli_credential.env_vars
} else {
let mut env = HashMap::new();
if let Some(auth) = request
.catalog
.provider(&request.provider_id)
.and_then(|provider| provider.auth.as_ref())
{
for credential_ref in &auth.credentials {
let CredentialRef::Env(name) = credential_ref else {
continue;
};
if let Some(value) = process_env_var(name) {
env.insert(name.clone(), value);
}
}
}
env
};
if let Some(provider) = request.tool_env {
if request.github_token_refresh_managed {
request.emitter.notice(
RunNoticeLevel::Info,
RunNoticeCode::GithubTokenRefreshLimited,
format!(
"{} agent stages receive GitHub tokens at process launch; stages running \
beyond token expiry may need to be retried.",
request.stage_label
),
);
}
let tool_env = provider.resolve().await.map_err(|err| {
Error::handler_with_anyhow(
format!("Failed to resolve {} agent env", request.stage_label),
err,
)
})?;
launch_env.extend(tool_env);
}
Ok(launch_env)
}

View file

@ -2,11 +2,10 @@ pub mod acp;
pub mod activation_lease;
pub mod api;
pub mod changed_files;
pub mod cli;
pub mod launch_env;
pub mod preamble;
pub mod router;
pub mod routing;
pub use acp::AgentAcpBackend;
pub use api::AgentApiBackend;
pub use cli::{AgentCliBackend, BackendRouter, parse_cli_response};
pub use router::BackendRouter;

View file

@ -0,0 +1,148 @@
use std::sync::Arc;
use async_trait::async_trait;
use fabro_graphviz::graph::Node;
use fabro_types::AgentBackend;
use super::super::agent::{CodergenBackend, CodergenResult, CodergenRunRequest, OneShotRequest};
use super::acp::AgentAcpBackend;
use super::routing;
use crate::error::Error;
use crate::event::Emitter;
/// Routes codergen invocations to API or ACP backends based on node attributes.
pub struct BackendRouter {
api: Box<dyn CodergenBackend>,
acp: AgentAcpBackend,
}
impl BackendRouter {
#[must_use]
pub fn new(api_backend: Box<dyn CodergenBackend>, acp_backend: AgentAcpBackend) -> Self {
Self {
api: api_backend,
acp: acp_backend,
}
}
fn select_backend(node: &Node) -> Result<AgentBackend, Error> {
routing::select_run_backend(node)
}
fn select_one_shot_backend(node: &Node) -> Result<AgentBackend, Error> {
routing::select_one_shot_backend(node)
}
}
#[async_trait]
impl CodergenBackend for BackendRouter {
async fn run(&self, request: CodergenRunRequest<'_>) -> Result<CodergenResult, Error> {
match Self::select_backend(request.node)? {
AgentBackend::Api => self.api.run(request).await,
AgentBackend::Acp => self.acp.run(request).await,
}
}
async fn one_shot(&self, request: OneShotRequest<'_>) -> Result<CodergenResult, Error> {
match Self::select_one_shot_backend(request.node)? {
AgentBackend::Api => self.api.one_shot(request).await,
AgentBackend::Acp => {
unreachable!("ACP one-shot is rejected by select_one_shot_backend")
}
}
}
async fn shutdown(&self, emitter: &Arc<Emitter>) {
self.api.shutdown(emitter).await;
}
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use async_trait::async_trait;
use fabro_agent::{LocalSandbox, Sandbox};
use fabro_graphviz::graph::{AttrValue, Node};
use tokio_util::sync::CancellationToken;
use super::*;
use crate::context::Context;
use crate::event::{Emitter, StageScope};
#[test]
fn router_uses_api_by_default() {
let node = Node::new("test");
assert_eq!(
BackendRouter::select_backend(&node).unwrap(),
AgentBackend::Api
);
}
#[test]
fn router_rejects_cli_backend() {
let mut node = Node::new("test");
node.attrs
.insert("backend".to_string(), AttrValue::String("cli".to_string()));
let err = BackendRouter::select_backend(&node).unwrap_err();
assert_eq!(
err.to_string(),
"Validation error: unsupported agent backend \"cli\"; expected one of: api, acp"
);
}
#[tokio::test]
async fn router_routes_one_shot_to_api_by_default() {
let node = Node::new("test");
let sandbox: Arc<dyn Sandbox> = Arc::new(LocalSandbox::new(
tempfile::tempdir().unwrap().path().to_path_buf(),
));
let context = Context::new();
let router = BackendRouter::new(Box::new(StubBackend), AgentAcpBackend::new());
let emitter = Arc::new(Emitter::default());
let stage_scope = StageScope::for_handler(&context, "test");
let result = router
.one_shot(OneShotRequest {
node: &node,
prompt: "prompt",
system_prompt: None,
emitter: &emitter,
stage_scope: &stage_scope,
sandbox: &sandbox,
cancel_token: CancellationToken::new(),
})
.await
.unwrap();
let CodergenResult::Text { text, .. } = result else {
panic!("expected text result");
};
assert_eq!(text, "api one-shot");
}
struct StubBackend;
#[async_trait]
impl CodergenBackend for StubBackend {
async fn run(&self, _request: CodergenRunRequest<'_>) -> Result<CodergenResult, Error> {
Ok(CodergenResult::Text {
text: "api run".to_string(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
})
}
async fn one_shot(&self, _request: OneShotRequest<'_>) -> Result<CodergenResult, Error> {
Ok(CodergenResult::Text {
text: "api one-shot".to_string(),
usage: None,
files_touched: Vec::new(),
last_file_touched: None,
})
}
}
}

View file

@ -1,19 +1,12 @@
use fabro_graphviz::graph::{self, Node};
use fabro_model::{AgentProfileKind, Catalog, ProviderId};
use fabro_types::LlmBackend;
use fabro_types::AgentBackend;
use super::cli::is_cli_only_model;
use crate::error::Error;
pub(crate) fn select_run_backend(node: &Node) -> Result<LlmBackend, Error> {
match node.llm_backend() {
None => {
if node.model().is_some_and(is_cli_only_model) {
Ok(LlmBackend::Cli)
} else {
Ok(LlmBackend::Api)
}
}
pub(crate) fn select_run_backend(node: &Node) -> Result<AgentBackend, Error> {
match node.agent_backend() {
None => Ok(AgentBackend::Api),
Some(Ok(backend)) => Ok(backend),
Some(Err(_)) => Err(unsupported_backend_error(
node.backend().unwrap_or_default(),
@ -21,10 +14,12 @@ pub(crate) fn select_run_backend(node: &Node) -> Result<LlmBackend, Error> {
}
}
pub(crate) fn select_one_shot_backend(node: &Node) -> Result<LlmBackend, Error> {
match node.llm_backend() {
Some(Ok(LlmBackend::Acp)) => Ok(LlmBackend::Acp),
Some(Ok(LlmBackend::Api | LlmBackend::Cli)) | None => Ok(LlmBackend::Api),
pub(crate) fn select_one_shot_backend(node: &Node) -> Result<AgentBackend, Error> {
match node.agent_backend() {
Some(Ok(AgentBackend::Acp)) => Err(Error::Validation(
"backend=\"acp\" is only valid on agent nodes; prompt nodes are API-only".to_string(),
)),
Some(Ok(AgentBackend::Api)) | None => Ok(AgentBackend::Api),
Some(Err(_)) => Err(unsupported_backend_error(
node.backend().unwrap_or_default(),
)),
@ -37,8 +32,8 @@ pub(crate) fn node_needs_api_backend(node: &Node) -> bool {
}
match node.handler_type() {
Some("prompt") => !matches!(select_one_shot_backend(node), Ok(LlmBackend::Acp)),
_ => matches!(select_run_backend(node), Ok(LlmBackend::Api)),
Some("prompt") => true,
_ => matches!(select_run_backend(node), Ok(AgentBackend::Api)),
}
}
@ -93,7 +88,7 @@ pub(crate) fn resolve_node_provider_context(
fn unsupported_backend_error(raw: &str) -> Error {
Error::Validation(format!(
"unsupported LLM backend \"{raw}\"; expected one of: {}",
LlmBackend::expected_values()
"unsupported agent backend \"{raw}\"; expected one of: {}",
AgentBackend::expected_values()
))
}

View file

@ -229,9 +229,6 @@ fn replay_event_for_fork_projection(body: &EventBody) -> bool {
| EventBody::InterviewTimeout(_)
| EventBody::InterviewInterrupted(_)
| EventBody::AgentSessionActivated(_)
| EventBody::AgentCliStarted(_)
| EventBody::AgentCliCancelled(_)
| EventBody::AgentCliTimedOut(_)
| EventBody::AgentAcpStarted(_)
| EventBody::AgentAcpCancelled(_)
| EventBody::AgentAcpTimedOut(_)
@ -324,11 +321,9 @@ mod tests {
fn fork_replay_preserves_agent_acp_projection_events() {
assert!(replay_event_for_fork_projection(
&EventBody::AgentAcpStarted(fabro_types::run_event::AgentAcpStartedProps {
visit: 1,
mode: "acp".to_string(),
provider: "openai".to_string(),
model: "fake-acp".to_string(),
command: "python fake_agent.py".to_string(),
visit: 1,
command: "python fake_agent.py".to_string(),
config_name: Some("fake".to_string()),
})
));
assert!(replay_event_for_fork_projection(

View file

@ -5,8 +5,7 @@ use std::time::Instant;
use fabro_agent::Sandbox;
use fabro_auth::{
CredentialResolver, CredentialSource, EnvCredentialSource, VaultCredentialSource,
auth_issue_message,
CredentialSource, EnvCredentialSource, VaultCredentialSource, auth_issue_message,
};
use fabro_graphviz::graph;
use fabro_hooks::{HookContext, HookDecision, HookEvent, HookRunner};
@ -29,9 +28,7 @@ use crate::devcontainer_bridge::{devcontainer_to_snapshot_config, run_devcontain
use crate::error::Error;
use crate::event::{Event, RunNoticeCode, RunNoticeLevel};
use crate::github_token_source::{AppIatMinter, GitHubTokenSource};
use crate::handler::llm::{
AgentAcpBackend, AgentApiBackend, AgentCliBackend, BackendRouter, routing,
};
use crate::handler::llm::{AgentAcpBackend, AgentApiBackend, BackendRouter, routing};
use crate::handler::{HandlerRegistry, default_registry};
use crate::run_metadata::{RunMetadataRuntime, build_metadata_writer, metadata_branch_name};
use crate::run_options::{GitCheckpointOptions, RunOptions};
@ -126,7 +123,6 @@ async fn build_registry(
graph: &graph::Graph,
llm_source: Arc<dyn CredentialSource>,
catalog: Arc<Catalog>,
cli_resolver: Option<CredentialResolver>,
) -> Result<(Arc<HandlerRegistry>, bool), Error> {
let no_backend_interviewer = Arc::clone(&interviewer);
let build_no_backend = move || {
@ -172,24 +168,9 @@ async fn build_registry(
.with_run_model_controls(model_controls.clone())
.with_tool_env_provider(tool_env_provider.clone())
.with_mcp_servers(mcp_servers.clone());
let cli = cli_resolver
.clone()
.map_or_else(
|| AgentCliBackend::new_from_env(model.clone(), provider_id.clone()),
|resolver| AgentCliBackend::new(model.clone(), provider_id.clone(), resolver),
)
.with_catalog(Arc::clone(&catalog_for_api))
.with_run_model_controls(model_controls.clone())
let acp = AgentAcpBackend::new()
.with_tool_env_provider(tool_env_provider.clone(), github_token_refresh_managed);
let acp = cli_resolver
.clone()
.map_or_else(
|| AgentAcpBackend::new_from_env(model.clone(), provider_id.clone()),
|resolver| AgentAcpBackend::new(model.clone(), provider_id.clone(), resolver),
)
.with_catalog(Arc::clone(&catalog_for_api))
.with_tool_env_provider(tool_env_provider.clone(), github_token_refresh_managed);
Some(Box::new(BackendRouter::new(Box::new(api), cli, acp)))
Some(Box::new(BackendRouter::new(Box::new(api), acp)))
}))
};
@ -350,7 +331,6 @@ pub async fn initialize(
let llm_source = build_llm_source(options.vault.clone());
let catalog = Arc::clone(&options.catalog);
let cli_resolver = options.vault.clone().map(CredentialResolver::new);
let sandbox_git = Arc::new(SandboxGitRuntime::new());
let metadata_runtime = Arc::new(RunMetadataRuntime::new());
@ -518,7 +498,6 @@ pub async fn initialize(
&graph,
Arc::clone(&llm_source),
Arc::clone(&catalog),
cli_resolver,
)
.await?
};
@ -1046,7 +1025,6 @@ mod tests {
&graph,
Arc::new(VaultCredentialSource::new(Arc::clone(&vault))),
test_catalog(),
Some(CredentialResolver::new(vault)),
)
.await
.unwrap();
@ -1065,7 +1043,7 @@ mod tests {
let source = format!(
r#"digraph test {{
start [shape=Mdiamond];
writer [type="agent", backend="acp", provider="openai", model="fake-acp", prompt="write hello", acp_command="python3 {}"];
writer [type="agent", backend="acp", prompt="write hello", acp.command="python3 {}"];
exit [shape=Msquare];
start -> writer;
writer -> exit;
@ -1085,20 +1063,12 @@ mod tests {
writer
.attrs
.insert("backend".to_string(), AttrValue::String("acp".to_string()));
writer.attrs.insert(
"provider".to_string(),
AttrValue::String("openai".to_string()),
);
writer.attrs.insert(
"model".to_string(),
AttrValue::String("fake-acp".to_string()),
);
writer.attrs.insert(
"prompt".to_string(),
AttrValue::String("write hello".to_string()),
);
writer.attrs.insert(
"acp_command".to_string(),
"acp.command".to_string(),
AttrValue::String(format!(
"python3 {}",
fabro_sandbox::shell_quote(&script_path.to_string_lossy())

View file

@ -598,7 +598,8 @@ impl ImportTransform {
| "reasoning_effort"
| "speed"
| "backend"
| "acp_command"
| "acp.command"
| "acp.config"
| "fidelity"
| "max_retries"
| "thread_id"
@ -1005,7 +1006,7 @@ mod tests {
let graph = apply_import(
r#"digraph Deploy {
start [shape=Mdiamond]
validate [import="./validate.fabro", model="haiku", backend="acp", acp_command="python fake_agent.py", class="fast, shared"]
validate [import="./validate.fabro", model="haiku", backend="acp", acp.command="python fake_agent.py", class="fast, shared"]
exit [shape=Msquare]
start -> validate -> exit
}"#,
@ -1055,7 +1056,7 @@ mod tests {
assert_eq!(
graph.nodes["validate.test"]
.attrs
.get("acp_command")
.get("acp.command")
.and_then(AttrValue::as_str),
Some("python fake_agent.py")
);

View file

@ -216,20 +216,20 @@ mod tests {
#[test]
fn apply_backend_property_via_stylesheet() {
let ss = parse_stylesheet("* { backend: cli; }").unwrap();
let ss = parse_stylesheet("* { backend: acp; }").unwrap();
let mut graph = Graph::new("test");
graph.nodes.insert("a".into(), Node::new("a"));
apply_stylesheet(&ss, &mut graph);
assert_eq!(
graph.nodes["a"].attrs.get("backend"),
Some(&AttrValue::String("cli".into()))
Some(&AttrValue::String("acp".into()))
);
}
#[test]
fn backend_property_not_overridden_by_stylesheet() {
let ss = parse_stylesheet("* { backend: cli; }").unwrap();
let ss = parse_stylesheet("* { backend: acp; }").unwrap();
let mut graph = Graph::new("test");
let mut node = Node::new("a");
node.attrs

View file

@ -24,7 +24,6 @@ use std::sync::Arc;
use fabro_agent::Sandbox;
use fabro_graphviz::graph::{AttrValue, Edge, Graph, Node};
use fabro_model::ProviderId;
use fabro_sandbox::daytona::{DaytonaConfig, DaytonaSandbox, DaytonaSnapshotConfig};
use fabro_static::EnvVars;
use fabro_store::{ArtifactKey, ArtifactStore, Database};
@ -991,151 +990,6 @@ async fn daytona_parallel_git_branching_e2e() {
env.cleanup().await.expect("Daytona cleanup should succeed");
}
// ---------------------------------------------------------------------------
// CLI Backend on Daytona — real CLI tools via exec_command
// ---------------------------------------------------------------------------
use fabro_workflow::handler::agent::{CodergenBackend, CodergenResult, CodergenRunRequest};
use fabro_workflow::handler::llm::AgentCliBackend;
/// Helper: run a real CLI backend test on Daytona.
///
/// Installs the CLI tool in the sandbox, then runs the AgentCliBackend against
/// it.
async fn run_daytona_cli_test(provider: ProviderId, model: &str, install_command: &str) {
let creds = load_github_app_credentials();
let config = DaytonaConfig {
snapshot: Some(DaytonaSnapshotConfig {
name: "daytona-medium".into(),
cpu: None,
memory: None,
disk: None,
dockerfile: None,
}),
..DaytonaConfig::default()
};
let env = DaytonaSandbox::new(config, Some(creds), None, None, None, None)
.await
.expect("Failed to create Daytona client — is DAYTONA_API_KEY set?");
env.initialize()
.await
.expect("Daytona sandbox should initialize");
let env: Arc<dyn Sandbox> = Arc::new(env);
// Install prerequisites (bash, curl, Node 20 via nodesource) if not available
let prereq_check = env
.exec_command(
"bash --version && curl --version && node --version && npm --version",
10_000,
None,
None,
None,
)
.await;
if prereq_check.as_ref().map_or(true, |r| !r.is_success()) {
let prereq = env
.exec_command(
"apt-get update -qq && apt-get install -y -qq bash curl ca-certificates gnupg >/dev/null 2>&1 \
&& curl -fsSL https://deb.nodesource.com/setup_20.x | bash - >/dev/null 2>&1 \
&& apt-get install -y -qq nodejs >/dev/null 2>&1",
180_000,
None,
None,
None,
)
.await
.expect("prerequisite install should not error");
assert_eq!(
prereq.exit_code,
Some(0),
"prerequisite install failed: {}",
prereq.stderr
);
}
// Install the CLI tool inside the Daytona sandbox
let install_result = env
.exec_command(install_command, 120_000, None, None, None)
.await
.expect("install command should not error");
assert_eq!(
install_result.exit_code,
Some(0),
"install command failed (exit {:?}): {}",
install_result.exit_code,
install_result.stdout
);
let backend = AgentCliBackend::new_from_env(model.to_string(), provider.clone());
let node = Node::new("daytona_cli_test");
let context = Context::new();
let emitter = Arc::new(Emitter::default());
let result = backend
.run(CodergenRunRequest {
node: &node,
prompt: "What is 2+2? Reply with just the number.",
context: &context,
thread_id: None,
emitter: &emitter,
sandbox: &env,
tool_hooks: None,
cancel_token: CancellationToken::new(),
})
.await;
match result {
Ok(CodergenResult::Text { text, usage, .. }) => {
assert!(
text.contains('4'),
"{provider}/{model} on Daytona: expected '4', got: {text}"
);
if let Some(u) = usage {
assert!(
u.tokens().input_tokens > 0,
"{provider}/{model}: input_tokens should be > 0"
);
}
}
Ok(CodergenResult::Full(_)) => panic!("expected Text result"),
Err(e) => panic!("{provider}/{model} on Daytona failed: {e}"),
}
env.cleanup()
.await
.expect("Daytona sandbox cleanup should succeed");
}
#[fabro_macros::e2e_test(live("DAYTONA_API_KEY"), live("GITHUB_APP_PRIVATE_KEY"))]
async fn daytona_cli_claude() {
run_daytona_cli_test(
ProviderId::anthropic(),
"haiku",
"curl -fsSL https://claude.ai/install.sh | bash",
)
.await;
}
#[fabro_macros::e2e_test(live("DAYTONA_API_KEY"), live("GITHUB_APP_PRIVATE_KEY"))]
async fn daytona_cli_codex() {
run_daytona_cli_test(
ProviderId::openai(),
"o4-mini",
"npm install -g @openai/codex",
)
.await;
}
#[fabro_macros::e2e_test(live("DAYTONA_API_KEY"), live("GITHUB_APP_PRIVATE_KEY"))]
async fn daytona_cli_gemini() {
run_daytona_cli_test(
ProviderId::gemini(),
"gemini-2.5-flash",
"npm install -g @google/gemini-cli",
)
.await;
}
// ---------------------------------------------------------------------------
// Daytona shadow commit E2E with sandbox-native metadata
// ---------------------------------------------------------------------------

File diff suppressed because it is too large Load diff