Skip to main content

supercode_harness/
session_activity.rs

1//! Protocol-neutral activity for persisted and live harness sessions.
2//!
3//! Activity is deliberately separate from transcript freshness and UI
4//! attention. A process receipt proves presence; only a harness lifecycle
5//! boundary or runtime state proves whether a turn is working.
6
7use 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/// Whether a durable session currently has a proven live owner.
20#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
21#[serde(rename_all = "snake_case")]
22pub enum SessionPresence {
23    /// The session is durable, but no live owner is proven.
24    Persisted,
25    /// A harness or Supercode runtime currently owns the session.
26    Running,
27    /// A Supercode-owned runtime is shutting down.
28    ShuttingDown,
29}
30
31/// Turn activity, independent of presence and frontend attention.
32#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
33#[serde(rename_all = "snake_case")]
34pub enum SessionTurnState {
35    /// Presence is known but the harness exposes no trustworthy turn state.
36    Unknown,
37    /// The runtime is ready for user input.
38    Idle,
39    /// A model, tool, or scheduler turn is active.
40    Working,
41    /// The runtime has issued a structured request that needs a response.
42    NeedsInput,
43}
44
45/// Provenance for one normalized activity observation.
46#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
47pub struct SessionActivityEvidence {
48    /// Stable, non-sensitive evidence source.
49    pub source: String,
50    /// Harness-native state token, when one was published.
51    pub native_state: Option<String>,
52    /// Wall-clock time at which Supercode sampled the evidence.
53    pub observed_at_ms: u64,
54    /// Harness version attached to the evidence, when available.
55    pub harness_version: Option<String>,
56}
57
58/// Normalized lifecycle state for one harness-native session.
59#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
60pub struct SessionActivity {
61    /// Owning harness.
62    pub harness: HarnessId,
63    /// Harness-native durable session id.
64    pub session_id: String,
65    /// Live ownership, independent of turn state.
66    pub presence: SessionPresence,
67    /// Current turn state, independent of unread/attention state.
68    pub turn: SessionTurnState,
69    /// Why this state is trustworthy. Never contains a pid, path, socket, or token.
70    pub evidence: SessionActivityEvidence,
71}
72
73impl SessionActivity {
74    /// Compare transition-bearing state while ignoring the observation clock.
75    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    /// Stable subscription identity without exposing a persistence path.
86    pub(crate) fn key(&self) -> (String, String) {
87        (self.harness.as_str().to_string(), self.session_id.clone())
88    }
89}
90
91/// Stateful activity resolver. It caches only expensive process-ownership
92/// discovery; every lifecycle boundary is still sampled on each poll.
93#[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
153/// Resolve stock-harness activity without consulting Supercode-owned runtime
154/// receipts. Discovery and the subscription lane share this exact mapping.
155pub(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
213/// Strongest proven status among live descendants, keyed by root session id.
214/// Codex currently defaults to depth one, but walking the live parent map also
215/// handles nested children without turning each rollout into a top-level row.
216fn 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/// Internal test fixture for the evidence precedence table.
392#[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}