Skip to content

Commit f37c7fe

Browse files
authored
Merge pull request #201 from cbusillo/fix/agent-archive-lifecycle-reset
Fix agent archive lifecycle reset
2 parents 1f13f5f + 8e70d1a commit f37c7fe

1 file changed

Lines changed: 49 additions & 12 deletions

File tree

code-rs/core/src/agent_tool.rs

Lines changed: 49 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -406,7 +406,6 @@ pub struct AgentManager {
406406
archived_terminal_agents: HashMap<String, Agent>,
407407
handles: HashMap<String, JoinHandle<()>>,
408408
event_senders: Vec<AgentStatusSender>,
409-
event_sender_lifecycle_started: bool,
410409
debug_log_root: Option<PathBuf>,
411410
watchdog_handle: Option<JoinHandle<()>>,
412411
inactivity_timeout: Duration,
@@ -521,7 +520,6 @@ impl AgentManager {
521520
archived_terminal_agents: HashMap::new(),
522521
handles: HashMap::new(),
523522
event_senders: Vec::new(),
524-
event_sender_lifecycle_started: false,
525523
debug_log_root: None,
526524
watchdog_handle: None,
527525
inactivity_timeout: Duration::minutes(30),
@@ -536,18 +534,20 @@ impl AgentManager {
536534
) {
537535
self.event_senders
538536
.retain(|registered| !registered.sender.is_closed());
539-
let first_session = !self.event_sender_lifecycle_started;
540-
self.event_sender_lifecycle_started = true;
537+
let fresh_sender_set = self.event_senders.is_empty();
541538
self.event_senders.push(AgentStatusSender {
542539
owner_session_id,
543540
sender,
544541
});
545-
// New manager lifecycle: keep only live agents and reset diagnostics
546-
// when the first UI connects, without clearing another live session's
547-
// archived terminal agents when a second UI registers.
548-
if first_session {
549-
self.archived_terminal_agents.clear();
542+
// A fresh sender set is either the first UI for this manager or a UI
543+
// reconnect after all previous senders closed. Keep archived results
544+
// for the reconnecting session, but drop archives from older sessions
545+
// so a long-lived manager does not retain stale terminal agents forever.
546+
if fresh_sender_set {
547+
self.archived_terminal_agents
548+
.retain(|_, agent| agent_belongs_to_session(agent, Some(owner_session_id)));
550549
self.diagnostics = AgentManagerDiagnostics::default();
550+
self.diagnostics.archived_terminal_agents = self.archived_terminal_agents.len() as u64;
551551
}
552552
self.start_watchdog();
553553
}
@@ -3341,27 +3341,64 @@ mod tests {
33413341
#[tokio::test]
33423342
async fn reconnect_after_sender_gap_keeps_existing_archived_agents() {
33433343
let mut manager = AgentManager::new();
3344+
let session_id = Uuid::new_v4();
33443345
let (tx_a, rx_a) = tokio::sync::mpsc::unbounded_channel();
3345-
manager.set_event_sender(Uuid::new_v4(), tx_a);
3346+
manager.set_event_sender(session_id, tx_a);
33463347
drop(rx_a);
33473348

33483349
let archived_id = "archived-after-gap".to_string();
33493350
manager.archived_terminal_agents.insert(
33503351
archived_id.clone(),
33513352
test_agent(
33523353
&archived_id,
3353-
Uuid::new_v4(),
3354+
session_id,
33543355
"batch-archived-gap",
33553356
AgentStatus::Completed,
33563357
),
33573358
);
33583359

33593360
let (tx_b, _rx_b) = tokio::sync::mpsc::unbounded_channel();
3360-
manager.set_event_sender(Uuid::new_v4(), tx_b);
3361+
manager.set_event_sender(session_id, tx_b);
33613362

33623363
assert!(manager.archived_terminal_agents.contains_key(&archived_id));
33633364
}
33643365

3366+
#[tokio::test]
3367+
async fn fresh_sender_set_drops_archives_from_other_sessions() {
3368+
let mut manager = AgentManager::new();
3369+
let old_session = Uuid::new_v4();
3370+
let new_session = Uuid::new_v4();
3371+
let (tx_old, rx_old) = tokio::sync::mpsc::unbounded_channel();
3372+
manager.set_event_sender(old_session, tx_old);
3373+
drop(rx_old);
3374+
3375+
manager.archived_terminal_agents.insert(
3376+
"old-archived".to_string(),
3377+
test_agent(
3378+
"old-archived",
3379+
old_session,
3380+
"batch-old-archived",
3381+
AgentStatus::Completed,
3382+
),
3383+
);
3384+
manager.archived_terminal_agents.insert(
3385+
"new-archived".to_string(),
3386+
test_agent(
3387+
"new-archived",
3388+
new_session,
3389+
"batch-new-archived",
3390+
AgentStatus::Completed,
3391+
),
3392+
);
3393+
3394+
let (tx_new, _rx_new) = tokio::sync::mpsc::unbounded_channel();
3395+
manager.set_event_sender(new_session, tx_new);
3396+
3397+
assert!(manager.archived_terminal_agents.contains_key("new-archived"));
3398+
assert!(!manager.archived_terminal_agents.contains_key("old-archived"));
3399+
assert_eq!(manager.diagnostics.archived_terminal_agents, 1);
3400+
}
3401+
33653402
#[tokio::test]
33663403
async fn read_only_agents_use_code_binary_path() {
33673404
let _lock = env_lock().lock().expect("env lock");

0 commit comments

Comments
 (0)