Skip to main content

lash_core/runtime/process/
observation.rs

1use std::collections::BTreeSet;
2use std::sync::Arc;
3
4use serde::{Deserialize, Serialize};
5
6use crate::plugin::PluginError;
7
8use super::events::{ProcessAwaitOutput, ProcessEvent};
9use super::model::{
10    AbandonRequest, ProcessExecutionEnvRef, ProcessExternalRef, ProcessHandleDescriptor, ProcessId,
11    ProcessIdentity, ProcessInput, ProcessLease, ProcessLifecycleStatus, ProcessListFilter,
12    ProcessOriginator, ProcessRecord, ProcessStarted, ProcessStatusFilter, RecoveryDisposition,
13    SessionScope, WaitState,
14};
15use super::registry::ProcessRegistry;
16use super::time::epoch_ms_from_system_time;
17
18#[derive(Clone)]
19pub struct ProcessWorkObserver {
20    registry: Arc<dyn ProcessRegistry>,
21}
22
23#[derive(Clone, Debug, Serialize, Deserialize)]
24pub struct ProcessWorkSnapshot {
25    pub session_id: String,
26    pub visible_process_ids: Vec<ProcessId>,
27    pub items: Vec<ObservedWorkItem>,
28}
29
30#[derive(Clone, Debug, Serialize, Deserialize)]
31pub struct ObservedWorkItem {
32    pub process: ObservedProcess,
33    pub descriptor: ProcessHandleDescriptor,
34    pub events: Vec<ObservedProcessEvent>,
35    pub kind: String,
36    pub label: String,
37}
38
39#[derive(Clone, Debug, Serialize, Deserialize)]
40pub struct ObservedProcess {
41    pub process_id: ProcessId,
42    pub graph_key: String,
43    pub kind: String,
44    pub lifecycle: ProcessLifecycleStatus,
45    pub identity: ProcessIdentity,
46    pub status_label: String,
47    pub terminal: bool,
48    /// Declared recovery contract (ADR 0019). Raw fact; hosts classify.
49    pub disposition: RecoveryDisposition,
50    #[serde(default, skip_serializing_if = "Option::is_none")]
51    pub error: Option<String>,
52    pub created_at_ms: u64,
53    pub updated_at_ms: u64,
54    /// Durable execution-started fact, if the row has begun executing.
55    #[serde(default, skip_serializing_if = "Option::is_none")]
56    pub first_started: Option<ProcessStarted>,
57    /// Current lease holder identity, if the row is leased (ADR 0019). Raw
58    /// fact for host-side staleness classification — no derived "stuck" verdict.
59    #[serde(default, skip_serializing_if = "Option::is_none")]
60    pub lease_holder: Option<crate::LeaseOwnerIdentity>,
61    /// Current lease expiry, paired with `lease_holder`.
62    #[serde(default, skip_serializing_if = "Option::is_none")]
63    pub lease_expires_at_ms: Option<u64>,
64    /// Pending Abandon Request the sweep reconciles once the lease lapses.
65    #[serde(default, skip_serializing_if = "Option::is_none")]
66    pub abandon_request: Option<AbandonRequest>,
67    pub input: ProcessInput,
68    pub originator: ProcessOriginator,
69    #[serde(default, skip_serializing_if = "Option::is_none")]
70    pub env_ref: Option<ProcessExecutionEnvRef>,
71    #[serde(default, skip_serializing_if = "Option::is_none")]
72    pub wake_target: Option<SessionScope>,
73    #[serde(default, skip_serializing_if = "Option::is_none")]
74    pub caused_by: Option<crate::CausalRef>,
75    #[serde(default, skip_serializing_if = "Option::is_none")]
76    pub external_ref: Option<ProcessExternalRef>,
77    #[serde(default, skip_serializing_if = "Option::is_none")]
78    pub wait: Option<WaitState>,
79    #[serde(default, skip_serializing_if = "Option::is_none")]
80    pub child_session_id: Option<String>,
81    pub label: String,
82}
83
84#[derive(Clone, Debug, Serialize, Deserialize)]
85pub struct ObservedProcessEvent {
86    pub sequence: u64,
87    pub event_type: String,
88    pub occurred_at_ms: u64,
89    pub payload: serde_json::Value,
90}
91
92/// Per-item event tail in session snapshots. Snapshots are polled by
93/// docks/UIs, so per-poll cost must stay bounded instead of growing with a
94/// process's full event history; detail views page through `events_after`
95/// with a cursor.
96pub const SNAPSHOT_EVENT_TAIL: usize = 32;
97
98impl ProcessWorkObserver {
99    pub fn new(registry: Arc<dyn ProcessRegistry>) -> Self {
100        Self { registry }
101    }
102
103    pub async fn snapshot_for_session(
104        &self,
105        session_id: impl Into<String>,
106    ) -> Result<ProcessWorkSnapshot, PluginError> {
107        let session_id = session_id.into();
108        let session_scope = SessionScope::new(session_id.clone());
109        let entries = self.registry.list_handle_grants(&session_scope).await?;
110        let mut items = Vec::new();
111        let mut seen_process_ids = BTreeSet::new();
112        for (grant, record) in entries {
113            seen_process_ids.insert(record.id.clone());
114            items.push(self.work_item_from_record(record, grant.descriptor).await?);
115        }
116        let visible_records = self
117            .registry
118            .list_processes(&ProcessListFilter {
119                status: ProcessStatusFilter::Any,
120                ..ProcessListFilter::default()
121            })
122            .await?;
123        for record in visible_records {
124            if seen_process_ids.contains(&record.id)
125                || !process_visible_to_session(&record, &session_id)
126            {
127                continue;
128            }
129            seen_process_ids.insert(record.id.clone());
130            let descriptor = descriptor_from_process_identity(&record.identity);
131            items.push(self.work_item_from_record(record, descriptor).await?);
132        }
133        items.sort_by(|left, right| {
134            right
135                .process
136                .updated_at_ms
137                .cmp(&left.process.updated_at_ms)
138                .then_with(|| right.process.created_at_ms.cmp(&left.process.created_at_ms))
139                .then_with(|| left.process.process_id.cmp(&right.process.process_id))
140        });
141        let visible_process_ids = items
142            .iter()
143            .map(|item| item.process.process_id.clone())
144            .collect();
145        Ok(ProcessWorkSnapshot {
146            session_id,
147            visible_process_ids,
148            items,
149        })
150    }
151
152    /// Snapshot every process matching `filter`, including the bounded event
153    /// tail used by host work rails. Unlike [`Self::snapshot_for_session`],
154    /// this is the runtime-wide observation surface: it does not depend on a
155    /// session handle grant and therefore continues to expose processes whose
156    /// originating session has been deleted.
157    /// Because grants are bypassed, the host must authorize access; routing
158    /// identity is not authorization.
159    pub async fn snapshot_all(
160        &self,
161        filter: &ProcessListFilter,
162    ) -> Result<Vec<ObservedWorkItem>, PluginError> {
163        let records = self.registry.list_processes(filter).await?;
164        let mut items = Vec::with_capacity(records.len());
165        for record in records {
166            let descriptor = descriptor_from_process_identity(&record.identity);
167            items.push(self.work_item_from_record(record, descriptor).await?);
168        }
169        items.sort_by(|left, right| {
170            right
171                .process
172                .updated_at_ms
173                .cmp(&left.process.updated_at_ms)
174                .then_with(|| right.process.created_at_ms.cmp(&left.process.created_at_ms))
175                .then_with(|| left.process.process_id.cmp(&right.process.process_id))
176        });
177        Ok(items)
178    }
179
180    async fn work_item_from_record(
181        &self,
182        record: ProcessRecord,
183        descriptor: ProcessHandleDescriptor,
184    ) -> Result<ObservedWorkItem, PluginError> {
185        let events = self
186            .registry
187            .recent_events(&record.id, SNAPSHOT_EVENT_TAIL)
188            .await?
189            .into_iter()
190            .map(ObservedProcessEvent::from)
191            .collect();
192        let lease = self.registry.get_process_lease(&record.id).await?;
193        let process = ObservedProcess::from_record(record, lease);
194        let kind = process.identity.kind.clone();
195        let label = process
196            .identity
197            .label
198            .clone()
199            .or_else(|| descriptor.label.clone())
200            .unwrap_or_else(|| kind.clone());
201        Ok(ObservedWorkItem {
202            process,
203            descriptor,
204            events,
205            kind,
206            label,
207        })
208    }
209
210    pub async fn process(&self, process_id: &str) -> Option<ObservedProcess> {
211        let record = self.registry.get_process(process_id).await?;
212        let lease = self
213            .registry
214            .get_process_lease(process_id)
215            .await
216            .ok()
217            .flatten();
218        Some(ObservedProcess::from_record(record, lease))
219    }
220
221    pub async fn list(
222        &self,
223        filter: &ProcessListFilter,
224    ) -> Result<Vec<ObservedProcess>, PluginError> {
225        let records = self.registry.list_processes(filter).await?;
226        self.observe_records(records).await
227    }
228
229    /// List processes a session may address — the grant filter (ADR 0019 /
230    /// process design grill). "Granted to" is the security lens: a process is
231    /// visible here only if `scope` holds a handle grant for it. This is the
232    /// single home for the grant-scoped view; the session facade sugar is a thin
233    /// caller of this method, never a parallel implementation.
234    pub async fn list_granted_to(
235        &self,
236        scope: &SessionScope,
237        filter: &ProcessListFilter,
238    ) -> Result<Vec<ObservedProcess>, PluginError> {
239        let entries = self.registry.list_handle_grants(scope).await?;
240        let records = entries
241            .into_iter()
242            .map(|(_, record)| record)
243            .filter(|record| filter.matches_record(record))
244            .collect::<Vec<_>>();
245        self.observe_records(records).await
246    }
247
248    /// List processes a session originated — the provenance filter (ADR 0019 /
249    /// process design grill). "Originated by" is the lineage lens, distinct from
250    /// the grant lens: a process matches when its recorded originator is a
251    /// session whose id equals `scope.session_id` (and its agent frame, when
252    /// `scope` names one), regardless of who currently holds a grant.
253    pub async fn list_originated_by(
254        &self,
255        scope: &SessionScope,
256        filter: &ProcessListFilter,
257    ) -> Result<Vec<ObservedProcess>, PluginError> {
258        let records = self
259            .registry
260            .list_processes(filter)
261            .await?
262            .into_iter()
263            .filter(|record| originator_matches(&record.provenance.originator, scope))
264            .collect::<Vec<_>>();
265        self.observe_records(records).await
266    }
267
268    async fn observe_records(
269        &self,
270        records: Vec<ProcessRecord>,
271    ) -> Result<Vec<ObservedProcess>, PluginError> {
272        let mut observed = Vec::with_capacity(records.len());
273        for record in records {
274            let lease = self.registry.get_process_lease(&record.id).await?;
275            observed.push(ObservedProcess::from_record(record, lease));
276        }
277        Ok(observed)
278    }
279
280    pub async fn events_after(
281        &self,
282        process_id: &str,
283        after_sequence: u64,
284    ) -> Result<Vec<ObservedProcessEvent>, PluginError> {
285        Ok(self
286            .registry
287            .events_after(process_id, after_sequence)
288            .await?
289            .into_iter()
290            .map(ObservedProcessEvent::from)
291            .collect())
292    }
293}
294
295impl ObservedProcess {
296    /// Build a read-side view of a process. `lease` is the current lease row (if
297    /// any), read separately so the observer exposes holder identity and expiry
298    /// as raw facts — no derived "stuck" classification (ADR 0019).
299    fn from_record(record: ProcessRecord, lease: Option<ProcessLease>) -> Self {
300        let lifecycle = ProcessLifecycleStatus::from(&record.status);
301        let input = record.input.as_ref().clone();
302        let identity = record.identity;
303        let kind = identity.kind.clone();
304        let label = identity.label.clone().unwrap_or_else(|| kind.clone());
305        let process_id = record.id;
306        let (lease_holder, lease_expires_at_ms) = match lease {
307            Some(lease) => (Some(lease.owner), Some(lease.expires_at_epoch_ms)),
308            None => (None, None),
309        };
310        Self {
311            graph_key: format!("process:{process_id}"),
312            process_id,
313            kind,
314            lifecycle,
315            identity,
316            status_label: lifecycle.label().to_string(),
317            terminal: lifecycle.is_terminal(),
318            disposition: record.disposition,
319            error: terminal_error(&record.status),
320            created_at_ms: record.created_at_ms,
321            updated_at_ms: record.updated_at_ms,
322            first_started: record.first_started.map(|started| *started),
323            lease_holder,
324            lease_expires_at_ms,
325            abandon_request: record.abandon_request.map(|request| *request),
326            originator: record.provenance.originator,
327            env_ref: record.env_ref,
328            wake_target: record.wake_target,
329            caused_by: record.provenance.caused_by,
330            external_ref: record.external_ref,
331            wait: record.wait,
332            child_session_id: child_session_id(&input),
333            input,
334            label,
335        }
336    }
337}
338
339impl From<ProcessEvent> for ObservedProcessEvent {
340    fn from(event: ProcessEvent) -> Self {
341        Self {
342            sequence: event.sequence,
343            event_type: event.event_type,
344            occurred_at_ms: epoch_ms_from_system_time(event.occurred_at),
345            payload: event.payload,
346        }
347    }
348}
349
350fn terminal_error(status: &super::model::ProcessStatus) -> Option<String> {
351    match status.await_output()? {
352        ProcessAwaitOutput::Failure { message, .. }
353        | ProcessAwaitOutput::Cancelled { message, .. } => Some(message.clone()),
354        // Abandonment is not a reported failure; the status label conveys it and
355        // the evidence rides the terminal event. No derived error string here.
356        ProcessAwaitOutput::Success { .. } | ProcessAwaitOutput::Abandoned { .. } => None,
357    }
358}
359
360fn child_session_id(input: &ProcessInput) -> Option<String> {
361    match input {
362        ProcessInput::SessionTurn { create_request, .. } => create_request.session_id.clone(),
363        ProcessInput::ToolCall { .. }
364        | ProcessInput::Engine { .. }
365        | ProcessInput::External { .. } => None,
366    }
367}
368
369/// Whether `originator` names the session (or session+frame) identified by
370/// `scope`. Frame is matched only when `scope` names one, so a session-level
371/// provenance filter captures every frame the session originated.
372fn originator_matches(originator: &ProcessOriginator, scope: &SessionScope) -> bool {
373    match originator {
374        ProcessOriginator::Host { .. } => false,
375        ProcessOriginator::Session {
376            scope: origin_scope,
377        } => {
378            origin_scope.session_id == scope.session_id
379                && (scope.agent_frame_id.is_none()
380                    || origin_scope.agent_frame_id == scope.agent_frame_id)
381        }
382    }
383}
384
385fn process_visible_to_session(record: &ProcessRecord, session_id: &str) -> bool {
386    record
387        .wake_target
388        .as_ref()
389        .is_some_and(|scope| scope.session_id == session_id)
390}
391
392fn descriptor_from_process_identity(identity: &ProcessIdentity) -> ProcessHandleDescriptor {
393    ProcessHandleDescriptor::new(Some(identity.kind.clone()), identity.label.clone())
394}
395
396#[cfg(test)]
397mod tests {
398    use std::sync::Arc;
399    use std::time::Duration;
400
401    use serde_json::json;
402
403    use super::*;
404    use crate::{
405        InputItem, PluginOptions, PreparedToolCall, ProcessEventAppendRequest,
406        ProcessExecutionEnvRef, ProcessIdentity, ProcessProvenance, ProcessRegistration,
407        SessionCreateRequest, SessionScope, SessionStartPoint, SubagentSessionContext,
408        ToolFailureClass, ToolOutputContract, TurnInput, WaitKind,
409    };
410
411    fn observer(registry: Arc<dyn ProcessRegistry>) -> ProcessWorkObserver {
412        ProcessWorkObserver::new(registry)
413    }
414
415    fn external_registration(process_id: &str, label: &str) -> ProcessRegistration {
416        ProcessRegistration::new(
417            process_id,
418            ProcessInput::External {
419                metadata: json!({ "label": label }),
420            },
421            RecoveryDisposition::ExternallyOwned,
422            ProcessProvenance::host(),
423        )
424    }
425
426    async fn register_visible(
427        registry: &Arc<dyn ProcessRegistry>,
428        scope: &SessionScope,
429        registration: ProcessRegistration,
430        descriptor: ProcessHandleDescriptor,
431    ) {
432        let process_id = registration.id.clone();
433        registry
434            .register_process(registration)
435            .await
436            .expect("register process");
437        registry
438            .grant_handle(scope, &process_id, descriptor)
439            .await
440            .expect("grant process handle");
441    }
442
443    #[tokio::test]
444    async fn snapshot_for_session_reads_visible_grants_and_events_as_epoch_ms() {
445        let registry =
446            Arc::new(super::super::TestLocalProcessRegistry::default()) as Arc<dyn ProcessRegistry>;
447        let visible_scope = SessionScope::new("visible");
448        register_visible(
449            &registry,
450            &visible_scope,
451            external_registration("visible-process", "Visible"),
452            ProcessHandleDescriptor::new(Some("visible-kind"), Some("Visible descriptor")),
453        )
454        .await;
455        register_visible(
456            &registry,
457            &SessionScope::new("other"),
458            external_registration("hidden-process", "Hidden"),
459            ProcessHandleDescriptor::new(Some("hidden-kind"), Some("Hidden")),
460        )
461        .await;
462        registry
463            .append_event(
464                "visible-process",
465                ProcessEventAppendRequest::new("process.cancel_requested", json!({"why": "test"}))
466                    .with_replay_key("visible-process:cancel-requested"),
467            )
468            .await
469            .expect("append event");
470
471        let snapshot = observer(Arc::clone(&registry))
472            .snapshot_for_session("visible")
473            .await
474            .expect("snapshot");
475
476        assert_eq!(snapshot.session_id, "visible");
477        assert_eq!(snapshot.visible_process_ids, vec!["visible-process"]);
478        assert_eq!(snapshot.items.len(), 1);
479        assert_eq!(snapshot.items[0].events.len(), 1);
480        assert_eq!(
481            snapshot.items[0].events[0].event_type,
482            "process.cancel_requested"
483        );
484        assert!(snapshot.items[0].events[0].occurred_at_ms > 0);
485    }
486
487    #[tokio::test]
488    async fn runtime_snapshot_keeps_orphaned_processes_after_session_deletion() {
489        let registry =
490            Arc::new(super::super::TestLocalProcessRegistry::default()) as Arc<dyn ProcessRegistry>;
491        register_visible(
492            &registry,
493            &SessionScope::new("deleted-session"),
494            external_registration("surviving-process", "Survivor"),
495            ProcessHandleDescriptor::new(Some("test"), Some("Survivor")),
496        )
497        .await;
498
499        let report = registry
500            .delete_session_process_state("deleted-session")
501            .await
502            .expect("delete session process edges");
503        assert_eq!(report.orphaned_process_ids, vec!["surviving-process"]);
504        assert!(
505            observer(Arc::clone(&registry))
506                .snapshot_for_session("deleted-session")
507                .await
508                .expect("deleted session snapshot")
509                .items
510                .is_empty()
511        );
512
513        let runtime_items = observer(registry)
514            .snapshot_all(&ProcessListFilter {
515                status: super::super::ProcessStatusFilter::Any,
516                ..ProcessListFilter::default()
517            })
518            .await
519            .expect("runtime process snapshot");
520        assert_eq!(runtime_items.len(), 1);
521        assert_eq!(runtime_items[0].process.process_id, "surviving-process");
522    }
523
524    #[tokio::test]
525    async fn snapshot_for_session_includes_frame_wake_targets_without_handle_grants() {
526        let registry =
527            Arc::new(super::super::TestLocalProcessRegistry::default()) as Arc<dyn ProcessRegistry>;
528        let frame_scope = SessionScope::for_agent_frame("visible", "frame-a");
529        registry
530            .register_process(ProcessRegistration::new(
531                "frame-originated",
532                ProcessInput::External {
533                    metadata: json!({ "label": "Frame originated" }),
534                },
535                RecoveryDisposition::ExternallyOwned,
536                ProcessProvenance::session(frame_scope.clone()),
537            ))
538            .await
539            .expect("register frame-originated process");
540        registry
541            .register_process(
542                external_registration("frame-wake-targeted", "Frame wake targeted")
543                    .with_wake_target(Some(frame_scope)),
544            )
545            .await
546            .expect("register frame wake-targeted process");
547        registry
548            .register_process(
549                external_registration("hidden-frame", "Hidden")
550                    .with_wake_target(Some(SessionScope::for_agent_frame("other", "frame-b"))),
551            )
552            .await
553            .expect("register hidden process");
554
555        let snapshot = observer(Arc::clone(&registry))
556            .snapshot_for_session("visible")
557            .await
558            .expect("snapshot");
559        let visible_process_ids = snapshot
560            .visible_process_ids
561            .iter()
562            .cloned()
563            .collect::<std::collections::BTreeSet<_>>();
564
565        assert_eq!(
566            visible_process_ids,
567            std::collections::BTreeSet::from(["frame-wake-targeted".to_string()])
568        );
569        assert_eq!(snapshot.items.len(), 1);
570    }
571
572    #[tokio::test]
573    async fn snapshot_for_session_labels_engine_wake_targets_from_identity_without_handle_grants() {
574        let registry =
575            Arc::new(super::super::TestLocalProcessRegistry::default()) as Arc<dyn ProcessRegistry>;
576        let scope = SessionScope::new("visible");
577        registry
578            .register_process(
579                ProcessRegistration::new(
580                    "engine-wake-targeted",
581                    ProcessInput::Engine {
582                        kind: "test-engine".to_string(),
583                        payload: json!({}),
584                    },
585                    RecoveryDisposition::Rerunnable,
586                    ProcessProvenance::host(),
587                )
588                .with_identity(
589                    ProcessIdentity::new("test-engine").with_label(Some("remember".to_string())),
590                )
591                .with_execution_env_ref(Some(ProcessExecutionEnvRef::new("process-env:test")))
592                .with_wake_target(Some(scope)),
593            )
594            .await
595            .expect("register engine wake-targeted process");
596
597        let snapshot = observer(Arc::clone(&registry))
598            .snapshot_for_session("visible")
599            .await
600            .expect("snapshot");
601
602        assert_eq!(snapshot.items.len(), 1);
603        assert_eq!(snapshot.items[0].kind, "test-engine");
604        assert_eq!(snapshot.items[0].label, "remember");
605        assert_eq!(
606            snapshot.items[0].descriptor.kind.as_deref(),
607            Some("test-engine")
608        );
609        assert_eq!(
610            snapshot.items[0].descriptor.label.as_deref(),
611            Some("remember")
612        );
613        assert_eq!(snapshot.items[0].process.kind, "test-engine");
614        assert_eq!(snapshot.items[0].process.label, "remember");
615    }
616
617    #[tokio::test]
618    async fn snapshot_for_session_sorts_work_by_updated_then_created_descending() {
619        let registry =
620            Arc::new(super::super::TestLocalProcessRegistry::default()) as Arc<dyn ProcessRegistry>;
621        let scope = SessionScope::new("sort");
622        register_visible(
623            &registry,
624            &scope,
625            external_registration("older", "Older"),
626            ProcessHandleDescriptor::new(None::<String>, None::<String>),
627        )
628        .await;
629        tokio::time::sleep(Duration::from_millis(2)).await;
630        register_visible(
631            &registry,
632            &scope,
633            external_registration("newer", "Newer"),
634            ProcessHandleDescriptor::new(None::<String>, None::<String>),
635        )
636        .await;
637        tokio::time::sleep(Duration::from_millis(2)).await;
638        registry
639            .append_event(
640                "older",
641                ProcessEventAppendRequest::new("process.cancel_requested", json!({}))
642                    .with_replay_key("older:cancel-requested"),
643            )
644            .await
645            .expect("update older process");
646
647        let snapshot = observer(Arc::clone(&registry))
648            .snapshot_for_session("sort")
649            .await
650            .expect("snapshot");
651
652        assert_eq!(snapshot.visible_process_ids, vec!["older", "newer"]);
653    }
654
655    #[tokio::test]
656    async fn observed_process_reports_terminal_status_and_error_messages() {
657        let registry =
658            Arc::new(super::super::TestLocalProcessRegistry::default()) as Arc<dyn ProcessRegistry>;
659        for process_id in ["failed", "cancelled"] {
660            registry
661                .register_process(external_registration(process_id, process_id))
662                .await
663                .expect("register");
664        }
665        registry
666            .complete_process(
667                "failed",
668                ProcessAwaitOutput::Failure {
669                    class: ToolFailureClass::External,
670                    code: "boom".to_string(),
671                    message: "failed loudly".to_string(),
672                    raw: None,
673                    control: None,
674                },
675                crate::ProcessCompletionAuthority::external_owner("test"),
676            )
677            .await
678            .expect("fail process");
679        registry
680            .complete_process(
681                "cancelled",
682                ProcessAwaitOutput::Cancelled {
683                    message: "cancelled intentionally".to_string(),
684                    raw: None,
685                    control: None,
686                },
687                crate::ProcessCompletionAuthority::external_owner("test"),
688            )
689            .await
690            .expect("cancel process");
691
692        let observer = observer(Arc::clone(&registry));
693        let failed = observer.process("failed").await.expect("failed process");
694        let cancelled = observer
695            .process("cancelled")
696            .await
697            .expect("cancelled process");
698
699        assert_eq!(failed.status_label, "failed");
700        assert!(failed.terminal);
701        assert_eq!(failed.error.as_deref(), Some("failed loudly"));
702        assert_eq!(cancelled.status_label, "cancelled");
703        assert!(cancelled.terminal);
704        assert_eq!(cancelled.error.as_deref(), Some("cancelled intentionally"));
705    }
706
707    #[tokio::test]
708    async fn observed_process_exposes_current_wait_state() {
709        let registry =
710            Arc::new(super::super::TestLocalProcessRegistry::default()) as Arc<dyn ProcessRegistry>;
711        let scope = SessionScope::new("wait");
712        register_visible(
713            &registry,
714            &scope,
715            external_registration("waiting-process", "Waiting"),
716            ProcessHandleDescriptor::new(Some("external"), Some("Waiting")),
717        )
718        .await;
719        let wait = WaitState {
720            since_ms: 1234,
721            kind: WaitKind::Signal {
722                name: "ready".to_string(),
723                event_type: "signal.ready".to_string(),
724                key: "process:waiting-process:signal.ready:1".to_string(),
725                ordinal: 1,
726            },
727        };
728        registry
729            .set_process_wait("waiting-process", wait.clone())
730            .await
731            .expect("set wait");
732
733        let observer = observer(Arc::clone(&registry));
734        let observed = observer
735            .process("waiting-process")
736            .await
737            .expect("waiting process");
738        let snapshot = observer
739            .snapshot_for_session("wait")
740            .await
741            .expect("snapshot");
742
743        assert_eq!(observed.wait, Some(wait.clone()));
744        assert_eq!(snapshot.items.len(), 1);
745        assert_eq!(snapshot.items[0].process.wait, Some(wait));
746    }
747
748    #[tokio::test]
749    async fn snapshot_for_session_prefers_typed_labels_and_extracts_child_session_id() {
750        let registry =
751            Arc::new(super::super::TestLocalProcessRegistry::default()) as Arc<dyn ProcessRegistry>;
752        let scope = SessionScope::new("labels");
753        let mut child_request = SessionCreateRequest::child_session(
754            "labels",
755            SessionStartPoint::Empty,
756            PluginOptions::default(),
757        )
758        .with_session_id("child-session");
759        child_request.subagent = Some(SubagentSessionContext {
760            parent_session_id: "labels".to_string(),
761            capability: "researcher".to_string(),
762            depth: 1,
763            max_depth: 4,
764        });
765        let cases = [
766            (
767                "tool",
768                ProcessInput::ToolCall {
769                    call: PreparedToolCall::from_parts(
770                        "call-1",
771                        "tool:shell.run",
772                        "shell.run",
773                        json!({}),
774                        None,
775                        serde_json::Value::Null,
776                    ),
777                },
778                "tool",
779                "shell.run",
780                None,
781            ),
782            (
783                "engine",
784                ProcessInput::Engine {
785                    kind: "test-engine".to_string(),
786                    payload: json!({}),
787                },
788                "test-engine",
789                "remember",
790                None,
791            ),
792            (
793                "session",
794                ProcessInput::SessionTurn {
795                    create_request: Box::new(child_request),
796                    turn_input: Box::new(TurnInput::items([InputItem::text("run child")])),
797                    output_contract: ToolOutputContract::Static,
798                },
799                "session_turn",
800                "researcher",
801                Some("child-session"),
802            ),
803            (
804                "external",
805                ProcessInput::External {
806                    metadata: json!({ "label": "external job" }),
807                },
808                "external",
809                "external job",
810                None,
811            ),
812        ];
813        for (process_id, input, kind, label, _child_session_id) in cases {
814            let needs_env = matches!(
815                input,
816                ProcessInput::ToolCall { .. } | ProcessInput::Engine { .. }
817            );
818            let disposition = match input {
819                ProcessInput::External { .. } => RecoveryDisposition::ExternallyOwned,
820                _ => RecoveryDisposition::Rerunnable,
821            };
822            let mut registration =
823                ProcessRegistration::new(process_id, input, disposition, ProcessProvenance::host())
824                    .with_identity(ProcessIdentity::new(kind).with_label(Some(label.to_string())));
825            if needs_env {
826                registration = registration.with_execution_env_ref(Some(
827                    ProcessExecutionEnvRef::new(format!("process-env:test:{process_id}")),
828                ));
829            }
830            register_visible(
831                &registry,
832                &scope,
833                registration,
834                ProcessHandleDescriptor::new(Some("descriptor-kind"), Some("Descriptor label")),
835            )
836            .await;
837        }
838
839        let snapshot = observer(Arc::clone(&registry))
840            .snapshot_for_session("labels")
841            .await
842            .expect("snapshot");
843        let by_id = snapshot
844            .items
845            .iter()
846            .map(|item| (item.process.process_id.as_str(), item))
847            .collect::<std::collections::BTreeMap<_, _>>();
848
849        assert_eq!(by_id["tool"].label, "shell.run");
850        assert_eq!(by_id["engine"].label, "remember");
851        assert_eq!(by_id["engine"].process.kind, "test-engine");
852        assert_eq!(by_id["session"].label, "researcher");
853        assert_eq!(
854            by_id["session"].process.child_session_id.as_deref(),
855            Some("child-session")
856        );
857        assert_eq!(by_id["external"].label, "external job");
858    }
859
860    #[tokio::test]
861    async fn observed_process_missing_lookup_returns_none() {
862        let registry =
863            Arc::new(super::super::TestLocalProcessRegistry::default()) as Arc<dyn ProcessRegistry>;
864
865        assert!(observer(registry).process("missing").await.is_none());
866    }
867}