diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 94158ea8a..e099f8d75 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -1595,28 +1595,29 @@ async fn delete_run_internal( } else { None }; + let durable_status = if managed_run.is_some() { + load_durable_run_status(state.as_ref(), &id).await + } else { + None + }; + let should_signal_cancel = !durable_status.is_some_and(RunStatus::is_terminal); if let Some(managed_run) = managed_run.as_mut() { - if let Some(token) = &managed_run.cancel_token { - token.cancel(); - } - if let Some(answer_transport) = managed_run.answer_transport.clone() { - let _ = answer_transport.cancel_run().await; - } - if let Some(cancel_tx) = managed_run.cancel_tx.take() { - let _ = cancel_tx.send(()); + if should_signal_cancel { + if let Some(token) = &managed_run.cancel_token { + token.cancel(); + } + if let Some(answer_transport) = managed_run.answer_transport.clone() { + let _ = answer_transport.cancel_run().await; + } + if let Some(cancel_tx) = managed_run.cancel_tx.take() { + let _ = cancel_tx.send(()); + } } // Terminal runs can still carry a stale worker PID briefly after their // completion events land, so avoid paying the full cancellation grace. - let delete_grace = if matches!( - managed_run.status, - RunStatus::Submitted - | RunStatus::Queued - | RunStatus::Starting - | RunStatus::Running - | RunStatus::Blocked { .. } - | RunStatus::Paused { .. } - ) { + let delete_grace = if should_signal_cancel && managed_run.status.requires_force_to_delete() + { WORKER_CANCEL_GRACE } else { TERMINAL_DELETE_WORKER_GRACE @@ -1658,6 +1659,12 @@ async fn delete_run_internal( Ok(delete_outcome) } +async fn load_durable_run_status(state: &AppState, id: &RunId) -> Option { + let run_store = state.store.open_run(id).await.ok()?; + let projection = run_store.state().await.ok()?; + Some(projection.status) +} + async fn delete_run_sandbox_resource( state: &Arc, id: RunId, @@ -2178,6 +2185,11 @@ async fn shutdown_active_workers_with_grace( async fn persist_cancelled_run_status(state: &AppState, run_id: RunId) -> anyhow::Result<()> { let run_store = state.store.open_run(&run_id).await?; + let run_state = run_store.state().await?; + if run_state.status.is_terminal() { + return Ok(()); + } + workflow_event::append_event( &run_store, &run_id, diff --git a/lib/crates/fabro-server/src/server/tests.rs b/lib/crates/fabro-server/src/server/tests.rs index 867be6ba6..0e0d7aa0a 100644 --- a/lib/crates/fabro-server/src/server/tests.rs +++ b/lib/crates/fabro-server/src/server/tests.rs @@ -2298,6 +2298,86 @@ async fn append_default_run_created(run_store: &fabro_store::RunDatabase, run_id .unwrap(); } +#[tokio::test] +async fn persist_cancelled_run_status_ignores_already_terminal_runs() { + let state = test_app_state(); + let run_id = fixtures::RUN_1; + create_durable_run_with_events(&state, run_id, &[ + workflow_event::Event::WorkflowRunCompleted { + duration_ms: 1000, + artifact_count: 0, + status: "succeeded".to_string(), + reason: SuccessReason::Completed, + total_usd_micros: None, + final_git_commit_sha: None, + final_patch: None, + diff_summary: None, + billing: None, + }, + ]) + .await; + + persist_cancelled_run_status(state.as_ref(), run_id) + .await + .unwrap(); + + let run_store = state.store.open_run(&run_id).await.unwrap(); + let projection = run_store.state().await.unwrap(); + assert_eq!(projection.status, RunStatus::Succeeded { + reason: SuccessReason::Completed, + }); + assert!(!run_store.list_events().await.unwrap().iter().any(|event| { + matches!( + event.event.body, + EventBody::RunFailed(ref props) if props.reason == FailureReason::Cancelled + ) + })); +} + +#[tokio::test] +async fn delete_terminal_managed_run_does_not_send_cancel_signal() { + let state = test_app_state(); + let run_id = fixtures::RUN_1; + create_durable_run_with_events(&state, run_id, &[ + workflow_event::Event::WorkflowRunCompleted { + duration_ms: 1000, + artifact_count: 0, + status: "succeeded".to_string(), + reason: SuccessReason::Completed, + total_usd_micros: None, + final_git_commit_sha: None, + final_patch: None, + diff_summary: None, + billing: None, + }, + ]) + .await; + + let temp = tempfile::tempdir().unwrap(); + let run_dir = temp.path().join("run"); + std::fs::create_dir_all(&run_dir).unwrap(); + let cancel_token = CancellationToken::new(); + let mut run = managed_run( + MINIMAL_DOT.to_string(), + RunStatus::Running, + Utc::now(), + run_dir, + RunExecutionMode::Start, + ); + run.cancel_token = Some(cancel_token.clone()); + let (cancel_tx, _cancel_rx) = oneshot::channel(); + run.cancel_tx = Some(cancel_tx); + state + .runs + .lock() + .expect("runs lock poisoned") + .insert(run_id, run); + + delete_run_internal(&state, run_id, true).await.unwrap(); + + assert!(!cancel_token.is_cancelled()); +} + /// Append a stage lifecycle event with an explicit `StageScope`, so the /// stored envelope carries the full `stage_id` (`node_id@visit`). The bare /// [`workflow_event::append_event`] helper only writes `node_id` because