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/// The harness session running beneath each of `roots` (a terminal pane's
92/// shell, say) with its activity: the nearest process below the root that
93/// owns a Claude Code session (its registry pid) or a Codex conversation (the
94/// rollout it holds open). Roots with no harness session beneath them are
95/// left out. This is how a pane's liveness comes from the harness's own
96/// lifecycle instead of a guess from the pane's foreground command.
97#[cfg(feature = "adapter-api")]
98pub(crate) async fn activity_under(
99    roots: &[u32],
100    homes: &HarnessHomes,
101) -> Result<Vec<(u32, SessionActivity)>, crate::SdkError> {
102    let table = process_table();
103    let mut children: HashMap<u32, Vec<u32>> = HashMap::new();
104    for (&pid, (parent, _)) in &table {
105        children.entry(*parent).or_default().push(pid);
106    }
107    let registry = crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(homes));
108    let mut found: Vec<(u32, SessionLocator)> = Vec::new();
109    for &root in roots {
110        let mut queue = std::collections::VecDeque::from([root]);
111        let mut visited = 0;
112        while let Some(pid) = queue.pop_front() {
113            visited += 1;
114            if visited > 256 {
115                break;
116            }
117            if let Some(session) = registry.iter().find(|session| session.pid == pid) {
118                found.push((
119                    root,
120                    SessionLocator {
121                        harness: HarnessId::new(HarnessId::CLAUDE_CODE),
122                        session_id: session.session_id.clone(),
123                        storage: crate::StorageLocator::File {
124                            path: PathBuf::new(),
125                        },
126                    },
127                ));
128                break;
129            }
130            let is_codex = table
131                .get(&pid)
132                .is_some_and(|(_, command)| command.rsplit('/').next() == Some("codex"));
133            if is_codex {
134                if let Some((session_id, path)) = crate::codex_peer::session_of_process(pid) {
135                    found.push((
136                        root,
137                        SessionLocator {
138                            harness: HarnessId::new(HarnessId::CODEX),
139                            session_id,
140                            storage: crate::StorageLocator::File { path },
141                        },
142                    ));
143                    break;
144                }
145            }
146            queue.extend(children.get(&pid).into_iter().flatten().copied());
147        }
148    }
149    if found.is_empty() {
150        return Ok(Vec::new());
151    }
152    let locators: Vec<SessionLocator> = found.iter().map(|(_, locator)| locator.clone()).collect();
153    let activities = SessionActivityMonitor::default()
154        .resolve(&locators, homes)
155        .await?;
156    Ok(found
157        .into_iter()
158        .filter_map(|(root, locator)| {
159            activities
160                .iter()
161                .find(|activity| {
162                    activity.harness == locator.harness && activity.session_id == locator.session_id
163                })
164                .cloned()
165                .map(|activity| (root, activity))
166        })
167        .collect())
168}
169
170/// Every process: pid → (parent pid, command), from `ps`.
171#[cfg(feature = "adapter-api")]
172fn process_table() -> HashMap<u32, (u32, String)> {
173    let Ok(output) = std::process::Command::new("ps")
174        .args(["-axo", "pid=,ppid=,comm="])
175        .output()
176    else {
177        return HashMap::new();
178    };
179    String::from_utf8_lossy(&output.stdout)
180        .lines()
181        .filter_map(|line| {
182            let mut fields = line.split_whitespace();
183            let pid = fields.next()?.parse().ok()?;
184            let parent = fields.next()?.parse().ok()?;
185            let command = fields.collect::<Vec<_>>().join(" ");
186            Some((pid, (parent, command)))
187        })
188        .collect()
189}
190
191/// Stateful activity resolver. It caches only expensive process-ownership
192/// discovery; every lifecycle boundary is still sampled on each poll.
193#[cfg(feature = "adapter-api")]
194#[derive(Debug, Default)]
195pub(crate) struct SessionActivityMonitor {
196    codex: CodexPeerTracker,
197}
198
199#[cfg(feature = "adapter-api")]
200impl SessionActivityMonitor {
201    pub(crate) async fn resolve(
202        &mut self,
203        locators: &[SessionLocator],
204        homes: &HarnessHomes,
205    ) -> Result<Vec<SessionActivity>, crate::SdkError> {
206        let authorization = crate::RuntimeAuthorization::observer();
207        let sources = locators
208            .iter()
209            .map(|locator| (locator.harness.0.clone(), locator.session_id.clone()))
210            .collect::<BTreeSet<_>>();
211        let owned = crate::LocalRuntimeRegistry::new()
212            .source_states(&sources, &authorization)
213            .await?;
214        let claude = read_claude(locators, homes);
215        let codex = if locators
216            .iter()
217            .any(|locator| locator.harness.as_str() == HarnessId::CODEX)
218        {
219            self.codex.sample(&homes.codex)
220        } else {
221            HashMap::new()
222        };
223        let hermes = read_hermes(locators, homes);
224        let mut activities = resolve_stock_with_evidence(locators, &claude, &codex, &hermes)
225            .into_iter()
226            .map(|activity| (activity.key(), activity))
227            .collect::<BTreeMap<_, _>>();
228        let observed_at_ms = now_ms();
229        for locator in locators {
230            let key = (
231                locator.harness.as_str().to_string(),
232                locator.session_id.clone(),
233            );
234            if let Some(state) = owned.get(&key).copied() {
235                activities.insert(key, owned_activity(locator, state, observed_at_ms));
236            }
237        }
238        Ok(locators
239            .iter()
240            .filter_map(|locator| {
241                activities.remove(&(
242                    locator.harness.as_str().to_string(),
243                    locator.session_id.clone(),
244                ))
245            })
246            .collect())
247    }
248}
249
250/// Resolve stock-harness activity without consulting Supercode-owned runtime
251/// receipts. Discovery and the subscription lane share this exact mapping.
252pub(crate) fn resolve_stock_session_activities(
253    locators: &[SessionLocator],
254    homes: &HarnessHomes,
255) -> Vec<SessionActivity> {
256    let claude = read_claude(locators, homes);
257    let codex = if locators
258        .iter()
259        .any(|locator| locator.harness.as_str() == HarnessId::CODEX)
260    {
261        live_rollouts(&homes.codex)
262    } else {
263        HashMap::new()
264    };
265    let hermes = read_hermes(locators, homes);
266    resolve_stock_with_evidence(locators, &claude, &codex, &hermes)
267}
268
269/// What Hermes's own execution ledger says about a session: the fire that
270/// opened it (joined by the orchestration codec, which also resolves
271/// compression chains) and the pid Hermes recorded for the attempt.
272#[derive(Debug, Clone, Copy, PartialEq, Eq)]
273pub(crate) struct HermesFire {
274    status: supercode_interchange::orchestration::FireStatus,
275    pid: Option<u32>,
276}
277
278/// Hermes keeps no process registry, but every scheduled or task fire records
279/// its attempt in `cron/executions.db` with the pid that ran it. A running fire
280/// whose pid is alive is a proven owner; a running fire whose pid is gone is an
281/// attempt its owner died under, and the ledger will never close it.
282fn read_hermes(locators: &[SessionLocator], homes: &HarnessHomes) -> HashMap<String, HermesFire> {
283    if !locators
284        .iter()
285        .any(|locator| locator.harness.as_str() == HarnessId::HERMES)
286    {
287        return HashMap::new();
288    }
289    // `HarnessHomes::hermes` addresses `state.db`; the ledger is its sibling.
290    let home = homes
291        .hermes
292        .parent()
293        .map_or_else(|| PathBuf::from("."), std::path::Path::to_path_buf);
294    let Ok(loaded) = supercode_interchange::orchestration::codec::from_hermes(&home) else {
295        return HashMap::new();
296    };
297    let mut fires = HashMap::new();
298    for profile in loaded.orchestration.profiles.values() {
299        for fire in &profile.fires {
300            let Some(session) = fire.session_id.clone() else {
301                continue;
302            };
303            let pid = fire.residue.0.get("pid").and_then(|value| {
304                value
305                    .as_u64()
306                    .or_else(|| value.as_str().and_then(|text| text.trim().parse().ok()))
307                    .and_then(|pid| u32::try_from(pid).ok())
308            });
309            fires.insert(
310                session,
311                HermesFire {
312                    status: fire.status,
313                    pid,
314                },
315            );
316        }
317    }
318    fires
319}
320
321fn hermes_activity(
322    locator: &SessionLocator,
323    fire: HermesFire,
324    live: impl Fn(u32) -> bool,
325    observed_at_ms: u64,
326) -> SessionActivity {
327    use supercode_interchange::orchestration::FireStatus;
328    let (presence, turn, native_state) = match fire.status {
329        FireStatus::Claimed | FireStatus::Running => match fire.pid {
330            Some(pid) if live(pid) => (
331                SessionPresence::Running,
332                SessionTurnState::Working,
333                fire.status.hermes_word(),
334            ),
335            // The attempt's owner is gone and the ledger still says running:
336            // Hermes will never write this fire's end.
337            _ => (
338                SessionPresence::Persisted,
339                SessionTurnState::Unknown,
340                "abandoned",
341            ),
342        },
343        status => (
344            SessionPresence::Persisted,
345            SessionTurnState::Unknown,
346            status.hermes_word(),
347        ),
348    };
349    activity(
350        locator,
351        presence,
352        turn,
353        "hermes_executions",
354        Some(native_state),
355        None,
356        observed_at_ms,
357    )
358}
359
360fn read_claude(
361    locators: &[SessionLocator],
362    homes: &HarnessHomes,
363) -> HashMap<String, ClaudePeerSession> {
364    if locators
365        .iter()
366        .any(|locator| locator.harness.as_str() == HarnessId::CLAUDE_CODE)
367    {
368        crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(homes))
369            .into_iter()
370            .map(|peer| (peer.session_id.clone(), peer))
371            .collect::<HashMap<_, _>>()
372    } else {
373        HashMap::new()
374    }
375}
376
377fn resolve_stock_with_evidence(
378    locators: &[SessionLocator],
379    claude: &HashMap<String, ClaudePeerSession>,
380    codex: &HashMap<PathBuf, CodexPeerStatus>,
381    hermes: &HashMap<String, HermesFire>,
382) -> Vec<SessionActivity> {
383    let observed_at_ms = now_ms();
384    let codex_descendants = codex_descendant_statuses(codex);
385
386    locators
387        .iter()
388        .map(|locator| {
389            if locator.harness.as_str() == HarnessId::CLAUDE_CODE {
390                if let Some(peer) = claude.get(&locator.session_id) {
391                    return claude_activity(locator, peer, observed_at_ms);
392                }
393            }
394            if locator.harness.as_str() == HarnessId::HERMES {
395                if let Some(fire) = hermes.get(&locator.session_id) {
396                    return hermes_activity(locator, *fire, process_is_live, observed_at_ms);
397                }
398            }
399            if locator.harness.as_str() == HarnessId::CODEX {
400                if let Some(status) = merge_codex_status(
401                    rollout_status(codex, locator.storage.path()),
402                    codex_descendants.get(&locator.session_id).copied(),
403                ) {
404                    return codex_activity(locator, status, observed_at_ms);
405                }
406            }
407            persisted_activity(locator, observed_at_ms)
408        })
409        .collect()
410}
411
412/// Strongest proven status among live descendants, keyed by root session id.
413/// Codex currently defaults to depth one, but walking the live parent map also
414/// handles nested children without turning each rollout into a top-level row.
415fn codex_descendant_statuses(
416    live: &HashMap<PathBuf, CodexPeerStatus>,
417) -> HashMap<String, CodexPeerStatus> {
418    let lineage = live
419        .iter()
420        .filter_map(|(path, status)| {
421            rollout_lineage(path).map(|(session_id, parent)| (session_id, parent, *status))
422        })
423        .collect::<Vec<_>>();
424    let parent_by_child = lineage
425        .iter()
426        .filter_map(|(child, parent, _)| {
427            parent
428                .as_ref()
429                .map(|parent| (child.clone(), parent.clone()))
430        })
431        .collect::<HashMap<_, _>>();
432    let mut statuses = HashMap::new();
433    for (_, parent, status) in lineage {
434        let Some(mut root) = parent else {
435            continue;
436        };
437        let mut visited = std::collections::HashSet::new();
438        while visited.insert(root.clone()) {
439            let Some(parent) = parent_by_child.get(&root) else {
440                break;
441            };
442            root = parent.clone();
443        }
444        statuses
445            .entry(root)
446            .and_modify(|current| *current = stronger_codex_status(*current, status))
447            .or_insert(status);
448    }
449    statuses
450}
451
452fn merge_codex_status(
453    own: Option<CodexPeerStatus>,
454    descendant: Option<CodexPeerStatus>,
455) -> Option<CodexPeerStatus> {
456    match (own, descendant) {
457        (Some(own), Some(descendant)) => Some(stronger_codex_status(own, descendant)),
458        (Some(status), None) | (None, Some(status)) => Some(status),
459        (None, None) => None,
460    }
461}
462
463fn stronger_codex_status(left: CodexPeerStatus, right: CodexPeerStatus) -> CodexPeerStatus {
464    fn rank(status: CodexPeerStatus) -> u8 {
465        match status {
466            CodexPeerStatus::Running => 0,
467            CodexPeerStatus::Idle => 1,
468            CodexPeerStatus::Busy => 2,
469        }
470    }
471    if rank(right) > rank(left) {
472        right
473    } else {
474        left
475    }
476}
477
478#[cfg(feature = "adapter-api")]
479fn owned_activity(
480    locator: &SessionLocator,
481    state: crate::RuntimeRegistryState,
482    observed_at_ms: u64,
483) -> SessionActivity {
484    use crate::RuntimeRegistryState;
485    let (presence, turn) = match state {
486        RuntimeRegistryState::Persisted => (SessionPresence::Persisted, SessionTurnState::Unknown),
487        RuntimeRegistryState::Idle => (SessionPresence::Running, SessionTurnState::Idle),
488        RuntimeRegistryState::Busy => (SessionPresence::Running, SessionTurnState::Working),
489        RuntimeRegistryState::ShuttingDown => {
490            (SessionPresence::ShuttingDown, SessionTurnState::Unknown)
491        }
492    };
493    activity(
494        locator,
495        presence,
496        turn,
497        "supercode_runtime",
498        Some(state.as_str()),
499        None,
500        observed_at_ms,
501    )
502}
503
504fn claude_activity(
505    locator: &SessionLocator,
506    peer: &ClaudePeerSession,
507    observed_at_ms: u64,
508) -> SessionActivity {
509    let turn = match peer.status {
510        Some(ClaudePeerStatus::Busy) => SessionTurnState::Working,
511        Some(ClaudePeerStatus::Idle) => SessionTurnState::Idle,
512        Some(ClaudePeerStatus::Waiting) => SessionTurnState::NeedsInput,
513        None => SessionTurnState::Unknown,
514    };
515    activity(
516        locator,
517        SessionPresence::Running,
518        turn,
519        "claude_registry",
520        peer.status.map(|status| status.as_str()),
521        peer.version.as_deref(),
522        observed_at_ms,
523    )
524}
525
526fn codex_activity(
527    locator: &SessionLocator,
528    status: CodexPeerStatus,
529    observed_at_ms: u64,
530) -> SessionActivity {
531    let turn = match status {
532        CodexPeerStatus::Running => SessionTurnState::Unknown,
533        CodexPeerStatus::Idle => SessionTurnState::Idle,
534        CodexPeerStatus::Busy => SessionTurnState::Working,
535    };
536    activity(
537        locator,
538        SessionPresence::Running,
539        turn,
540        "codex_rollout",
541        Some(status.as_str()),
542        None,
543        observed_at_ms,
544    )
545}
546
547fn persisted_activity(locator: &SessionLocator, observed_at_ms: u64) -> SessionActivity {
548    activity(
549        locator,
550        SessionPresence::Persisted,
551        SessionTurnState::Unknown,
552        "persisted_store",
553        None,
554        None,
555        observed_at_ms,
556    )
557}
558
559fn activity(
560    locator: &SessionLocator,
561    presence: SessionPresence,
562    turn: SessionTurnState,
563    source: &str,
564    native_state: Option<&str>,
565    harness_version: Option<&str>,
566    observed_at_ms: u64,
567) -> SessionActivity {
568    SessionActivity {
569        harness: locator.harness.clone(),
570        session_id: locator.session_id.clone(),
571        presence,
572        turn,
573        evidence: SessionActivityEvidence {
574            source: source.to_string(),
575            native_state: native_state.map(str::to_string),
576            observed_at_ms,
577            harness_version: harness_version.map(str::to_string),
578        },
579    }
580}
581
582fn now_ms() -> u64 {
583    SystemTime::now()
584        .duration_since(UNIX_EPOCH)
585        .unwrap_or_default()
586        .as_millis()
587        .try_into()
588        .unwrap_or(u64::MAX)
589}