1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
//! Global registry of running agents with cancellation support.
use std::collections::HashMap;
use std::sync::LazyLock;
use std::sync::Mutex;
use std::sync::atomic::AtomicU64;
use crate::util::UnwrapPoison;
use chrono::{DateTime, Utc};
use serde::Serialize;
use tokio_util::sync::CancellationToken;
/// Monotonically increasing generation counter for registry entries.
/// Used by [`deregister`](AgentRegistry::deregister) to detect stale entries
/// — when a new agent is registered with the same `run_id` (e.g. the Manager
/// interrupt-and-resume pattern), the old entry's generation will not match
/// the new entry, so `deregister` will not incorrectly remove the replacement.
static NEXT_ENTRY_GENERATION: AtomicU64 = AtomicU64::new(1);
/// Public handle returned by `list()` — serializable, no cancel_token exposed.
#[derive(Clone, Debug, Serialize)]
pub struct AgentHandle {
pub run_id: String,
pub role: String,
pub ticket_id: Option<String>,
/// Filesystem path of the workspace (not the name) — this is used for
/// agent display/location and is intentionally distinct from the
/// workspace_name identifier used in the board database.
pub workspace_path: String,
pub started_at: DateTime<Utc>,
pub label: String,
}
struct AgentEntry {
generation: u64,
handle: AgentHandle,
cancel_token: CancellationToken,
}
#[derive(Default)]
pub struct AgentRegistry {
inner: Mutex<HashMap<String, AgentEntry>>,
}
impl AgentRegistry {
/// Register an agent entry and return the generation counter.
///
/// Used by [`crate::Agent::new`] where deregistration is handled by [`crate::Agent::drop`]
/// instead of a guard.
pub fn register_unguarded(
&self,
run_id: String,
role: String,
ticket_id: Option<String>,
ws: &crate::Workspace,
label: String,
cancel_token: CancellationToken,
) -> u64 {
let generation = NEXT_ENTRY_GENERATION.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let handle = AgentHandle {
run_id: run_id.clone(),
role,
ticket_id,
workspace_path: ws.path.clone(),
started_at: Utc::now(),
label,
};
let mut map = self.inner.lock().unwrap_poison();
if let Some(old) = map.remove(&run_id) {
old.cancel_token.cancel();
}
map.insert(
run_id,
AgentEntry {
generation,
handle,
cancel_token,
},
);
generation
}
/// Cancel a specific agent by run_id. Removes it from the registry.
pub fn cancel(&self, run_id: &str) {
let mut map = self.inner.lock().unwrap_poison();
if let Some(entry) = map.remove(run_id) {
entry.cancel_token.cancel();
}
}
/// Cancel all agents running for a specific `ticket_id`.
/// Used on ticket status transitions — stops any agent currently working on it.
pub fn cancel_by_ticket_id(&self, ticket_id: &str) {
let to_cancel: Vec<String> = {
let map = self.inner.lock().unwrap_poison();
map.iter()
.filter(|(_, entry)| entry.handle.ticket_id.as_deref() == Some(ticket_id))
.map(|(id, _)| id.clone())
.collect()
};
for run_id in to_cancel {
self.cancel(&run_id);
}
}
/// Snapshot of all currently running agents (serializable).
#[must_use]
pub fn list(&self) -> Vec<AgentHandle> {
self.inner
.lock()
.unwrap_poison()
.values()
.map(|e| e.handle.clone())
.collect()
}
/// Cancel all running agents. Used during daemon shutdown.
pub fn shutdown_all(&self) {
let entries: Vec<(String, CancellationToken)> = self
.inner
.lock()
.unwrap_poison()
.drain()
.map(|(id, entry)| (id, entry.cancel_token))
.collect();
for (_id, token) in entries {
token.cancel();
}
}
/// Remove a registry entry only if its generation still matches.
/// Used by [`crate::Agent::drop`] to safely deregister without stale-removal risk.
pub fn deregister(&self, run_id: &str, generation: u64) {
let mut map = self.inner.lock().unwrap_poison();
if let Some(entry) = map.get(run_id)
&& entry.generation == generation
{
map.remove(run_id);
}
}
}
/// Global static registry.
pub static AGENT_REGISTRY: LazyLock<AgentRegistry> = LazyLock::new(AgentRegistry::default);