mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-09-06 08:18:58 +00:00
Merge remote-tracking branch 'origin/main' into add-fabro-mcp-server
This commit is contained in:
commit
e0be041c1e
2 changed files with 109 additions and 17 deletions
|
|
@ -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<RunStatus> {
|
||||
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<AppState>,
|
||||
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,
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue