fabro(01KQRFAT7326ADECC8AAWBXK6M): implement (succeeded)

Fabro-Run: 01KQRFAT7326ADECC8AAWBXK6M
Fabro-Completed: 5

⚒️ Generated with [Fabro](https://fabro.sh)
This commit is contained in:
Fabro 2026-05-04 03:32:03 +00:00
parent ee24291f54
commit 62f6032e12
9 changed files with 334 additions and 19 deletions

View file

@ -282,14 +282,18 @@ impl Sandbox for LocalSandbox {
let stdout_task = tokio::spawn(async move {
let mut buf = String::new();
if let Some(ref mut r) = stdout_pipe {
let _ = r.read_to_string(&mut buf).await;
if let Err(err) = r.read_to_string(&mut buf).await {
tracing::warn!(error = %err, stream = "stdout", "Failed to drain child stdout");
}
}
buf
});
let stderr_task = tokio::spawn(async move {
let mut buf = String::new();
if let Some(ref mut r) = stderr_pipe {
let _ = r.read_to_string(&mut buf).await;
if let Err(err) = r.read_to_string(&mut buf).await {
tracing::warn!(error = %err, stream = "stderr", "Failed to drain child stderr");
}
}
buf
});

View file

@ -52,6 +52,8 @@ pub trait CodergenBackend: Send + Sync {
_node: &Node,
_prompt: &str,
_system_prompt: Option<&str>,
_emitter: &Arc<Emitter>,
_stage_scope: &StageScope,
) -> Result<CodergenResult, Error> {
Err(Error::Validation(
"one_shot mode not supported by this backend".into(),

View file

@ -285,6 +285,8 @@ impl CodergenBackend for AgentApiBackend {
node: &Node,
prompt: &str,
system_prompt: Option<&str>,
emitter: &Arc<Emitter>,
stage_scope: &StageScope,
) -> Result<CodergenResult, Error> {
let client = Client::from_source(self.source.as_ref())
.await
@ -358,14 +360,16 @@ impl CodergenBackend for AgentApiBackend {
let mut found = None;
for target in fallback_chain {
tracing::warn!(
stage = node.id.as_str(),
from_provider = from_provider.as_str(),
from_model = from_model.as_str(),
to_provider = target.provider.as_str(),
to_model = target.model.as_str(),
error = error_msg.as_str(),
"LLM provider failover (prompt)"
emitter.emit_scoped(
&Event::Failover {
stage: node.id.clone(),
from_provider: from_provider.clone(),
from_model: from_model.clone(),
to_provider: target.provider.clone(),
to_model: target.model.clone(),
error: error_msg.clone(),
},
stage_scope,
);
let max_tokens = node.max_tokens().or_else(|| {

View file

@ -810,9 +810,13 @@ impl CodergenBackend for BackendRouter {
node: &Node,
prompt: &str,
system_prompt: Option<&str>,
emitter: &Arc<Emitter>,
stage_scope: &StageScope,
) -> Result<CodergenResult, Error> {
// CLI backend doesn't support one_shot, always route to API
self.api_backend.one_shot(node, prompt, system_prompt).await
self.api_backend
.one_shot(node, prompt, system_prompt, emitter, stage_scope)
.await
}
}

View file

@ -12,7 +12,7 @@ use tokio::sync::Semaphore;
use super::{EngineServices, Handler};
use crate::context::{Context, WorkflowContext, keys};
use crate::error::Error;
use crate::event::{Event, StageScope};
use crate::event::{Event, RunNoticeLevel, StageScope};
use crate::git::sanitize_ref_component;
use crate::hook_context::set_hook_node;
use crate::millis_u64;
@ -207,6 +207,11 @@ impl Handler for ParallelHandler {
error = %fabro_sandbox::display_for_log(&e),
"parallel base checkpoint failed"
);
services.run.emitter.notice(
RunNoticeLevel::Warn,
"parallel_base_checkpoint_failed",
format!("Could not checkpoint base state before parallel branches: {e}"),
);
None
}
}
@ -875,4 +880,107 @@ mod tests {
let branch_count = context.get(keys::PARALLEL_BRANCH_COUNT);
assert_eq!(branch_count, Some(serde_json::json!(2)));
}
#[expect(
clippy::disallowed_methods,
reason = "Test sets up a temporary git repo to exercise the parallel checkpoint failure path."
)]
#[tokio::test]
async fn parallel_base_checkpoint_failure_emits_notice() {
// Set up a git repo then lock the index to make `git add` fail.
let repo_dir = tempfile::tempdir().unwrap();
let init = std::process::Command::new("git")
.args(["init", "-b", "main"])
.current_dir(repo_dir.path())
.output()
.unwrap();
assert!(init.status.success());
for (key, value) in [("user.name", "Test"), ("user.email", "test@test.com")] {
let config = std::process::Command::new("git")
.args(["config", key, value])
.current_dir(repo_dir.path())
.output()
.unwrap();
assert!(config.status.success());
}
let commit = std::process::Command::new("git")
.args(["commit", "--allow-empty", "-m", "initial"])
.current_dir(repo_dir.path())
.output()
.unwrap();
assert!(commit.status.success());
// Lock the index to make `git add` fail
let lock_path = repo_dir.path().join(".git/index.lock");
std::fs::write(&lock_path, b"locked").unwrap();
let store = test_store();
let run_store = store.create_run(&fixtures::RUN_1).await.unwrap();
let mut services = EngineServices::test_default();
let emitter = Arc::new(crate::event::Emitter::new(fixtures::RUN_1));
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())
});
services.run = services
.run
.with_emitter(emitter)
.with_run_store(run_store.clone().into())
.with_sandbox(Arc::new(fabro_agent::LocalSandbox::new(
repo_dir.path().to_path_buf(),
)));
// Supply git_state so the parallel handler enters the git path
services.set_git_state(Some(Arc::new(crate::sandbox_git::GitState {
run_id: fixtures::RUN_1,
base_sha: "abc123".to_string(),
run_branch: Some("fabro/run/test".to_string()),
meta_branch: None,
checkpoint_exclude_globs: Vec::new(),
git_author: crate::git::GitAuthor::default(),
})));
let logger = crate::event::StoreProgressLogger::new(run_store.clone());
logger.register(services.run.emitter.as_ref());
let mut node = Node::new("par");
node.attrs.insert(
"shape".to_string(),
AttrValue::String("component".to_string()),
);
let context = test_context();
let mut graph = Graph::new("test");
graph.nodes.insert("par".to_string(), node.clone());
graph
.nodes
.insert("branch_a".to_string(), Node::new("branch_a"));
graph.edges.push(Edge::new("par", "branch_a"));
let outcome = ParallelHandler
.execute(&node, &context, &graph, repo_dir.path(), &services)
.await
.unwrap();
logger.flush().await;
// Clean up the lock so tempdir cleanup doesn't leave artifacts
let _ = std::fs::remove_file(&lock_path);
// Should still succeed (the checkpoint failure is non-fatal)
assert_eq!(outcome.status, StageOutcome::Succeeded);
let events = seen.lock().unwrap();
let notice = events.iter().find(|e| {
matches!(
&e.body,
fabro_types::EventBody::RunNotice(props)
if props.code == "parallel_base_checkpoint_failed"
)
});
assert!(
notice.is_some(),
"expected parallel_base_checkpoint_failed notice, got events: {:?}",
events
.iter()
.map(fabro_types::RunEvent::event_name)
.collect::<Vec<_>>()
);
}
}

View file

@ -105,7 +105,13 @@ impl Handler for PromptHandler {
let (response_text, stage_usage, backend_files_touched) =
if let Some(backend) = &self.backend {
let result = backend
.one_shot(node, &prompt, system_prompt.as_deref())
.one_shot(
node,
&prompt,
system_prompt.as_deref(),
&services.run.emitter,
&stage_scope,
)
.await;
match result {
Ok(CodergenResult::Full(outcome)) => return Ok(outcome),
@ -279,6 +285,8 @@ mod tests {
_node: &Node,
_prompt: &str,
_system_prompt: Option<&str>,
_emitter: &Arc<crate::event::Emitter>,
_stage_scope: &StageScope,
) -> Result<CodergenResult, Error> {
Ok(CodergenResult::Text {
text: "one-shot response".to_string(),
@ -339,6 +347,8 @@ mod tests {
_node: &Node,
_prompt: &str,
_system_prompt: Option<&str>,
_emitter: &Arc<crate::event::Emitter>,
_stage_scope: &StageScope,
) -> Result<CodergenResult, Error> {
Ok(CodergenResult::Text {
text: "one-shot response".to_string(),
@ -396,6 +406,8 @@ mod tests {
_node: &Node,
prompt: &str,
system_prompt: Option<&str>,
_emitter: &Arc<crate::event::Emitter>,
_stage_scope: &StageScope,
) -> Result<CodergenResult, Error> {
*self.captured_prompt.lock().unwrap() = Some(prompt.to_string());
*self.captured_system_prompt.lock().unwrap() = Some(system_prompt.map(String::from));

View file

@ -292,6 +292,14 @@ impl RunLifecycle<WorkflowGraph> for GitLifecycle {
error = %fabro_sandbox::display_for_log(&err),
"git push from run lifecycle failed"
);
self.emitter.emit(&Event::RunNotice {
level: RunNoticeLevel::Warn,
code: "git_push_failed".to_string(),
message: format!(
"Failed to push run branch {branch}: {err}"
),
exec_output_tail: exec_output_tail.clone(),
});
(false, exec_output_tail)
}
};
@ -1062,4 +1070,76 @@ mod tests {
Ok(None)
}
}
#[expect(
clippy::disallowed_methods,
reason = "Test sets up a temporary git repo with a bogus remote to verify push-failure notice."
)]
#[tokio::test]
async fn checkpoint_push_failure_emits_git_push_failed_notice() {
let repo_dir = tempfile::tempdir().unwrap();
init_git_repo(repo_dir.path());
// Add a bogus remote so git_push_ref actually attempts the push.
let add_remote = std::process::Command::new("git")
.args(["remote", "add", "origin", "https://127.0.0.1:1/nope.git"])
.current_dir(repo_dir.path())
.output()
.unwrap();
assert!(add_remote.status.success());
let branch = "fabro/metadata/run";
let run_branch = "fabro/run/test";
let emitter = Arc::new(Emitter::new(fixtures::RUN_1));
let events = record_events(&emitter);
let run_store = run_store(fixtures::RUN_1).await;
let handle = RunStoreHandle::local(run_store.clone());
// run_options with a run_branch so the push path is entered
let options = Arc::new(RunOptions {
settings: WorkflowSettings::default(),
run_dir: repo_dir.path().to_path_buf(),
cancel_token: None,
run_id: fixtures::RUN_1,
labels: HashMap::new(),
workflow_slug: Some("metadata".to_string()),
github_app: None,
pre_run_git: None,
fork_source_ref: None,
base_branch: None,
display_base_sha: None,
git: Some(GitCheckpointOptions {
base_sha: None,
run_branch: Some(run_branch.to_string()),
meta_branch: Some(branch.to_string()),
}),
});
let lifecycle = git_lifecycle(
repo_dir.path(),
emitter,
handle,
options,
Arc::new(RunMetadataRuntime::new()),
);
let graph = workflow_graph();
let node = graph.get_node("build").unwrap();
let mut state = ExecutionState::new(&graph).unwrap();
state.increment_visits("build");
let result = WfNodeResult::new(Outcome::success(), Duration::from_millis(10), 1, 1);
lifecycle
.on_checkpoint(&node, &result, Some("exit"), &state)
.await
.unwrap();
let events = events.lock().unwrap();
let notice = events.iter().find(|e| {
matches!(
&e.body,
EventBody::RunNotice(props) if props.code == "git_push_failed"
)
});
assert!(
notice.is_some(),
"expected git_push_failed notice, got events: {:?}",
events.iter().map(RunEvent::event_name).collect::<Vec<_>>()
);
}
}

View file

@ -239,11 +239,14 @@ async fn build_sandbox_env(
Ok(token) => {
env.insert("GITHUB_TOKEN".to_string(), token);
}
Err(e) => emitter.notice(
RunNoticeLevel::Warn,
"github_token_failed",
format!("Failed to mint GitHub token: {e}"),
),
Err(e) => {
tracing::warn!(error = %e, "Failed to mint GitHub token");
emitter.notice(
RunNoticeLevel::Warn,
"github_token_failed",
format!("Failed to mint GitHub token: {e}"),
);
}
}
}
}
@ -504,6 +507,15 @@ pub async fn initialize(
}
}
} else {
tracing::warn!(
worktree_mode = ?options.worktree_mode,
"worktree requested but cwd is not a git repository; running without a worktree"
);
options.emitter.notice(
RunNoticeLevel::Warn,
"worktree_skipped_no_git",
"Worktree mode requested but no Git repository was found; running without a worktree.",
);
Arc::new(ReadBeforeWriteSandbox::new(inner))
}
} else {
@ -619,7 +631,15 @@ pub async fn initialize(
options.run_options.base_branch = info.base_branch;
}
}
Ok(None) => {}
Ok(None) => {
if sandbox.origin_url().is_some() {
options.emitter.notice(
RunNoticeLevel::Warn,
"sandbox_git_unavailable",
"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));
}
@ -1373,4 +1393,82 @@ mod tests {
assert!(matches!(result, Err(Error::Cancelled)));
}
#[tokio::test]
async fn initialize_worktree_skipped_no_git_emits_notice() {
// Use a non-git temp dir so resolve_worktree_base_sha returns Ok(None)
let temp = tempfile::tempdir().unwrap();
let non_git_dir = temp.path().to_path_buf();
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 _initialized = initialize(persisted, InitOptions {
run_id: test_run_id(),
run_store: {
let store = memory_store();
let inner = store.create_run(&test_run_id()).await.unwrap();
inner.into()
},
dry_run: false,
emitter,
sandbox: SandboxSpec::Local {
working_directory: non_git_dir,
},
llm: LlmSpec {
model: "test-model".to_string(),
provider: fabro_llm::Provider::Anthropic,
fallback_chain: Vec::new(),
mcp_servers: Vec::new(),
dry_run: true,
},
interviewer: Arc::new(AutoApproveInterviewer::engine()),
lifecycle: crate::run_options::LifecycleOptions {
setup_commands: vec![],
setup_command_timeout_ms: 1_000,
devcontainer_phases: vec![],
},
run_options: test_settings(&run_dir),
workflow_path: None,
workflow_bundle: None,
hooks: fabro_hooks::HookSettings { hooks: vec![] },
sandbox_env: SandboxEnvSpec {
devcontainer_env: HashMap::new(),
toml_env: HashMap::new(),
github_permissions: None,
origin_url: None,
},
vault: None,
devcontainer: None,
git: None,
worktree_mode: Some(WorktreeMode::Always),
run_control: None,
registry_override: None,
artifact_sink: None,
checkpoint: None,
seed_context: None,
})
.await
.unwrap();
let events = seen.lock().unwrap();
let notice = events.iter().find(|e| {
matches!(
&e.body,
EventBody::RunNotice(props) if props.code == "worktree_skipped_no_git"
)
});
assert!(
notice.is_some(),
"expected worktree_skipped_no_git notice, got events: {:?}",
events.iter().map(RunEvent::event_name).collect::<Vec<_>>()
);
}
}

View file

@ -6193,6 +6193,7 @@ mod real_llm {
use fabro_types::WorkflowSettings;
use fabro_workflow::context::Context;
use fabro_workflow::error::Error;
use fabro_workflow::event::StageScope;
use fabro_workflow::handler::agent::{AgentHandler, CodergenBackend, CodergenResult};
struct LlmCodergenBackend {
@ -6221,6 +6222,8 @@ mod real_llm {
_node: &Node,
prompt: &str,
_system_prompt: Option<&str>,
_emitter: &Arc<Emitter>,
_stage_scope: &StageScope,
) -> Result<CodergenResult, Error> {
self.complete(prompt).await
}