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 = 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
156/// Resolve stock-harness activity without consulting Supercode-owned runtime
157/// receipts. Discovery and the subscription lane share this exact mapping.
158pub(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
220/// Strongest proven status among live descendants, keyed by root session id.
221/// Codex currently defaults to depth one, but walking the live parent map also
222/// handles nested children without turning each rollout into a top-level row.
223fn 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/// Internal test fixture for the evidence precedence table.
399#[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}