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