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, 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/// 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 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
150/// Resolve stock-harness activity without consulting Supercode-owned runtime
151/// receipts. Discovery and the subscription lane share this exact mapping.
152pub(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/// What Hermes's own execution ledger says about a session: the fire that
170/// opened it (joined by the orchestration codec, which also resolves
171/// compression chains) and the pid Hermes recorded for the attempt.
172#[derive(Debug, Clone, Copy, PartialEq, Eq)]
173pub(crate) struct HermesFire {
174    status: supercode_interchange::orchestration::FireStatus,
175    pid: Option<u32>,
176}
177
178/// Hermes keeps no process registry, but every scheduled or task fire records
179/// its attempt in `cron/executions.db` with the pid that ran it. A running fire
180/// whose pid is alive is a proven owner; a running fire whose pid is gone is an
181/// attempt its owner died under, and the ledger will never close it.
182fn 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    // `HarnessHomes::hermes` addresses `state.db`; the ledger is its sibling.
190    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            // The attempt's owner is gone and the ledger still says running:
236            // Hermes will never write this fire's end.
237            _ => (
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
312/// Strongest proven status among live descendants, keyed by root session id.
313/// Codex currently defaults to depth one, but walking the live parent map also
314/// handles nested children without turning each rollout into a top-level row.
315fn 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/// Internal test fixture for the evidence precedence table.
491#[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}