Skip to content

Commit 99488f8

Browse files
committed
fix(agent-core): drain archived background jobs before delete
Archive receipts could become quiesced after a job was marked killed even though its process tree, replay pipeline, or subagent task was still executing. Team Delete also left retained background-job registry state after the durable hierarchy was removed. Add exact Session indexes and execution-finality barriers, make Archive wait for runtime, memory, and background owners, and make Delete fail closed before commit and purge only after a successful transaction. Add deterministic race coverage, rendered lifecycle coverage, and read-only WebDriver evidence. Verification: - cargo test -p agent_core: 3230 passed, 2 ignored - cargo clippy -p agent_core --all-targets -- -D warnings: passed - cargo fmt --all -- --check: passed - pnpm test: 8777 passed - pnpm run lint: passed with 0 errors and 5 pre-existing warnings - pnpm run check:circular: blocked by the existing Madge raw-import resolver failure - WEBDRIVER=1 pnpm run tauri:build:fast: passed - BuildFast real-provider Archive x3, Pause/Resume, Task smoke, restart, and UI Delete: passed Pre-commit hook ran. Total eslint: 5, total circular: 0
1 parent 3b29828 commit 99488f8

13 files changed

Lines changed: 1435 additions & 365 deletions

File tree

src-tauri/crates/agent-core/src/core/coordination/agent_org_archive.rs

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -544,6 +544,18 @@ pub fn summary_for_run(run_id: &str) -> Result<Option<ArchiveTeardownSummary>, S
544544
summary_for_run_with_connection(&conn, run_id)
545545
}
546546

