mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-06 02:48:25 +00:00
fix(run): propagate worker cancellation into command stages
Pass a real cancel signal from __run-worker through the workflow engine into sandbox command execution so cancelled runs reap gated shell loops instead of leaking slow.gate waiters.
This commit is contained in:
parent
099362cd42
commit
4470eb21ad
7 changed files with 129 additions and 11 deletions
|
|
@ -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<RunControlState>) -> Result<()> {
|
||||
fn install_signal_handlers(
|
||||
run_control: Arc<RunControlState>,
|
||||
cancel_token: Arc<AtomicBool>,
|
||||
) -> Result<()> {
|
||||
#[cfg(unix)]
|
||||
{
|
||||
let mut pause = signal(SignalKind::user_defined1())?;
|
||||
|
|
@ -220,6 +225,21 @@ fn install_signal_handlers(run_control: Arc<RunControlState>) -> 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(())
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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<Option<String>>,
|
||||
captured_env_vars: std::sync::Mutex<Option<std::collections::HashMap<String, String>>>,
|
||||
captured_cancel_token: std::sync::Mutex<Option<bool>>,
|
||||
}
|
||||
|
||||
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<String, String>>,
|
||||
_cancel_token: Option<tokio_util::sync::CancellationToken>,
|
||||
cancel_token: Option<tokio_util::sync::CancellationToken>,
|
||||
) -> Result<fabro_agent::sandbox::ExecResult, String> {
|
||||
*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;
|
||||
|
|
|
|||
|
|
@ -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<String, String>,
|
||||
/// 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<Arc<AtomicBool>>,
|
||||
/// Logical path of the current workflow when running from a bundle.
|
||||
pub workflow_path: Option<PathBuf>,
|
||||
/// 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<CancellationToken> {
|
||||
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,
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
};
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
});
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue