fabro/lib/components/fabro-workflow/src/pipeline/initialize.rs
Scott Werner d65785d888 Simplify blob activation and share the test store fixture
Blob activation cleanups:
- Reuse fabro-db's append_to_path, remove_file_if_exists, and
  set_private_permissions instead of local duplicates.
- Return the store directly from activate_blob_storage; the report
  wrapper existed only to be logged internally and then discarded.
- Collapse compute_disk_preflight to return the required free bytes
  instead of echoing its inputs back through a struct.
- Deduplicate the "exactly one ok row" PRAGMA integrity_check protocol
  into one executor-generic helper used by the backup and live checks.
- Skip re-validating a freshly published backup; the staging copy was
  validated immediately before the atomic rename, so only a
  concurrently published file needs its own validation.
- Replace the manual anyhow wrapping plus duplicate error log in
  serve.rs with a plain .context(), matching other startup errors.
- Extract the disk-candidate enumeration in resource_sampler.rs that
  available_space_for_path had copy-pasted from sample_disk_resources.

Test fixture cleanups:
- Route all hand-assembled Database::new(..., test_blob_store()) test
  fixtures (32 sites) through fabro_store::test_support::test_database,
  and make that helper infallible instead of returning an unconditional
  Ok.
- Install the test blob schema from fabro_db::BLOBS_MIGRATION_SQL via a
  test-support-gated optional dependency instead of a four-level
  relative include_str! into fabro-db's migrations directory.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-24 14:02:35 -04:00

1799 lines
66 KiB
Rust

