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