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")]
94#[derive(Debug, Default)]
95pub(crate) struct SessionActivityMonitor {
96 codex: CodexPeerTracker,
97}
98
99#[cfg(feature = "adapter-api")]
100impl SessionActivityMonitor {
101 pub(crate) async fn resolve(
102 &mut self,
103 locators: &[SessionLocator],
104 homes: &HarnessHomes,
105 ) -> Result<Vec<SessionActivity>, crate::SdkError> {
106 let authorization = crate::RuntimeAuthorization::observer();
107 let sources = locators
108 .iter()
109 .map(|locator| (locator.harness.0.clone(), locator.session_id.clone()))
110 .collect::<BTreeSet<_>>();
111 let owned = crate::LocalRuntimeRegistry::new()
112 .source_states(&sources, &authorization)
113 .await?;
114 let claude = read_claude(locators, homes);
115 let codex = if locators
116 .iter()
117 .any(|locator| locator.harness.as_str() == HarnessId::CODEX)
118 {
119 self.codex.sample(&homes.codex)
120 } else {
121 HashMap::new()
122 };
123 let hermes = read_hermes(locators, homes);
124 let mut activities = resolve_stock_with_evidence(locators, &claude, &codex, &hermes)
125 .into_iter()
126 .map(|activity| (activity.key(), activity))
127 .collect::<BTreeMap<_, _>>();
128 let observed_at_ms = now_ms();
129 for locator in locators {
130 let key = (
131 locator.harness.as_str().to_string(),
132 locator.session_id.clone(),
133 );
134 if let Some(state) = owned.get(&key).copied() {
135 activities.insert(key, owned_activity(locator, state, observed_at_ms));
136 }
137 }
138 Ok(locators
139 .iter()
140 .filter_map(|locator| {
141 activities.remove(&(
142 locator.harness.as_str().to_string(),
143 locator.session_id.clone(),
144 ))
145 })
146 .collect())
147 }
148}
149
150pub(crate) fn resolve_stock_session_activities(
153 locators: &[SessionLocator],
154 homes: &HarnessHomes,
155) -> Vec<SessionActivity> {
156 let claude = read_claude(locators, homes);
157 let codex = if locators
158 .iter()
159 .any(|locator| locator.harness.as_str() == HarnessId::CODEX)
160 {
161 live_rollouts(&homes.codex)
162 } else {
163 HashMap::new()
164 };
165 let hermes = read_hermes(locators, homes);
166 resolve_stock_with_evidence(locators, &claude, &codex, &hermes)
167}
168
169#[derive(Debug, Clone, Copy, PartialEq, Eq)]
173pub(crate) struct HermesFire {
174 status: supercode_interchange::orchestration::FireStatus,
175 pid: Option<u32>,
176}
177
178fn read_hermes(locators: &[SessionLocator], homes: &HarnessHomes) -> HashMap<String, HermesFire> {
183 if !locators
184 .iter()
185 .any(|locator| locator.harness.as_str() == HarnessId::HERMES)
186 {
187 return HashMap::new();
188 }
189 let home = homes
191 .hermes
192 .parent()
193 .map_or_else(|| PathBuf::from("."), std::path::Path::to_path_buf);
194 let Ok(loaded) = supercode_interchange::orchestration::codec::from_hermes(&home) else {
195 return HashMap::new();
196 };
197 let mut fires = HashMap::new();
198 for profile in loaded.orchestration.profiles.values() {
199 for fire in &profile.fires {
200 let Some(session) = fire.session_id.clone() else {
201 continue;
202 };
203 let pid = fire.residue.0.get("pid").and_then(|value| {
204 value
205 .as_u64()
206 .or_else(|| value.as_str().and_then(|text| text.trim().parse().ok()))
207 .and_then(|pid| u32::try_from(pid).ok())
208 });
209 fires.insert(
210 session,
211 HermesFire {
212 status: fire.status,
213 pid,
214 },
215 );
216 }
217 }
218 fires
219}
220
221fn hermes_activity(
222 locator: &SessionLocator,
223 fire: HermesFire,
224 live: impl Fn(u32) -> bool,
225 observed_at_ms: u64,
226) -> SessionActivity {
227 use supercode_interchange::orchestration::FireStatus;
228 let (presence, turn, native_state) = match fire.status {
229 FireStatus::Claimed | FireStatus::Running => match fire.pid {
230 Some(pid) if live(pid) => (
231 SessionPresence::Running,
232 SessionTurnState::Working,
233 fire.status.hermes_word(),
234 ),
235 _ => (
238 SessionPresence::Persisted,
239 SessionTurnState::Unknown,
240 "abandoned",
241 ),
242 },
243 status => (
244 SessionPresence::Persisted,
245 SessionTurnState::Unknown,
246 status.hermes_word(),
247 ),
248 };
249 activity(
250 locator,
251 presence,
252 turn,
253 "hermes_executions",
254 Some(native_state),
255 None,
256 observed_at_ms,
257 )
258}
259
260fn read_claude(
261 locators: &[SessionLocator],
262 homes: &HarnessHomes,
263) -> HashMap<String, ClaudePeerSession> {
264 if locators
265 .iter()
266 .any(|locator| locator.harness.as_str() == HarnessId::CLAUDE_CODE)
267 {
268 crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(homes))
269 .into_iter()
270 .map(|peer| (peer.session_id.clone(), peer))
271 .collect::<HashMap<_, _>>()
272 } else {
273 HashMap::new()
274 }
275}
276
277fn resolve_stock_with_evidence(
278 locators: &[SessionLocator],
279 claude: &HashMap<String, ClaudePeerSession>,
280 codex: &HashMap<PathBuf, CodexPeerStatus>,
281 hermes: &HashMap<String, HermesFire>,
282) -> Vec<SessionActivity> {
283 let observed_at_ms = now_ms();
284 let codex_descendants = codex_descendant_statuses(codex);
285
286 locators
287 .iter()
288 .map(|locator| {
289 if locator.harness.as_str() == HarnessId::CLAUDE_CODE {
290 if let Some(peer) = claude.get(&locator.session_id) {
291 return claude_activity(locator, peer, observed_at_ms);
292 }
293 }
294 if locator.harness.as_str() == HarnessId::HERMES {
295 if let Some(fire) = hermes.get(&locator.session_id) {
296 return hermes_activity(locator, *fire, process_is_live, observed_at_ms);
297 }
298 }
299 if locator.harness.as_str() == HarnessId::CODEX {
300 if let Some(status) = merge_codex_status(
301 rollout_status(codex, locator.storage.path()),
302 codex_descendants.get(&locator.session_id).copied(),
303 ) {
304 return codex_activity(locator, status, observed_at_ms);
305 }
306 }
307 persisted_activity(locator, observed_at_ms)
308 })
309 .collect()
310}
311
312fn codex_descendant_statuses(
316 live: &HashMap<PathBuf, CodexPeerStatus>,
317) -> HashMap<String, CodexPeerStatus> {
318 let lineage = live
319 .iter()
320 .filter_map(|(path, status)| {
321 rollout_lineage(path).map(|(session_id, parent)| (session_id, parent, *status))
322 })
323 .collect::<Vec<_>>();
324 let parent_by_child = lineage
325 .iter()
326 .filter_map(|(child, parent, _)| {
327 parent
328 .as_ref()
329 .map(|parent| (child.clone(), parent.clone()))
330 })
331 .collect::<HashMap<_, _>>();
332 let mut statuses = HashMap::new();
333 for (_, parent, status) in lineage {
334 let Some(mut root) = parent else {
335 continue;
336 };
337 let mut visited = std::collections::HashSet::new();
338 while visited.insert(root.clone()) {
339 let Some(parent) = parent_by_child.get(&root) else {
340 break;
341 };
342 root = parent.clone();
343 }
344 statuses
345 .entry(root)
346 .and_modify(|current| *current = stronger_codex_status(*current, status))
347 .or_insert(status);
348 }
349 statuses
350}
351
352fn merge_codex_status(
353 own: Option<CodexPeerStatus>,
354 descendant: Option<CodexPeerStatus>,
355) -> Option<CodexPeerStatus> {
356 match (own, descendant) {
357 (Some(own), Some(descendant)) => Some(stronger_codex_status(own, descendant)),
358 (Some(status), None) | (None, Some(status)) => Some(status),
359 (None, None) => None,
360 }
361}
362
363fn stronger_codex_status(left: CodexPeerStatus, right: CodexPeerStatus) -> CodexPeerStatus {
364 fn rank(status: CodexPeerStatus) -> u8 {
365 match status {
366 CodexPeerStatus::Running => 0,
367 CodexPeerStatus::Idle => 1,
368 CodexPeerStatus::Busy => 2,
369 }
370 }
371 if rank(right) > rank(left) {
372 right
373 } else {
374 left
375 }
376}
377
378#[cfg(feature = "adapter-api")]
379fn owned_activity(
380 locator: &SessionLocator,
381 state: crate::RuntimeRegistryState,
382 observed_at_ms: u64,
383) -> SessionActivity {
384 use crate::RuntimeRegistryState;
385 let (presence, turn) = match state {
386 RuntimeRegistryState::Persisted => (SessionPresence::Persisted, SessionTurnState::Unknown),
387 RuntimeRegistryState::Idle => (SessionPresence::Running, SessionTurnState::Idle),
388 RuntimeRegistryState::Busy => (SessionPresence::Running, SessionTurnState::Working),
389 RuntimeRegistryState::ShuttingDown => {
390 (SessionPresence::ShuttingDown, SessionTurnState::Unknown)
391 }
392 };
393 activity(
394 locator,
395 presence,
396 turn,
397 "supercode_runtime",
398 Some(state.as_str()),
399 None,
400 observed_at_ms,
401 )
402}
403
404fn claude_activity(
405 locator: &SessionLocator,
406 peer: &ClaudePeerSession,
407 observed_at_ms: u64,
408) -> SessionActivity {
409 let turn = match peer.status {
410 Some(ClaudePeerStatus::Busy) => SessionTurnState::Working,
411 Some(ClaudePeerStatus::Idle) => SessionTurnState::Idle,
412 None => SessionTurnState::Unknown,
413 };
414 activity(
415 locator,
416 SessionPresence::Running,
417 turn,
418 "claude_registry",
419 peer.status.map(|status| status.as_str()),
420 peer.version.as_deref(),
421 observed_at_ms,
422 )
423}
424
425fn codex_activity(
426 locator: &SessionLocator,
427 status: CodexPeerStatus,
428 observed_at_ms: u64,
429) -> SessionActivity {
430 let turn = match status {
431 CodexPeerStatus::Running => SessionTurnState::Unknown,
432 CodexPeerStatus::Idle => SessionTurnState::Idle,
433 CodexPeerStatus::Busy => SessionTurnState::Working,
434 };
435 activity(
436 locator,
437 SessionPresence::Running,
438 turn,
439 "codex_rollout",
440 Some(status.as_str()),
441 None,
442 observed_at_ms,
443 )
444}
445
446fn persisted_activity(locator: &SessionLocator, observed_at_ms: u64) -> SessionActivity {
447 activity(
448 locator,
449 SessionPresence::Persisted,
450 SessionTurnState::Unknown,
451 "persisted_store",
452 None,
453 None,
454 observed_at_ms,
455 )
456}
457
458fn activity(
459 locator: &SessionLocator,
460 presence: SessionPresence,
461 turn: SessionTurnState,
462 source: &str,
463 native_state: Option<&str>,
464 harness_version: Option<&str>,
465 observed_at_ms: u64,
466) -> SessionActivity {
467 SessionActivity {
468 harness: locator.harness.clone(),
469 session_id: locator.session_id.clone(),
470 presence,
471 turn,
472 evidence: SessionActivityEvidence {
473 source: source.to_string(),
474 native_state: native_state.map(str::to_string),
475 observed_at_ms,
476 harness_version: harness_version.map(str::to_string),
477 },
478 }
479}
480
481fn now_ms() -> u64 {
482 SystemTime::now()
483 .duration_since(UNIX_EPOCH)
484 .unwrap_or_default()
485 .as_millis()
486 .try_into()
487 .unwrap_or(u64::MAX)
488}
489
490#[cfg(all(test, feature = "adapter-api"))]
492pub(crate) fn resolve_fixture(
493 locator: &SessionLocator,
494 owned: Option<crate::RuntimeRegistryState>,
495) -> SessionActivity {
496 owned.map_or_else(
497 || persisted_activity(locator, 0),
498 |state| owned_activity(locator, state, 0),
499 )
500}
501
502#[cfg(all(test, feature = "adapter-api"))]
503mod tests {
504 use std::fs;
505 use std::path::PathBuf;
506
507 use super::*;
508 use crate::{RuntimeRegistryState, StorageLocator};
509
510 #[test]
511 fn normalized_activity_keeps_presence_and_turn_orthogonal() {
512 let locator = SessionLocator {
513 harness: HarnessId("fixture".into()),
514 session_id: "session-1".into(),
515 storage: StorageLocator::File {
516 path: PathBuf::from("/not-read"),
517 },
518 };
519 let cases = [
520 (None, SessionPresence::Persisted, SessionTurnState::Unknown),
521 (
522 Some(RuntimeRegistryState::Idle),
523 SessionPresence::Running,
524 SessionTurnState::Idle,
525 ),
526 (
527 Some(RuntimeRegistryState::Busy),
528 SessionPresence::Running,
529 SessionTurnState::Working,
530 ),
531 (
532 Some(RuntimeRegistryState::ShuttingDown),
533 SessionPresence::ShuttingDown,
534 SessionTurnState::Unknown,
535 ),
536 ];
537 for (native, presence, turn) in cases {
538 let activity = resolve_fixture(&locator, native);
539 assert_eq!((activity.presence, activity.turn), (presence, turn));
540 assert_eq!(activity.harness, locator.harness);
541 assert_eq!(activity.session_id, locator.session_id);
542 assert!(!activity.evidence.source.contains('/'));
543 }
544 }
545
546 #[test]
547 fn busy_codex_child_makes_the_root_conversation_working() {
548 let child_path = std::env::temp_dir().join(format!(
549 "supercode-codex-child-{}-{}.jsonl",
550 std::process::id(),
551 std::thread::current().name().unwrap_or("test")
552 ));
553 fs::write(
554 &child_path,
555 serde_json::json!({
556 "timestamp": "2026-01-01T00:00:00Z",
557 "type": "session_meta",
558 "payload": {
559 "id": "child-session",
560 "cwd": "/project",
561 "parent_thread_id": "root-session",
562 "source": {"subagent":{"thread_spawn":{
563 "parent_thread_id":"root-session",
564 "depth":1,
565 "agent_path":"/root/reviewer"
566 }}}
567 }
568 })
569 .to_string()
570 + "\n",
571 )
572 .unwrap();
573 let root = SessionLocator {
574 harness: HarnessId::from(HarnessId::CODEX),
575 session_id: "root-session".into(),
576 storage: StorageLocator::File {
577 path: PathBuf::from("/not-open/root.jsonl"),
578 },
579 };
580 let activities = resolve_stock_with_evidence(
581 &[root],
582 &HashMap::new(),
583 &HashMap::from([(child_path.clone(), CodexPeerStatus::Busy)]),
584 &HashMap::new(),
585 );
586 assert_eq!(activities.len(), 1);
587 assert_eq!(activities[0].presence, SessionPresence::Running);
588 assert_eq!(activities[0].turn, SessionTurnState::Working);
589 assert_eq!(activities[0].session_id, "root-session");
590 fs::remove_file(child_path).ok();
591 }
592}