use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::Arc;
use std::time::Instant;
use fabro_agent::{Sandbox, ToolSecrets};
use fabro_auth::{
CredentialSource, ExtraHeadersCredentialSource, VaultCredentialSource, auth_issue_message,
};
use fabro_github::token_source::InstallationTokenSource;
use fabro_graphviz::graph;
use fabro_hooks::{HookContext, HookDecision, HookEvent, HookExecutionContext, HookRunner};
use fabro_model::Catalog;
use fabro_sandbox::{
GitSetupIntent, SandboxEventCallback, SandboxSpec, reconnect_for_run_with_callback, shell_quote,
};
use fabro_static::EnvVars;
use fabro_types::RunSandboxKind;
use fabro_vault::Vault;
use tokio::runtime::Handle;
use tokio::sync::RwLock as AsyncRwLock;
use super::types::{InitOptions, Initialized, LlmSpec, Persisted, SandboxEnvSpec};
use crate::error::Error;
use crate::event::{Event, RunNoticeCode, RunNoticeLevel};
use crate::git::GitAuthor;
use crate::git_bridge;
use crate::handler::llm::{AgentAcpBackend, AgentApiBackend, BackendRouter, routing};
use crate::handler::{HandlerRegistry, default_registry};
#[cfg(test)]
use crate::model_fallback::ModelFallbackPolicy;
use crate::run_metadata::{RunMetadataRuntime, build_metadata_writer, metadata_branch_name};
use crate::run_options::{GitCheckpointOptions, RunOptions};
use crate::sandbox_git_runtime::SandboxGitRuntime;
use crate::services::{
EngineServices, FabroRunToolServices, RunLocations, RunServices, WorkflowToolEnvProvider,
};
use crate::stage_execution::{StageExecutionSeed, StageExecutionTracker};
use crate::steering_hub::SteeringHub;
struct BuiltSandboxEnv {
env: HashMap<String, String>,
github_token: Option<Arc<InstallationTokenSource>>,
/// The validated effective repository set behind `github_token`.
/// Present only in App mode or when additional repositories are
/// declared; drives the eager access validation at initialization.
github_access: Option<fabro_github::GitHubRepositoryAccess>,
}
async fn run_hooks(
hook_runner: Option<&HookRunner>,
hook_context: &HookContext,
sandbox: Arc<dyn Sandbox>,
execution_context: HookExecutionContext,
) -> HookDecision {
let Some(runner) = hook_runner else {
return HookDecision::Proceed;
};
runner.run(hook_context, sandbox, execution_context).await
}
fn git_setup_intent(run_options: &RunOptions) -> GitSetupIntent {
if let Some(source) = run_options.fork_source_ref.as_ref() {
GitSetupIntent::ForkFromCheckpoint {
new_run_id: run_options.run_id.to_string(),
source_run_id: source.source_run_id.to_string(),
checkpoint_sha: source.checkpoint_sha.clone(),
}
} else {
GitSetupIntent::NewRun {
run_id: run_options.run_id.to_string(),
}
}
}
async fn configure_sandbox_git_identity(
sandbox: &dyn Sandbox,
author: &GitAuthor,
) -> Result<(), Error> {
let command = format!(
"git config --local user.name {} && git config --local user.email {}",
shell_quote(&author.name),
shell_quote(&author.email)
);
sandbox
.exec_command(&command, 10_000, None, None, None)
.await
.map_err(|err| Error::engine_with_source("Sandbox git identity setup failed", err))?
.into_result("git config user identity")
.map_err(|err| Error::engine_with_source("Sandbox git identity setup failed", err))?;
Ok(())
}
fn build_sandbox_env(
spec: &SandboxEnvSpec,
github_app: Option<&fabro_github::GitHubCredentials>,
) -> Result<BuiltSandboxEnv, Error> {
let mut env = spec.toml_env.clone();
let no_token = |env| BuiltSandboxEnv {
env,
github_token: None,
github_access: None,
};
let Some(integration) = spec
.github_integration
.as_ref()
.filter(|integration| integration.is_token_requested())
else {
return Ok(no_token(env));
};
let declares_additional = integration.has_additional_repositories();
let Some(creds) = github_app else {
if declares_additional {
// Legacy permissions-only configuration stays best-effort, but a
// declared additional set is an explicit access requirement.
return Err(Error::Precondition(
"run.integrations.github.additional_repositories requires GitHub credentials, \
but none are configured"
.to_string(),
));
}
return Ok(no_token(env));
};
// Validate the effective repository set whenever it matters: App mode
// scopes the mint to it, and any declared additional set must hold its
// invariants regardless of credential kind. Legacy PAT/static
// permissions-only runs skip it to preserve their origin-agnostic
// behavior.
let github_access =
if declares_additional || matches!(creds, fabro_github::GitHubCredentials::App(_)) {
fabro_github::GitHubRepositoryAccess::new(
spec.origin_url.as_deref(),
&integration.additional_repositories,
integration.permissions.clone(),
)
.map_err(|err| {
Error::engine_with_anyhow("Failed to validate GitHub repository access", err)
})?
} else {
None
};
let github_token = match github_access.as_ref() {
Some(access) => Some(InstallationTokenSource::for_access(creds, access).map_err(
|err| Error::engine_with_anyhow("Failed to build GitHub token source", err),
)?),
None => match creds {
fabro_github::GitHubCredentials::Pat(token) => {
Some(InstallationTokenSource::pat(token.clone()))
}
fabro_github::GitHubCredentials::Installation(token) => {
Some(InstallationTokenSource::installation(token.clone()))
}
// No origin URL and nothing declared: keep the legacy App-mode
// best-effort skip.
fabro_github::GitHubCredentials::App(_) => None,
},
};
if declares_additional {
let access = github_access
.as_ref()
.expect("access is always constructed when additional repositories are declared");
git_bridge::merge_git_bridge_env(&mut env, &access.targets())?;
}
Ok(BuiltSandboxEnv {
env,
github_token,
github_access,
})
}
/// When additional repositories are declared, resolve their token before the
/// first workflow stage. App-backed sources first check that every target is
/// on one installation. Static credentials resolve locally; the first Git
/// operation remains their access check. Legacy permissions-only runs skip
/// eager resolution.
async fn resolve_declared_repository_token(built: &BuiltSandboxEnv) -> Result<(), Error> {
let Some(_) = built
.github_access
.as_ref()
.filter(|access| access.has_additional_repositories())
else {
return Ok(());
};
// `build_sandbox_env` guarantees a token source whenever additional
// repositories are declared; fail closed if that ever breaks.
let Some(source) = built.github_token.as_ref() else {
return Err(Error::Precondition(
"run.integrations.github.additional_repositories requires GitHub credentials, but \
none are configured"
.to_string(),
));
};
source.resolve().await.map_err(|err| {
Error::engine_with_anyhow(
"Failed to resolve the GitHub token for the declared repository set",
err,
)
})?;
Ok(())
}
async fn build_registry(
spec: &LlmSpec,
interviewer: Arc<dyn fabro_interview::Interviewer>,
steering_hub: Arc<SteeringHub>,
tool_env_provider: Arc<WorkflowToolEnvProvider>,
github_token_refresh_managed: bool,
graph: &graph::Graph,
llm_source: Arc<dyn CredentialSource>,
catalog: Arc<Catalog>,
tool_secrets: ToolSecrets,
fabro_run_tools: Option<FabroRunToolServices>,
) -> Result<(Arc<HandlerRegistry>, bool), Error> {
let no_backend_interviewer = Arc::clone(&interviewer);
let build_no_backend = move || {
Arc::new(default_registry(
Arc::clone(&no_backend_interviewer),
|| None,
))
};
if spec.dry_run {
return Ok((build_no_backend(), true));
}
let graph_needs_llm = graph
.nodes
.values()
.any(|n| graph::is_llm_handler_type(n.handler_type()));
if !graph_needs_llm {
return Ok((build_no_backend(), false));
}
let build_llm_registry = || {
let model = spec.model.clone();
let provider_id = spec.provider_id.clone();
let fallbacks = spec.fallbacks.clone();
let mcp_servers = spec.mcp_servers.clone();
let model_controls = spec.model_controls.clone();
let tool_secrets_for_api = tool_secrets.clone();
let llm_source_for_api = Arc::clone(&llm_source);
let catalog_for_api = Arc::clone(&catalog);
let steering_hub_for_api = Arc::clone(&steering_hub);
let tool_env_provider_for_backend = Arc::clone(&tool_env_provider);
let fabro_run_tools_for_api = fabro_run_tools.clone();
Arc::new(default_registry(interviewer, move || {
let tool_env_provider = Arc::clone(&tool_env_provider_for_backend);
let mut api = AgentApiBackend::new_with_catalog(
model.clone(),
provider_id.clone(),
fallbacks.clone(),
Arc::clone(&llm_source_for_api),
Arc::clone(&steering_hub_for_api),
Arc::clone(&catalog_for_api),
)
.with_run_model_controls(model_controls.clone())
.with_tool_env_provider(tool_env_provider.clone())
.with_tool_secrets(tool_secrets_for_api.clone())
.with_mcp_servers(mcp_servers.clone());
if let Some(services) = fabro_run_tools_for_api.clone() {
api = api.with_fabro_run_tools(services);
}
let acp = AgentAcpBackend::new()
.with_tool_env_provider(tool_env_provider.clone(), github_token_refresh_managed)
.with_steering_hub(Arc::clone(&steering_hub));
Some(Box::new(BackendRouter::new(Box::new(api), acp)))
}))
};
if !graph_needs_api_backend(graph) {
return Ok((build_llm_registry(), false));
}
match llm_source.resolve(catalog.as_ref()).await {
Ok(result) if result.credentials.is_empty() => {
if graph_needs_llm {
let detail = (!result.auth_issues.is_empty()).then(|| {
result
.auth_issues
.iter()
.map(|(provider, issue)| auth_issue_message(provider, issue))
.collect::<Vec<_>>()
.join("; ")
});
let prefix = detail.map_or_else(
|| "No LLM providers configured".to_string(),
|detail| format!("No usable LLM providers configured: {detail}"),
);
return Err(Error::Precondition(format!(
"{prefix}. Set ANTHROPIC_API_KEY or OPENAI_API_KEY, or pass --dry-run to simulate."
)));
}
Ok((build_no_backend(), false))
}
Ok(_result) => Ok((build_llm_registry(), false)),
Err(e) => {
if graph_needs_llm {
return Err(Error::Precondition(format!(
"Failed to initialize LLM client: {e}. Set ANTHROPIC_API_KEY or OPENAI_API_KEY, or pass --dry-run to simulate.",
)));
}
Ok((build_no_backend(), false))
}
}
}
async fn tool_secrets_from_configured_sources(vault: &Arc<AsyncRwLock<Vault>>) -> ToolSecrets {
let vault = vault.read().await;
ToolSecrets {
brave_search_api_key: vault.get(EnvVars::BRAVE_SEARCH_API_KEY).map(str::to_string),
venice_api_key: vault.get(EnvVars::VENICE_API_KEY).map(str::to_string),
}
}
fn graph_needs_api_backend(graph: &graph::Graph) -> bool {
graph.nodes.values().any(routing::node_needs_api_backend)
}
/// Trace header attached to every LLM request in a run so gateways that
/// understand it (e.g. OpenRouter broadcast) can group the run's requests
/// into one session. Explicit `extra_headers` provider configuration wins.
const SESSION_ID_HEADER: &str = "x-session-id";
fn build_llm_source(
vault: Arc<AsyncRwLock<Vault>>,
run_id: fabro_types::RunId,
) -> Arc<dyn CredentialSource> {
Arc::new(ExtraHeadersCredentialSource::new(
Arc::new(VaultCredentialSource::new(vault)),
HashMap::from([(SESSION_ID_HEADER.to_string(), run_id.to_string())]),
))
}
/// INITIALIZE phase: prepare the sandbox, env, and handlers for execution.
pub async fn initialize(
persisted: Persisted,
mut options: InitOptions,
) -> Result<Initialized, Error> {
let (graph, source, _diagnostics, run_dir, run_spec) = persisted.into_parts();
let (checkpoint, stage_executions) = options.resume.take().map_or_else(
|| (None, StageExecutionSeed::default()),
|resume| {
let (checkpoint, stage_executions) = resume.into_parts();
(Some(checkpoint), stage_executions)
},
);
let host_source_dir = run_spec.source_directory.as_deref().map(PathBuf::from);
options.run_options.run_dir = run_dir.clone();
options.run_options.git = options.git.clone();
let llm_source = build_llm_source(options.vault.clone(), options.run_options.run_id);
let tool_secrets = tool_secrets_from_configured_sources(&options.vault).await;
let catalog = Arc::clone(&options.catalog);
let sandbox_git = Arc::new(SandboxGitRuntime::new());
let metadata_runtime = Arc::new(RunMetadataRuntime::new());
let hook_runner = if options.hooks.hooks.is_empty() {
None
} else {
Some(Arc::new(HookRunner::new(
options.hooks.clone(),
Arc::clone(&llm_source),
Arc::clone(&catalog),
)))
};
let is_resume = checkpoint.is_some();
options.run_options.display_base_sha = options
.run_options
.pre_run_git
.as_ref()
.and_then(|git| git.sha.clone());
if !is_resume
&& !matches!(options.sandbox, SandboxSpec::Local { .. })
&& matches!(
options
.run_options
.pre_run_git
.as_ref()
.map(|git| git.dirty),
Some(fabro_types::DirtyStatus::Dirty)
)
{
options.emitter.notice(
RunNoticeLevel::Warn,
RunNoticeCode::DirtyWorktree,
"Uncommitted changes will not be included in the remote sandbox.",
);
}
let sandbox_event_callback: SandboxEventCallback = {
let emitter = Arc::clone(&options.emitter);
Arc::new(move |event| {
emitter.emit(&Event::Sandbox { event });
})
};
let attach_instance = if is_resume {
let record = options
.run_store
.state()
.await
.map_err(|err| Error::engine(err.to_string()))?
.sandbox
.ok_or_else(|| {
Error::Precondition("cannot resume run: run sandbox is missing".to_string())
})?;
// A fork carries a checkpoint from its source run, but its first
// `run.created` event contains only a sandbox plan. Materialize that
// sandbox before resuming. Later fork resumes reconnect the ready
// instance.
let fork_needs_materialization = options.run_options.fork_source_ref.is_some()
&& record.kind() == RunSandboxKind::Planned;
if fork_needs_materialization {
None
} else {
Some(record.into_instance().ok_or_else(|| {
Error::Precondition(
"cannot resume run: run sandbox was not initialized".to_string(),
)
})?)
}
} else {
None
};
let attach_existing = attach_instance.is_some();
let sandbox: Arc<dyn Sandbox> = if let Some(instance) = attach_instance {
let daytona_api_key = options
.vault
.read()
.await
.get(EnvVars::DAYTONA_API_KEY)
.map(str::to_string);
let sandbox = reconnect_for_run_with_callback(
&instance,
daytona_api_key,
Some(options.run_options.run_id),
Some(Arc::clone(&sandbox_event_callback)),
)
.await
.map_err(|err| Error::engine_with_anyhow("Failed to reconnect sandbox for resume", err))?;
Arc::from(sandbox)
} else {
options
.sandbox
.build(Some(Arc::clone(&sandbox_event_callback)))
.await
.map_err(|e| Error::engine_with_anyhow("Failed to build sandbox", e))?
};
let cleanup_guard = (!attach_existing).then(|| {
scopeguard::guard(Arc::clone(&sandbox), |sandbox| {
if let Ok(handle) = Handle::try_current() {
handle.spawn(async move {
let _ = sandbox.delete().await;
});
}
})
});
if attach_existing {
// Resume needs the full provider health check. `activate()` is the
// lighter access-time operation used after a run is already active.
sandbox
.start()
.await
.map_err(|e| Error::engine_with_source("Failed to start sandbox", e))?;
} else {
sandbox
.initialize()
.await
.map_err(|e| Error::engine_with_source("Failed to initialize sandbox", e))?;
}
let locations = RunLocations::for_sandbox(host_source_dir, sandbox.as_ref(), run_dir.clone());
let hook_ctx = HookContext::new(
HookEvent::SandboxReady,
options.run_options.run_id,
graph.name.clone(),
);
let decision = run_hooks(
hook_runner.as_deref(),
&hook_ctx,
Arc::clone(&sandbox),
locations.hook_execution_context(),
)
.await;
if let HookDecision::Block { reason } = decision {
let msg = reason.unwrap_or_else(|| "blocked by SandboxReady hook".into());
return Err(Error::engine(msg));
}
if !attach_existing {
let run_sandbox = options
.sandbox
.to_run_sandbox_instance(&*sandbox, options.run_options.run_id);
let runtime = &run_sandbox.runtime;
options.emitter.emit(&Event::SandboxInitialized {
working_directory: runtime.working_directory.clone(),
provider: run_sandbox.provider,
id: runtime.id.clone(),
image: run_sandbox.image.clone(),
snapshot: run_sandbox.snapshot.clone(),
repo_cloned: runtime.repo_cloned,
clone_origin_url: runtime.clone_origin_url.clone(),
clone_branch: runtime.clone_branch.clone(),
workspace_root: runtime.workspace_root.clone(),
repos_root: runtime.repos_root.clone(),
primary_repo_path: runtime.primary_repo_path.clone(),
primary_repo_link: runtime.primary_repo_link.clone(),
});
}
let built_env = build_sandbox_env(
&options.sandbox_env,
options.run_options.github_app.as_ref(),
)?;
resolve_declared_repository_token(&built_env).await?;
let BuiltSandboxEnv {
env: base_env,
github_token,
github_access: _,
} = built_env;
let tool_env_provider = Arc::new(WorkflowToolEnvProvider {
base_env: base_env.clone(),
github_token: github_token.clone(),
});
let github_token_refresh_managed = github_token
.as_deref()
.is_some_and(InstallationTokenSource::mints_installation_tokens);
let (registry, effective_dry_run) = if let Some(registry) = options.registry_override.clone() {
// A caller-supplied registry owns execution behavior for its handlers.
(registry, options.dry_run)
} else {
build_registry(
&options.llm,
Arc::clone(&options.interviewer),
Arc::clone(&options.steering_hub),
Arc::clone(&tool_env_provider),
github_token_refresh_managed,
&graph,
Arc::clone(&llm_source),
Arc::clone(&catalog),
tool_secrets.clone(),
options.fabro_run_tools.clone(),
)
.await?
};
if effective_dry_run {
use fabro_types::settings::run::RunMode;
options.dry_run = true;
options.run_options.settings.run.execution.mode = RunMode::DryRun;
}
let has_run_branch = options
.run_options
.git
.as_ref()
.and_then(|g| g.run_branch.as_ref())
.is_some();
if options.run_options.settings.run.run_branch.enabled && !has_run_branch {
let intent = git_setup_intent(&options.run_options);
let sandbox_has_origin = sandbox.origin_url().is_some();
if sandbox_has_origin {
sandbox_git
.ensure_git_available(&*sandbox)
.await
.map_err(|err| Error::engine_with_source("sandbox git unavailable", err))?;
}
match sandbox.setup_git(&intent).await {
Ok(Some(info)) => {
let base_sha = options
.run_options
.git
.as_ref()
.and_then(|g| g.base_sha.clone())
.or(Some(info.base_sha.clone()));
options.run_options.display_base_sha.clone_from(&base_sha);
options.run_options.git = Some(GitCheckpointOptions {
base_sha,
run_branch: Some(info.run_branch.clone()),
meta_branch: options
.run_options
.settings
.run
.meta_branch
.enabled
.then(|| metadata_branch_name(&options.run_options.run_id.to_string())),
});
if options.run_options.base_branch.is_none() {
options.run_options.base_branch = info.base_branch;
}
}
Ok(None) => {
if sandbox_has_origin {
options.emitter.notice(
RunNoticeLevel::Warn,
RunNoticeCode::SandboxGitUnavailable,
"Sandbox could not set up Git despite a configured origin; running \
without checkpointing or PR support.",
);
}
}
Err(e) => {
return Err(Error::engine_with_source("Sandbox git setup failed", e));
}
}
}
if sandbox.origin_url().is_some() {
let git_author = options.run_options.git_author();
configure_sandbox_git_identity(sandbox.as_ref(), &git_author).await?;
}
if !options.lifecycle.setup_commands.is_empty() {
options.emitter.emit(&Event::SetupStarted {
command_count: options.lifecycle.setup_commands.len(),
});
let setup_start = Instant::now();
for (index, setup) in options.lifecycle.setup_commands.iter().enumerate() {
let command = &setup.command;
options.emitter.emit(&Event::SetupCommandStarted {
command: command.clone(),
index,
});
let cmd_start = Instant::now();
let cancel_token = options.run_options.cancel_token.child_token();
let step_env = (!setup.env.is_empty()).then_some(&setup.env);
let result = sandbox
.exec_command(
command,
options.lifecycle.setup_command_timeout_ms,
None,
step_env,
Some(cancel_token.clone()),
)
.await
.map_err(|e| Error::engine_with_source("Setup command failed", e))?;
if options.run_options.cancel_token.is_cancelled() {
return Err(Error::Cancelled);
}
cancel_token.cancel();
let duration_ms = crate::millis_u64(cmd_start.elapsed());
if !result.is_success() {
let exit_code = result.display_exit_code();
let exec_output_tail = result.default_redacted_output_tail();
options.emitter.emit(&Event::SetupFailed {
command: command.clone(),
index,
exit_code,
stderr: result.stderr.clone(),
exec_output_tail,
});
return Err(Error::engine(format!(
"Setup command failed (exit code {}): {command}\n{}",
exit_code, result.stderr,
)));
}
let exit_code = result.exit_code.unwrap_or(0);
options.emitter.emit(&Event::SetupCommandCompleted {
command: command.clone(),
index,
exit_code,
duration_ms,
});
}
options.emitter.emit(&Event::SetupCompleted {
duration_ms: crate::millis_u64(setup_start.elapsed()),
});
}
let metadata_writer =
match build_metadata_writer(&options.run_options, sandbox.push_token_source()) {
Ok(writer) => writer,
Err(err) => {
let message = format!("failed to initialize checkpoint metadata writer: {err}");
if metadata_runtime.mark_metadata_degraded(false) {
options.emitter.notice(
RunNoticeLevel::Warn,
RunNoticeCode::CheckpointMetadataWriteFailed,
message,
);
}
None
}
};
let run_services = RunServices::new(
options.run_store.clone(),
Arc::clone(&options.emitter),
Arc::clone(&sandbox),
hook_runner.clone(),
locations,
options.run_options.cancel_token.clone(),
options.llm.provider_id.clone(),
options.llm.model.clone(),
Arc::clone(&llm_source),
catalog,
sandbox_git,
metadata_runtime,
metadata_writer,
StageExecutionTracker::seeded(stage_executions),
);
let engine = Arc::new(EngineServices {
run: Arc::clone(&run_services),
registry,
interviewer: Arc::clone(&options.interviewer),
base_env,
github_token,
inputs: options.run_options.settings.run.inputs.clone(),
dry_run: options.dry_run,
workflow_path: options.workflow_path.clone(),
workflow_bundle: options.workflow_bundle.clone(),
});
if let Some(cleanup_guard) = cleanup_guard {
scopeguard::ScopeGuard::into_inner(cleanup_guard);
}
Ok(Initialized {
graph,
source,
run_options: options.run_options,
checkpoint,
seed_context: options.seed_context,
on_node: None,
artifact_sink: options.artifact_sink,
run_control: options.run_control,
engine,
model: options.llm.model,
})
}
#[cfg(test)]
mod tests {
use std::collections::{BTreeMap, HashMap};
use std::sync::Arc;
use std::time::Duration;
use fabro_acp::test_support::fake_acp_agent_script;
use fabro_auth::test_support as auth_test_support;
use fabro_graphviz::graph::{AttrValue, Edge, Graph, Node};
use fabro_interview::AutoApproveInterviewer;
use fabro_sandbox::SandboxSpec;
use fabro_store::{Database, RunDatabase};
use fabro_types::settings::run::RunModelControls;
use fabro_types::{
EventBody, ForkSourceRef, RunEvent, RunId, WorkflowSettings, fixtures, test_support,
};
use fabro_vault::{SecretType, Vault};
use object_store::memory::InMemory;
use tokio::fs::{create_dir_all, write};
use tokio::sync::RwLock as AsyncRwLock;
use super::*;
use crate::context::{Context, keys};
use crate::event::StoreProgressLogger;
use crate::pipeline::ResumeState;
use crate::pipeline::types::InitOptions;
use crate::records::{Checkpoint, CheckpointExt, RunSpec};
use crate::run_options::RunOptions;
use crate::stage_execution::StageExecutionSeed;
const CHECKPOINT_SHA: &str = "abc123";
fn test_run_id() -> RunId {
fixtures::RUN_1
}
fn setup_cmd(command: &str) -> crate::run_options::SetupCommand {
crate::run_options::SetupCommand {
command: command.to_string(),
env: HashMap::new(),
}
}
fn test_catalog() -> Arc<Catalog> {
Arc::new(Catalog::from_builtin().expect("default catalog should build"))
}
fn memory_store() -> Arc<Database> {
Arc::new(fabro_store::test_support::test_database(
Arc::new(InMemory::new()),
"",
Duration::from_millis(1),
None,
))
}
async fn seed_run_created(
run_store: &RunDatabase,
settings: serde_json::Value,
graph: serde_json::Value,
source_directory: Option<String>,
fork_source_ref: Option<ForkSourceRef>,
) {
crate::event::append_event(run_store, &test_run_id(), &Event::RunCreated {
run_id: test_run_id(),
title: None,
settings,
graph,
workflow_source: None,
labels: BTreeMap::new(),
source_directory,
workflow_slug: Some("test".to_string()),
workflow_version_id: None,
target: None,
automation: None,
provenance: test_support::test_run_provenance(),
manifest_blob: None,
spec_blob: None,
git: None,
fork_source_ref,
retried_from: None,
parent_id: None,
web_url: None,
})
.await
.unwrap();
}
fn simple_graph() -> (Graph, String) {
let source = r"digraph test {
start [shape=Mdiamond];
exit [shape=Msquare];
start -> exit;
}"
.to_string();
let mut graph = Graph::new("test");
let mut start = Node::new("start");
start.attrs.insert(
"shape".to_string(),
AttrValue::String("Mdiamond".to_string()),
);
let mut exit = Node::new("exit");
exit.attrs.insert(
"shape".to_string(),
AttrValue::String("Msquare".to_string()),
);
graph.nodes.insert("start".to_string(), start);
graph.nodes.insert("exit".to_string(), exit);
graph.edges.push(Edge::new("start", "exit"));
(graph, source)
}
fn llm_graph() -> (Graph, String) {
let source = r"digraph test {
start [shape=Mdiamond];
writer [shape=box];
exit [shape=Msquare];
start -> writer;
writer -> exit;
}"
.to_string();
let mut graph = Graph::new("test");
let mut start = Node::new("start");
start.attrs.insert(
"shape".to_string(),
AttrValue::String("Mdiamond".to_string()),
);
let mut writer = Node::new("writer");
writer
.attrs
.insert("shape".to_string(), AttrValue::String("box".to_string()));
let mut exit = Node::new("exit");
exit.attrs.insert(
"shape".to_string(),
AttrValue::String("Msquare".to_string()),
);
graph.nodes.insert("start".to_string(), start);
graph.nodes.insert("writer".to_string(), writer);
graph.nodes.insert("exit".to_string(), exit);
graph.edges.push(Edge::new("start", "writer"));
graph.edges.push(Edge::new("writer", "exit"));
(graph, source)
}
fn test_settings(run_dir: &std::path::Path) -> RunOptions {
RunOptions {
settings: WorkflowSettings::default(),
run_dir: run_dir.to_path_buf(),
cancel_token: tokio_util::sync::CancellationToken::new(),
run_id: test_run_id(),
labels: HashMap::new(),
workflow_slug: None,
github_app: None,
pre_run_git: None,
fork_source_ref: None,
base_branch: None,
display_base_sha: None,
git: None,
}
}
fn test_init_options(
run_store: crate::runtime_store::RunStoreHandle,
emitter: Arc<crate::event::Emitter>,
working_directory: std::path::PathBuf,
run_options: RunOptions,
) -> InitOptions {
InitOptions {
run_store,
dry_run: false,
emitter: Arc::clone(&emitter),
sandbox: SandboxSpec::Local { working_directory },
llm: LlmSpec {
model: "test-model".to_string(),
provider_id: fabro_model::ProviderId::anthropic(),
fallbacks: ModelFallbackPolicy::default(),
mcp_servers: Vec::new(),
model_controls: RunModelControls::default(),
dry_run: true,
},
interviewer: Arc::new(AutoApproveInterviewer::engine()),
steering_hub: Arc::new(crate::steering_hub::SteeringHub::new(emitter)),
catalog: test_catalog(),
lifecycle: crate::run_options::LifecycleOptions {
setup_commands: vec![],
setup_command_timeout_ms: 1_000,
},
run_options,
workflow_path: None,
workflow_bundle: None,
hooks: fabro_hooks::HookSettings { hooks: vec![] },
sandbox_env: SandboxEnvSpec {
toml_env: HashMap::new(),
github_integration: None,
origin_url: None,
},
vault: auth_test_support::empty_vault(),
git: None,
run_control: None,
registry_override: None,
artifact_sink: None,
resume: None,
seed_context: None,
fabro_run_tools: None,
}
}
fn test_persisted(graph: Graph, source: String, run_dir: &std::path::Path) -> Persisted {
test_persisted_run(graph, source, run_dir, WorkflowSettings::default(), None)
}
fn test_persisted_run(
graph: Graph,
source: String,
run_dir: &std::path::Path,
settings: WorkflowSettings,
fork_source_ref: Option<ForkSourceRef>,
) -> Persisted {
Persisted::new(
graph.clone(),
source,
vec![],
run_dir.to_path_buf(),
RunSpec {
run_id: test_run_id(),
settings,
graph,
graph_source: None,
workflow_slug: Some("test".to_string()),
workflow_version_id: None,
target: None,
automation: None,
source_directory: Some(std::env::current_dir().unwrap().display().to_string()),
git: Some(fabro_types::GitContext {
origin_url: String::new(),
branch: "main".to_string(),
sha: None,
dirty: fabro_types::DirtyStatus::Clean,
}),
labels: HashMap::new(),
provenance: test_support::test_run_provenance(),
manifest_blob: None,
definition_blob: None,
spec_blob: None,
fork_source_ref,
},
)
}
async fn initialize_with_setup_command(
command: &str,
) -> (crate::error::Result<Initialized>, Vec<RunEvent>) {
initialize_with_setup_step(setup_cmd(command)).await
}
async fn initialize_with_setup_step(
setup: crate::run_options::SetupCommand,
) -> (crate::error::Result<Initialized>, Vec<RunEvent>) {
let temp = tempfile::tempdir().unwrap();
let run_dir = temp.path().join("run");
std::fs::create_dir_all(&run_dir).unwrap();
let (graph, source) = simple_graph();
let persisted = test_persisted(graph, source, &run_dir);
let emitter = Arc::new(crate::event::Emitter::new(test_run_id()));
let seen = Arc::new(std::sync::Mutex::new(Vec::new()));
emitter.on_event({
let seen = Arc::clone(&seen);
move |event| seen.lock().unwrap().push(event.clone())
});
let run_store = memory_store().create_run(&test_run_id()).await.unwrap();
let result = initialize(persisted, InitOptions {
lifecycle: crate::run_options::LifecycleOptions {
setup_commands: vec![setup],
setup_command_timeout_ms: 1_000,
},
..test_init_options(
run_store.into(),
emitter,
std::env::current_dir().unwrap(),
test_settings(&run_dir),
)
})
.await;
let events = seen.lock().unwrap().clone();
(result, events)
}
#[tokio::test]
async fn configure_sandbox_git_identity_uses_run_author() {
let sandbox = fabro_sandbox::test_support::MockSandbox::linux();
let author = GitAuthor::from_options(
Some("Fabro Bot".to_string()),
Some("fabro-bot@example.com".to_string()),
);
configure_sandbox_git_identity(&sandbox, &author)
.await
.expect("git identity should configure");
let commands = sandbox
.captured_commands
.lock()
.expect("captured_commands lock poisoned")
.clone();
assert_eq!(commands, vec![
"git config --local user.name 'Fabro Bot' && git config --local user.email \
fabro-bot@example.com"
]);
}
#[tokio::test]
async fn initialize_prepares_sandbox_and_uses_persisted_run_dir() {
let temp = tempfile::tempdir().unwrap();
let run_dir = temp.path().join("run");
std::fs::create_dir_all(&run_dir).unwrap();
let (graph, source) = simple_graph();
let persisted = test_persisted(graph, source.clone(), &run_dir);
let emitter = Arc::new(crate::event::Emitter::new(test_run_id()));
let run_store = memory_store().create_run(&test_run_id()).await.unwrap();
let initialized = initialize(persisted, InitOptions {
sandbox_env: SandboxEnvSpec {
toml_env: HashMap::from([("TEST_KEY".to_string(), "value".to_string())]),
github_integration: None,
origin_url: None,
},
..test_init_options(
run_store.into(),
emitter,
std::env::current_dir().unwrap(),
test_settings(&run_dir),
)
})
.await
.unwrap();
assert_eq!(initialized.run_options.run_dir, run_dir);
assert_eq!(initialized.source, source);
assert!(initialized.engine.run.hook_runner.is_none());
assert_eq!(
initialized.engine.run.locations.host_source_dir.as_deref(),
Some(std::env::current_dir().unwrap().as_path())
);
assert_eq!(
initialized.engine.run.locations.sandbox_work_dir.as_deref(),
Some(std::env::current_dir().unwrap().as_path())
);
assert_eq!(
initialized.engine.run.locations.run_scratch_dir.as_path(),
run_dir.as_path()
);
assert_eq!(
initialized
.engine
.base_env
.get("TEST_KEY")
.map(String::as_str),
Some("value")
);
assert!(initialized.engine.dry_run);
assert_eq!(initialized.model, "test-model");
assert_eq!(
initialized.engine.run.provider_id,
fabro_model::ProviderId::anthropic()
);
assert!(
initialized
.engine
.run
.llm_source
.resolve(&initialized.engine.run.catalog)
.await
.unwrap()
.credentials
.is_empty()
);
}
async fn initialize_resume_with_planned_sandbox(
temp: &tempfile::TempDir,
fork_source_ref: Option<ForkSourceRef>,
) -> Result<Initialized, Error> {
let run_dir = temp.path().join("run");
let workspace = temp.path().join("workspace");
std::fs::create_dir_all(&run_dir).unwrap();
std::fs::create_dir_all(&workspace).unwrap();
let (graph, source) = simple_graph();
let mut settings = WorkflowSettings::default();
settings.run.run_branch.enabled = false;
let persisted = test_persisted_run(
graph.clone(),
source,
&run_dir,
settings.clone(),
fork_source_ref.clone(),
);
let emitter = Arc::new(crate::event::Emitter::new(test_run_id()));
let store = memory_store();
let run_store = store.create_run(&test_run_id()).await.unwrap();
let mut checkpoint = Checkpoint::from_context(
&Context::new(),
"start",
vec!["start".to_string()],
HashMap::new(),
HashMap::new(),
Some("exit".to_string()),
HashMap::new(),
HashMap::new(),
HashMap::new(),
);
checkpoint.git_commit_sha = Some(CHECKPOINT_SHA.to_string());
let mut run_options = test_settings(&run_dir);
run_options.settings = settings;
run_options.fork_source_ref = fork_source_ref;
seed_run_created(
&run_store,
serde_json::to_value(&run_options.settings).unwrap(),
serde_json::to_value(&graph).unwrap(),
Some(workspace.display().to_string()),
run_options.fork_source_ref.clone(),
)
.await;
initialize(persisted, InitOptions {
resume: Some(ResumeState::for_test(
checkpoint,
StageExecutionSeed::default(),
)),
..test_init_options(run_store.into(), emitter, workspace, run_options)
})
.await
}
#[tokio::test]
async fn forked_run_resume_materializes_fresh_sandbox() {
let temp = tempfile::tempdir().unwrap();
let workspace = temp.path().join("workspace");
let fork_source_ref = ForkSourceRef {
source_run_id: fixtures::RUN_64,
checkpoint_sha: CHECKPOINT_SHA.to_string(),
};
let initialized = initialize_resume_with_planned_sandbox(&temp, Some(fork_source_ref))
.await
.expect("a forked run should materialize a fresh sandbox before resuming");
assert_eq!(
initialized.engine.run.sandbox.working_directory(),
workspace.to_string_lossy().as_ref()
);
}
#[tokio::test]
async fn same_run_resume_does_not_recreate_uninitialized_sandbox() {
let temp = tempfile::tempdir().unwrap();
match initialize_resume_with_planned_sandbox(&temp, None).await {
Err(Error::Precondition(message)) => {
assert!(
message.contains("was not initialized"),
"unexpected precondition message: {message}"
);
}
Err(error) => panic!("expected sandbox precondition error, got {error}"),
Ok(_) => panic!("same-run resume should not recreate an uninitialized sandbox"),
}
}
#[tokio::test]
async fn build_registry_accepts_vault_only_llm_provider() {
let dir = tempfile::tempdir().unwrap();
let mut vault = Vault::load(dir.path().join("secrets.json")).unwrap();
vault
.set(
"ANTHROPIC_API_KEY",
"anthropic-key",
SecretType::Token,
None,
)
.unwrap();
let (graph, _) = llm_graph();
let vault = Arc::new(AsyncRwLock::new(vault));
let test_emitter = Arc::new(crate::event::Emitter::new(test_run_id()));
let tool_env_provider = Arc::new(WorkflowToolEnvProvider {
base_env: HashMap::new(),
github_token: None,
});
let (_registry, effective_dry_run) = build_registry(
&LlmSpec {
model: "claude-opus-4-6".to_string(),
provider_id: fabro_model::ProviderId::anthropic(),
fallbacks: ModelFallbackPolicy::default(),
mcp_servers: Vec::new(),
model_controls: RunModelControls::default(),
dry_run: false,
},
Arc::new(AutoApproveInterviewer::engine()),
Arc::new(crate::steering_hub::SteeringHub::new(test_emitter)),
tool_env_provider,
false,
&graph,
Arc::new(VaultCredentialSource::new(Arc::clone(&vault))),
test_catalog(),
ToolSecrets::default(),
None,
)
.await
.unwrap();
assert!(!effective_dry_run);
}
#[tokio::test]
async fn build_llm_source_appends_run_session_trace_header() {
let mut vault = Vault::from_entries(HashMap::new());
fabro_auth::vault_set_token(&mut vault, EnvVars::ANTHROPIC_API_KEY, "anthropic-key")
.unwrap();
let vault = Arc::new(AsyncRwLock::new(vault));
let run_id = test_run_id();
let expected_session_id = run_id.to_string();
let source = build_llm_source(vault, run_id);
let resolved = source.resolve(test_catalog().as_ref()).await.unwrap();
assert!(!resolved.credentials.is_empty());
for credential in &resolved.credentials {
assert_eq!(
credential
.extra_headers
.get(SESSION_ID_HEADER)
.map(String::as_str),
Some(expected_session_id.as_str())
);
}
}
#[tokio::test]
async fn initialize_executes_acp_backend_node_from_registry() {
let temp = tempfile::tempdir().unwrap();
let run_dir = temp.path().join("run");
create_dir_all(&run_dir).await.unwrap();
let script_path = temp.path().join("fake_acp_agent.py");
write(&script_path, fake_acp_agent_script()).await.unwrap();
let source = format!(
r#"digraph test {{
start [shape=Mdiamond];
writer [type="agent", backend="acp", prompt="write hello", acp.command="python3 {}"];
exit [shape=Msquare];
start -> writer;
writer -> exit;
}}"#,
script_path.display()
);
let mut graph = Graph::new("test");
let mut start = Node::new("start");
start.attrs.insert(
"shape".to_string(),
AttrValue::String("Mdiamond".to_string()),
);
let mut writer = Node::new("writer");
writer
.attrs
.insert("type".to_string(), AttrValue::String("agent".to_string()));
writer
.attrs
.insert("backend".to_string(), AttrValue::String("acp".to_string()));
writer.attrs.insert(
"prompt".to_string(),
AttrValue::String("write hello".to_string()),
);
writer.attrs.insert(
"acp.command".to_string(),
AttrValue::String(format!(
"python3 {}",
fabro_sandbox::shell_quote(&script_path.to_string_lossy())
)),
);
let mut exit = Node::new("exit");
exit.attrs.insert(
"shape".to_string(),
AttrValue::String("Msquare".to_string()),
);
graph.nodes.insert("start".to_string(), start);
graph.nodes.insert("writer".to_string(), writer);
graph.nodes.insert("exit".to_string(), exit);
graph.edges.push(Edge::new("start", "writer"));
graph.edges.push(Edge::new("writer", "exit"));
let mut vault = Vault::load(temp.path().join("secrets.json")).unwrap();
vault
.set("OPENAI_API_KEY", "openai-key", SecretType::Token, None)
.unwrap();
let vault = Arc::new(AsyncRwLock::new(vault));
let emitter = Arc::new(crate::event::Emitter::new(test_run_id()));
let seen = Arc::new(std::sync::Mutex::new(Vec::new()));
emitter.on_event({
let seen = Arc::clone(&seen);
move |event| seen.lock().unwrap().push(event.event_name().to_string())
});
let store = memory_store();
let run_store = store.create_run(&test_run_id()).await.unwrap();
let initialized = initialize(test_persisted(graph, source, &run_dir), InitOptions {
run_store: run_store.into(),
dry_run: false,
emitter: emitter.clone(),
sandbox: SandboxSpec::Local {
working_directory: temp.path().to_path_buf(),
},
llm: LlmSpec {
model: "fake-acp".to_string(),
provider_id: fabro_model::ProviderId::openai(),
fallbacks: ModelFallbackPolicy::default(),
mcp_servers: Vec::new(),
model_controls: RunModelControls::default(),
dry_run: false,
},
interviewer: Arc::new(AutoApproveInterviewer::engine()),
steering_hub: Arc::new(crate::steering_hub::SteeringHub::new(emitter)),
catalog: test_catalog(),
lifecycle: crate::run_options::LifecycleOptions {
setup_commands: Vec::new(),
setup_command_timeout_ms: 1_000,
},
run_options: test_settings(&run_dir),
workflow_path: None,
workflow_bundle: None,
hooks: fabro_hooks::HookSettings { hooks: vec![] },
sandbox_env: SandboxEnvSpec {
toml_env: HashMap::new(),
github_integration: None,
origin_url: None,
},
vault,
git: None,
run_control: None,
registry_override: None,
artifact_sink: None,
resume: None,
seed_context: None,
fabro_run_tools: None,
})
.await
.unwrap();
let node = initialized.graph.nodes.get("writer").unwrap().clone();
let handler = initialized.engine.registry.resolve(&node);
let context = Context::new();
context.set(
keys::INTERNAL_RUN_ID,
serde_json::json!(test_run_id().to_string()),
);
let outcome = handler
.execute(
&node,
&context,
&initialized.graph,
&initialized.run_options.run_dir,
&initialized.engine,
)
.await
.unwrap();
assert_eq!(
outcome.context_updates.get(&keys::response_key("writer")),
Some(&serde_json::json!("hello from acp"))
);
assert!(
seen.lock()
.unwrap()
.contains(&"agent.acp.started".to_string())
);
assert!(
seen.lock()
.unwrap()
.contains(&"agent.acp.completed".to_string())
);
}
#[tokio::test]
async fn initialize_runs_setup_commands() {
let temp = tempfile::tempdir().unwrap();
let run_dir = temp.path().join("run");
std::fs::create_dir_all(&run_dir).unwrap();
let (graph, source) = simple_graph();
let persisted = test_persisted(graph.clone(), source, &run_dir);
let emitter = Arc::new(crate::event::Emitter::new(test_run_id()));
let store = memory_store();
let run_store = store.create_run(&test_run_id()).await.unwrap();
seed_run_created(
&run_store,
serde_json::to_value(WorkflowSettings::default()).unwrap(),
serde_json::to_value(graph).unwrap(),
None,
None,
)
.await;
let store_logger = StoreProgressLogger::new(run_store.clone());
let seen = Arc::new(std::sync::Mutex::new(Vec::new()));
emitter.on_event({
let seen = Arc::clone(&seen);
move |event| seen.lock().unwrap().push(event.event_name().to_string())
});
store_logger.register(&emitter);
let initialized = initialize(persisted, InitOptions {
run_store: run_store.into(),
dry_run: false,
emitter: emitter.clone(),
sandbox: SandboxSpec::Local {
working_directory: std::env::current_dir().unwrap(),
},
llm: LlmSpec {
model: "test-model".to_string(),
provider_id: fabro_model::ProviderId::anthropic(),
fallbacks: ModelFallbackPolicy::default(),
mcp_servers: Vec::new(),
model_controls: RunModelControls::default(),
dry_run: true,
},
interviewer: Arc::new(AutoApproveInterviewer::engine()),
steering_hub: Arc::new(crate::steering_hub::SteeringHub::new(emitter.clone())),
catalog: test_catalog(),
lifecycle: crate::run_options::LifecycleOptions {
setup_commands: vec![setup_cmd("true")],
setup_command_timeout_ms: 1_000,
},
run_options: test_settings(&run_dir),
workflow_path: None,
workflow_bundle: None,
hooks: fabro_hooks::HookSettings { hooks: vec![] },
sandbox_env: SandboxEnvSpec {
toml_env: HashMap::new(),
github_integration: None,
origin_url: None,
},
vault: auth_test_support::empty_vault(),
git: None,
run_control: None,
registry_override: None,
artifact_sink: None,
resume: None,
seed_context: None,
fabro_run_tools: None,
})
.await
.unwrap();
store_logger.flush().await.unwrap();
assert_eq!(initialized.run_options.run_dir, run_dir);
assert!(
seen.lock()
.unwrap()
.iter()
.any(|event| event == "sandbox.initialized")
);
}
#[tokio::test]
async fn initialize_passes_per_step_env_to_setup_command() {
// The command only succeeds when the per-step env var is visible to the
// shell, so a green run proves the env reached `exec_command`.
let setup = crate::run_options::SetupCommand {
command: "test \"$PREPARE_STAGE\" = build".to_string(),
env: HashMap::from([("PREPARE_STAGE".to_string(), "build".to_string())]),
};
let (result, events) = initialize_with_setup_step(setup).await;
assert!(result.is_ok(), "setup with per-step env should succeed");
assert!(
events
.iter()
.any(|event| event.event_name() == "setup.completed")
);
}
#[tokio::test]
async fn initialize_setup_command_without_step_env_does_not_see_it() {
// Negative control: the same command without the per-step env fails,
// confirming the success above is attributable to the per-step env.
let (result, _events) =
initialize_with_setup_command("test \"$PREPARE_STAGE\" = build").await;
assert!(result.is_err(), "setup should fail without per-step env");
}
#[tokio::test]
async fn initialize_setup_failure_preserves_stderr_and_adds_exec_tail() {
let (result, events) =
initialize_with_setup_command("printf setup-out; printf setup-err >&2; exit 7").await;
assert!(result.is_err());
let failed = events
.iter()
.find(|event| event.event_name() == "setup.failed")
.expect("setup failed event");
match &failed.body {
EventBody::SetupFailed(props) => {
assert_eq!(props.exit_code, 7);
assert_eq!(props.stderr, "setup-err");
let tail = props.exec_output_tail.as_ref().expect("exec output tail");
assert_eq!(tail.stdout.as_deref(), Some("setup-out"));
assert_eq!(tail.stderr.as_deref(), Some("setup-err"));
}
other => panic!("expected setup failed body, got {other:?}"),
}
}
#[tokio::test]
async fn initialize_setup_failure_with_stdout_only_adds_stdout_tail() {
let (result, events) = initialize_with_setup_command("printf setup-out; exit 5").await;
assert!(result.is_err());
let failed = events
.iter()
.find(|event| event.event_name() == "setup.failed")
.expect("setup failed event");
match &failed.body {
EventBody::SetupFailed(props) => {
assert_eq!(props.exit_code, 5);
assert!(props.stderr.is_empty());
let tail = props.exec_output_tail.as_ref().expect("exec output tail");
assert_eq!(tail.stdout.as_deref(), Some("setup-out"));
assert!(tail.stderr.is_none());
}
other => panic!("expected setup failed body, got {other:?}"),
}
}
#[tokio::test]
async fn initialize_cancelled_setup_command_returns_cancelled() {
let temp = tempfile::tempdir().unwrap();
let run_dir = temp.path().join("run");
std::fs::create_dir_all(&run_dir).unwrap();
let (graph, source) = simple_graph();
let persisted = test_persisted(graph, source, &run_dir);
let cancel_token = tokio_util::sync::CancellationToken::new();
cancel_token.cancel();
let mut run_options = test_settings(&run_dir);
run_options.cancel_token = cancel_token;
let emitter = Arc::new(crate::event::Emitter::new(test_run_id()));
let result = initialize(persisted, InitOptions {
run_store: {
let store = memory_store();
let inner = store.create_run(&test_run_id()).await.unwrap();
inner.into()
},
dry_run: false,
emitter: emitter.clone(),
sandbox: SandboxSpec::Local {
working_directory: std::env::current_dir().unwrap(),
},
llm: LlmSpec {
model: "test-model".to_string(),
provider_id: fabro_model::ProviderId::anthropic(),
fallbacks: ModelFallbackPolicy::default(),
mcp_servers: Vec::new(),
model_controls: RunModelControls::default(),
dry_run: true,
},
interviewer: Arc::new(AutoApproveInterviewer::engine()),
steering_hub: Arc::new(crate::steering_hub::SteeringHub::new(emitter.clone())),
catalog: test_catalog(),
lifecycle: crate::run_options::LifecycleOptions {
setup_commands: vec![setup_cmd("sleep 5")],
setup_command_timeout_ms: 5_000,
},
run_options,
workflow_path: None,
workflow_bundle: None,
hooks: fabro_hooks::HookSettings { hooks: vec![] },
sandbox_env: SandboxEnvSpec {
toml_env: HashMap::new(),
github_integration: None,
origin_url: None,
},
vault: auth_test_support::empty_vault(),
git: None,
run_control: None,
registry_override: None,
artifact_sink: None,
resume: None,
seed_context: None,
fabro_run_tools: None,
})
.await;
assert!(matches!(result, Err(Error::Cancelled)));
}
mod github_integration_env {
//! Focused tests for `build_sandbox_env` /
//! `resolve_declared_repository_token` around declared additional
//! repositories. Installation-resolution failure naming is covered
//! by `fabro_github::access` tests; these prove the initialization
//! wiring: hard errors for declared sets, best-effort behavior for
//! legacy permissions-only configuration.
use fabro_github::test_support::{InstallationTokenMinter, installation_token_source};
use fabro_github::{GitHubAppCredentials, GitHubCredentials, InstallationToken};
use fabro_types::settings::run::ResolvedGithubIntegration;
use super::*;
fn integration(additional: &[&str]) -> ResolvedGithubIntegration {
ResolvedGithubIntegration {
permissions: HashMap::from([(
"contents".to_string(),
"read".to_string(),
)]),
additional_repositories: additional
.iter()
.map(|value| value.parse().expect("test slug should parse"))
.collect(),
}
}
fn spec(
origin: Option<&str>,
github_integration: Option<ResolvedGithubIntegration>,
) -> SandboxEnvSpec {
SandboxEnvSpec {
toml_env: HashMap::new(),
github_integration,
origin_url: origin.map(str::to_string),
}
}
#[test]
fn declared_additional_repositories_require_credentials() {
let spec = spec(
Some("https://github.com/fabro-sh/fabro"),
Some(integration(&["fabro-sh/keystone"])),
);
let Err(err) = build_sandbox_env(&spec, None) else {
panic!("declared additional repositories without credentials must fail");
};
assert!(
err.to_string().contains("requires GitHub credentials"),
"{err}"
);
}
#[test]
fn declared_additional_repositories_require_an_origin() {
let spec = spec(None, Some(integration(&["fabro-sh/keystone"])));
let creds = GitHubCredentials::Pat("ghp_x".to_string());
let Err(err) = build_sandbox_env(&spec, Some(&creds)) else {
panic!("declared additional repositories without an origin must fail");
};
assert!(
err.to_string().contains("GitHub repository access"),
"{err}"
);
}
#[test]
fn declared_repositories_inject_bridge_entries_and_keep_the_pat_source() {
let spec = spec(
Some("https://github.com/fabro-sh/fabro"),
Some(integration(&["fabro-sh/keystone"])),
);
let creds = GitHubCredentials::Pat("ghp_x".to_string());
let built = build_sandbox_env(&spec, Some(&creds)).unwrap();
assert!(built.github_token.is_some());
let access = built.github_access.expect("access should be constructed");
assert!(access.has_additional_repositories());
// Helper entry plus two SSH rewrites for each of the two
// effective repositories (origin + declared additional).
assert_eq!(
built.env.get("GIT_CONFIG_COUNT").map(String::as_str),
Some("5")
);
assert_eq!(
built.env.get("GIT_CONFIG_KEY_0").map(String::as_str),
Some("credential.https://github.com.helper")
);
assert_eq!(
built.env.get("GIT_TERMINAL_PROMPT").map(String::as_str),
Some("0")
);
}
#[test]
fn legacy_permissions_only_configuration_stays_best_effort() {
// No credentials: no error, no token source, no bridge entries.
let no_creds = spec(
Some("https://github.com/fabro-sh/fabro"),
Some(integration(&[])),
);
let built = build_sandbox_env(&no_creds, None).unwrap();
assert!(built.github_token.is_none());
assert!(!built.env.contains_key("GIT_CONFIG_COUNT"));
// App credentials without an origin: legacy best-effort skip.
let creds = GitHubCredentials::App(GitHubAppCredentials {
app_id: "1".to_string(),
private_key_pem: "unused".to_string(),
slug: None,
});
let no_origin = spec(None, Some(integration(&[])));
let built = build_sandbox_env(&no_origin, Some(&creds)).unwrap();
assert!(built.github_token.is_none());
assert!(built.github_access.is_none());
}
struct FailingMinter;
#[async_trait::async_trait]
impl InstallationTokenMinter for FailingMinter {
async fn mint(&self) -> anyhow::Result<InstallationToken> {
Err(anyhow::anyhow!("scripted mint failure"))
}
}
#[tokio::test]
async fn eager_validation_fails_when_the_declared_token_cannot_resolve() {
let access = fabro_github::GitHubRepositoryAccess::new(
Some("https://github.com/fabro-sh/fabro"),
&["fabro-sh/keystone".parse().unwrap()].into_iter().collect(),
HashMap::from([("contents".to_string(), "read".to_string())]),
)
.unwrap();
let built = BuiltSandboxEnv {
env: HashMap::new(),
github_token: Some(installation_token_source(
"fabro-sh/fabro (+1 additional)",
Arc::new(FailingMinter),
)),
github_access: access,
};
let err = resolve_declared_repository_token(&built).await.unwrap_err();
let message = err.to_string();
assert!(message.contains("declared repository set"), "{message}");
}
#[tokio::test]
async fn eager_validation_skips_legacy_permissions_only_runs() {
let built = BuiltSandboxEnv {
env: HashMap::new(),
github_token: Some(installation_token_source(
"fabro-sh/fabro",
Arc::new(FailingMinter),
)),
github_access: None,
};
resolve_declared_repository_token(&built)
.await
.expect("legacy permissions-only runs must not resolve eagerly");
}
}
}