diff --git a/lib/crates/fabro-workflows/src/operations/start.rs b/lib/crates/fabro-workflows/src/operations/start.rs index 5f4205903..11e9b8f28 100644 --- a/lib/crates/fabro-workflows/src/operations/start.rs +++ b/lib/crates/fabro-workflows/src/operations/start.rs @@ -115,3 +115,376 @@ pub async fn start(validated: Validated, options: StartOptions) -> Result exit + }"#; + + const EMIT_DOT: &str = r#"digraph Test { + graph [goal="Ship feature"] + start [shape=Mdiamond] + work [type="emit"] + exit [shape=Msquare] + start -> work -> exit + }"#; + + struct CleanupCountingSandbox { + inner: Arc, + cleanup_count: Arc, + } + + #[async_trait] + impl Sandbox for CleanupCountingSandbox { + async fn read_file( + &self, + path: &str, + offset: Option, + limit: Option, + ) -> Result { + self.inner.read_file(path, offset, limit).await + } + + async fn write_file(&self, path: &str, content: &str) -> Result<(), String> { + self.inner.write_file(path, content).await + } + + async fn delete_file(&self, path: &str) -> Result<(), String> { + self.inner.delete_file(path).await + } + + async fn file_exists(&self, path: &str) -> Result { + self.inner.file_exists(path).await + } + + async fn list_directory( + &self, + path: &str, + depth: Option, + ) -> Result, String> { + self.inner.list_directory(path, depth).await + } + + async fn exec_command( + &self, + command: &str, + timeout_ms: u64, + working_dir: Option<&str>, + env_vars: Option<&HashMap>, + cancel_token: Option, + ) -> Result { + self.inner + .exec_command(command, timeout_ms, working_dir, env_vars, cancel_token) + .await + } + + async fn grep( + &self, + pattern: &str, + path: &str, + options: &GrepOptions, + ) -> Result, String> { + self.inner.grep(pattern, path, options).await + } + + async fn glob(&self, pattern: &str, path: Option<&str>) -> Result, String> { + self.inner.glob(pattern, path).await + } + + async fn download_file_to_local( + &self, + remote_path: &str, + local_path: &Path, + ) -> Result<(), String> { + self.inner + .download_file_to_local(remote_path, local_path) + .await + } + + async fn upload_file_from_local( + &self, + local_path: &Path, + remote_path: &str, + ) -> Result<(), String> { + self.inner + .upload_file_from_local(local_path, remote_path) + .await + } + + async fn initialize(&self) -> Result<(), String> { + self.inner.initialize().await + } + + async fn cleanup(&self) -> Result<(), String> { + self.cleanup_count.fetch_add(1, Ordering::SeqCst); + self.inner.cleanup().await + } + + fn working_directory(&self) -> &str { + self.inner.working_directory() + } + + fn platform(&self) -> &str { + self.inner.platform() + } + + fn os_version(&self) -> String { + self.inner.os_version() + } + + fn sandbox_info(&self) -> String { + self.inner.sandbox_info() + } + + async fn refresh_push_credentials(&self) -> Result<(), String> { + self.inner.refresh_push_credentials().await + } + + async fn set_autostop_interval(&self, minutes: i32) -> Result<(), String> { + self.inner.set_autostop_interval(minutes).await + } + + async fn setup_git_for_run( + &self, + run_id: &str, + ) -> Result, String> { + self.inner.setup_git_for_run(run_id).await + } + + fn resume_setup_commands(&self, run_branch: &str) -> Vec { + self.inner.resume_setup_commands(run_branch) + } + + async fn git_push_branch(&self, branch: &str) -> bool { + self.inner.git_push_branch(branch).await + } + + fn host_git_dir(&self) -> Option<&str> { + self.inner.host_git_dir() + } + + fn parallel_worktree_path( + &self, + run_dir: &Path, + run_id: &str, + node_id: &str, + key: &str, + ) -> String { + self.inner + .parallel_worktree_path(run_dir, run_id, node_id, key) + } + + async fn ssh_access_command(&self) -> Result, String> { + self.inner.ssh_access_command().await + } + + fn origin_url(&self) -> Option<&str> { + self.inner.origin_url() + } + + async fn get_preview_url( + &self, + port: u16, + ) -> Result)>, String> { + self.inner.get_preview_url(port).await + } + + fn mark_agent_read(&self, path: &str) { + self.inner.mark_agent_read(path); + } + } + + struct EmitCheckpointHandler; + + #[async_trait] + impl Handler for EmitCheckpointHandler { + async fn execute( + &self, + node: &Node, + _context: &Context, + _graph: &Graph, + _run_dir: &Path, + services: &crate::handler::EngineServices, + ) -> Result { + services.emitter.emit(&WorkflowRunEvent::CheckpointCompleted { + node_id: node.id.clone(), + status: "success".to_string(), + git_commit_sha: Some("sha-test".to_string()), + }); + Ok(Outcome::success()) + } + } + + fn validated_workflow(dot: &str) -> Validated { + let validated = crate::operations::create(dot, crate::operations::CreateOptions::default()) + .unwrap(); + validated.raise_on_errors().unwrap(); + validated + } + + fn test_settings(run_dir: &std::path::Path) -> RunSettings { + RunSettings { + config: FabroConfig::default(), + run_dir: run_dir.to_path_buf(), + cancel_token: None, + dry_run: false, + run_id: "run-test".to_string(), + labels: HashMap::new(), + git_author: crate::git::GitAuthor::default(), + workflow_slug: None, + github_app: None, + host_repo_path: None, + base_branch: None, + git: None, + } + } + + fn test_registry() -> HandlerRegistry { + let mut registry = HandlerRegistry::new(Box::new(StartHandler)); + registry.register("start", Box::new(StartHandler)); + registry.register("exit", Box::new(ExitHandler)); + registry.register("emit", Box::new(EmitCheckpointHandler)); + registry + } + + fn test_start_options( + run_dir: &std::path::Path, + sandbox: Arc, + emitter: Arc, + registry: Arc, + lifecycle: LifecycleConfig, + preserve_sandbox: bool, + ) -> StartOptions { + StartOptions { + init: InitOptions { + run_id: "run-test".to_string(), + run_dir: run_dir.to_path_buf(), + dry_run: false, + emitter, + sandbox, + registry, + lifecycle, + run_settings: test_settings(run_dir), + hooks: fabro_hooks::HookConfig { hooks: vec![] }, + sandbox_env: HashMap::new(), + checkpoint: None, + seed_context: None, + }, + retro: StartRetroConfig { + enabled: false, + dry_run: false, + llm_client: None, + provider: fabro_llm::Provider::Anthropic, + model: "test-model".to_string(), + }, + finalize: StartFinalizeConfig { + preserve_sandbox, + pr_config: None, + github_app: None, + origin_url: None, + model: "test-model".to_string(), + }, + } + } + + fn counting_sandbox() -> (Arc, Arc) { + let cleanup_count = Arc::new(AtomicUsize::new(0)); + let inner: Arc = Arc::new(LocalSandbox::new(std::env::current_dir().unwrap())); + ( + Arc::new(CleanupCountingSandbox { + inner, + cleanup_count: Arc::clone(&cleanup_count), + }), + cleanup_count, + ) + } + + #[tokio::test] + async fn start_cleans_up_sandbox_when_initialize_fails() { + let temp = tempfile::tempdir().unwrap(); + let run_dir = temp.path().join("run"); + let emitter = Arc::new(EventEmitter::new()); + let registry = Arc::new(test_registry()); + let (sandbox, cleanup_count) = counting_sandbox(); + + let result = start( + validated_workflow(MINIMAL_DOT), + test_start_options( + &run_dir, + sandbox, + emitter, + registry, + LifecycleConfig { + setup_commands: vec!["false".to_string()], + setup_command_timeout_ms: 1_000, + devcontainer_phases: vec![], + }, + false, + ), + ) + .await; + + assert!(result.is_err()); + tokio::task::yield_now().await; + tokio::time::sleep(Duration::from_millis(20)).await; + assert_eq!(cleanup_count.load(Ordering::SeqCst), 1); + } + + #[tokio::test] + async fn start_captures_checkpoint_git_sha_in_conclusion() { + let temp = tempfile::tempdir().unwrap(); + let run_dir = temp.path().join("run"); + let emitter = Arc::new(EventEmitter::new()); + let registry = Arc::new(test_registry()); + let sandbox: Arc = + Arc::new(LocalSandbox::new(std::env::current_dir().unwrap())); + + let started = start( + validated_workflow(EMIT_DOT), + test_start_options( + &run_dir, + sandbox, + emitter, + registry, + LifecycleConfig { + setup_commands: vec![], + setup_command_timeout_ms: 1_000, + devcontainer_phases: vec![], + }, + true, + ), + ) + .await + .unwrap(); + + assert_eq!( + started.finalized.conclusion.final_git_commit_sha.as_deref(), + Some("sha-test") + ); + assert_eq!(started.finalized.conclusion.status, StageStatus::Success); + assert!(started.retro.is_none()); + } +} diff --git a/lib/crates/fabro-workflows/src/pipeline/initialize.rs b/lib/crates/fabro-workflows/src/pipeline/initialize.rs index 4d7ae3d62..2802a4eca 100644 --- a/lib/crates/fabro-workflows/src/pipeline/initialize.rs +++ b/lib/crates/fabro-workflows/src/pipeline/initialize.rs @@ -35,8 +35,10 @@ pub async fn initialize( // Create run directory and write graph std::fs::create_dir_all(&options.run_dir)?; - let graph_path = options.run_dir.join("graph.fabro"); - std::fs::write(&graph_path, &source)?; + if !source.is_empty() { + let graph_path = options.run_dir.join("graph.fabro"); + std::fs::write(&graph_path, &source)?; + } let hook_runner = if options.hooks.hooks.is_empty() { None @@ -291,4 +293,44 @@ mod tests { Some("value") ); } + + #[tokio::test] + async fn initialize_skips_empty_graph_source() { + let temp = tempfile::tempdir().unwrap(); + let run_dir = temp.path().join("run"); + let (graph, _source) = simple_graph(); + let validated = Validated::new(graph, String::new(), vec![]); + let emitter = Arc::new(crate::event::EventEmitter::new()); + let sandbox = Arc::new(fabro_agent::LocalSandbox::new( + std::env::current_dir().unwrap(), + )); + let registry = Arc::new(default_registry(Arc::new(AutoApproveInterviewer), || None)); + + let initialized = initialize( + validated, + InitOptions { + run_id: "run-test".to_string(), + run_dir: run_dir.clone(), + dry_run: false, + emitter, + sandbox, + registry, + lifecycle: crate::run_settings::LifecycleConfig { + setup_commands: vec![], + setup_command_timeout_ms: 1_000, + devcontainer_phases: vec![], + }, + run_settings: test_settings(&run_dir), + hooks: fabro_hooks::HookConfig { hooks: vec![] }, + sandbox_env: HashMap::new(), + checkpoint: None, + seed_context: None, + }, + ) + .await + .unwrap(); + + assert!(!run_dir.join("graph.fabro").exists()); + assert!(initialized.source.is_empty()); + } }