diff --git a/lib/crates/fabro-api/src/server.rs b/lib/crates/fabro-api/src/server.rs index 507bbf27f..563f4b305 100644 --- a/lib/crates/fabro-api/src/server.rs +++ b/lib/crates/fabro-api/src/server.rs @@ -533,6 +533,7 @@ async fn start_run( workflow_slug: None, }, settings, + cwd: std::env::current_dir().unwrap_or_else(|_| std::env::temp_dir()), run_dir: Some(run_dir.clone()), run_id: Some(run_id.clone()), host_repo_path: None, diff --git a/lib/crates/fabro-cli/src/commands/preflight.rs b/lib/crates/fabro-cli/src/commands/preflight.rs index aa119355c..b2d37312d 100644 --- a/lib/crates/fabro-cli/src/commands/preflight.rs +++ b/lib/crates/fabro-cli/src/commands/preflight.rs @@ -29,6 +29,7 @@ pub async fn execute(mut args: PreflightArgs) -> anyhow::Result<()> { fabro_workflows::operations::ResolveWorkflowRequest { workflow: fabro_workflows::operations::WorkflowInput::Path(args.workflow.clone()), settings: settings.clone(), + cwd: std::env::current_dir()?, }, )?; diff --git a/lib/crates/fabro-cli/src/commands/run/attach.rs b/lib/crates/fabro-cli/src/commands/run/attach.rs index c8898358d..ca015b81f 100644 --- a/lib/crates/fabro-cli/src/commands/run/attach.rs +++ b/lib/crates/fabro-cli/src/commands/run/attach.rs @@ -271,7 +271,7 @@ fn read_status_record(path: &Path) -> Option { } fn read_launcher_pid(run_dir: &Path) -> Option { - super::launcher::launcher_record_for_run(run_dir) + super::launcher::active_launcher_record_for_run(run_dir) .map(|record| record.pid) .or_else(|| { std::fs::read_to_string(run_dir.join("run.pid")) diff --git a/lib/crates/fabro-cli/src/commands/run/create.rs b/lib/crates/fabro-cli/src/commands/run/create.rs index 2a214d19f..c2d0ca882 100644 --- a/lib/crates/fabro-cli/src/commands/run/create.rs +++ b/lib/crates/fabro-cli/src/commands/run/create.rs @@ -27,19 +27,17 @@ pub async fn create_run( cli_args_config, true, )?; - let working_directory = std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")); - let base_branch = fabro_sandbox::daytona::detect_repo_info(&working_directory) - .ok() - .and_then(|(_, branch)| branch); + let cwd = std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")); let created = match fabro_workflows::operations::create(fabro_workflows::operations::CreateRequest { workflow: fabro_workflows::operations::WorkflowInput::Path(workflow_path.clone()), settings, + cwd, run_dir: None, run_id: args.run_id.clone(), - base_branch, - host_repo_path: Some(working_directory.to_string_lossy().to_string()), + base_branch: None, + host_repo_path: None, }) { Ok(created) => created, Err(fabro_workflows::error::FabroError::ValidationFailed { diagnostics }) => { diff --git a/lib/crates/fabro-cli/src/commands/run/launcher.rs b/lib/crates/fabro-cli/src/commands/run/launcher.rs index 91f1540c6..990fe9e22 100644 --- a/lib/crates/fabro-cli/src/commands/run/launcher.rs +++ b/lib/crates/fabro-cli/src/commands/run/launcher.rs @@ -43,8 +43,99 @@ pub(crate) fn remove_launcher_record(path: &Path) { let _ = std::fs::remove_file(path); } -pub(crate) fn launcher_record_for_run(run_dir: &Path) -> Option { - let record = fabro_workflows::records::RunRecord::load(run_dir).ok()?; - let path = launcher_record_path(&record.settings.storage_dir(), &record.run_id); - read_launcher_record(&path) +pub(crate) fn active_launcher_record_for_run(run_dir: &Path) -> Option { + let run_record = fabro_workflows::records::RunRecord::load(run_dir).ok()?; + let path = launcher_record_path(&run_record.settings.storage_dir(), &run_record.run_id); + let launcher = read_launcher_record(&path)?; + if launcher_record_is_running(&launcher) { + Some(launcher) + } else { + remove_launcher_record(&path); + None + } +} + +pub(crate) fn launcher_record_is_running(record: &LauncherRecord) -> bool { + process_alive(record.pid) && launcher_process_matches(record) +} + +#[cfg(unix)] +fn process_alive(pid: u32) -> bool { + unsafe { libc::kill(pid as i32, 0) == 0 } +} + +#[cfg(not(unix))] +fn process_alive(_pid: u32) -> bool { + true +} + +#[cfg(unix)] +fn launcher_process_matches(record: &LauncherRecord) -> bool { + let output = match std::process::Command::new("ps") + .args(["-ww", "-o", "command=", "-p", &record.pid.to_string()]) + .output() + { + Ok(output) if output.status.success() => output, + _ => return false, + }; + + let command = String::from_utf8_lossy(&output.stdout); + let run_dir = record.run_dir.to_string_lossy(); + command.contains("__detached") && command.contains(run_dir.as_ref()) +} + +#[cfg(not(unix))] +fn launcher_process_matches(_record: &LauncherRecord) -> bool { + true +} + +#[cfg(test)] +mod tests { + use super::*; + use chrono::Utc; + use fabro_config::FabroSettings; + use fabro_graphviz::graph::Graph; + use fabro_workflows::records::RunRecord; + + #[test] + fn active_launcher_record_for_run_removes_stale_record() { + let dir = tempfile::tempdir().unwrap(); + let storage_dir = dir.path().join("storage"); + let run_dir = dir.path().join("run"); + std::fs::create_dir_all(&run_dir).unwrap(); + + RunRecord { + run_id: "run-test".to_string(), + created_at: Utc::now(), + settings: FabroSettings { + storage_dir: Some(storage_dir.clone()), + ..Default::default() + }, + graph: Graph::default(), + workflow_slug: None, + working_directory: dir.path().to_path_buf(), + host_repo_path: None, + base_branch: None, + labels: std::collections::HashMap::new(), + } + .save(&run_dir) + .unwrap(); + + let launcher_path = launcher_record_path(&storage_dir, "run-test"); + write_launcher_record( + &launcher_path, + &LauncherRecord { + run_id: "run-test".to_string(), + run_dir: run_dir.clone(), + pid: u32::MAX, + resume: false, + log_path: dir.path().join("launcher.log"), + started_at: Utc::now(), + }, + ) + .unwrap(); + + assert!(active_launcher_record_for_run(&run_dir).is_none()); + assert!(!launcher_path.exists()); + } } diff --git a/lib/crates/fabro-cli/src/commands/run/resume.rs b/lib/crates/fabro-cli/src/commands/run/resume.rs index cf7a3ebd2..3851e3c92 100644 --- a/lib/crates/fabro-cli/src/commands/run/resume.rs +++ b/lib/crates/fabro-cli/src/commands/run/resume.rs @@ -39,7 +39,7 @@ pub async fn resume_command(args: ResumeArgs, styles: &'static Styles) -> anyhow } fn launcher_pid_alive(run_dir: &std::path::Path) -> bool { - super::launcher::launcher_record_for_run(run_dir) + super::launcher::active_launcher_record_for_run(run_dir) .map(|record| process_alive(record.pid)) .or_else(|| { std::fs::read_to_string(run_dir.join("run.pid")) diff --git a/lib/crates/fabro-cli/src/commands/run/start.rs b/lib/crates/fabro-cli/src/commands/run/start.rs index c2a57f08d..2f58145e1 100644 --- a/lib/crates/fabro-cli/src/commands/run/start.rs +++ b/lib/crates/fabro-cli/src/commands/run/start.rs @@ -4,7 +4,8 @@ use anyhow::{anyhow, Result}; use chrono::Utc; use super::launcher::{ - launcher_log_path, launcher_record_path, write_launcher_record, LauncherRecord, + launcher_log_path, launcher_record_path, remove_launcher_record, write_launcher_record, + LauncherRecord, }; /// Spawn a detached engine process for the given run directory. @@ -67,6 +68,10 @@ pub fn start_run(run_dir: &Path, resume: bool) -> Result { return Err(err); } + if matches!(child.try_wait(), Ok(Some(_))) { + remove_launcher_record(&launcher_path); + } + Ok(child) } diff --git a/lib/crates/fabro-workflows/src/operations/create.rs b/lib/crates/fabro-workflows/src/operations/create.rs index 690b8230b..bef8886b5 100644 --- a/lib/crates/fabro-workflows/src/operations/create.rs +++ b/lib/crates/fabro-workflows/src/operations/create.rs @@ -29,29 +29,13 @@ pub struct ValidateOptions { pub struct CreateRequest { pub workflow: WorkflowInput, pub settings: FabroSettings, + pub cwd: PathBuf, pub run_dir: Option, pub run_id: Option, pub host_repo_path: Option, pub base_branch: Option, } -impl Default for CreateRequest { - fn default() -> Self { - Self { - workflow: WorkflowInput::DotSource { - source: String::new(), - base_dir: None, - workflow_slug: None, - }, - settings: FabroSettings::default(), - run_dir: None, - run_id: None, - host_repo_path: None, - base_branch: None, - } - } -} - #[derive(Debug)] pub struct CreatedRun { pub persisted: Persisted, @@ -67,7 +51,7 @@ struct PersistCreateOptions { workflow_slug: Option, labels: HashMap, base_branch: Option, - working_directory: Option, + working_directory: PathBuf, host_repo_path: Option, goal_override: Option, base_dir: Option, @@ -106,6 +90,7 @@ pub fn create(request: CreateRequest) -> Result { let resolved = resolve_workflow(ResolveWorkflowRequest { workflow: request.workflow, settings: request.settings, + cwd: request.cwd, }) .map_err(|err| FabroError::Parse(err.to_string()))?; @@ -116,6 +101,7 @@ pub fn create(request: CreateRequest) -> Result { let CreateRequest { workflow: _, settings: _, + cwd: _, run_dir, run_id, host_repo_path, @@ -150,7 +136,7 @@ pub fn create(request: CreateRequest) -> Result { workflow_slug: resolved.workflow_slug.clone(), labels: resolved.settings.labels.clone(), base_branch, - working_directory: Some(working_directory), + working_directory, host_repo_path, goal_override: resolved.goal_override.clone(), base_dir: resolved.base_dir.clone(), @@ -276,8 +262,6 @@ fn persist_validated( let run_id = run_id.unwrap_or_else(|| ulid::Ulid::new().to_string()); let run_dir = run_dir.unwrap_or_else(|| default_run_dir(&run_id, settings.dry_run_enabled())); - let working_directory = working_directory - .unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."))); let run_record = RunRecord { run_id, @@ -529,14 +513,19 @@ mod tests { graph [goal="Test"] work [label="Work"] }"#; + let dir = tempfile::tempdir().unwrap(); let err = create(CreateRequest { workflow: WorkflowInput::DotSource { source: dot.to_string(), base_dir: None, workflow_slug: None, }, - run_dir: Some(tempfile::tempdir().unwrap().path().join("run")), - ..Default::default() + settings: FabroSettings::default(), + cwd: dir.path().to_path_buf(), + run_dir: Some(dir.path().join("run")), + run_id: None, + host_repo_path: None, + base_branch: None, }) .unwrap_err(); @@ -572,11 +561,11 @@ mod tests { labels: HashMap::from([("env".to_string(), "test".to_string())]), ..Default::default() }, + cwd: dir.path().to_path_buf(), run_dir: Some(dir.path().join("run")), run_id: Some("run-123".to_string()), host_repo_path: Some(dir.path().display().to_string()), base_branch: Some("main".to_string()), - ..Default::default() }) .unwrap(); @@ -644,7 +633,11 @@ mod tests { dry_run: Some(true), ..Default::default() }, - ..Default::default() + cwd: dir.path().to_path_buf(), + run_dir: None, + run_id: None, + host_repo_path: None, + base_branch: None, }) .unwrap(); @@ -653,4 +646,43 @@ mod tests { "version = 1\ngraph = \"workflow.fabro\"\n" ); } + + #[test] + fn create_resolves_working_directory_and_repo_path_from_request_cwd() { + let dir = tempfile::tempdir().unwrap(); + let workspace = dir.path().join("workspace"); + std::fs::create_dir_all(&workspace).unwrap(); + + let created = create(CreateRequest { + workflow: WorkflowInput::DotSource { + source: MINIMAL_DOT.to_string(), + base_dir: None, + workflow_slug: None, + }, + settings: FabroSettings { + work_dir: Some("workspace".to_string()), + dry_run: Some(true), + ..Default::default() + }, + cwd: dir.path().to_path_buf(), + run_dir: Some(dir.path().join("run")), + run_id: Some("run-cwd".to_string()), + host_repo_path: None, + base_branch: None, + }) + .unwrap(); + + assert_eq!(created.persisted.run_record().working_directory, workspace); + assert_eq!( + created.persisted.run_record().host_repo_path.as_deref(), + Some( + created + .persisted + .run_record() + .working_directory + .to_string_lossy() + .as_ref() + ) + ); + } } diff --git a/lib/crates/fabro-workflows/src/operations/source.rs b/lib/crates/fabro-workflows/src/operations/source.rs index 074674b7d..19d3d455c 100644 --- a/lib/crates/fabro-workflows/src/operations/source.rs +++ b/lib/crates/fabro-workflows/src/operations/source.rs @@ -29,6 +29,7 @@ pub struct WorkflowPathResolution { pub struct ResolveWorkflowRequest { pub workflow: WorkflowInput, pub settings: FabroSettings, + pub cwd: PathBuf, } #[derive(Clone, Debug)] @@ -93,6 +94,24 @@ pub fn workflow_slug_from_path(workflow_path: &Path) -> Option { Some(file_stem.into_owned()) } +fn cached_workflow_graph_path(path: &Path) -> Option { + if path.file_name().and_then(|name| name.to_str()) != Some("workflow.toml") { + return None; + } + + let canonical = path.with_file_name(RUN_GRAPH_FILE); + if canonical.exists() { + return Some(canonical); + } + + let legacy = path.with_file_name(LEGACY_RUN_GRAPH_FILE); + if legacy.exists() { + return Some(legacy); + } + + None +} + pub fn resolve_workflow_path(workflow_path: &Path) -> anyhow::Result { let path = project_config::resolve_workflow_arg(workflow_path)?; let workflow_slug = workflow_slug_from_path(&path); @@ -111,15 +130,9 @@ pub fn resolve_workflow_path(workflow_path: &Path) -> anyhow::Result - { - let canonical = path.with_file_name(RUN_GRAPH_FILE); - let legacy = path.with_file_name(LEGACY_RUN_GRAPH_FILE); - let dot_path = if canonical.exists() || !legacy.exists() { - canonical - } else { - legacy + Err(_) if !path.exists() => { + let Some(dot_path) = cached_workflow_graph_path(&path) else { + anyhow::bail!("Workflow not found: {}", path.display()); }; Ok(WorkflowPathResolution { resolved_workflow_path: path, @@ -177,12 +190,11 @@ pub fn resolve_settings_for_path( } pub fn resolve_workflow(request: ResolveWorkflowRequest) -> anyhow::Result { - let caller_cwd = std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")); match request.workflow { WorkflowInput::Path(workflow_path) => { let resolution = resolve_workflow_path(&workflow_path)?; let settings = request.settings; - let working_directory = resolve_working_directory(&settings, &caller_cwd); + let working_directory = resolve_working_directory(&settings, &request.cwd); let raw_source = std::fs::read_to_string(&resolution.dot_path) .with_context(|| format!("Failed to read {}", resolution.dot_path.display()))?; let goal_override = settings.goal.clone().or(resolve_goal_file( @@ -214,7 +226,7 @@ pub fn resolve_workflow(request: ResolveWorkflowRequest) -> anyhow::Result { let settings = request.settings; - let working_directory = resolve_working_directory(&settings, &caller_cwd); + let working_directory = resolve_working_directory(&settings, &request.cwd); let goal_override = settings.goal.clone().or(resolve_goal_file( settings.goal_file.as_deref(), &working_directory, @@ -233,3 +245,50 @@ pub fn resolve_workflow(request: ResolveWorkflowRequest) -> anyhow::Result exit }", + ) + .unwrap(); + + let resolution = resolve_workflow_path(&run_dir.join("workflow.toml")).unwrap(); + + assert_eq!(resolution.dot_path, run_dir.join("workflow.fabro")); + assert!(resolution.workflow_config.is_none()); + assert!(resolution.workflow_toml_path.is_none()); + } + + #[test] + fn resolve_workflow_uses_explicit_cwd_for_relative_work_dir() { + let dir = tempfile::tempdir().unwrap(); + let resolved = resolve_workflow(ResolveWorkflowRequest { + workflow: WorkflowInput::DotSource { + source: "digraph Test { start -> exit }".to_string(), + base_dir: None, + workflow_slug: None, + }, + settings: FabroSettings { + work_dir: Some("workspace".to_string()), + ..Default::default() + }, + cwd: dir.path().to_path_buf(), + }) + .unwrap(); + + assert_eq!(resolved.working_directory, dir.path().join("workspace")); + } +} diff --git a/lib/crates/fabro-workflows/src/operations/start.rs b/lib/crates/fabro-workflows/src/operations/start.rs index 1833ebe09..5b70df10a 100644 --- a/lib/crates/fabro-workflows/src/operations/start.rs +++ b/lib/crates/fabro-workflows/src/operations/start.rs @@ -144,22 +144,6 @@ async fn execute_persisted_run( } }; - let original_cwd = std::env::current_dir().ok(); - if let Err(err) = std::env::set_current_dir(&persisted.run_record().working_directory) { - let err = FabroError::Io(format!( - "Failed to set working directory to {}: {err}", - persisted.run_record().working_directory.display() - )); - let _ = persist_detached_failure(run_dir, "bootstrap", StatusReason::BootstrapFailed, &err); - bootstrap_guard.defuse(); - return Err(err); - } - let _cwd_guard = scopeguard::guard(original_cwd, |cwd| { - if let Some(cwd) = cwd { - let _ = std::env::set_current_dir(cwd); - } - }); - let options = match derive_start_options(&persisted, services) { Ok(options) => options, Err(err) => { @@ -801,11 +785,14 @@ mod tests { dry_run: Some(true), ..Default::default() }, + cwd: run_dir + .parent() + .unwrap_or_else(|| Path::new(".")) + .to_path_buf(), run_dir: Some(run_dir.to_path_buf()), run_id: Some("run-test".to_string()), - host_repo_path: Some(std::env::current_dir().unwrap().display().to_string()), - base_branch: Some("main".to_string()), - ..Default::default() + host_repo_path: None, + base_branch: None, }) .unwrap() .persisted