Tighten workflows request context and launcher cleanup

This commit is contained in:
Bryan Helmkamp 2026-03-27 17:04:45 -04:00
parent 3c3cefb452
commit f8cc95e97d
10 changed files with 243 additions and 69 deletions

View file

@ -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,

View file

@ -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()?,
},
)?;

View file

@ -271,7 +271,7 @@ fn read_status_record(path: &Path) -> Option<RunStatusRecord> {
}
fn read_launcher_pid(run_dir: &Path) -> Option<u32> {
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"))

View file

@ -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 }) => {

View file

@ -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<LauncherRecord> {
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<LauncherRecord> {
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());
}
}

View file

@ -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"))

View file

@ -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<std::process::Child> {
return Err(err);
}
if matches!(child.try_wait(), Ok(Some(_))) {
remove_launcher_record(&launcher_path);
}
Ok(child)
}

View file

@ -29,29 +29,13 @@ pub struct ValidateOptions {
pub struct CreateRequest {
pub workflow: WorkflowInput,
pub settings: FabroSettings,
pub cwd: PathBuf,
pub run_dir: Option<PathBuf>,
pub run_id: Option<String>,
pub host_repo_path: Option<String>,
pub base_branch: Option<String>,
}
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<String>,
labels: HashMap<String, String>,
base_branch: Option<String>,
working_directory: Option<PathBuf>,
working_directory: PathBuf,
host_repo_path: Option<String>,
goal_override: Option<String>,
base_dir: Option<PathBuf>,
@ -106,6 +90,7 @@ pub fn create(request: CreateRequest) -> Result<CreatedRun, FabroError> {
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<CreatedRun, FabroError> {
let CreateRequest {
workflow: _,
settings: _,
cwd: _,
run_dir,
run_id,
host_repo_path,
@ -150,7 +136,7 @@ pub fn create(request: CreateRequest) -> Result<CreatedRun, FabroError> {
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()
)
);
}
}

View file

@ -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<String> {
Some(file_stem.into_owned())
}
fn cached_workflow_graph_path(path: &Path) -> Option<PathBuf> {
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<WorkflowPathResolution> {
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<WorkflowPat
workflow_slug,
})
}
Err(_)
if !path.exists() && path.starts_with(crate::run_lookup::default_runs_base()) =>
{
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<ResolvedWorkflow> {
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<Resol
workflow_slug,
} => {
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<Resol
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn resolve_workflow_path_uses_cached_graph_sibling_for_missing_workflow_toml() {
let dir = tempfile::tempdir().unwrap();
let run_dir = dir
.path()
.join("custom-storage")
.join("runs")
.join("run-123");
std::fs::create_dir_all(&run_dir).unwrap();
std::fs::write(
run_dir.join("workflow.fabro"),
"digraph Test { start -> 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"));
}
}

View file

@ -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