Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
}
Expand Down Expand Up @@ -403,7 +403,7 @@ impl WorkAttemptProcessRegistryV1 {
.processes
.get(&Self::key(worktree_id, identity))
{
notify.cancellation.notify_waiters();
notify.cancellation.notify_one();
}
}

Expand All @@ -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 =
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -801,6 +801,14 @@ fn settle_unstarted<S>(
?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 {
Expand Down Expand Up @@ -828,6 +836,53 @@ fn settle_unstarted<S>(
}
}

/// 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<S>(
attempts: &tracedecay_contracts::WorkAttemptService<S>,
context: &RequestContext,
attempt: &WorkAttemptV1,
actual_route: Option<WorkProviderRouteV1>,
provider_fallback: Option<WorkProviderFallbackRecordV1>,
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,
Expand Down Expand Up @@ -1006,7 +1061,16 @@ async fn execute_provider_with_environment<S>(
);
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;
}
};
Expand Down Expand Up @@ -1188,7 +1252,16 @@ async fn execute_app_server<S>(
?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;
}
};
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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! {
Expand All @@ -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
Expand All @@ -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)
// ---------------------------------------------------------------------------
Expand Down Expand Up @@ -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();
Expand Down
7 changes: 5 additions & 2 deletions crates/tracedecay/benches/coverage/admin.rs
Original file line number Diff line number Diff line change
Expand Up @@ -262,7 +262,7 @@ pub(crate) fn groups(ctx: &QueryContext, out: &mut Vec<ToolGroup>) {
out.push(ToolGroup {
tool: "tracedecay_worktree_cleanup_inspect",
queries: five(|_i| {
eqn(
eq(
"tracedecay_worktree_cleanup_inspect",
"wt_cleanup_inspect",
wt_claim(json!({})),
Expand All @@ -281,7 +281,7 @@ pub(crate) fn groups(ctx: &QueryContext, out: &mut Vec<ToolGroup>) {
out.push(ToolGroup {
tool,
queries: five(|_i| {
eqn(
eq(
tool,
label,
wt_claim(json!({
Expand Down Expand Up @@ -579,6 +579,9 @@ fn wt_cleanup_primes(ctx: &QueryContext, _iter: u64) -> Vec<PrimeStep> {
.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![
Expand Down
21 changes: 0 additions & 21 deletions crates/tracedecay/benches/coverage/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}

Expand Down
Loading