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 if let Some((session_id, path)) = crate::codex_peer::session_of_process(pid) {
135 found.push((
136 root,
137 SessionLocator {
138 harness: HarnessId::new(HarnessId::CODEX),
139 session_id,
140 storage: crate::StorageLocator::File { path },
141 },
142 ));
143 break;
144 }
145 }
146 queue.extend(children.get(&pid).into_iter().flatten().copied());
147 }
148 }
149 if found.is_empty() {
150 return Ok(Vec::new());
151 }
152 let locators: Vec<SessionLocator> = found.iter().map(|(_, locator)| locator.clone()).collect();
153 let activities = SessionActivityMonitor::default()
154 .resolve(&locators, homes)
155 .await?;
156 Ok(found
157 .into_iter()
158 .filter_map(|(root, locator)| {
159 activities
160 .iter()
161 .find(|activity| {
162 activity.harness == locator.harness && activity.session_id == locator.session_id
163 })
164 .cloned()
165 .map(|activity| (root, activity))
166 })
167 .collect())
168}
169
170#[cfg(feature = "adapter-api")]
172fn process_table() -> HashMap<u32, (u32, String)> {
173 let Ok(output) = std::process::Command::new("ps")
174 .args(["-axo", "pid=,ppid=,comm="])
175 .output()
176 else {
177 return HashMap::new();
178 };
179 String::from_utf8_lossy(&output.stdout)
180 .lines()
181 .filter_map(|line| {
182 let mut fields = line.split_whitespace();
183 let pid = fields.next()?.parse().ok()?;
184 let parent = fields.next()?.parse().ok()?;
185 let command = fields.collect::<Vec<_>>().join(" ");
186 Some((pid, (parent, command)))
187 })
188 .collect()
189}
190
191#[cfg(feature = "adapter-api")]
194#[derive(Debug, Default)]
195pub(crate) struct SessionActivityMonitor {
196 codex: CodexPeerTracker,
197}
198
199#[cfg(feature = "adapter-api")]
200impl SessionActivityMonitor {
201 pub(crate) async fn resolve(
202 &mut self,
203 locators: &[SessionLocator],
204 homes: &HarnessHomes,
205 ) -> Result<Vec<SessionActivity>, crate::SdkError> {
206 let authorization = crate::RuntimeAuthorization::observer();
207 let sources = locators
208 .iter()
209 .map(|locator| (locator.harness.0.clone(), locator.session_id.clone()))
210 .collect::<BTreeSet<_>>();
211 let owned = crate::LocalRuntimeRegistry::new()
212 .source_states(&sources, &authorization)
213 .await?;
214 let claude = read_claude(locators, homes);
215 let codex = if locators
216 .iter()
217 .any(|locator| locator.harness.as_str() == HarnessId::CODEX)
218 {
219 self.codex.sample(&homes.codex)
220 } else {
221 HashMap::new()
222 };
223 let hermes = read_hermes(locators, homes);
224 let mut activities = resolve_stock_with_evidence(locators, &claude, &codex, &hermes)
225 .into_iter()
226 .map(|activity| (activity.key(), activity))
227 .collect::<BTreeMap<_, _>>();
228 let observed_at_ms = now_ms();
229 for locator in locators {
230 let key = (
231 locator.harness.as_str().to_string(),
232 locator.session_id.clone(),
233 );
234 if let Some(state) = owned.get(&key).copied() {
235 activities.insert(key, owned_activity(locator, state, observed_at_ms));
236 }
237 }
238 Ok(locators
239 .iter()
240 .filter_map(|locator| {
241 activities.remove(&(
242 locator.harness.as_str().to_string(),
243 locator.session_id.clone(),
244 ))
245 })
246 .collect())
247 }
248}
249
250pub(crate) fn resolve_stock_session_activities(
253 locators: &[SessionLocator],
254 homes: &HarnessHomes,
255) -> Vec<SessionActivity> {
256 let claude = read_claude(locators, homes);
257 let codex = if locators
258 .iter()
259 .any(|locator| locator.harness.as_str() == HarnessId::CODEX)
260 {
261 live_rollouts(&homes.codex)
262 } else {
263 HashMap::new()
264 };
265 let hermes = read_hermes(locators, homes);
266 resolve_stock_with_evidence(locators, &claude, &codex, &hermes)
267}
268
269#[derive(Debug, Clone, Copy, PartialEq, Eq)]
273pub(crate) struct HermesFire {
274 status: supercode_interchange::orchestration::FireStatus,
275 pid: Option<u32>,
276}
277
278fn read_hermes(locators: &[SessionLocator], homes: &HarnessHomes) -> HashMap<String, HermesFire> {
283 if !locators
284 .iter()
285 .any(|locator| locator.harness.as_str() == HarnessId::HERMES)
286 {
287 return HashMap::new();
288 }
289 let home = homes
291 .hermes
292 .parent()
293 .map_or_else(|| PathBuf::from("."), std::path::Path::to_path_buf);
294 let Ok(loaded) = supercode_interchange::orchestration::codec::from_hermes(&home) else {
295 return HashMap::new();
296 };
297 let mut fires = HashMap::new();
298 for profile in loaded.orchestration.profiles.values() {
299 for fire in &profile.fires {
300 let Some(session) = fire.session_id.clone() else {
301 continue;
302 };
303 let pid = fire.residue.0.get("pid").and_then(|value| {
304 value
305 .as_u64()
306 .or_else(|| value.as_str().and_then(|text| text.trim().parse().ok()))
307 .and_then(|pid| u32::try_from(pid).ok())
308 });
309 fires.insert(
310 session,
311 HermesFire {
312 status: fire.status,
313 pid,
314 },
315 );
316 }
317 }
318 fires
319}
320
321fn hermes_activity(
322 locator: &SessionLocator,
323 fire: HermesFire,
324 live: impl Fn(u32) -> bool,
325 observed_at_ms: u64,
326) -> SessionActivity {
327 use supercode_interchange::orchestration::FireStatus;
328 let (presence, turn, native_state) = match fire.status {
329 FireStatus::Claimed | FireStatus::Running => match fire.pid {
330 Some(pid) if live(pid) => (
331 SessionPresence::Running,
332 SessionTurnState::Working,
333 fire.status.hermes_word(),
334 ),
335 _ => (
338 SessionPresence::Persisted,
339 SessionTurnState::Unknown,
340 "abandoned",
341 ),
342 },
343 status => (
344 SessionPresence::Persisted,
345 SessionTurnState::Unknown,
346 status.hermes_word(),
347 ),
348 };
349 activity(
350 locator,
351 presence,
352 turn,
353 "hermes_executions",
354 Some(native_state),
355 None,
356 observed_at_ms,
357 )
358}
359
360fn read_claude(
361 locators: &[SessionLocator],
362 homes: &HarnessHomes,
363) -> HashMap<String, ClaudePeerSession> {
364 if locators
365 .iter()
366 .any(|locator| locator.harness.as_str() == HarnessId::CLAUDE_CODE)
367 {
368 crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(homes))
369 .into_iter()
370 .map(|peer| (peer.session_id.clone(), peer))
371 .collect::<HashMap<_, _>>()
372 } else {
373 HashMap::new()
374 }
375}
376
377fn resolve_stock_with_evidence(
378 locators: &[SessionLocator],
379 claude: &HashMap<String, ClaudePeerSession>,
380 codex: &HashMap<PathBuf, CodexPeerStatus>,
381 hermes: &HashMap<String, HermesFire>,
382) -> Vec<SessionActivity> {
383 let observed_at_ms = now_ms();
384 let codex_descendants = codex_descendant_statuses(codex);
385
386 locators
387 .iter()
388 .map(|locator| {
389 if locator.harness.as_str() == HarnessId::CLAUDE_CODE {
390 if let Some(peer) = claude.get(&locator.session_id) {
391 return claude_activity(locator, peer, observed_at_ms);
392 }
393 }
394 if locator.harness.as_str() == HarnessId::HERMES {
395 if let Some(fire) = hermes.get(&locator.session_id) {
396 return hermes_activity(locator, *fire, process_is_live, observed_at_ms);
397 }
398 }
399 if locator.harness.as_str() == HarnessId::CODEX {
400 if let Some(status) = merge_codex_status(
401 rollout_status(codex, locator.storage.path()),
402 codex_descendants.get(&locator.session_id).copied(),
403 ) {
404 return codex_activity(locator, status, observed_at_ms);
405 }
406 }
407 persisted_activity(locator, observed_at_ms)
408 })
409 .collect()
410}
411
412fn codex_descendant_statuses(
416 live: &HashMap<PathBuf, CodexPeerStatus>,
417) -> HashMap<String, CodexPeerStatus> {
418 let lineage = live
419 .iter()
420 .filter_map(|(path, status)| {
421 rollout_lineage(path).map(|(session_id, parent)| (session_id, parent, *status))
422 })
423 .collect::<Vec<_>>();
424 let parent_by_child = lineage
425 .iter()
426 .filter_map(|(child, parent, _)| {
427 parent
428 .as_ref()
429 .map(|parent| (child.clone(), parent.clone()))
430 })
431 .collect::<HashMap<_, _>>();
432 let mut statuses = HashMap::new();
433 for (_, parent, status) in lineage {
434 let Some(mut root) = parent else {
435 continue;
436 };
437 let mut visited = std::collections::HashSet::new();
438 while visited.insert(root.clone()) {
439 let Some(parent) = parent_by_child.get(&root) else {
440 break;
441 };
442 root = parent.clone();
443 }
444 statuses
445 .entry(root)
446 .and_modify(|current| *current = stronger_codex_status(*current, status))
447 .or_insert(status);
448 }
449 statuses
450}
451
452fn merge_codex_status(
453 own: Option<CodexPeerStatus>,
454 descendant: Option<CodexPeerStatus>,
455) -> Option<CodexPeerStatus> {
456 match (own, descendant) {
457 (Some(own), Some(descendant)) => Some(stronger_codex_status(own, descendant)),
458 (Some(status), None) | (None, Some(status)) => Some(status),
459 (None, None) => None,
460 }
461}
462
463fn stronger_codex_status(left: CodexPeerStatus, right: CodexPeerStatus) -> CodexPeerStatus {
464 fn rank(status: CodexPeerStatus) -> u8 {
465 match status {
466 CodexPeerStatus::Running => 0,
467 CodexPeerStatus::Idle => 1,
468 CodexPeerStatus::Busy => 2,
469 }
470 }
471 if rank(right) > rank(left) {
472 right
473 } else {
474 left
475 }
476}
477
478#[cfg(feature = "adapter-api")]
479fn owned_activity(
480 locator: &SessionLocator,
481 state: crate::RuntimeRegistryState,
482 observed_at_ms: u64,
483) -> SessionActivity {
484 use crate::RuntimeRegistryState;
485 let (presence, turn) = match state {
486 RuntimeRegistryState::Persisted => (SessionPresence::Persisted, SessionTurnState::Unknown),
487 RuntimeRegistryState::Idle => (SessionPresence::Running, SessionTurnState::Idle),
488 RuntimeRegistryState::Busy => (SessionPresence::Running, SessionTurnState::Working),
489 RuntimeRegistryState::ShuttingDown => {
490 (SessionPresence::ShuttingDown, SessionTurnState::Unknown)
491 }
492 };
493 activity(
494 locator,
495 presence,
496 turn,
497 "supercode_runtime",
498 Some(state.as_str()),
499 None,
500 observed_at_ms,
501 )
502}
503
504fn claude_activity(
505 locator: &SessionLocator,
506 peer: &ClaudePeerSession,
507 observed_at_ms: u64,
508) -> SessionActivity {
509 let turn = match peer.status {
510 Some(ClaudePeerStatus::Busy) => SessionTurnState::Working,
511 Some(ClaudePeerStatus::Idle) => SessionTurnState::Idle,
512 Some(ClaudePeerStatus::Waiting) => SessionTurnState::NeedsInput,
513 None => SessionTurnState::Unknown,
514 };
515 activity(
516 locator,
517 SessionPresence::Running,
518 turn,
519 "claude_registry",
520 peer.status.map(|status| status.as_str()),
521 peer.version.as_deref(),
522 observed_at_ms,
523 )
524}
525
526fn codex_activity(
527 locator: &SessionLocator,
528 status: CodexPeerStatus,
529 observed_at_ms: u64,
530) -> SessionActivity {
531 let turn = match status {
532 CodexPeerStatus::Running => SessionTurnState::Unknown,
533 CodexPeerStatus::Idle => SessionTurnState::Idle,
534 CodexPeerStatus::Busy => SessionTurnState::Working,
535 };
536 activity(
537 locator,
538 SessionPresence::Running,
539 turn,
540 "codex_rollout",
541 Some(status.as_str()),
542 None,
543 observed_at_ms,
544 )
545}
546
547fn persisted_activity(locator: &SessionLocator, observed_at_ms: u64) -> SessionActivity {
548 activity(
549 locator,
550 SessionPresence::Persisted,
551 SessionTurnState::Unknown,
552 "persisted_store",
553 None,
554 None,
555 observed_at_ms,
556 )
557}
558
559fn activity(
560 locator: &SessionLocator,
561 presence: SessionPresence,
562 turn: SessionTurnState,
563 source: &str,
564 native_state: Option<&str>,
565 harness_version: Option<&str>,
566 observed_at_ms: u64,
567) -> SessionActivity {
568 SessionActivity {
569 harness: locator.harness.clone(),
570 session_id: locator.session_id.clone(),
571 presence,
572 turn,
573 evidence: SessionActivityEvidence {
574 source: source.to_string(),
575 native_state: native_state.map(str::to_string),
576 observed_at_ms,
577 harness_version: harness_version.map(str::to_string),
578 },
579 }
580}
581
582fn now_ms() -> u64 {
583 SystemTime::now()
584 .duration_since(UNIX_EPOCH)
585 .unwrap_or_default()
586 .as_millis()
587 .try_into()
588 .unwrap_or(u64::MAX)
589}
590
591#[cfg(all(test, feature = "adapter-api"))]
593pub(crate) fn resolve_fixture(
594 locator: &SessionLocator,
595 owned: Option<crate::RuntimeRegistryState>,
596) -> SessionActivity {
597 owned.map_or_else(
598 || persisted_activity(locator, 0),
599 |state| owned_activity(locator, state, 0),
600 )
601}
602
603#[cfg(all(test, feature = "adapter-api"))]
604mod tests {
605 use std::fs;
606 use std::path::PathBuf;
607
608 use super::*;
609 use crate::{RuntimeRegistryState, StorageLocator};
610
611 #[test]
612 fn normalized_activity_keeps_presence_and_turn_orthogonal() {
613 let locator = SessionLocator {
614 harness: HarnessId("fixture".into()),
615 session_id: "session-1".into(),
616 storage: StorageLocator::File {
617 path: PathBuf::from("/not-read"),
618 },
619 };
620 let cases = [
621 (None, SessionPresence::Persisted, SessionTurnState::Unknown),
622 (
623 Some(RuntimeRegistryState::Idle),
624 SessionPresence::Running,
625 SessionTurnState::Idle,
626 ),
627 (
628 Some(RuntimeRegistryState::Busy),
629 SessionPresence::Running,
630 SessionTurnState::Working,
631 ),
632 (
633 Some(RuntimeRegistryState::ShuttingDown),
634 SessionPresence::ShuttingDown,
635 SessionTurnState::Unknown,
636 ),
637 ];
638 for (native, presence, turn) in cases {
639 let activity = resolve_fixture(&locator, native);
640 assert_eq!((activity.presence, activity.turn), (presence, turn));
641 assert_eq!(activity.harness, locator.harness);
642 assert_eq!(activity.session_id, locator.session_id);
643 assert!(!activity.evidence.source.contains('/'));
644 }
645 }
646
647 #[test]
648 fn busy_codex_child_makes_the_root_conversation_working() {
649 let child_path = std::env::temp_dir().join(format!(
650 "supercode-codex-child-{}-{}.jsonl",
651 std::process::id(),
652 std::thread::current().name().unwrap_or("test")
653 ));
654 fs::write(
655 &child_path,
656 serde_json::json!({
657 "timestamp": "2026-01-01T00:00:00Z",
658 "type": "session_meta",
659 "payload": {
660 "id": "child-session",
661 "cwd": "/project",
662 "parent_thread_id": "root-session",
663 "source": {"subagent":{"thread_spawn":{
664 "parent_thread_id":"root-session",
665 "depth":1,
666 "agent_path":"/root/reviewer"
667 }}}
668 }
669 })
670 .to_string()
671 + "\n",
672 )
673 .unwrap();
674 let root = SessionLocator {
675 harness: HarnessId::from(HarnessId::CODEX),
676 session_id: "root-session".into(),
677 storage: StorageLocator::File {
678 path: PathBuf::from("/not-open/root.jsonl"),
679 },
680 };
681 let activities = resolve_stock_with_evidence(
682 &[root],
683 &HashMap::new(),
684 &HashMap::from([(child_path.clone(), CodexPeerStatus::Busy)]),
685 &HashMap::new(),
686 );
687 assert_eq!(activities.len(), 1);
688 assert_eq!(activities[0].presence, SessionPresence::Running);
689 assert_eq!(activities[0].turn, SessionTurnState::Working);
690 assert_eq!(activities[0].session_id, "root-session");
691 fs::remove_file(child_path).ok();
692 }
693}