1use std::collections::{BTreeMap, HashMap};
8use std::path::PathBuf;
9use std::time::{SystemTime, UNIX_EPOCH};
10
11use serde::{Deserialize, Serialize};
12
13use crate::claude_peer::{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 entries = crate::LocalRuntimeRegistry::new()
108 .list(
109 &crate::RuntimeRegistryQuery {
110 persisted: Default::default(),
111 include_live: true,
112 include_persisted: false,
113 },
114 &authorization,
115 )
116 .await?;
117 let owned = entries
118 .into_iter()
119 .map(|entry| ((entry.source_harness, entry.source_session_id), entry.state))
120 .collect::<BTreeMap<_, _>>();
121 let claude = read_claude(locators, homes);
122 let codex = if locators
123 .iter()
124 .any(|locator| locator.harness.as_str() == HarnessId::CODEX)
125 {
126 self.codex.sample(&homes.codex)
127 } else {
128 HashMap::new()
129 };
130 let mut activities = resolve_stock_with_evidence(locators, &claude, &codex)
131 .into_iter()
132 .map(|activity| (activity.key(), activity))
133 .collect::<BTreeMap<_, _>>();
134 let observed_at_ms = now_ms();
135 for locator in locators {
136 let key = (
137 locator.harness.as_str().to_string(),
138 locator.session_id.clone(),
139 );
140 if let Some(state) = owned.get(&key).copied() {
141 activities.insert(key, owned_activity(locator, state, observed_at_ms));
142 }
143 }
144 Ok(locators
145 .iter()
146 .filter_map(|locator| {
147 activities.remove(&(
148 locator.harness.as_str().to_string(),
149 locator.session_id.clone(),
150 ))
151 })
152 .collect())
153 }
154}
155
156pub(crate) fn resolve_stock_session_activities(
159 locators: &[SessionLocator],
160 homes: &HarnessHomes,
161) -> Vec<SessionActivity> {
162 let claude = read_claude(locators, homes);
163 let codex = if locators
164 .iter()
165 .any(|locator| locator.harness.as_str() == HarnessId::CODEX)
166 {
167 live_rollouts(&homes.codex)
168 } else {
169 HashMap::new()
170 };
171 resolve_stock_with_evidence(locators, &claude, &codex)
172}
173
174fn read_claude(
175 locators: &[SessionLocator],
176 homes: &HarnessHomes,
177) -> HashMap<String, ClaudePeerSession> {
178 if locators
179 .iter()
180 .any(|locator| locator.harness.as_str() == HarnessId::CLAUDE_CODE)
181 {
182 crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(homes))
183 .into_iter()
184 .map(|peer| (peer.session_id.clone(), peer))
185 .collect::<HashMap<_, _>>()
186 } else {
187 HashMap::new()
188 }
189}
190
191fn resolve_stock_with_evidence(
192 locators: &[SessionLocator],
193 claude: &HashMap<String, ClaudePeerSession>,
194 codex: &HashMap<PathBuf, CodexPeerStatus>,
195) -> Vec<SessionActivity> {
196 let observed_at_ms = now_ms();
197 let codex_descendants = codex_descendant_statuses(codex);
198
199 locators
200 .iter()
201 .map(|locator| {
202 if locator.harness.as_str() == HarnessId::CLAUDE_CODE {
203 if let Some(peer) = claude.get(&locator.session_id) {
204 return claude_activity(locator, peer, observed_at_ms);
205 }
206 }
207 if locator.harness.as_str() == HarnessId::CODEX {
208 if let Some(status) = merge_codex_status(
209 rollout_status(codex, locator.storage.path()),
210 codex_descendants.get(&locator.session_id).copied(),
211 ) {
212 return codex_activity(locator, status, observed_at_ms);
213 }
214 }
215 persisted_activity(locator, observed_at_ms)
216 })
217 .collect()
218}
219
220fn codex_descendant_statuses(
224 live: &HashMap<PathBuf, CodexPeerStatus>,
225) -> HashMap<String, CodexPeerStatus> {
226 let lineage = live
227 .iter()
228 .filter_map(|(path, status)| {
229 rollout_lineage(path).map(|(session_id, parent)| (session_id, parent, *status))
230 })
231 .collect::<Vec<_>>();
232 let parent_by_child = lineage
233 .iter()
234 .filter_map(|(child, parent, _)| {
235 parent
236 .as_ref()
237 .map(|parent| (child.clone(), parent.clone()))
238 })
239 .collect::<HashMap<_, _>>();
240 let mut statuses = HashMap::new();
241 for (_, parent, status) in lineage {
242 let Some(mut root) = parent else {
243 continue;
244 };
245 let mut visited = std::collections::HashSet::new();
246 while visited.insert(root.clone()) {
247 let Some(parent) = parent_by_child.get(&root) else {
248 break;
249 };
250 root = parent.clone();
251 }
252 statuses
253 .entry(root)
254 .and_modify(|current| *current = stronger_codex_status(*current, status))
255 .or_insert(status);
256 }
257 statuses
258}
259
260fn merge_codex_status(
261 own: Option<CodexPeerStatus>,
262 descendant: Option<CodexPeerStatus>,
263) -> Option<CodexPeerStatus> {
264 match (own, descendant) {
265 (Some(own), Some(descendant)) => Some(stronger_codex_status(own, descendant)),
266 (Some(status), None) | (None, Some(status)) => Some(status),
267 (None, None) => None,
268 }
269}
270
271fn stronger_codex_status(left: CodexPeerStatus, right: CodexPeerStatus) -> CodexPeerStatus {
272 fn rank(status: CodexPeerStatus) -> u8 {
273 match status {
274 CodexPeerStatus::Running => 0,
275 CodexPeerStatus::Idle => 1,
276 CodexPeerStatus::Busy => 2,
277 }
278 }
279 if rank(right) > rank(left) {
280 right
281 } else {
282 left
283 }
284}
285
286#[cfg(feature = "adapter-api")]
287fn owned_activity(
288 locator: &SessionLocator,
289 state: crate::RuntimeRegistryState,
290 observed_at_ms: u64,
291) -> SessionActivity {
292 use crate::RuntimeRegistryState;
293 let (presence, turn) = match state {
294 RuntimeRegistryState::Persisted => (SessionPresence::Persisted, SessionTurnState::Unknown),
295 RuntimeRegistryState::Idle => (SessionPresence::Running, SessionTurnState::Idle),
296 RuntimeRegistryState::Busy => (SessionPresence::Running, SessionTurnState::Working),
297 RuntimeRegistryState::ShuttingDown => {
298 (SessionPresence::ShuttingDown, SessionTurnState::Unknown)
299 }
300 };
301 activity(
302 locator,
303 presence,
304 turn,
305 "supercode_runtime",
306 Some(state.as_str()),
307 None,
308 observed_at_ms,
309 )
310}
311
312fn claude_activity(
313 locator: &SessionLocator,
314 peer: &ClaudePeerSession,
315 observed_at_ms: u64,
316) -> SessionActivity {
317 let turn = match peer.status {
318 Some(ClaudePeerStatus::Busy) => SessionTurnState::Working,
319 Some(ClaudePeerStatus::Idle) => SessionTurnState::Idle,
320 None => SessionTurnState::Unknown,
321 };
322 activity(
323 locator,
324 SessionPresence::Running,
325 turn,
326 "claude_registry",
327 peer.status.map(|status| status.as_str()),
328 peer.version.as_deref(),
329 observed_at_ms,
330 )
331}
332
333fn codex_activity(
334 locator: &SessionLocator,
335 status: CodexPeerStatus,
336 observed_at_ms: u64,
337) -> SessionActivity {
338 let turn = match status {
339 CodexPeerStatus::Running => SessionTurnState::Unknown,
340 CodexPeerStatus::Idle => SessionTurnState::Idle,
341 CodexPeerStatus::Busy => SessionTurnState::Working,
342 };
343 activity(
344 locator,
345 SessionPresence::Running,
346 turn,
347 "codex_rollout",
348 Some(status.as_str()),
349 None,
350 observed_at_ms,
351 )
352}
353
354fn persisted_activity(locator: &SessionLocator, observed_at_ms: u64) -> SessionActivity {
355 activity(
356 locator,
357 SessionPresence::Persisted,
358 SessionTurnState::Unknown,
359 "persisted_store",
360 None,
361 None,
362 observed_at_ms,
363 )
364}
365
366fn activity(
367 locator: &SessionLocator,
368 presence: SessionPresence,
369 turn: SessionTurnState,
370 source: &str,
371 native_state: Option<&str>,
372 harness_version: Option<&str>,
373 observed_at_ms: u64,
374) -> SessionActivity {
375 SessionActivity {
376 harness: locator.harness.clone(),
377 session_id: locator.session_id.clone(),
378 presence,
379 turn,
380 evidence: SessionActivityEvidence {
381 source: source.to_string(),
382 native_state: native_state.map(str::to_string),
383 observed_at_ms,
384 harness_version: harness_version.map(str::to_string),
385 },
386 }
387}
388
389fn now_ms() -> u64 {
390 SystemTime::now()
391 .duration_since(UNIX_EPOCH)
392 .unwrap_or_default()
393 .as_millis()
394 .try_into()
395 .unwrap_or(u64::MAX)
396}
397
398#[cfg(all(test, feature = "adapter-api"))]
400pub(crate) fn resolve_fixture(
401 locator: &SessionLocator,
402 owned: Option<crate::RuntimeRegistryState>,
403) -> SessionActivity {
404 owned.map_or_else(
405 || persisted_activity(locator, 0),
406 |state| owned_activity(locator, state, 0),
407 )
408}
409
410#[cfg(all(test, feature = "adapter-api"))]
411mod tests {
412 use std::fs;
413 use std::path::PathBuf;
414
415 use super::*;
416 use crate::{RuntimeRegistryState, StorageLocator};
417
418 #[test]
419 fn normalized_activity_keeps_presence_and_turn_orthogonal() {
420 let locator = SessionLocator {
421 harness: HarnessId("fixture".into()),
422 session_id: "session-1".into(),
423 storage: StorageLocator::File {
424 path: PathBuf::from("/not-read"),
425 },
426 };
427 let cases = [
428 (None, SessionPresence::Persisted, SessionTurnState::Unknown),
429 (
430 Some(RuntimeRegistryState::Idle),
431 SessionPresence::Running,
432 SessionTurnState::Idle,
433 ),
434 (
435 Some(RuntimeRegistryState::Busy),
436 SessionPresence::Running,
437 SessionTurnState::Working,
438 ),
439 (
440 Some(RuntimeRegistryState::ShuttingDown),
441 SessionPresence::ShuttingDown,
442 SessionTurnState::Unknown,
443 ),
444 ];
445 for (native, presence, turn) in cases {
446 let activity = resolve_fixture(&locator, native);
447 assert_eq!((activity.presence, activity.turn), (presence, turn));
448 assert_eq!(activity.harness, locator.harness);
449 assert_eq!(activity.session_id, locator.session_id);
450 assert!(!activity.evidence.source.contains('/'));
451 }
452 }
453
454 #[test]
455 fn busy_codex_child_makes_the_root_conversation_working() {
456 let child_path = std::env::temp_dir().join(format!(
457 "supercode-codex-child-{}-{}.jsonl",
458 std::process::id(),
459 std::thread::current().name().unwrap_or("test")
460 ));
461 fs::write(
462 &child_path,
463 serde_json::json!({
464 "timestamp": "2026-01-01T00:00:00Z",
465 "type": "session_meta",
466 "payload": {
467 "id": "child-session",
468 "cwd": "/project",
469 "parent_thread_id": "root-session",
470 "source": {"subagent":{"thread_spawn":{
471 "parent_thread_id":"root-session",
472 "depth":1,
473 "agent_path":"/root/reviewer"
474 }}}
475 }
476 })
477 .to_string()
478 + "\n",
479 )
480 .unwrap();
481 let root = SessionLocator {
482 harness: HarnessId::from(HarnessId::CODEX),
483 session_id: "root-session".into(),
484 storage: StorageLocator::File {
485 path: PathBuf::from("/not-open/root.jsonl"),
486 },
487 };
488 let activities = resolve_stock_with_evidence(
489 &[root],
490 &HashMap::new(),
491 &HashMap::from([(child_path.clone(), CodexPeerStatus::Busy)]),
492 );
493 assert_eq!(activities.len(), 1);
494 assert_eq!(activities[0].presence, SessionPresence::Running);
495 assert_eq!(activities[0].turn, SessionTurnState::Working);
496 assert_eq!(activities[0].session_id, "root-session");
497 fs::remove_file(child_path).ok();
498 }
499}