547+
/// Canonical Team Session scope for debug/WebDriver runtime evidence. Kept
548+
/// out of release builds so the production API surface remains unchanged.
549+
#[cfg(debug_assertions)]
550+
pub fn debug_owned_session_ids_for_run(run_id: &str) -> Result<Vec<String>, String> {
551+
let conn = database::db::get_connection().map_err(|error| error.to_string())?;
552+
Ok(load_team_for_run(&conn, run_id)?
553+
.sessions
554+
.into_iter()
555+
.map(|session| session.session_id)
556+
.collect())
557+
}
558+
547559
pub(crate) fn summary_for_run_with_connection(
548560
conn: &Connection,
549561
run_id: &str,

src-tauri/crates/agent-core/src/core/tools/impls/coding/exec/registry.rs

Lines changed: 101 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,14 @@ use tokio_util::sync::CancellationToken;
1717

1818
use crate::tools::call_context::{TurnProcessControl, TurnProcessOwner};
1919

20+
mod session_lifecycle;
21+
22+
pub use session_lifecycle::{
23+
execution_blockers_for_sessions, purge_deleted_sessions, request_cancel_for_session,
24+
retained_tombstone_count, session_runtime_evidence, wait_for_session_finality,
25+
PurgedSessionJobs, SessionJobEvidence,
26+
};
27+
2028
/// Status of a background job.
2129
#[derive(Debug, Clone)]
2230
pub enum JobStatus {
@@ -129,6 +137,10 @@ pub struct BackgroundJob {
129137
/// Cancellation was requested, but the monitor has not yet proved the
130138
/// process group and replay pipeline are terminal.
131139
shell_kill_requested: bool,
140+
/// Archive has installed its short, Session-scoped escalation for this
141+
/// subagent. Separate from `Killed`: status is user-visible, while this
142+
/// flag prevents repeated finality polls from spawning duplicate timers.
143+
session_cancel_escalation_requested: bool,
132144
/// Set to `true` once the agent has read the completed job's output via
133145
/// `AwaitTool` (monitor/wait_for). Acknowledged completed jobs are excluded
134146
/// from the per-turn system reminder to avoid the stale-reminder
@@ -369,6 +381,7 @@ const TOMBSTONE_TTL: std::time::Duration = std::time::Duration::from_secs(10 * 6
369381
/// mistyped it), instead of synthesising a guess from the handle string.
370382
#[derive(Clone)]
371383
struct Tombstone {
384+
session_id: String,
372385
status: JobStatus,
373386
kind: JobKind,
374387
created_at: Instant,
@@ -475,6 +488,7 @@ fn register_shell_inner(registration: ShellRegistration) -> broadcast::Sender<St
475488
let handle = pid.to_string();
476489
let (tx, _) = broadcast::channel(BROADCAST_CAPACITY);
477490
let sender = tx.clone();
491+
let indexed_session_id = session_id.clone();
478492
let job = BackgroundJob {
479493
handle: handle.clone(),
480494
label: command,
@@ -498,22 +512,32 @@ fn register_shell_inner(registration: ShellRegistration) -> broadcast::Sender<St
498512
shell_cancel,
499513
shell_completion,
500514
shell_kill_requested: false,
515+
session_cancel_escalation_requested: false,
501516
output_acknowledged: false,
502517
wake_dispatched: false,
503518
output_seq: 0,
504519
stalled_waiting_input: false,
505520
stall_delivered: false,
506521
};
507522
let mut reg = REGISTRY.lock().unwrap_or_else(|e| e.into_inner());
508-
reg.insert(handle.clone(), job);
523+
let replaced = reg.insert(handle.clone(), job);
524+
let mut owner_index = OWNER_INDEX.lock().unwrap_or_else(|e| e.into_inner());
525+
if let Some(previous_owner) = replaced
526+
.as_ref()
527+
.and_then(|previous| previous.turn_owner.as_ref())
528+
{
529+
remove_indexed_handle(&mut owner_index, previous_owner, &handle);
530+
}
509531
if let Some(owner) = indexed_owner {
510-
OWNER_INDEX
511-
.lock()
512-
.unwrap_or_else(|e| e.into_inner())
513-
.entry(owner)
514-
.or_default()
515-
.insert(handle);
532+
owner_index.entry(owner).or_default().insert(handle.clone());
516533
}
534+
session_lifecycle::replace_live_index(
535+
replaced
536+
.as_ref()
537+
.map(|previous| previous.session_id.as_str()),
538+
&indexed_session_id,
539+
&handle,
540+
);
517541
sender
518542
}
519543

@@ -622,22 +646,32 @@ fn register_subagent_inner(
622646
shell_cancel: None,
623647
shell_completion: None,
624648
shell_kill_requested: false,
649+
session_cancel_escalation_requested: false,
625650
output_acknowledged: false,
626651
wake_dispatched: false,
627652
output_seq: 0,
628653
stalled_waiting_input: false,
629654
stall_delivered: false,
630655
};
631656
let mut reg = REGISTRY.lock().unwrap_or_else(|e| e.into_inner());
632-
reg.insert(handle.clone(), job);
657+
let replaced = reg.insert(handle.clone(), job);
658+
let mut owner_index = OWNER_INDEX.lock().unwrap_or_else(|e| e.into_inner());
659+
if let Some(previous_owner) = replaced
660+
.as_ref()
661+
.and_then(|previous| previous.turn_owner.as_ref())
662+
{
663+
remove_indexed_handle(&mut owner_index, previous_owner, &handle);
664+
}
633665
if let Some(owner) = indexed_owner {
634-
OWNER_INDEX
635-
.lock()
636-
.unwrap_or_else(|e| e.into_inner())
637-
.entry(owner)
638-
.or_default()
639-
.insert(handle.clone());
666+
owner_index.entry(owner).or_default().insert(handle.clone());
640667
}
668+
session_lifecycle::replace_live_index(
669+
replaced
670+
.as_ref()
671+
.map(|previous| previous.session_id.as_str()),
672+
&session_id,
673+
&handle,
674+
);
641675
drop(reg);
642676
broadcast_subagent_job_changed(&session_id, &handle, &agent_name, &subagent_type, "running");
643677
sender
@@ -783,29 +817,56 @@ pub fn remove(handle: &str) {
783817
let removed = {
784818
let mut reg = REGISTRY.lock().unwrap_or_else(|e| e.into_inner());
785819
let removed = reg.remove(handle);
820+
let mut index = OWNER_INDEX.lock().unwrap_or_else(|e| e.into_inner());
786821
if let Some(owner) = removed.as_ref().and_then(|job| job.turn_owner.as_ref()) {
787-
let mut index = OWNER_INDEX.lock().unwrap_or_else(|e| e.into_inner());
788-
if let Some(handles) = index.get_mut(owner) {
789-
handles.remove(handle);
790-
if handles.is_empty() {
791-
index.remove(owner);
792-
}
793-
}
822+
remove_indexed_handle(&mut index, owner, handle);
823+
}
824+
if let Some(job) = removed.as_ref() {
825+
session_lifecycle::remove_live_index(&job.session_id, handle);
794826
}
795827
removed
796828
};
797829
if let Some(job) = removed {
798830
let mut tombs = TOMBSTONES.lock().unwrap_or_else(|e| e.into_inner());
799831
let now = Instant::now();
800-
tombs.retain(|_, t| now.duration_since(t.created_at) < TOMBSTONE_TTL);
801-
tombs.insert(
832+
let expired = tombs
833+
.iter()
834+
.filter(|(_, tombstone)| now.duration_since(tombstone.created_at) >= TOMBSTONE_TTL)
835+
.map(|(expired_handle, tombstone)| {
836+
(tombstone.session_id.clone(), expired_handle.clone())
837+
})
838+
.collect::<Vec<_>>();
839+
tombs.retain(|_, tombstone| now.duration_since(tombstone.created_at) < TOMBSTONE_TTL);
840+
session_lifecycle::remove_expired_tombstone_indexes(&expired);
841+
let replaced = tombs.insert(
802842
handle.to_string(),
803843
Tombstone {
844+
session_id: job.session_id.clone(),
804845
status: job.status.clone(),
805846
kind: job.kind.clone(),
806847
created_at: now,
807848
},
808849
);
850+
session_lifecycle::replace_tombstone_index(
851+
replaced
852+
.as_ref()
853+
.map(|previous| previous.session_id.as_str()),
854+
&job.session_id,
855+
handle,
856+
);
857+
}
858+
}
859+
860+
fn remove_indexed_handle(
861+
index: &mut HashMap<TurnProcessOwner, HashSet<String>>,
862+
owner: &TurnProcessOwner,
863+
handle: &str,
864+
) {
865+
if let Some(handles) = index.get_mut(owner) {
866+
handles.remove(handle);
867+
if handles.is_empty() {
868+
index.remove(owner);
869+
}
809870
}
810871
}
811872

@@ -833,14 +894,15 @@ pub fn resolve_status_with_tombstone(handle: &str) -> Option<(JobStatus, JobKind
833894
if let Some(found) = get_status(handle) {
834895
return Some(found);
835896
}
836-
let tombs = TOMBSTONES.lock().unwrap_or_else(|e| e.into_inner());
837-
tombs.get(handle).and_then(|t| {
838-
if Instant::now().duration_since(t.created_at) < TOMBSTONE_TTL {
839-
Some((t.status.clone(), t.kind.clone()))
840-
} else {
841-
None
842-
}
843-
})
897+
let mut tombs = TOMBSTONES.lock().unwrap_or_else(|e| e.into_inner());
898+
let tombstone = tombs.get(handle).cloned()?;
899+
if Instant::now().duration_since(tombstone.created_at) < TOMBSTONE_TTL {
900+
Some((tombstone.status, tombstone.kind))
901+
} else {
902+
tombs.remove(handle);
903+
session_lifecycle::remove_tombstone_index(&tombstone.session_id, handle);
904+
None
905+
}
844906
}
845907

846908
/// Get the final result text for a job.
@@ -939,13 +1001,14 @@ pub async fn await_shells_terminated_for_owner(
9391001
/// scope, `None` for global scope.
9401002
pub fn list_jobs(session_id: Option<&str>) -> Vec<JobSnapshot> {
9411003
let reg = REGISTRY.lock().unwrap_or_else(|e| e.into_inner());
942-
reg.values()
943-
.filter(|job| match session_id {
944-
Some(sid) => job.session_id == sid,
945-
None => true,
946-
})
947-
.map(|job| job.snapshot())
948-
.collect()
1004+
match session_id {
1005+
Some(session_id) => session_lifecycle::live_handles(session_id)
1006+
.into_iter()
1007+
.filter_map(|handle| reg.get(&handle))
1008+
.map(BackgroundJob::snapshot)
1009+
.collect(),
1010+
None => reg.values().map(BackgroundJob::snapshot).collect(),
1011+
}
9491012
}
9501013

9511014
/// Mark a completed job's output as acknowledged. Once acknowledged, the job

0 commit comments

Comments
 (0)