From e80c43007650215430d894f2d56e2cd4cbcfacc4 Mon Sep 17 00:00:00 2001 From: "zackary.l.jackson" Date: Sun, 4 Oct 2026 15:29:06 +0000 Subject: [PATCH] fix(work): seal cancellations that land before running A cancellation persisted before mark_running parked the attempt in cancellation_requested: the refused transition killed the provider and returned without settling. Seal it as cancelled through the normal acknowledge/settle ladder, and store a notify permit so an early registry signal is not lost. The bench cleanup chain now requests JSON from worktree_cleanup_*, and its resume-sweep workaround for parked cancels is removed. Refs #3053 Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- .../src/invocation/work_attempt_exec.rs | 85 +++++++++++- .../src/invocation/work_attempt_exec/tests.rs | 126 +++++++++++++++++- crates/tracedecay/benches/coverage/admin.rs | 7 +- crates/tracedecay/benches/coverage/mod.rs | 21 --- 4 files changed, 203 insertions(+), 36 deletions(-) diff --git a/crates/tracedecay-daemon-service/src/invocation/work_attempt_exec.rs b/crates/tracedecay-daemon-service/src/invocation/work_attempt_exec.rs index 854d2eedde..3001fd53c3 100644 --- a/crates/tracedecay-daemon-service/src/invocation/work_attempt_exec.rs +++ b/crates/tracedecay-daemon-service/src/invocation/work_attempt_exec.rs @@ -292,7 +292,7 @@ impl WorkAttemptProcessRegistryV1 { tokio::select! { () = &mut future => {} () = task_lifecycle.cancelled() => { - task_cancellation.notify_waiters(); + task_cancellation.notify_one(); let _ = tokio::time::timeout(CANCELLATION_GRACE, &mut future).await; } } @@ -403,7 +403,7 @@ impl WorkAttemptProcessRegistryV1 { .processes .get(&Self::key(worktree_id, identity)) { - notify.cancellation.notify_waiters(); + notify.cancellation.notify_one(); } } @@ -429,7 +429,7 @@ impl WorkAttemptProcessRegistryV1 { std::mem::take(&mut state.processes) }; for process in processes.values() { - process.cancellation.notify_waiters(); + process.cancellation.notify_one(); process.lifecycle.cancel(); } let deadline = @@ -469,7 +469,7 @@ impl Drop for WorkAttemptProcessRegistryV1 { .unwrap_or_else(std::sync::PoisonError::into_inner); state.accepting = false; for process in state.processes.values() { - process.cancellation.notify_waiters(); + process.cancellation.notify_one(); process.lifecycle.cancel(); if let Some(handle) = &process.handle { handle.abort(); @@ -801,6 +801,14 @@ fn settle_unstarted( ?problem, "work attempt could not be fenced for provider unavailability" ); + seal_cancellation_before_running( + attempts, + context, + attempt, + None, + provider_fallback, + observability_producer, + ); return; } let evidence = WorkAttemptEvidenceRecordV1 { @@ -828,6 +836,53 @@ fn settle_unstarted( } } +/// A cancellation requested before the provider was durably `Running` wins +/// the race: the row is `CancellationRequested`, which cannot become +/// `Running`, and no provider is left to exit. Seal it here as `Cancelled`. +/// Any other refusal leaves the durable row as the authority. +fn seal_cancellation_before_running( + attempts: &tracedecay_contracts::WorkAttemptService, + context: &RequestContext, + attempt: &WorkAttemptV1, + actual_route: Option, + provider_fallback: Option, + observability_producer: Option<&BoundedObservabilityProducerV1>, +) where + S: tracedecay_contracts::WorkAttemptStoragePort, +{ + let identity = attempt.identity(); + let observed_at = now_micros(); + if attempts + .acknowledge_cancellation(context, identity, observed_at) + .is_err() + { + return; + } + let evidence = WorkAttemptEvidenceRecordV1 { + identity: identity.clone(), + requested_route: attempt.requested_route().clone(), + actual_route, + outcome: WorkAttemptProviderOutcomeV1::Cancelled, + stdout: None, + stderr: None, + provider_session: None, + provider_fallback, + observed_at, + }; + match attempts.settle_with_artifacts(context, identity, &evidence, Vec::new()) { + Ok(settled) => { + let _ = record_terminal_attempt_product_views(observability_producer, &settled); + } + Err(problem) => { + tracing::warn!( + task = identity.task_id().as_str(), + ?problem, + "work attempt cancellation could not be sealed before running" + ); + } + } +} + enum EffectDispatchAdmission { Recorded, Replayed, @@ -1006,7 +1061,16 @@ async fn execute_provider_with_environment( ); terminate(&mut child, TerminationSignal::Kill); let _ = child.wait().await; - settle_effect_dispatch(attempt_effects, context, attempt, true); + if settle_effect_dispatch(attempt_effects, context, attempt, true) { + seal_cancellation_before_running( + attempts, + context, + attempt, + Some(selection.actual_route.clone()), + selection.fallback.clone(), + observability_producer, + ); + } return; } }; @@ -1188,7 +1252,16 @@ async fn execute_app_server( ?problem, "work attempt could not be marked running; app-server was not started" ); - settle_effect_dispatch(attempt_effects, context, attempt, false); + if settle_effect_dispatch(attempt_effects, context, attempt, false) { + seal_cancellation_before_running( + attempts, + context, + attempt, + None, + selection.fallback.clone(), + observability_producer, + ); + } return; } }; diff --git a/crates/tracedecay-daemon-service/src/invocation/work_attempt_exec/tests.rs b/crates/tracedecay-daemon-service/src/invocation/work_attempt_exec/tests.rs index a36588ad44..f343abd12b 100644 --- a/crates/tracedecay-daemon-service/src/invocation/work_attempt_exec/tests.rs +++ b/crates/tracedecay-daemon-service/src/invocation/work_attempt_exec/tests.rs @@ -1234,12 +1234,8 @@ async fn a_provider_that_ignores_interrupt_is_escalated_to_a_kill_on_the_record( }, ) .unwrap(); - // `notify_waiters` only wakes waiters already parked on the channel, - // so keep signalling until the execution arm observes it. - loop { - cancel.notify_waiters(); - tokio::time::sleep(std::time::Duration::from_millis(20)).await; - } + cancel.notify_one(); + std::future::pending::<()>().await; }; tokio::select! { @@ -1256,7 +1252,7 @@ async fn a_provider_that_ignores_interrupt_is_escalated_to_a_kill_on_the_record( None, AttemptAdmissionTimingV1::for_test(), ) => {} - _ = driver => unreachable!("the driver loops until execution settles"), + () = driver => unreachable!("the driver never settles"), } // The graceful rung really was delivered to the child, and really was @@ -1281,6 +1277,80 @@ async fn a_provider_that_ignores_interrupt_is_escalated_to_a_kill_on_the_record( ); } +/// A cancellation requested after the provider was launched but before the +/// row became `Running` must still end `Cancelled`: `CancellationRequested` +/// cannot become `Running`, and no later provider exit will seal it. +#[cfg(unix)] +#[tokio::test] +async fn a_cancellation_requested_before_mark_running_seals_cancelled() { + let (_artifact_dir, workflow_artifacts) = provider_artifact_store().await; + let directory = tempfile::TempDir::new().unwrap(); + let root = directory.path(); + let executable = fake_executable(root, "claude-code", "#!/bin/sh\nexit 0\n"); + let fixture = leased_attempt(root, "Cancelled before running.", &SnapshotShape::default()); + let identity = fixture.identity().clone(); + fixture + .attempts + .request_cancellation( + &fixture.context, + CancelWorkAttemptCommand { + task_id: identity.task_id().clone(), + run_id: identity.run_id().clone(), + attempt_id: identity.attempt_id().clone(), + request_id: id("cancellation.work-attempt-exec.before-running"), + occurred_at: now_micros(), + }, + ) + .unwrap(); + let admitted_environment = + admitted_provider_environment(fixture.attempt.execution().execution_snapshot()); + let route = requested_route(WorkProviderBackendV1::ClaudeCodeCli); + + execute_provider_with_environment( + &fixture.attempts, + &fixture.effects, + Some(&workflow_artifacts), + &fixture.context, + &fixture.attempt, + &preferred( + executable, + WorkProviderProtocol::ClaudeStreamJson, + &CLAUDE_STREAM_JSON_ARGV, + route.clone(), + ), + &admitted_environment, + Arc::new(Notify::new()), + None, + None, + AttemptAdmissionTimingV1::for_test(), + ) + .await; + + assert_eq!(fixture.state(), WorkAttemptStateV1::Cancelled); + assert_eq!( + fixture.rows.observed_states(), + vec![ + WorkAttemptStateV1::Leased, + WorkAttemptStateV1::CancellationRequested, + WorkAttemptStateV1::CancellationAcknowledged, + WorkAttemptStateV1::Cancelled, + ] + ); + let evidence = fixture.sealed_evidence(); + assert_eq!(evidence.outcome, WorkAttemptProviderOutcomeV1::Cancelled); + assert_eq!(evidence.requested_route, route); + assert_eq!(evidence.actual_route, Some(route)); + assert_eq!(evidence.stdout, None); + assert_eq!( + fixture + .effects + .load(&fixture.context, fixture.identity()) + .unwrap() + .and_then(|holder| holder.resolution()), + Some(WorkAttemptEffectResolutionV1::NoEffect) + ); +} + // --------------------------------------------------------------------------- // 3b. Wall exhaustion (no-progress terminal) // --------------------------------------------------------------------------- @@ -1936,6 +2006,48 @@ async fn the_process_registry_admits_one_live_owner_per_attempt() { ); } +/// The durable request is persisted before the live owner is signalled, so +/// the owner may not be waiting yet; the signal must still reach it. +#[tokio::test] +async fn a_cancellation_signalled_before_the_owner_waits_still_reaches_it() { + let directory = tempfile::TempDir::new().unwrap(); + let fixture = leased_attempt( + directory.path(), + "Early cancellation signal.", + &SnapshotShape::default(), + ); + let worktree_id: WorktreeId = id("worktree.registry.early-signal"); + let registry = Arc::new(WorkAttemptProcessRegistryV1::default()); + let signalled = Arc::new(Notify::new()); + let observed = Arc::new(std::sync::atomic::AtomicBool::new(false)); + let task_signalled = Arc::clone(&signalled); + let task_observed = Arc::clone(&observed); + + assert!(registry.spawn_for_worktree( + fixture.identity(), + &worktree_id, + move |cancellation| async move { + task_signalled.notified().await; + if tokio::time::timeout(std::time::Duration::from_secs(5), cancellation.notified()) + .await + .is_ok() + { + task_observed.store(true, std::sync::atomic::Ordering::SeqCst); + } + }, + )); + registry.signal_cancellation(&worktree_id, fixture.identity()); + signalled.notify_one(); + while registry.holds_attempt(&worktree_id, fixture.identity()) { + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + + assert!( + observed.load(std::sync::atomic::Ordering::SeqCst), + "a cancellation signalled before the owner waited was lost" + ); +} + #[tokio::test] async fn attempt_process_registry_shutdown_fences_and_joins_cooperative_owner() { let directory = tempfile::TempDir::new().unwrap(); diff --git a/crates/tracedecay/benches/coverage/admin.rs b/crates/tracedecay/benches/coverage/admin.rs index 004d11842a..88d02f0fc6 100644 --- a/crates/tracedecay/benches/coverage/admin.rs +++ b/crates/tracedecay/benches/coverage/admin.rs @@ -262,7 +262,7 @@ pub(crate) fn groups(ctx: &QueryContext, out: &mut Vec) { out.push(ToolGroup { tool: "tracedecay_worktree_cleanup_inspect", queries: five(|_i| { - eqn( + eq( "tracedecay_worktree_cleanup_inspect", "wt_cleanup_inspect", wt_claim(json!({})), @@ -281,7 +281,7 @@ pub(crate) fn groups(ctx: &QueryContext, out: &mut Vec) { out.push(ToolGroup { tool, queries: five(|_i| { - eqn( + eq( tool, label, wt_claim(json!({ @@ -579,6 +579,9 @@ fn wt_cleanup_primes(ctx: &QueryContext, _iter: u64) -> Vec { .unwrap_or_else(|| "worktree.bench.missing".into()), }, }); + let mut claim = claim; + // Application surfaces default to markdown; the prime chain parses JSON. + claim["format"] = json!("json"); let mut confirm_args = claim.clone(); confirm_args["inspection_digest"] = json!("{{wt_inspection_digest}}"); vec![ diff --git a/crates/tracedecay/benches/coverage/mod.rs b/crates/tracedecay/benches/coverage/mod.rs index 2b1c2a1d07..6649271a07 100644 --- a/crates/tracedecay/benches/coverage/mod.rs +++ b/crates/tracedecay/benches/coverage/mod.rs @@ -1847,27 +1847,6 @@ async fn settle_attempt( tokio::time::sleep(Duration::from_millis(100)).await; } } - // A cancellation observed while the provider is still spawning can leave - // the row parked in `cancellation_requested`: the orphaned mark_running - // transition already lost its fence, so only the recovery sweep seals the - // lost cancellation to `cancelled`. - if !terminal(&state) { - let _ = call( - "tracedecay_work_resume_attempts", - json!({"occurred_at": now_micros()}), - ) - .await; - for _ in 0..60 { - if let Ok(s) = call("tracedecay_work_attempt_status", status(attempt_id)).await - && let Some(st) = dig_str(&s, "state") - && terminal(st) - { - state = st.to_owned(); - break; - } - tokio::time::sleep(Duration::from_millis(100)).await; - } - } (started, identity, state) }