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 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 #[serde(default, skip_serializing_if = "Option::is_none")]
56 pub first_started: Option<ProcessStarted>,
57 #[serde(default, skip_serializing_if = "Option::is_none")]
60 pub lease_holder: Option<crate::LeaseOwnerIdentity>,
61 #[serde(default, skip_serializing_if = "Option::is_none")]
63 pub lease_expires_at_ms: Option<u64>,
64 #[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
92pub 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 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 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 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 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 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
369fn 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 ®istry,
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 ®istry,
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(®istry))
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 ®istry,
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(®istry))
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(®istry))
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(®istry))
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 ®istry,
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 ®istry,
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(®istry))
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(®istry));
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 ®istry,
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(®istry));
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 ®istry,
832 &scope,
833 registration,
834 ProcessHandleDescriptor::new(Some("descriptor-kind"), Some("Descriptor label")),
835 )
836 .await;
837 }
838
839 let snapshot = observer(Arc::clone(®istry))
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}