1use std::collections::{BTreeMap, BTreeSet, HashMap};
8use std::path::PathBuf;
9use std::time::{SystemTime, UNIX_EPOCH};
10
11use serde::{Deserialize, Serialize};
12
13use crate::claude_peer::{process_is_live, ClaudePeerSession, ClaudePeerStatus};
14#[cfg(feature = "adapter-api")]
15use crate::codex_peer::CodexPeerTracker;
16use crate::codex_peer::{live_rollouts, rollout_lineage, rollout_status, CodexPeerStatus};
17use crate::{HarnessHomes, HarnessId, SessionLocator};
18
19#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
21#[serde(rename_all = "snake_case")]
22pub enum SessionPresence {
23 Persisted,
25 Running,
27 ShuttingDown,
29}
30
31#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
33#[serde(rename_all = "snake_case")]
34pub enum SessionTurnState {
35 Unknown,
37 Idle,
39 Working,
41 NeedsInput,
43}
44
45#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
47pub struct SessionActivityEvidence {
48 pub source: String,
50 pub native_state: Option<String>,
52 pub observed_at_ms: u64,
54 pub harness_version: Option<String>,
56}
57
58#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
60pub struct SessionActivity {
61 pub harness: HarnessId,
63 pub session_id: String,
65 pub presence: SessionPresence,
67 pub turn: SessionTurnState,
69 pub evidence: SessionActivityEvidence,
71}
72
73impl SessionActivity {
74 pub(crate) fn same_state(&self, other: &Self) -> bool {
76 self.harness == other.harness
77 && self.session_id == other.session_id
78 && self.presence == other.presence
79 && self.turn == other.turn
80 && self.evidence.source == other.evidence.source
81 && self.evidence.native_state == other.evidence.native_state
82 && self.evidence.harness_version == other.evidence.harness_version
83 }
84
85 pub(crate) fn key(&self) -> (String, String) {
87 (self.harness.as_str().to_string(), self.session_id.clone())
88 }
89}
90
91#[cfg(feature = "adapter-api")]
98pub(crate) async fn activity_under(
99 roots: &[u32],
100 homes: &HarnessHomes,
101) -> Result<Vec<(u32, SessionActivity)>, crate::SdkError> {
102 let table = process_table();
103 let mut children: HashMap<u32, Vec<u32>> = HashMap::new();
104 for (&pid, (parent, _)) in &table {
105 children.entry(*parent).or_default().push(pid);
106 }
107 let registry = crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(homes));
108 let mut found: Vec<(u32, SessionLocator)> = Vec::new();
109 for &root in roots {
110 let mut queue = std::collections::VecDeque::from([root]);
111 let mut visited = 0;
112 while let Some(pid) = queue.pop_front() {
113 visited += 1;
114 if visited > 256 {
115 break;
116 }
117 if let Some(session) = registry.iter().find(|session| session.pid == pid) {
118 found.push((
119 root,
120 SessionLocator {
121 harness: HarnessId::new(HarnessId::CLAUDE_CODE),
122 session_id: session.session_id.clone(),
123 storage: crate::StorageLocator::File {
124 path: PathBuf::new(),
125 },
126 },
127 ));
128 break;
129 }
130 let is_codex = table
131 .get(&pid)
132 .is_some_and(|(_, command)| command.rsplit('/').next() == Some("codex"));
133 if is_codex {
134 let resolved = crate::codex_peer::session_of_process(pid).or_else(|| {
136 let session_id = crate::codex_peer::resumed_session_of_process(pid)?;
137 let path = crate::codex_peer::rollout_of_session(&homes.codex, &session_id)?;
138 Some((session_id, path))
139 });
140 if let Some((session_id, path)) = resolved {
141 found.push((
142 root,
143 SessionLocator {
144 harness: HarnessId::new(HarnessId::CODEX),
145 session_id,
146 storage: crate::StorageLocator::File { path },
147 },
148 ));
149 break;
150 }
151 }
152 queue.extend(children.get(&pid).into_iter().flatten().copied());
153 }
154 }
155 if found.is_empty() {
156 return Ok(Vec::new());
157 }
158 let locators: Vec<SessionLocator> = found.iter().map(|(_, locator)| locator.clone()).collect();
159 let activities = SessionActivityMonitor::default()
160 .resolve(&locators, homes)
161 .await?;
162 Ok(found
163 .into_iter()
164 .filter_map(|(root, locator)| {
165 activities
166 .iter()
167 .find(|activity| {
168 activity.harness == locator.harness && activity.session_id == locator.session_id
169 })
170 .cloned()
171 .map(|activity| (root, activity))
172 })
173 .collect())
174}
175
176#[cfg(feature = "adapter-api")]
178fn process_table() -> HashMap<u32, (u32, String)> {
179 let Ok(output) = std::process::Command::new("ps")
180 .args(["-axo", "pid=,ppid=,comm="])
181 .output()
182 else {
183 return HashMap::new();
184 };
185 String::from_utf8_lossy(&output.stdout)
186 .lines()
187 .filter_map(|line| {
188 let mut fields = line.split_whitespace();
189 let pid = fields.next()?.parse().ok()?;
190 let parent = fields.next()?.parse().ok()?;
191 let command = fields.collect::<Vec<_>>().join(" ");
192 Some((pid, (parent, command)))
193 })
194 .collect()
195}
196
197#[cfg(feature = "adapter-api")]
200#[derive(Debug, Default)]
201pub(crate) struct SessionActivityMonitor {
202 codex: CodexPeerTracker,
203}
204
205#[cfg(feature = "adapter-api")]
206impl SessionActivityMonitor {
207 pub(crate) async fn resolve(
208 &mut self,
209 locators: &[SessionLocator],
210 homes: &HarnessHomes,
211 ) -> Result<Vec<SessionActivity>, crate::SdkError> {
212 let authorization = crate::RuntimeAuthorization::observer();
213 let sources = locators
214 .iter()
215 .map(|locator| (locator.harness.0.clone(), locator.session_id.clone()))
216 .collect::<BTreeSet<_>>();
217 let owned = crate::LocalRuntimeRegistry::new()
218 .source_states(&sources, &authorization)
219 .await?;
220 let claude = read_claude(locators, homes);
221 let codex = if locators
222 .iter()
223 .any(|locator| locator.harness.as_str() == HarnessId::CODEX)
224 {
225 self.codex.sample(&homes.codex)
226 } else {
227 HashMap::new()
228 };
229 let hermes = read_hermes(locators, homes);
230 let mut activities = resolve_stock_with_evidence(locators, &claude, &codex, &hermes)
231 .into_iter()
232 .map(|activity| (activity.key(), activity))
233 .collect::<BTreeMap<_, _>>();
234 let observed_at_ms = now_ms();
235 for locator in locators {
236 let key = (
237 locator.harness.as_str().to_string(),
238 locator.session_id.clone(),
239 );
240 if let Some(state) = owned.get(&key).copied() {
241 activities.insert(key, owned_activity(locator, state, observed_at_ms));
242 }
243 }
244 Ok(locators
245 .iter()
246 .filter_map(|locator| {
247 activities.remove(&(
248 locator.harness.as_str().to_string(),
249 locator.session_id.clone(),
250 ))
251 })
252 .collect())
253 }
254}
255
256pub(crate) fn resolve_stock_session_activities(
259 locators: &[SessionLocator],
260 homes: &HarnessHomes,
261) -> Vec<SessionActivity> {
262 let claude = read_claude(locators, homes);
263 let codex = if locators
264 .iter()
265 .any(|locator| locator.harness.as_str() == HarnessId::CODEX)
266 {
267 live_rollouts(&homes.codex)
268 } else {
269 HashMap::new()
270 };
271 let hermes = read_hermes(locators, homes);
272 resolve_stock_with_evidence(locators, &claude, &codex, &hermes)
273}
274
275#[derive(Debug, Clone, Copy, PartialEq, Eq)]
279pub(crate) struct HermesFire {
280 status: supercode_interchange::orchestration::FireStatus,
281 pid: Option<u32>,
282}
283
284fn read_hermes(locators: &[SessionLocator], homes: &HarnessHomes) -> HashMap<String, HermesFire> {
289 if !locators
290 .iter()
291 .any(|locator| locator.harness.as_str() == HarnessId::HERMES)
292 {
293 return HashMap::new();
294 }
295 let home = homes
297 .hermes
298 .parent()
299 .map_or_else(|| PathBuf::from("."), std::path::Path::to_path_buf);
300 let Ok(loaded) = supercode_interchange::orchestration::codec::from_hermes(&home) else {
301 return HashMap::new();
302 };
303 let mut fires = HashMap::new();
304 for profile in loaded.orchestration.profiles.values() {
305 for fire in &profile.fires {
306 let Some(session) = fire.session_id.clone() else {
307 continue;
308 };
309 let pid = fire.residue.0.get("pid").and_then(|value| {
310 value
311 .as_u64()
312 .or_else(|| value.as_str().and_then(|text| text.trim().parse().ok()))
313 .and_then(|pid| u32::try_from(pid).ok())
314 });
315 fires.insert(
316 session,
317 HermesFire {
318 status: fire.status,
319 pid,
320 },
321 );
322 }
323 }
324 fires
325}
326
327fn hermes_activity(
328 locator: &SessionLocator,
329 fire: HermesFire,
330 live: impl Fn(u32) -> bool,
331 observed_at_ms: u64,
332) -> SessionActivity {
333 use supercode_interchange::orchestration::FireStatus;
334 let (presence, turn, native_state) = match fire.status {
335 FireStatus::Claimed | FireStatus::Running => match fire.pid {
336 Some(pid) if live(pid) => (
337 SessionPresence::Running,
338 SessionTurnState::Working,
339 fire.status.hermes_word(),
340 ),
341 _ => (
344 SessionPresence::Persisted,
345 SessionTurnState::Unknown,
346 "abandoned",
347 ),
348 },
349 status => (
350 SessionPresence::Persisted,
351 SessionTurnState::Unknown,
352 status.hermes_word(),
353 ),
354 };
355 activity(
356 locator,
357 presence,
358 turn,
359 "hermes_executions",
360 Some(native_state),
361 None,
362 observed_at_ms,
363 )
364}
365
366fn read_claude(
367 locators: &[SessionLocator],
368 homes: &HarnessHomes,
369) -> HashMap<String, ClaudePeerSession> {
370 if locators
371 .iter()
372 .any(|locator| locator.harness.as_str() == HarnessId::CLAUDE_CODE)
373 {
374 crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(homes))
375 .into_iter()
376 .map(|peer| (peer.session_id.clone(), peer))
377 .collect::<HashMap<_, _>>()
378 } else {
379 HashMap::new()
380 }
381}
382
383fn resolve_stock_with_evidence(
384 locators: &[SessionLocator],
385 claude: &HashMap<String, ClaudePeerSession>,
386 codex: &HashMap<PathBuf, CodexPeerStatus>,
387 hermes: &HashMap<String, HermesFire>,
388) -> Vec<SessionActivity> {
389 let observed_at_ms = now_ms();
390 let codex_descendants = codex_descendant_statuses(codex);
391
392 locators
393 .iter()
394 .map(|locator| {
395 if locator.harness.as_str() == HarnessId::CLAUDE_CODE {
396 if let Some(peer) = claude.get(&locator.session_id) {
397 return claude_activity(locator, peer, observed_at_ms);
398 }
399 }
400 if locator.harness.as_str() == HarnessId::HERMES {
401 if let Some(fire) = hermes.get(&locator.session_id) {
402 return hermes_activity(locator, *fire, process_is_live, observed_at_ms);
403 }
404 }
405 if locator.harness.as_str() == HarnessId::CODEX {
406 if let Some(status) = merge_codex_status(
407 rollout_status(codex, locator.storage.path()),
408 codex_descendants.get(&locator.session_id).copied(),
409 ) {
410 return codex_activity(locator, status, observed_at_ms);
411 }
412 }
413 persisted_activity(locator, observed_at_ms)
414 })
415 .collect()
416}
417
418fn codex_descendant_statuses(
422 live: &HashMap<PathBuf, CodexPeerStatus>,
423) -> HashMap<String, CodexPeerStatus> {
424 let lineage = live
425 .iter()
426 .filter_map(|(path, status)| {
427 rollout_lineage(path).map(|(session_id, parent)| (session_id, parent, *status))
428 })
429 .collect::<Vec<_>>();
430 let parent_by_child = lineage
431 .iter()
432 .filter_map(|(child, parent, _)| {
433 parent
434 .as_ref()
435 .map(|parent| (child.clone(), parent.clone()))
436 })
437 .collect::<HashMap<_, _>>();
438 let mut statuses = HashMap::new();
439 for (_, parent, status) in lineage {
440 let Some(mut root) = parent else {
441 continue;
442 };
443 let mut visited = std::collections::HashSet::new();
444 while visited.insert(root.clone()) {
445 let Some(parent) = parent_by_child.get(&root) else {
446 break;
447 };
448 root = parent.clone();
449 }
450 statuses
451 .entry(root)
452 .and_modify(|current| *current = stronger_codex_status(*current, status))
453 .or_insert(status);
454 }
455 statuses
456}
457
458fn merge_codex_status(
459 own: Option<CodexPeerStatus>,
460 descendant: Option<CodexPeerStatus>,
461) -> Option<CodexPeerStatus> {
462 match (own, descendant) {
463 (Some(own), Some(descendant)) => Some(stronger_codex_status(own, descendant)),
464 (Some(status), None) | (None, Some(status)) => Some(status),
465 (None, None) => None,
466 }
467}
468
469fn stronger_codex_status(left: CodexPeerStatus, right: CodexPeerStatus) -> CodexPeerStatus {
470 fn rank(status: CodexPeerStatus) -> u8 {
471 match status {
472 CodexPeerStatus::Running => 0,
473 CodexPeerStatus::Idle => 1,
474 CodexPeerStatus::Busy => 2,
475 CodexPeerStatus::Waiting => 3,
476 }
477 }
478 if rank(right) > rank(left) {
479 right
480 } else {
481 left
482 }
483}
484
485#[cfg(feature = "adapter-api")]
486fn owned_activity(
487 locator: &SessionLocator,
488 state: crate::RuntimeRegistryState,
489 observed_at_ms: u64,
490) -> SessionActivity {
491 use crate::RuntimeRegistryState;
492 let (presence, turn) = match state {
493 RuntimeRegistryState::Persisted => (SessionPresence::Persisted, SessionTurnState::Unknown),
494 RuntimeRegistryState::Idle => (SessionPresence::Running, SessionTurnState::Idle),
495 RuntimeRegistryState::Busy => (SessionPresence::Running, SessionTurnState::Working),
496 RuntimeRegistryState::ShuttingDown => {
497 (SessionPresence::ShuttingDown, SessionTurnState::Unknown)
498 }
499 };
500 activity(
501 locator,
502 presence,
503 turn,
504 "supercode_runtime",
505 Some(state.as_str()),
506 None,
507 observed_at_ms,
508 )
509}
510
511fn claude_activity(
512 locator: &SessionLocator,
513 peer: &ClaudePeerSession,
514 observed_at_ms: u64,
515) -> SessionActivity {
516 let turn = match peer.status {
517 Some(ClaudePeerStatus::Busy) => SessionTurnState::Working,
518 Some(ClaudePeerStatus::Idle) => SessionTurnState::Idle,
519 Some(ClaudePeerStatus::Waiting) => SessionTurnState::NeedsInput,
520 None => SessionTurnState::Unknown,
521 };
522 activity(
523 locator,
524 SessionPresence::Running,
525 turn,
526 "claude_registry",
527 peer.status.map(|status| status.as_str()),
528 peer.version.as_deref(),
529 observed_at_ms,
530 )
531}
532
533fn codex_activity(
534 locator: &SessionLocator,
535 status: CodexPeerStatus,
536 observed_at_ms: u64,
537) -> SessionActivity {
538 let turn = match status {
539 CodexPeerStatus::Running => SessionTurnState::Unknown,
540 CodexPeerStatus::Idle => SessionTurnState::Idle,
541 CodexPeerStatus::Busy => SessionTurnState::Working,
542 CodexPeerStatus::Waiting => SessionTurnState::NeedsInput,
543 };
544 activity(
545 locator,
546 SessionPresence::Running,
547 turn,
548 "codex_rollout",
549 Some(status.as_str()),
550 None,
551 observed_at_ms,
552 )
553}
554
555fn persisted_activity(locator: &SessionLocator, observed_at_ms: u64) -> SessionActivity {
556 activity(
557 locator,
558 SessionPresence::Persisted,
559 SessionTurnState::Unknown,
560 "persisted_store",
561 None,
562 None,
563 observed_at_ms,
564 )
565}
566
567fn activity(
568 locator: &SessionLocator,
569 presence: SessionPresence,
570 turn: SessionTurnState,
571 source: &str,
572 native_state: Option<&str>,
573 harness_version: Option<&str>,
574 observed_at_ms: u64,
575) -> SessionActivity {
576 SessionActivity {
577 harness: locator.harness.clone(),
578 session_id: locator.session_id.clone(),
579 presence,
580 turn,
581 evidence: SessionActivityEvidence {
582 source: source.to_string(),
583 native_state: native_state.map(str::to_string),
584 observed_at_ms,
585 harness_version: harness_version.map(str::to_string),
586 },
587 }
588}
589
590fn now_ms() -> u64 {
591 SystemTime::now()
592 .duration_since(UNIX_EPOCH)
593 .unwrap_or_default()
594 .as_millis()
595 .try_into()
596 .unwrap_or(u64::MAX)
597}