diff --git a/lib/crates/fabro-cli/src/commands/run/runner.rs b/lib/crates/fabro-cli/src/commands/run/runner.rs index 50ab14a4a..c61c56da6 100644 --- a/lib/crates/fabro-cli/src/commands/run/runner.rs +++ b/lib/crates/fabro-cli/src/commands/run/runner.rs @@ -1,5 +1,6 @@ use std::path::PathBuf; use std::sync::Arc; +use std::sync::atomic::{AtomicBool, Ordering}; use std::time::Duration; use anyhow::{Context, Result, anyhow}; @@ -58,12 +59,13 @@ pub(crate) async fn execute( scratch.interview_claim_path(), )); let run_control = RunControlState::new(); - install_signal_handlers(Arc::clone(&run_control))?; + let cancel_token = Arc::new(AtomicBool::new(false)); + install_signal_handlers(Arc::clone(&run_control), Arc::clone(&cancel_token))?; let github_app = maybe_build_github_app_credentials(&run_record.settings)?; let event_client = client.clone_for_reuse(); let services = fabro_workflow::operations::StartServices { run_id, - cancel_token: None, + cancel_token: Some(Arc::clone(&cancel_token)), emitter: Arc::new(Emitter::new(run_id)), interviewer, run_store: run_store.clone(), @@ -203,7 +205,10 @@ fn maybe_build_github_app_credentials( } } -fn install_signal_handlers(run_control: Arc) -> Result<()> { +fn install_signal_handlers( + run_control: Arc, + cancel_token: Arc, +) -> Result<()> { #[cfg(unix)] { let mut pause = signal(SignalKind::user_defined1())?; @@ -220,6 +225,21 @@ fn install_signal_handlers(run_control: Arc) -> Result<()> { run_control.request_unpause(); } }); + + let mut terminate = signal(SignalKind::terminate())?; + let terminate_cancel = Arc::clone(&cancel_token); + tokio::spawn(async move { + while terminate.recv().await.is_some() { + terminate_cancel.store(true, Ordering::SeqCst); + } + }); + + let mut interrupt = signal(SignalKind::interrupt())?; + tokio::spawn(async move { + while interrupt.recv().await.is_some() { + cancel_token.store(true, Ordering::SeqCst); + } + }); } Ok(()) diff --git a/lib/crates/fabro-cli/tests/it/cmd/run.rs b/lib/crates/fabro-cli/tests/it/cmd/run.rs index c0fcd2b73..33293f675 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/run.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/run.rs @@ -4,8 +4,8 @@ use httpmock::MockServer; use serde_json::Value; use super::support::{ - only_run, output_stderr, run_count_for_test_case, run_state, wait_for_status, - write_gated_workflow, + only_run, output_stderr, run_count_for_test_case, run_state, wait_for_no_process_match, + wait_for_status, write_gated_workflow, }; use crate::support::{example_fixture, fabro_json_snapshot, run_output_filters, unique_run_id}; @@ -1357,7 +1357,7 @@ fn detach_creates_run_dir_with_detach_log() { #[test] fn ctrl_c_cancels_active_run_via_server() { let context = test_context!(); - let _gate = write_gated_workflow(&context.temp_dir.join("slow.fabro"), "slow", "Run slowly"); + let gate = write_gated_workflow(&context.temp_dir.join("slow.fabro"), "slow", "Run slowly"); let mut run_cmd = std::process::Command::new(env!("CARGO_BIN_EXE_fabro")); run_cmd.current_dir(&context.temp_dir); @@ -1418,4 +1418,6 @@ fn ctrl_c_cancels_active_run_via_server() { .and_then(|record| record.reason), Some(StatusReason::Cancelled) ); + let gate_pattern = gate.gate_path().to_string_lossy().into_owned(); + wait_for_no_process_match(&gate_pattern); } diff --git a/lib/crates/fabro-cli/tests/it/cmd/support.rs b/lib/crates/fabro-cli/tests/it/cmd/support.rs index 0b695bcf4..73fef9765 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/support.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/support.rs @@ -341,11 +341,34 @@ pub(crate) fn setup_project_fixture(context: &TestContext) -> ProjectFixture { } impl WorkflowGate { + pub(crate) fn gate_path(&self) -> &Path { + &self.gate_path + } + pub(crate) fn release(&self) { write_text_file(&self.gate_path, "open\n"); } } +pub(crate) fn wait_for_no_process_match(pattern: &str) { + let deadline = Instant::now() + COMMAND_TIMEOUT; + loop { + let output = std::process::Command::new("pgrep") + .args(["-f", pattern]) + .output() + .expect("pgrep should execute"); + if !output.status.success() { + return; + } + assert!( + Instant::now() < deadline, + "timed out waiting for processes matching {pattern:?} to exit: {}", + String::from_utf8_lossy(&output.stdout) + ); + std::thread::sleep(Duration::from_millis(50)); + } +} + pub(crate) fn setup_artifact_run(context: &TestContext) -> WorkspaceRunSetup { let workspace_dir = context.temp_dir.join("artifact-run"); std::fs::create_dir_all(&workspace_dir) diff --git a/lib/crates/fabro-workflow/src/handler/command.rs b/lib/crates/fabro-workflow/src/handler/command.rs index 74bc3066d..0034aad9f 100644 --- a/lib/crates/fabro-workflow/src/handler/command.rs +++ b/lib/crates/fabro-workflow/src/handler/command.rs @@ -106,12 +106,17 @@ impl Handler for CommandHandler { } else { Some(&services.env) }; + let cancel_token = services.sandbox_cancel_token(); let result = services .sandbox - .exec_command(&command, timeout_ms, None, env_vars, None) - .await - .map_err(|e| FabroError::handler(format!("Failed to spawn script: {e}")))?; + .exec_command(&command, timeout_ms, None, env_vars, cancel_token.clone()) + .await; + if let Some(token) = cancel_token { + token.cancel(); + } + let result = + result.map_err(|e| FabroError::handler(format!("Failed to spawn script: {e}")))?; services.emitter.emit(&Event::CommandCompleted { node_id: node.id.clone(), @@ -173,6 +178,7 @@ mod tests { use fabro_types::fixtures; use object_store::memory::InMemory; use std::sync::Arc; + use std::sync::atomic::AtomicBool; use std::time::Duration; fn make_services() -> EngineServices { @@ -702,6 +708,7 @@ mod tests { exec_result: fabro_agent::sandbox::ExecResult, captured_command: std::sync::Mutex>, captured_env_vars: std::sync::Mutex>>, + captured_cancel_token: std::sync::Mutex>, } impl SpySandbox { @@ -710,6 +717,7 @@ mod tests { exec_result, captured_command: std::sync::Mutex::new(None), captured_env_vars: std::sync::Mutex::new(None), + captured_cancel_token: std::sync::Mutex::new(None), } } @@ -750,10 +758,11 @@ mod tests { _timeout_ms: u64, _working_dir: Option<&str>, env_vars: Option<&std::collections::HashMap>, - _cancel_token: Option, + cancel_token: Option, ) -> Result { *self.captured_command.lock().unwrap() = Some(command.to_string()); *self.captured_env_vars.lock().unwrap() = env_vars.cloned(); + *self.captured_cancel_token.lock().unwrap() = Some(cancel_token.is_some()); Ok(self.exec_result.clone()) } async fn grep( @@ -919,6 +928,35 @@ mod tests { ); } + #[tokio::test] + async fn passes_run_cancellation_to_sandbox() { + let spy = std::sync::Arc::new(SpySandbox::new(fabro_agent::sandbox::ExecResult { + stdout: String::new(), + stderr: String::new(), + exit_code: 0, + timed_out: false, + duration_ms: 5, + })); + + let handler = CommandHandler; + let mut node = Node::new("script_node"); + node.attrs + .insert("script".to_string(), AttrValue::String("true".to_string())); + let context = Context::new(); + let graph = Graph::new("test"); + let run_dir = tempfile::tempdir().unwrap(); + + let mut services = make_spy_services(spy.clone()); + services.cancel_requested = Some(Arc::new(AtomicBool::new(false))); + + handler + .execute(&node, &context, &graph, run_dir.path(), &services) + .await + .unwrap(); + + assert_eq!(*spy.captured_cancel_token.lock().unwrap(), Some(true)); + } + #[tokio::test] async fn tool_output_context_key_not_emitted() { let handler = CommandHandler; diff --git a/lib/crates/fabro-workflow/src/handler/mod.rs b/lib/crates/fabro-workflow/src/handler/mod.rs index 2f21df16a..01220b71d 100644 --- a/lib/crates/fabro-workflow/src/handler/mod.rs +++ b/lib/crates/fabro-workflow/src/handler/mod.rs @@ -15,7 +15,7 @@ use std::any::Any; use std::collections::HashMap; use std::path::{Path, PathBuf}; use std::sync::Arc; -#[cfg(test)] +use std::sync::atomic::{AtomicBool, Ordering}; use std::time::Duration; use async_trait::async_trait; @@ -35,6 +35,8 @@ use crate::workflow_bundle::WorkflowBundle; use fabro_graphviz::graph::{Graph, Node, shape_to_handler_type}; use fabro_hooks::{HookContext, HookDecision, HookRunner}; use fabro_interview::Interviewer; +use tokio::time; +use tokio_util::sync::CancellationToken; /// Shared services available to all handlers during execution. pub struct EngineServices { @@ -51,6 +53,8 @@ pub struct EngineServices { pub env: HashMap, /// When true, handlers should skip real execution and return simulated results. pub dry_run: bool, + /// Optional run-scoped cancellation flag from the core executor. + pub cancel_requested: Option>, /// Logical path of the current workflow when running from a bundle. pub workflow_path: Option, /// Bundled workflows available for child-workflow resolution. @@ -68,6 +72,33 @@ impl EngineServices { *self.git_state.write().unwrap() = state; } + /// Bridge the core executor's atomic cancel flag to sandbox command cancellation. + pub fn sandbox_cancel_token(&self) -> Option { + let cancel_requested = self.cancel_requested.clone()?; + let token = CancellationToken::new(); + + if cancel_requested.load(Ordering::Relaxed) { + token.cancel(); + return Some(token); + } + + let token_clone = token.clone(); + tokio::spawn(async move { + loop { + if token_clone.is_cancelled() { + return; + } + if cancel_requested.load(Ordering::Relaxed) { + token_clone.cancel(); + return; + } + time::sleep(Duration::from_millis(10)).await; + } + }); + + Some(token) + } + /// Run lifecycle hooks and return the merged decision. /// Returns `Proceed` if no hook runner is configured. pub async fn run_hooks(&self, hook_context: &HookContext) -> HookDecision { @@ -111,6 +142,7 @@ impl EngineServices { hook_runner: None, env: HashMap::new(), dry_run: false, + cancel_requested: None, workflow_path: None, workflow_bundle: None, } diff --git a/lib/crates/fabro-workflow/src/handler/parallel.rs b/lib/crates/fabro-workflow/src/handler/parallel.rs index 92341ea9b..6266042ce 100644 --- a/lib/crates/fabro-workflow/src/handler/parallel.rs +++ b/lib/crates/fabro-workflow/src/handler/parallel.rs @@ -271,6 +271,7 @@ impl Handler for ParallelHandler { let run_store = services.run_store.clone(); let env = services.env.clone(); let dry_run = services.dry_run; + let cancel_requested = services.cancel_requested.clone(); let workflow_path = services.workflow_path.clone(); let workflow_bundle = services.workflow_bundle.clone(); let graph = graph.clone(); @@ -324,6 +325,7 @@ impl Handler for ParallelHandler { hook_runner: hook_runner.clone(), env: env.clone(), dry_run, + cancel_requested, workflow_path, workflow_bundle, }; diff --git a/lib/crates/fabro-workflow/src/pipeline/execute.rs b/lib/crates/fabro-workflow/src/pipeline/execute.rs index 3c89039be..1c785c22f 100644 --- a/lib/crates/fabro-workflow/src/pipeline/execute.rs +++ b/lib/crates/fabro-workflow/src/pipeline/execute.rs @@ -81,6 +81,7 @@ pub async fn execute(init: Initialized) -> Executed { hook_runner: hook_runner.clone(), env, dry_run, + cancel_requested: run_options.cancel_token.clone(), workflow_path, workflow_bundle, });