1use super::*;
2use meerkat_core::ToolName;
3
4#[derive(Default)]
15pub struct MachineToolVisibilityOwner {
16 pub state: StdRwLock<SessionToolVisibilityState>,
17 dsl_authority: StdRwLock<Option<Arc<std::sync::Mutex<super::dsl::MeerkatMachineAuthority>>>>,
27}
28
29impl std::fmt::Debug for MachineToolVisibilityOwner {
30 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
31 f.debug_struct("MachineToolVisibilityOwner")
32 .field("state", &"<StdRwLock<SessionToolVisibilityState>>")
33 .field(
34 "dsl_authority",
35 &self
36 .dsl_authority
37 .read()
38 .ok()
39 .and_then(|slot| slot.as_ref().map(|_| "bound"))
40 .unwrap_or("unbound"),
41 )
42 .finish()
43 }
44}
45
46impl MachineToolVisibilityOwner {
47 pub fn new() -> Self {
48 Self::default()
49 }
50
51 pub fn bind_dsl_authority(
56 &self,
57 authority: Arc<std::sync::Mutex<super::dsl::MeerkatMachineAuthority>>,
58 ) {
59 let mut slot = self
60 .dsl_authority
61 .write()
62 .unwrap_or_else(std::sync::PoisonError::into_inner);
63 *slot = Some(authority);
64 }
65
66 fn dsl_authority_for_stage(
67 &self,
68 ) -> Result<Arc<std::sync::Mutex<super::dsl::MeerkatMachineAuthority>>, ToolScopeStageError>
69 {
70 let slot = self
71 .dsl_authority
72 .read()
73 .map_err(|_| ToolScopeStageError::Owner {
74 message: "machine visibility DSL authority slot lock poisoned".to_string(),
75 })?;
76 let authority = slot
77 .as_ref()
78 .cloned()
79 .ok_or_else(|| ToolScopeStageError::Owner {
80 message:
81 "machine visibility DSL authority not bound — staging call before session wiring"
82 .to_string(),
83 })?;
84 Ok(authority)
85 }
86
87 fn dsl_authority_for_apply(
88 &self,
89 ) -> Result<Arc<std::sync::Mutex<super::dsl::MeerkatMachineAuthority>>, ToolScopeApplyError>
90 {
91 let slot = self
92 .dsl_authority
93 .read()
94 .map_err(|_| ToolScopeApplyError::Owner {
95 message: "machine visibility DSL authority slot lock poisoned".to_string(),
96 })?;
97 let authority = slot
98 .as_ref()
99 .cloned()
100 .ok_or_else(|| ToolScopeApplyError::Owner {
101 message:
102 "machine visibility DSL authority not bound — apply call before session wiring"
103 .to_string(),
104 })?;
105 Ok(authority)
106 }
107}
108
109pub fn formal_projection_value<T: serde::Serialize>(value: &T) -> String {
110 serde_json::to_string(value).unwrap_or_else(|err| format!("\"<serialization error: {err}>\""))
111}
112
113fn authority_witnesses_for_names(
114 names: &std::collections::BTreeSet<ToolName>,
115 witnesses: &std::collections::BTreeMap<ToolName, ToolVisibilityWitness>,
116) -> std::collections::BTreeMap<ToolName, ToolVisibilityWitness> {
117 names
118 .iter()
119 .filter_map(|name| {
120 witnesses
121 .get(name)
122 .map(|witness| (name.clone(), witness.clone()))
123 })
124 .collect()
125}
126
127fn deferred_load_authority_map(
128 authorities: &[DeferredToolLoadAuthority],
129) -> Result<std::collections::BTreeMap<ToolName, ToolVisibilityWitness>, ToolScopeStageError> {
130 let mut by_name = std::collections::BTreeMap::new();
131 let mut invalid = Vec::new();
132
133 for authority in authorities {
134 match by_name.insert(authority.name.clone(), authority.witness.clone()) {
135 Some(existing) if existing != authority.witness => invalid.push(authority.name.clone()),
136 _ => {}
137 }
138 }
139
140 if invalid.is_empty() {
141 return Ok(by_name);
142 }
143
144 invalid.sort_unstable();
145 invalid.dedup();
146 Err(ToolScopeStageError::InvalidWitnesses { names: invalid })
147}
148
149fn dsl_witnesses(
150 witnesses: &std::collections::BTreeMap<ToolName, ToolVisibilityWitness>,
151) -> std::collections::BTreeMap<ToolName, super::dsl::ToolVisibilityWitness> {
152 witnesses
153 .iter()
154 .map(|(name, witness)| {
155 (
156 name.clone(),
157 super::dsl::ToolVisibilityWitness::from(witness),
158 )
159 })
160 .collect()
161}
162
163fn core_tool_source_kind(kind: super::dsl::ToolSourceKind) -> meerkat_core::types::ToolSourceKind {
164 match kind {
165 super::dsl::ToolSourceKind::Builtin => meerkat_core::types::ToolSourceKind::Builtin,
166 super::dsl::ToolSourceKind::Shell => meerkat_core::types::ToolSourceKind::Shell,
167 super::dsl::ToolSourceKind::Comms => meerkat_core::types::ToolSourceKind::Comms,
168 super::dsl::ToolSourceKind::Memory => meerkat_core::types::ToolSourceKind::Memory,
169 super::dsl::ToolSourceKind::Schedule => meerkat_core::types::ToolSourceKind::Schedule,
170 super::dsl::ToolSourceKind::WorkGraph => meerkat_core::types::ToolSourceKind::WorkGraph,
171 super::dsl::ToolSourceKind::Mob => meerkat_core::types::ToolSourceKind::Mob,
172 super::dsl::ToolSourceKind::Callback => meerkat_core::types::ToolSourceKind::Callback,
173 super::dsl::ToolSourceKind::Mcp => meerkat_core::types::ToolSourceKind::Mcp,
174 super::dsl::ToolSourceKind::RustBundle => meerkat_core::types::ToolSourceKind::RustBundle,
175 }
176}
177
178fn core_tool_visibility_witness(
179 witness: &super::dsl::ToolVisibilityWitness,
180) -> ToolVisibilityWitness {
181 ToolVisibilityWitness {
182 last_seen_provenance: witness.last_seen_provenance.as_ref().map(|provenance| {
183 meerkat_core::types::ToolProvenance {
184 kind: core_tool_source_kind(provenance.kind),
185 source_id: meerkat_core::types::ToolSourceId::new(provenance.source_id.clone()),
186 }
187 }),
188 }
189}
190
191fn core_witnesses(
192 witnesses: &std::collections::BTreeMap<ToolName, super::dsl::ToolVisibilityWitness>,
193) -> std::collections::BTreeMap<ToolName, ToolVisibilityWitness> {
194 witnesses
195 .iter()
196 .map(|(name, witness)| (name.clone(), core_tool_visibility_witness(witness)))
197 .collect()
198}
199
200fn mirror_visibility_projection_from_authority(
201 projection: &mut SessionToolVisibilityState,
202 authority: &super::dsl::MeerkatMachineAuthority,
203) {
204 let authority_state = authority.state();
205 projection.capability_base_filter = meerkat_core::ToolFilter::from(
206 authority_state
207 .current_session_capability_base_filter
208 .clone(),
209 );
210 projection.inherited_base_filter =
211 meerkat_core::ToolFilter::from(authority_state.inherited_base_filter.clone());
212 projection.active_filter =
213 meerkat_core::ToolFilter::from(authority_state.active_filter.clone());
214 projection.staged_filter =
215 meerkat_core::ToolFilter::from(authority_state.staged_filter.clone());
216 projection.active_requested_deferred_names = authority_state.active_deferred_names.clone();
217 projection.staged_requested_deferred_names = authority_state.staged_deferred_names.clone();
218 projection.active_revision = authority_state.active_visibility_revision;
219 projection.staged_revision = authority_state.staged_visibility_revision;
220 projection.requested_witnesses =
221 core_witnesses(&authority_state.requested_visibility_witnesses);
222 projection.filter_witnesses = core_witnesses(&authority_state.filter_visibility_witnesses);
223}
224
225impl ToolVisibilityOwner for MachineToolVisibilityOwner {
226 fn visibility_state(&self) -> Result<SessionToolVisibilityState, ToolScopeApplyError> {
227 self.state
228 .read()
229 .map(|state| state.clone())
230 .map_err(|_| ToolScopeApplyError::Owner {
231 message: "machine visibility state lock poisoned".to_string(),
232 })
233 }
234
235 fn replace_visibility_state(
236 &self,
237 visibility_state: SessionToolVisibilityState,
238 ) -> Result<(), ToolScopeApplyError> {
239 let active_deferred_authorities = dsl_witnesses(&authority_witnesses_for_names(
240 &visibility_state.active_requested_deferred_names,
241 &visibility_state.requested_witnesses,
242 ));
243 let staged_deferred_authorities = dsl_witnesses(&authority_witnesses_for_names(
244 &visibility_state.staged_requested_deferred_names,
245 &visibility_state.requested_witnesses,
246 ));
247 let dsl_visibility_state =
248 super::dsl::SessionToolVisibilityState::from_domain(&visibility_state);
249 let authority = self.dsl_authority_for_apply()?;
250 let mut state = self
251 .state
252 .write()
253 .unwrap_or_else(std::sync::PoisonError::into_inner);
254 let mut guard = authority
255 .lock()
256 .unwrap_or_else(std::sync::PoisonError::into_inner);
257 super::dsl::MeerkatMachineMutator::apply(
258 &mut *guard,
259 super::dsl::MeerkatMachineInput::ReplaceVisibilityState {
260 capability_base_filter: super::dsl::ToolFilter::from(
261 &visibility_state.capability_base_filter,
262 ),
263 inherited_base_filter: super::dsl::ToolFilter::from(
264 &visibility_state.inherited_base_filter,
265 ),
266 active_filter: super::dsl::ToolFilter::from(&visibility_state.active_filter),
267 staged_filter: super::dsl::ToolFilter::from(&visibility_state.staged_filter),
268 requested_witnesses: dsl_visibility_state.requested_witnesses.clone(),
269 filter_witnesses: dsl_visibility_state.filter_witnesses,
270 active_deferred_names: visibility_state.active_requested_deferred_names.clone(),
271 staged_deferred_names: visibility_state.staged_requested_deferred_names.clone(),
272 active_deferred_authorities,
273 staged_deferred_authorities,
274 active_revision: visibility_state.active_revision,
275 staged_revision: visibility_state.staged_revision,
276 },
277 )
278 .map_err(|err| ToolScopeApplyError::Owner {
279 message: super::dsl_authority::map_error(err, "ReplaceVisibilityState"),
280 })?;
281 mirror_visibility_projection_from_authority(&mut state, &guard);
282 Ok(())
283 }
284
285 fn stage_persistent_filter(
286 &self,
287 filter: ToolFilter,
288 witnesses: std::collections::BTreeMap<ToolName, ToolVisibilityWitness>,
289 ) -> Result<ToolScopeRevision, ToolScopeStageError> {
290 let authority = self.dsl_authority_for_stage()?;
291 let mut state = self.state.write().map_err(|_| ToolScopeStageError::Owner {
292 message: "machine visibility state lock poisoned".to_string(),
293 })?;
294 let mut guard = authority
295 .lock()
296 .unwrap_or_else(std::sync::PoisonError::into_inner);
297 let mut next_filter_witnesses = guard.state().filter_visibility_witnesses.clone();
298 next_filter_witnesses.extend(dsl_witnesses(&witnesses));
299 super::dsl::MeerkatMachineMutator::apply(
300 &mut *guard,
301 super::dsl::MeerkatMachineInput::StageVisibilityFilter {
302 filter: super::dsl::ToolFilter::from(&filter),
303 witnesses: next_filter_witnesses,
304 },
305 )
306 .map_err(|err| ToolScopeStageError::Owner {
307 message: super::dsl_authority::map_error(err, "StageVisibilityFilter"),
308 })?;
309 let revision = ToolScopeRevision(guard.state().staged_visibility_revision);
310 mirror_visibility_projection_from_authority(&mut state, &guard);
311 Ok(revision)
312 }
313
314 fn stage_requested_deferred_names(
315 &self,
316 names: std::collections::BTreeSet<ToolName>,
317 ) -> Result<ToolScopeRevision, ToolScopeStageError> {
318 if !names.is_empty() {
319 return Err(ToolScopeStageError::MissingWitnesses {
320 names: names.into_iter().collect(),
321 });
322 }
323 let authority = self.dsl_authority_for_stage()?;
324 let mut state = self.state.write().map_err(|_| ToolScopeStageError::Owner {
325 message: "machine visibility state lock poisoned".to_string(),
326 })?;
327 let mut guard = authority
328 .lock()
329 .unwrap_or_else(std::sync::PoisonError::into_inner);
330 super::dsl::MeerkatMachineMutator::apply(
331 &mut *guard,
332 super::dsl::MeerkatMachineInput::StageDeferredNames { names },
333 )
334 .map_err(|err| ToolScopeStageError::Owner {
335 message: super::dsl_authority::map_error(err, "StageDeferredNames"),
336 })?;
337 let revision = ToolScopeRevision(guard.state().staged_visibility_revision);
338 mirror_visibility_projection_from_authority(&mut state, &guard);
339 Ok(revision)
340 }
341
342 fn request_deferred_tools(
343 &self,
344 authorities: Vec<DeferredToolLoadAuthority>,
345 ) -> Result<ToolScopeRevision, ToolScopeStageError> {
346 let authorities = deferred_load_authority_map(&authorities)?;
347 if authorities.is_empty() {
348 return Err(ToolScopeStageError::Owner {
349 message: "deferred tool request requires at least one authority".to_string(),
350 });
351 }
352 let authority = self.dsl_authority_for_stage()?;
353 let mut state = self.state.write().map_err(|_| ToolScopeStageError::Owner {
354 message: "machine visibility state lock poisoned".to_string(),
355 })?;
356 let mut guard = authority
357 .lock()
358 .unwrap_or_else(std::sync::PoisonError::into_inner);
359 let mut target_authorities = guard.state().staged_deferred_authorities.clone();
360 target_authorities.extend(dsl_witnesses(&authorities));
361 super::dsl::MeerkatMachineMutator::apply(
362 &mut *guard,
363 super::dsl::MeerkatMachineInput::RequestDeferredTools {
364 authorities: target_authorities,
365 },
366 )
367 .map_err(|err| ToolScopeStageError::Owner {
368 message: super::dsl_authority::map_error(err, "RequestDeferredTools"),
369 })?;
370 let revision = ToolScopeRevision(guard.state().staged_visibility_revision);
371 mirror_visibility_projection_from_authority(&mut state, &guard);
372 Ok(revision)
373 }
374
375 fn replace_deferred_tool_authority_catalog(
376 &self,
377 catalog: std::collections::BTreeMap<ToolName, ToolVisibilityWitness>,
378 ) -> Result<(), ToolScopeApplyError> {
379 let authority = self.dsl_authority_for_apply()?;
380 let mut state = self.state.write().map_err(|_| ToolScopeApplyError::Owner {
381 message: "machine visibility state lock poisoned".to_string(),
382 })?;
383 let mut guard = authority
384 .lock()
385 .unwrap_or_else(std::sync::PoisonError::into_inner);
386 super::dsl::MeerkatMachineMutator::apply(
387 &mut *guard,
388 super::dsl::MeerkatMachineInput::ReplaceDeferredToolAuthorityCatalog {
389 catalog: dsl_witnesses(&catalog),
390 },
391 )
392 .map_err(|err| ToolScopeApplyError::Owner {
393 message: super::dsl_authority::map_error(err, "ReplaceDeferredToolAuthorityCatalog"),
394 })?;
395 mirror_visibility_projection_from_authority(&mut state, &guard);
396 Ok(())
397 }
398
399 fn replace_filter_tool_authority_catalog(
400 &self,
401 catalog: std::collections::BTreeMap<ToolName, ToolVisibilityWitness>,
402 ) -> Result<(), ToolScopeApplyError> {
403 let authority = self.dsl_authority_for_apply()?;
404 let mut state = self.state.write().map_err(|_| ToolScopeApplyError::Owner {
405 message: "machine visibility state lock poisoned".to_string(),
406 })?;
407 let mut guard = authority
408 .lock()
409 .unwrap_or_else(std::sync::PoisonError::into_inner);
410 super::dsl::MeerkatMachineMutator::apply(
411 &mut *guard,
412 super::dsl::MeerkatMachineInput::ReplaceFilterToolAuthorityCatalog {
413 catalog: dsl_witnesses(&catalog),
414 },
415 )
416 .map_err(|err| ToolScopeApplyError::Owner {
417 message: super::dsl_authority::map_error(err, "ReplaceFilterToolAuthorityCatalog"),
418 })?;
419 mirror_visibility_projection_from_authority(&mut state, &guard);
420 Ok(())
421 }
422
423 fn set_turn_overlay(
424 &self,
425 allow: Option<std::collections::BTreeSet<ToolName>>,
426 deny: std::collections::BTreeSet<ToolName>,
427 ) -> Result<ToolScopeTurnOverlay, ToolScopeStageError> {
428 let authority = self.dsl_authority_for_stage()?;
429 let allow_active = allow.is_some();
430 let allow_names = allow.unwrap_or_default();
431 let mut guard = authority
432 .lock()
433 .unwrap_or_else(std::sync::PoisonError::into_inner);
434 super::dsl::MeerkatMachineMutator::apply(
435 &mut *guard,
436 super::dsl::MeerkatMachineInput::SetTurnToolOverlay {
437 allow_active,
438 allow_names,
439 deny_names: deny,
440 },
441 )
442 .map_err(|err| ToolScopeStageError::Owner {
443 message: super::dsl_authority::map_error(err, "SetTurnToolOverlay"),
444 })?;
445 let state = guard.state();
446 Ok(ToolScopeTurnOverlay::from_string_sets(
447 state
448 .turn_tool_overlay_allow_active
449 .then(|| state.turn_tool_overlay_allow_names.clone()),
450 state.turn_tool_overlay_deny_names.clone(),
451 ))
452 }
453
454 fn clear_turn_overlay(&self) -> Result<ToolScopeTurnOverlay, ToolScopeStageError> {
455 let authority = self.dsl_authority_for_stage()?;
456 let mut guard = authority
457 .lock()
458 .unwrap_or_else(std::sync::PoisonError::into_inner);
459 super::dsl::MeerkatMachineMutator::apply(
460 &mut *guard,
461 super::dsl::MeerkatMachineInput::ClearTurnToolOverlay,
462 )
463 .map_err(|err| ToolScopeStageError::Owner {
464 message: super::dsl_authority::map_error(err, "ClearTurnToolOverlay"),
465 })?;
466 Ok(ToolScopeTurnOverlay::cleared())
467 }
468
469 fn requires_filter_witnesses(&self) -> bool {
470 true
471 }
472
473 fn boundary_applied(&self) -> Result<SessionToolVisibilityState, ToolScopeApplyError> {
474 let authority = self.dsl_authority_for_apply()?;
475 let mut state = self
476 .state
477 .write()
478 .unwrap_or_else(std::sync::PoisonError::into_inner);
479 let mut guard = authority
480 .lock()
481 .unwrap_or_else(std::sync::PoisonError::into_inner);
482 let staged_filter = guard.state().staged_filter.clone();
483 let staged_revision = guard.state().staged_visibility_revision;
484 let staged_deferred_authorities = guard.state().staged_deferred_authorities.clone();
485 super::dsl::MeerkatMachineMutator::apply(
486 &mut *guard,
487 super::dsl::MeerkatMachineInput::CommitVisibilityFilter {
488 filter: staged_filter,
489 revision: staged_revision,
490 },
491 )
492 .map_err(|err| ToolScopeApplyError::Owner {
493 message: super::dsl_authority::map_error(err, "CommitVisibilityFilter"),
494 })?;
495 super::dsl::MeerkatMachineMutator::apply(
496 &mut *guard,
497 super::dsl::MeerkatMachineInput::CommitDeferredNames {
498 authorities: staged_deferred_authorities,
499 },
500 )
501 .map_err(|err| ToolScopeApplyError::Owner {
502 message: super::dsl_authority::map_error(err, "CommitDeferredNames"),
503 })?;
504
505 mirror_visibility_projection_from_authority(&mut state, &guard);
506 Ok(state.clone())
507 }
508}
509
510impl MeerkatMachine {
511 pub async fn meerkat_machine_archive_snapshot(
512 &self,
513 session_id: &SessionId,
514 ) -> Option<MeerkatArchiveSnapshot> {
515 let (driver_handle, control_snapshot, completions_handle) = {
516 let sessions = self.sessions.read().await;
517 let entry = sessions.get(session_id)?;
518 let control = entry.control_snapshot();
519 let mut snapshot = control.clone();
520 let authority = entry
521 .dsl_authority
522 .lock()
523 .unwrap_or_else(std::sync::PoisonError::into_inner);
524 let dsl_phase =
525 crate::meerkat_machine::dsl_authority::runtime_phase_from_authority(&authority);
526 let dsl_current_run_id =
527 crate::meerkat_machine::dsl_authority::current_run_id_from_authority(&authority);
528 let dsl_pre_run_phase =
529 crate::meerkat_machine::dsl_authority::pre_run_phase_from_authority(&authority);
530 let plan = crate::meerkat_machine::resolve_visible_runtime_phase(
536 dsl_phase,
537 dsl_pre_run_phase,
538 control.phase,
539 control.pre_run_phase,
540 self.has_runtime_persistence(),
541 )
542 .unwrap_or(crate::meerkat_machine::VisibleRuntimePhasePlan {
543 publish_control: true,
544 selected_raw_phase: control.phase,
545 visible_phase: control.phase,
546 });
547 if !plan.publish_control {
548 snapshot.phase = plan.visible_phase;
549 snapshot.current_run_id = dsl_current_run_id;
550 snapshot.pre_run_phase = dsl_pre_run_phase;
551 }
552 (
553 Arc::clone(&entry.driver),
554 snapshot,
555 Arc::clone(&entry.completions),
556 )
557 };
558
559 let completion_waiters = {
560 let completions = completions_handle.lock().await;
561 let snapshot = completions.diagnostic_snapshot();
562 MeerkatCompletionWaitersSnapshot {
563 input_count: snapshot.input_count,
564 waiter_count: snapshot.waiter_count,
565 waiting_inputs: snapshot
566 .waiting_inputs
567 .into_iter()
568 .map(|entry| MeerkatCompletionWaiterSnapshot {
569 input_id: entry.input_id,
570 waiter_count: entry.waiter_count,
571 })
572 .collect(),
573 }
574 };
575
576 let (queue, steer_queue) = {
577 let driver = driver_handle.lock().await;
578 let ingress = driver.driver_ingress();
579 (ingress.queue(), ingress.steer_queue())
580 };
581
582 Some(MeerkatArchiveSnapshot {
583 control: MeerkatControlSnapshot {
584 phase: control_snapshot.phase,
585 current_run_id: control_snapshot.current_run_id,
586 pre_run_phase: control_snapshot.pre_run_phase,
587 },
588 queue,
589 steer_queue,
590 completion_waiters,
591 })
592 }
593
594 pub async fn meerkat_machine_spine_snapshot(
595 &self,
596 session_id: &SessionId,
597 ) -> Option<MeerkatMachineSpineSnapshot> {
598 tracing::info!(%session_id, "meerkat_machine_spine_snapshot start");
599 let (
600 driver_handle,
601 control_snapshot,
602 completions_handle,
603 ops_lifecycle,
604 completions_present,
605 ops_registry_present,
606 epoch_id,
607 _visibility_state,
608 formal_pre_run_phase,
609 formal_visibility_authority_catalogs,
610 ) = {
611 tracing::info!(%session_id, "meerkat_machine_spine_snapshot reading session entry");
612 let sessions = self.sessions.read().await;
613 let entry = sessions.get(session_id)?;
614 tracing::info!(%session_id, "meerkat_machine_spine_snapshot locking dsl authority");
615 let control = entry.control_snapshot();
621 let mut snapshot = control.clone();
622 let authority = entry
623 .dsl_authority
624 .lock()
625 .unwrap_or_else(std::sync::PoisonError::into_inner);
626 tracing::info!(%session_id, "meerkat_machine_spine_snapshot locked dsl authority");
627 let dsl_phase =
628 crate::meerkat_machine::dsl_authority::runtime_phase_from_authority(&authority);
629 let dsl_current_run_id =
630 crate::meerkat_machine::dsl_authority::current_run_id_from_authority(&authority);
631 let dsl_pre_run_phase =
632 crate::meerkat_machine::dsl_authority::pre_run_phase_from_authority(&authority);
633 let plan = crate::meerkat_machine::resolve_visible_runtime_phase(
636 dsl_phase,
637 dsl_pre_run_phase,
638 control.phase,
639 control.pre_run_phase,
640 self.has_runtime_persistence(),
641 )
642 .unwrap_or(crate::meerkat_machine::VisibleRuntimePhasePlan {
643 publish_control: true,
644 selected_raw_phase: control.phase,
645 visible_phase: control.phase,
646 });
647 if !plan.publish_control {
648 snapshot.phase = plan.visible_phase;
649 snapshot.current_run_id = dsl_current_run_id;
650 snapshot.pre_run_phase = dsl_pre_run_phase;
651 }
652 let formal_pre_run_phase = snapshot
653 .pre_run_phase
654 .and_then(crate::meerkat_machine::dsl_authority::pre_run_phase_from_runtime_state)
655 .map(|phase| format!("{phase:?}"));
656 let formal_visibility_authority_catalogs = {
657 let state = authority.state();
658 (
659 core_witnesses(&state.deferred_visibility_authority_catalog),
660 core_witnesses(&state.filter_visibility_authority_catalog),
661 )
662 };
663 (
664 Arc::clone(&entry.driver),
665 snapshot,
666 Arc::clone(&entry.completions),
667 Arc::clone(&entry.ops_lifecycle),
668 true,
669 true,
670 entry.epoch_id.clone(),
671 entry.tool_visibility_owner.visibility_state().ok()?,
672 formal_pre_run_phase,
673 formal_visibility_authority_catalogs,
674 )
675 };
676 tracing::info!(%session_id, "meerkat_machine_spine_snapshot locking completions");
677 let completion_waiters = {
678 let completions = completions_handle.lock().await;
679 tracing::info!(%session_id, "meerkat_machine_spine_snapshot locked completions");
680 let snapshot = completions.diagnostic_snapshot();
681 MeerkatCompletionWaitersSnapshot {
682 input_count: snapshot.input_count,
683 waiter_count: snapshot.waiter_count,
684 waiting_inputs: snapshot
685 .waiting_inputs
686 .into_iter()
687 .map(|entry| MeerkatCompletionWaiterSnapshot {
688 input_id: entry.input_id,
689 waiter_count: entry.waiter_count,
690 })
691 .collect(),
692 }
693 };
694 tracing::info!(%session_id, "meerkat_machine_spine_snapshot reading drain");
695 let drain = {
696 let sessions = self.sessions.read().await;
697 if let Some(entry) = sessions.get(session_id) {
698 let authority = entry
699 .dsl_authority
700 .lock()
701 .unwrap_or_else(std::sync::PoisonError::into_inner);
702 let state = authority.state();
703 let phase = crate::meerkat_machine::CommsDrainPhase::from(state.drain_phase);
704 let mode = state
705 .drain_mode
706 .map(crate::meerkat_machine::CommsDrainMode::from);
707 let handle_present = entry.drain_slot.handle_present();
708 let activated = phase != crate::meerkat_machine::CommsDrainPhase::Inactive
709 || mode.is_some()
710 || handle_present;
711 if activated {
712 MeerkatDrainSnapshot {
713 slot_present: true,
714 phase: Some(phase),
715 mode,
716 handle_present,
717 }
718 } else {
719 MeerkatDrainSnapshot {
720 slot_present: false,
721 phase: None,
722 mode: None,
723 handle_present: false,
724 }
725 }
726 } else {
727 MeerkatDrainSnapshot {
728 slot_present: false,
729 phase: None,
730 mode: None,
731 handle_present: false,
732 }
733 }
734 };
735 tracing::info!(%session_id, "meerkat_machine_spine_snapshot locking driver");
736 let driver = driver_handle.lock().await;
737 tracing::info!(%session_id, "meerkat_machine_spine_snapshot locked driver");
738 let driver_kind = match &*driver {
739 DriverEntry::Ephemeral(_) => MeerkatDriverKind::Ephemeral,
740 DriverEntry::Persistent(_) => MeerkatDriverKind::Persistent,
741 };
742 let ingress = driver.driver_ingress();
743
744 let binding = MeerkatBindingSnapshot {
745 session_id: session_id.clone(),
746 runtime_id: driver.runtime_id().clone(),
747 driver_kind,
748 driver_present: true,
749 completions_present,
750 ops_registry_present,
751 epoch_id,
752 cursor_state: {
753 let cursor_state = ops_lifecycle.completion_cursor_snapshot();
754 MeerkatCursorSnapshot {
755 agent_applied_cursor: cursor_state.agent_applied_cursor,
756 runtime_observed_seq: cursor_state.runtime_observed_seq,
757 runtime_last_injected_seq: cursor_state.runtime_last_injected_seq,
758 }
759 },
760 };
761
762 let control = MeerkatControlSnapshot {
763 phase: control_snapshot.phase,
764 current_run_id: control_snapshot.current_run_id,
765 pre_run_phase: control_snapshot.pre_run_phase,
766 };
767
768 let admission_order: Vec<MeerkatAdmittedInputSnapshot> = ingress
769 .admission_order()
770 .into_iter()
771 .map(|input_id| MeerkatAdmittedInputSnapshot {
772 content_shape: ingress.content_shape(&input_id),
773 request_id: ingress.request_id(&input_id),
774 reservation_key: ingress.reservation_key(&input_id),
775 handling_mode: ingress.handling_mode(&input_id),
776 live_interrupt_required: ingress.live_interrupt_required(&input_id),
777 lifecycle: driver.input_phase(&input_id),
778 terminal_outcome: driver.input_terminal_outcome(&input_id),
779 last_run_id: driver.input_last_run_id(&input_id),
780 last_boundary_sequence: driver.input_last_boundary_sequence(&input_id),
781 is_prompt: ingress.is_prompt(&input_id),
782 input_id,
783 })
784 .collect();
785
786 let current_run_contributors = if let Some(control_run_id) = &control.current_run_id {
787 admission_order
788 .iter()
789 .filter(|snapshot| {
790 snapshot.last_run_id.as_ref() == Some(control_run_id)
791 && matches!(
792 snapshot.lifecycle,
793 Some(
794 crate::input_state::InputLifecycleState::Staged
795 | crate::input_state::InputLifecycleState::Applied
796 | crate::input_state::InputLifecycleState::AppliedPendingConsumption
797 )
798 )
799 })
800 .map(|snapshot| snapshot.input_id.clone())
801 .collect()
802 } else {
803 Vec::new()
804 };
805
806 let inputs = MeerkatInputsSnapshot {
807 admission_order,
808 queue: ingress.queue(),
809 steer_queue: ingress.steer_queue(),
810 current_run_id: control.current_run_id.clone(),
811 current_run_contributors,
812 post_admission_signal: format!("{:?}", driver.post_admission_signal()),
813 silent_intent_overrides: driver.silent_comms_intents().into_iter().collect(),
814 };
815 let ledger = {
816 let mut snapshot = MeerkatLedgerSnapshot {
817 input_count: 0,
818 non_terminal_count: 0,
819 accepted_count: 0,
820 queued_count: 0,
821 staged_count: 0,
822 applied_count: 0,
823 applied_pending_consumption_count: 0,
824 consumed_count: 0,
825 superseded_count: 0,
826 coalesced_count: 0,
827 abandoned_count: 0,
828 };
829
830 for (input_id, _state) in driver.ledger().iter() {
831 snapshot.input_count += 1;
832 let Some(lifecycle) = driver.input_phase(input_id) else {
833 tracing::error!(
834 input_id = %input_id,
835 "missing generated input lifecycle authority for ledger snapshot"
836 );
837 continue;
838 };
839 let terminal = crate::meerkat_machine::input_phase_behavioral_terminality_via_authority(
840 input_id,
841 lifecycle,
842 driver.input_terminal_outcome(input_id),
843 )
844 .unwrap_or_else(|err| {
845 tracing::error!(
846 input_id = %input_id,
847 error = %err,
848 "generated input terminality authority rejected ledger snapshot classification"
849 );
850 true
851 });
852 if !terminal {
853 snapshot.non_terminal_count += 1;
854 }
855 match lifecycle {
856 InputLifecycleState::Accepted => snapshot.accepted_count += 1,
857 InputLifecycleState::Queued => snapshot.queued_count += 1,
858 InputLifecycleState::Staged => snapshot.staged_count += 1,
859 InputLifecycleState::Applied => snapshot.applied_count += 1,
860 InputLifecycleState::AppliedPendingConsumption => {
861 snapshot.applied_pending_consumption_count += 1;
862 }
863 InputLifecycleState::Consumed => snapshot.consumed_count += 1,
864 InputLifecycleState::Superseded => snapshot.superseded_count += 1,
865 InputLifecycleState::Coalesced => snapshot.coalesced_count += 1,
866 InputLifecycleState::Abandoned => snapshot.abandoned_count += 1,
867 }
868 }
869
870 snapshot
871 };
872 let ops_snapshot = match ops_lifecycle.diagnostic_snapshot() {
873 Ok(snapshot) => snapshot,
874 Err(error) => {
875 tracing::error!(
876 session_id = %session_id,
877 error = %error,
878 "generated ops lifecycle authority rejected spine projection"
879 );
880 return None;
881 }
882 };
883 let ops = MeerkatOpsSnapshot {
884 operation_count: ops_snapshot.operation_count,
885 active_count: ops_snapshot.active_count,
886 wait_request_id: ops_snapshot.wait_request_id,
887 pending_wait_present: ops_snapshot.pending_wait_present,
888 pending_wait_request_id: ops_snapshot.pending_wait_request_id,
889 wait_operation_ids: ops_snapshot.wait_operation_ids,
890 operations: ops_snapshot.operations,
891 };
892 let formal_state = {
893 let mut available_fields = std::collections::BTreeMap::new();
894 available_fields.insert(
895 "session_id".into(),
896 formal_projection_value(&Some(session_id.to_string())),
897 );
898 available_fields.insert(
899 "active_runtime_id".into(),
900 formal_projection_value(&Some(driver.runtime_id().to_string())),
901 );
902 available_fields.insert(
903 "current_run_id".into(),
904 formal_projection_value(&control.current_run_id.as_ref().map(ToString::to_string)),
905 );
906 available_fields.insert(
907 "pre_run_phase".into(),
908 formal_projection_value(&formal_pre_run_phase),
909 );
910 available_fields.insert(
911 "silent_intent_overrides".into(),
912 formal_projection_value(
913 &driver
914 .silent_comms_intents()
915 .into_iter()
916 .collect::<BTreeSet<_>>(),
917 ),
918 );
919 available_fields.insert(
920 "deferred_visibility_authority_catalog".into(),
921 formal_projection_value(&formal_visibility_authority_catalogs.0),
922 );
923 available_fields.insert(
924 "filter_visibility_authority_catalog".into(),
925 formal_projection_value(&formal_visibility_authority_catalogs.1),
926 );
927 MeerkatFormalStateProjection {
928 available_fields,
929 unavailable_fields: vec!["active_fence_token".into()],
930 }
931 };
932
933 Some(MeerkatMachineSpineSnapshot {
934 binding,
935 control,
936 inputs,
937 ledger,
938 completion_waiters,
939 ops,
940 drain,
941 formal_state,
942 })
943 }
944}