1use std::collections::{BTreeMap, BTreeSet};
4use std::path::{Path, PathBuf};
5use std::sync::Arc;
6
7use async_trait::async_trait;
8use futures::StreamExt;
9use meerkat::{AgentFactory, Config, FactoryAgentBuilder, SessionStore};
10use meerkat_client::types::LlmStream;
11use meerkat_client::{LlmClient, LlmRequest};
12use meerkat_core::agent::CommsRuntime;
13use meerkat_core::service::{
14 CreateSessionRequest, SessionError, SessionHistoryPage, SessionHistoryQuery,
15 SessionServiceHistoryExt,
16};
17use meerkat_core::{AgentSessionStore, AssistantBlock, Message, Provider};
18use meerkat_mob::{
19 MobBuilder, MobDefinition, MobError, MobHandle, MobSessionService, MobStorage, Profile,
20 ProfileName, SpawnMemberSpec,
21};
22use meerkat_runtime::input_state::{InputStatePersistenceRecord, StoredInputState};
23use meerkat_runtime::store::MachineLifecycleCommit;
24use meerkat_store::StoreAdapter;
25use serde::Serialize;
26use serde_json::Value;
27
28use crate::blob_store::{
29 Base64BlobStoreAdapter, BinaryBlobStore, BinaryBlobStoreAdapter, ObjectStoreBlobStore,
30};
31use crate::console_spawn::{
32 ConsoleSpawnSink, SharedConsoleSpawnSinkSlot, new_console_spawn_sink_slot,
33};
34
35pub(crate) const DELEGATE_IDLE_RETIRE_SECS_LABEL: &str = "implicit_delegate_idle_retire_secs";
36pub(crate) const DELEGATE_IDLE_RETIRE_DISABLED_LABEL: &str = "disabled";
37
38pub(crate) fn is_previous_member_cleanup_ambiguous_error(error: &str) -> bool {
39 error.contains("previous member cleanup ambiguous for member ")
40}
41
42pub(crate) fn is_recoverable_lifecycle_cleanup_error(error: &str) -> bool {
43 is_previous_member_cleanup_ambiguous_error(error)
44 || (error.contains("disposal completed but ArchiveSession failed")
45 && (
46 (error.contains("cancel-before-retire failed")
49 && error.contains("Runtime not ready: running"))
50 || (error.contains("guard rejected transition from Stopped")
57 && error.contains("input::Retire"))
58 || (error.contains("continuity save")
64 && (error.contains("continuity record not found")
65 || error.contains("stale fencing token")))
66 ))
67}
68
69fn is_session_owned_archive_absent_cleanup_error(error: &str) -> bool {
87 error.contains("disposal completed but ArchiveSession failed")
88 && error.contains("NotFound for registered runtime session")
89}
90
91pub(crate) fn is_recoverable_session_owned_retire_cleanup_error(error: &str) -> bool {
103 is_recoverable_lifecycle_cleanup_error(error)
104 || is_session_owned_archive_absent_cleanup_error(error)
105}
106
107pub(crate) fn is_stopped_session_archive_retire_rejection(error: &str) -> bool {
121 error.contains("machine archive retire failed")
122 && error.contains("guard rejected transition from Stopped")
123 && error.contains("input::Retire")
124}
125
126pub(crate) fn topology_restore_failed_peer_ids(
127 error: &meerkat_mob::MobRespawnError,
128) -> Option<Vec<String>> {
129 match error {
130 meerkat_mob::MobRespawnError::TopologyRestoreFailed {
131 receipt: _,
132 failed_peer_ids,
133 } => Some(failed_peer_ids.iter().map(ToString::to_string).collect()),
134 _ => None,
135 }
136}
137
138pub(crate) fn topology_restore_warning_json(failed_peer_ids: &[String]) -> Value {
139 serde_json::json!({
140 "kind": "topology_restore_degraded",
141 "failed_peer_ids": failed_peer_ids,
142 })
143}
144
145#[cfg(test)]
146use std::sync::{Mutex, OnceLock};
147
148pub const MEMBER_STATE_ACTIVE: &str = "active";
150pub const MEMBER_STATE_RETIRING: &str = "retiring";
152
153pub(crate) fn member_status_state_string(status: meerkat_mob::MobMemberStatus) -> String {
162 match status {
163 meerkat_mob::MobMemberStatus::Active => MEMBER_STATE_ACTIVE.to_string(),
164 meerkat_mob::MobMemberStatus::Retiring => MEMBER_STATE_RETIRING.to_string(),
165 other => format!("{other:?}").to_ascii_lowercase(),
166 }
167}
168
169#[derive(Clone, Default)]
171pub struct MobBootstrapOptions {
172 pub allow_ephemeral_sessions: bool,
173 pub notify_orchestrator_on_resume: bool,
174 pub default_llm_client: Option<Arc<dyn LlmClient>>,
175}
176
177pub struct ReplaySanitizingLlmClient {
181 inner: Arc<dyn LlmClient>,
182}
183
184type SharedDefaultLlmClientSlot = Arc<std::sync::RwLock<Option<Arc<dyn LlmClient>>>>;
185
186impl ReplaySanitizingLlmClient {
187 pub fn new(inner: Arc<dyn LlmClient>) -> Self {
188 Self { inner }
189 }
190
191 pub fn wrap(inner: Arc<dyn LlmClient>) -> Arc<dyn LlmClient> {
192 Arc::new(Self::new(inner))
193 }
194}
195
196pub struct ReplaySanitizingAgentLlmClient {
204 inner: Arc<dyn meerkat_core::AgentLlmClient>,
205}
206
207impl ReplaySanitizingAgentLlmClient {
208 pub fn new(inner: Arc<dyn meerkat_core::AgentLlmClient>) -> Self {
209 Self { inner }
210 }
211
212 pub fn wrap(
213 inner: Arc<dyn meerkat_core::AgentLlmClient>,
214 ) -> Arc<dyn meerkat_core::AgentLlmClient> {
215 Arc::new(Self::new(inner))
216 }
217}
218
219#[async_trait]
220impl meerkat_core::AgentLlmClient for ReplaySanitizingAgentLlmClient {
221 async fn stream_response(
222 &self,
223 messages: &[Message],
224 tools: &[Arc<meerkat_core::ToolDef>],
225 max_tokens: u32,
226 temperature: Option<f32>,
227 provider_params: Option<&meerkat_core::lifecycle::run_primitive::ProviderParamsOverride>,
228 ) -> Result<meerkat_core::agent::LlmStreamResult, meerkat_core::AgentError> {
229 let sanitized: Vec<Message> = messages
230 .iter()
231 .cloned()
232 .map(sanitize_message_for_stateless_replay)
233 .collect();
234 self.inner
235 .stream_response(&sanitized, tools, max_tokens, temperature, provider_params)
236 .await
237 }
238
239 fn provider(&self) -> meerkat_core::Provider {
240 self.inner.provider()
241 }
242
243 fn model(&self) -> &str {
244 self.inner.model()
245 }
246
247 fn compile_schema(
248 &self,
249 output_schema: &meerkat_core::OutputSchema,
250 ) -> Result<meerkat_core::schema::CompiledSchema, meerkat_core::schema::SchemaError> {
251 self.inner.compile_schema(output_schema)
252 }
253}
254
255#[async_trait]
256impl LlmClient for ReplaySanitizingLlmClient {
257 fn project_replay_messages(
258 &self,
259 messages: &[Message],
260 ) -> Result<Vec<Message>, meerkat_client::LlmError> {
261 let sanitized: Vec<Message> = messages
262 .iter()
263 .cloned()
264 .map(sanitize_message_for_stateless_replay)
265 .collect();
266 self.inner.project_replay_messages(&sanitized)
267 }
268
269 fn stream<'a>(&'a self, request: &'a LlmRequest) -> LlmStream<'a> {
270 let inner = Arc::clone(&self.inner);
271 let sanitized = sanitize_llm_request_for_stateless_replay(request);
272 Box::pin(async_stream::stream! {
273 let mut stream = inner.stream(&sanitized);
274 while let Some(event) = stream.next().await {
275 if runtime_turn_diagnostics_enabled()
276 && let Err(error) = &event
277 {
278 tracing::error!(
279 error = %error,
280 error_debug = ?error,
281 "mobkit llm client stream error"
282 );
283 }
284 yield event;
285 }
286 })
287 }
288
289 fn provider(&self) -> meerkat_core::Provider {
290 self.inner.provider()
291 }
292
293 async fn health_check(&self) -> Result<(), meerkat_client::LlmError> {
294 self.inner.health_check().await
295 }
296
297 fn compile_schema(
298 &self,
299 output_schema: &meerkat_core::OutputSchema,
300 ) -> Result<meerkat_core::schema::CompiledSchema, meerkat_core::schema::SchemaError> {
301 self.inner.compile_schema(output_schema)
302 }
303}
304
305pub(crate) type PreBuildHook = Arc<
333 dyn Fn(
334 &mut CreateSessionRequest,
335 ) -> std::pin::Pin<
336 Box<dyn std::future::Future<Output = Result<(), SessionError>> + Send + '_>,
337 > + Send
338 + Sync,
339>;
340
341pub type AfterCreateHook = Arc<
343 dyn Fn(
344 meerkat_core::types::SessionId,
345 SessionCreatedContext,
346 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send>>
347 + Send
348 + Sync,
349>;
350
351struct PreBuildMobSessionService {
358 inner: Arc<dyn MobSessionService>,
359 hook: PreBuildHook,
360 after_create_hook: Option<AfterCreateHook>,
361 runtime_adapter_override: Option<Arc<meerkat_runtime::MeerkatMachine>>,
362}
363
364fn no_op_pre_build_hook() -> PreBuildHook {
365 Arc::new(|_req: &mut CreateSessionRequest| Box::pin(async { Ok(()) }))
366}
367
368#[derive(Debug, Clone, Copy, PartialEq, Eq)]
369pub(crate) enum DelegateIdleRetireOverride {
370 Disabled,
371 Seconds(u64),
372}
373
374#[derive(Clone, Default)]
375pub(crate) struct ImplicitDelegateRetirementOverrides {
376 inner: Arc<tokio::sync::RwLock<BTreeMap<(String, String), DelegateIdleRetireOverride>>>,
377}
378
379impl ImplicitDelegateRetirementOverrides {
380 pub(crate) async fn set(
381 &self,
382 mob_id: impl Into<String>,
383 member_id: impl Into<String>,
384 override_policy: DelegateIdleRetireOverride,
385 ) {
386 self.inner
387 .write()
388 .await
389 .insert((mob_id.into(), member_id.into()), override_policy);
390 }
391
392 pub(crate) async fn get(
393 &self,
394 mob_id: &str,
395 member_id: &str,
396 ) -> Option<DelegateIdleRetireOverride> {
397 self.inner
398 .read()
399 .await
400 .get(&(mob_id.to_string(), member_id.to_string()))
401 .copied()
402 }
403}
404
405struct AutoWireParentMobToolsFactory {
406 inner: Arc<dyn meerkat_core::service::MobToolsFactory>,
407 implicit_delegate_retirement_overrides: ImplicitDelegateRetirementOverrides,
408 console_spawn_sink: SharedConsoleSpawnSinkSlot,
409}
410
411#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
412#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
413impl meerkat_core::service::MobToolsFactory for AutoWireParentMobToolsFactory {
414 async fn build_mob_tools(
415 &self,
416 args: meerkat_core::service::MobToolsBuildArgs,
417 ) -> Result<Arc<dyn meerkat_core::AgentToolDispatcher>, Box<dyn std::error::Error + Send + Sync>>
418 {
419 let spawner_comms_name = args.comms_name.clone();
420 let inner = self.inner.build_mob_tools(args).await?;
421 Ok(Arc::new(AutoWireParentMobToolDispatcher {
422 inner,
423 implicit_delegate_retirement_overrides: self
424 .implicit_delegate_retirement_overrides
425 .clone(),
426 console_spawn_sink: Arc::clone(&self.console_spawn_sink),
427 spawner_comms_name,
428 }))
429 }
430}
431
432struct AutoWireParentMobToolDispatcher {
433 inner: Arc<dyn meerkat_core::AgentToolDispatcher>,
434 implicit_delegate_retirement_overrides: ImplicitDelegateRetirementOverrides,
435 console_spawn_sink: SharedConsoleSpawnSinkSlot,
438 spawner_comms_name: Option<String>,
441}
442
443#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
444#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
445impl meerkat_core::AgentToolDispatcher for AutoWireParentMobToolDispatcher {
446 fn tools(&self) -> Arc<[Arc<meerkat_core::types::ToolDef>]> {
447 self.inner
448 .tools()
449 .iter()
450 .map(|tool| {
451 if tool.name == "delegate" {
452 Arc::new(delegate_tool_def_with_idle_retire_secs(tool))
453 } else if tool.name == "mob_spawn_member" {
454 Arc::new(mob_spawn_tool_def_with_idle_retire_secs(tool))
455 } else {
456 Arc::clone(tool)
457 }
458 })
459 .collect::<Vec<_>>()
460 .into()
461 }
462
463 async fn dispatch(
464 &self,
465 call: meerkat_core::types::ToolCallView<'_>,
466 ) -> Result<meerkat_core::ToolDispatchOutcome, meerkat_core::ToolError> {
467 if call.name == "delegate" {
468 return self.dispatch_delegate(call).await;
469 }
470 if call.name == "mob_spawn_member" {
471 return self.dispatch_mob_spawn_member(call).await;
472 }
473 if crate::console_spawn::is_console_spawn_tool(call.name) {
474 let args = serde_json::from_str::<Value>(call.args.get()).ok();
478 let name = call.name.to_string();
479 let outcome = self.inner.dispatch(call).await?;
480 if let Some(args) = args {
481 self.project_spawn_to_console(&name, &args, &outcome).await;
482 }
483 return Ok(outcome);
484 }
485 self.inner.dispatch(call).await
486 }
487
488 fn capabilities(&self) -> meerkat_core::agent::DispatcherCapabilities {
489 self.inner.capabilities()
490 }
491
492 fn bind_ops_lifecycle(
493 self: Arc<Self>,
494 registry: Arc<dyn meerkat_core::ops_lifecycle::OpsLifecycleRegistry>,
495 owner_bridge_session_id: meerkat_core::types::SessionId,
496 ) -> Result<meerkat_core::agent::BindOutcome, meerkat_core::agent::OpsLifecycleBindError> {
497 let owned = Arc::try_unwrap(self)
498 .map_err(|_| meerkat_core::agent::OpsLifecycleBindError::SharedOwnership)?;
499 let outcome = owned
500 .inner
501 .bind_ops_lifecycle(registry, owner_bridge_session_id)?;
502 let was_bound = outcome.was_bound();
503 let dispatcher = Arc::new(Self {
504 inner: outcome.into_dispatcher(),
505 implicit_delegate_retirement_overrides: owned.implicit_delegate_retirement_overrides,
506 console_spawn_sink: owned.console_spawn_sink,
507 spawner_comms_name: owned.spawner_comms_name,
508 });
509 Ok(if was_bound {
510 meerkat_core::agent::BindOutcome::Bound(dispatcher)
511 } else {
512 meerkat_core::agent::BindOutcome::Skipped(dispatcher)
513 })
514 }
515}
516
517impl AutoWireParentMobToolDispatcher {
518 async fn dispatch_mob_spawn_member(
519 &self,
520 call: meerkat_core::types::ToolCallView<'_>,
521 ) -> Result<meerkat_core::ToolDispatchOutcome, meerkat_core::ToolError> {
522 let mut args = serde_json::from_str::<Value>(call.args.get()).map_err(|error| {
523 meerkat_core::ToolError::invalid_arguments(call.name, error.to_string())
524 })?;
525 let idle_retire_override = delegate_idle_retire_override_from_args(call.name, &mut args)?;
526 let idle_retire_targets = idle_retire_targets_from_spawn_args(&args);
527 if let Some(object) = args.as_object_mut() {
528 object
529 .entry("auto_wire_parent".to_string())
530 .or_insert(Value::Bool(true));
531 }
532 let name = call.name.to_string();
533 let args_for_console = args.clone();
534 let args = serde_json::value::RawValue::from_string(args.to_string()).map_err(|error| {
535 meerkat_core::ToolError::invalid_arguments(call.name, error.to_string())
536 })?;
537 let call = meerkat_core::types::ToolCallView {
538 id: call.id,
539 name: call.name,
540 args: &args,
541 };
542 let outcome = self.inner.dispatch(call).await?;
543 self.register_idle_retire_override_from_outcome(
544 &outcome,
545 idle_retire_override,
546 &idle_retire_targets,
547 )
548 .await;
549 self.project_spawn_to_console(&name, &args_for_console, &outcome)
550 .await;
551
552 Ok(outcome)
553 }
554
555 async fn dispatch_delegate(
556 &self,
557 call: meerkat_core::types::ToolCallView<'_>,
558 ) -> Result<meerkat_core::ToolDispatchOutcome, meerkat_core::ToolError> {
559 let mut args = serde_json::from_str::<Value>(call.args.get()).map_err(|error| {
560 meerkat_core::ToolError::invalid_arguments(call.name, error.to_string())
561 })?;
562 let idle_retire_override = delegate_idle_retire_override_from_args(call.name, &mut args)?;
563 let name = call.name.to_string();
564 let args_for_console = args.clone();
565 let args = serde_json::value::RawValue::from_string(args.to_string()).map_err(|error| {
566 meerkat_core::ToolError::invalid_arguments(call.name, error.to_string())
567 })?;
568 let call = meerkat_core::types::ToolCallView {
569 id: call.id,
570 name: call.name,
571 args: &args,
572 };
573 let outcome = self.inner.dispatch(call).await?;
574
575 self.register_idle_retire_override_from_outcome(&outcome, idle_retire_override, &[])
576 .await;
577 self.project_spawn_to_console(&name, &args_for_console, &outcome)
578 .await;
579
580 Ok(outcome)
581 }
582
583 async fn project_spawn_to_console(
588 &self,
589 tool_name: &str,
590 args: &Value,
591 outcome: &meerkat_core::ToolDispatchOutcome,
592 ) {
593 if outcome.result.is_error {
594 return;
595 }
596 let sink = self
597 .console_spawn_sink
598 .read()
599 .unwrap_or_else(std::sync::PoisonError::into_inner)
600 .clone();
601 let Some(sink) = sink else {
602 return;
603 };
604 let seeds = crate::console_spawn::console_spawn_seeds(
605 tool_name,
606 args,
607 &outcome.result.text_content(),
608 self.spawner_comms_name.as_deref(),
609 );
610 for seed in &seeds {
611 sink.project_spawned_member(seed).await;
612 }
613 }
614
615 async fn register_idle_retire_override_from_outcome(
616 &self,
617 outcome: &meerkat_core::ToolDispatchOutcome,
618 idle_retire_override: Option<DelegateIdleRetireOverride>,
619 fallback_targets: &[IdleRetireTarget],
620 ) {
621 if outcome.result.is_error {
622 return;
623 }
624 let Some(override_policy) = idle_retire_override else {
625 return;
626 };
627 for target in
628 idle_retire_targets_from_outcome_text(&outcome.result.text_content(), fallback_targets)
629 {
630 self.implicit_delegate_retirement_overrides
631 .set(&target.mob_id, &target.member_id, override_policy)
632 .await;
633 }
634 }
635}
636
637#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
638struct IdleRetireTarget {
639 mob_id: String,
640 member_id: String,
641}
642
643fn text_field<'a>(value: &'a Value, key: &str) -> Option<&'a str> {
644 value
645 .get(key)
646 .and_then(Value::as_str)
647 .filter(|text| !text.is_empty())
648}
649
650fn member_identity_field(value: &Value) -> Option<&str> {
651 text_field(value, "agent_identity")
652 .or_else(|| text_field(value, "member_id"))
653 .or_else(|| text_field(value, "identity"))
654}
655
656fn target_from_value(value: &Value, default_mob_id: Option<&str>) -> Option<IdleRetireTarget> {
657 let mob_id = text_field(value, "mob_id").or(default_mob_id)?;
658 let member_id = member_identity_field(value)?;
659 Some(IdleRetireTarget {
660 mob_id: mob_id.to_string(),
661 member_id: member_id.to_string(),
662 })
663}
664
665fn idle_retire_targets_from_spawn_args(args: &Value) -> Vec<IdleRetireTarget> {
666 let default_mob_id = text_field(args, "mob_id");
667 let mut targets = BTreeSet::new();
668 if let Some(target) = target_from_value(args, default_mob_id) {
669 targets.insert(target);
670 }
671 for key in ["specs", "members"] {
672 let Some(values) = args.get(key).and_then(Value::as_array) else {
673 continue;
674 };
675 for value in values {
676 if let Some(target) = target_from_value(value, default_mob_id) {
677 targets.insert(target);
678 }
679 }
680 }
681 targets.into_iter().collect()
682}
683
684fn target_from_result_value(
685 value: &Value,
686 fallback_targets: &[IdleRetireTarget],
687) -> Option<IdleRetireTarget> {
688 if let Some(target) = target_from_value(value, None) {
689 return Some(target);
690 }
691 let member_id = member_identity_field(value)?;
692 let mut matches = fallback_targets
693 .iter()
694 .filter(|target| target.member_id == member_id);
695 let target = matches.next()?;
696 if matches.next().is_some() {
697 return None;
698 }
699 Some(target.clone())
700}
701
702fn collect_idle_retire_result_targets(
703 value: &Value,
704 fallback_targets: &[IdleRetireTarget],
705 targets: &mut BTreeSet<IdleRetireTarget>,
706) {
707 if let Some(target) = target_from_result_value(value, fallback_targets) {
708 targets.insert(target);
709 }
710 for key in ["members", "specs", "spawned", "results"] {
711 let Some(values) = value.get(key).and_then(Value::as_array) else {
712 continue;
713 };
714 for value in values {
715 collect_idle_retire_result_targets(value, fallback_targets, targets);
716 }
717 }
718}
719
720fn idle_retire_targets_from_outcome_text(
721 text: &str,
722 fallback_targets: &[IdleRetireTarget],
723) -> Vec<IdleRetireTarget> {
724 let Ok(payload) = serde_json::from_str::<Value>(text) else {
725 return fallback_targets.to_vec();
726 };
727 let mut targets = BTreeSet::new();
728 collect_idle_retire_result_targets(&payload, fallback_targets, &mut targets);
729 if targets.is_empty() {
730 targets.extend(fallback_targets.iter().cloned());
731 }
732 targets.into_iter().collect()
733}
734
735struct DefinitionSeededRealmProfileStore {
736 inner: Arc<dyn meerkat_mob::RealmProfileStore>,
737 profiles: BTreeMap<String, Profile>,
738}
739
740impl DefinitionSeededRealmProfileStore {
741 fn new(
742 definition: &MobDefinition,
743 inner: Arc<dyn meerkat_mob::RealmProfileStore>,
744 ) -> Option<Self> {
745 let profiles = definition
746 .profiles
747 .iter()
748 .filter_map(|(name, binding)| {
749 binding
750 .as_inline()
751 .cloned()
752 .map(|profile| (name.to_string(), profile))
753 })
754 .collect::<BTreeMap<_, _>>();
755
756 (!profiles.is_empty()).then_some(Self { inner, profiles })
757 }
758
759 fn stored(&self, name: &str, profile: &Profile) -> meerkat_mob::StoredRealmProfile {
760 let now = chrono::Utc::now();
761 meerkat_mob::StoredRealmProfile {
762 name: name.to_string(),
763 profile: profile.clone(),
764 revision: 0,
765 created_at: now,
766 updated_at: now,
767 }
768 }
769}
770
771#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
772#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
773impl meerkat_mob::RealmProfileStore for DefinitionSeededRealmProfileStore {
774 async fn create(
775 &self,
776 name: &str,
777 profile: &Profile,
778 ) -> Result<meerkat_mob::StoredRealmProfile, meerkat_mob::MobStoreError> {
779 if self.profiles.contains_key(name) {
780 return Err(meerkat_mob::MobStoreError::CasConflict(format!(
781 "realm profile already exists: {name}"
782 )));
783 }
784 self.inner.create(name, profile).await
785 }
786
787 async fn get(
788 &self,
789 name: &str,
790 ) -> Result<Option<meerkat_mob::StoredRealmProfile>, meerkat_mob::MobStoreError> {
791 if let Some(profile) = self.profiles.get(name) {
792 return Ok(Some(self.stored(name, profile)));
793 }
794 self.inner.get(name).await
795 }
796
797 async fn list(
798 &self,
799 ) -> Result<Vec<meerkat_mob::StoredRealmProfile>, meerkat_mob::MobStoreError> {
800 let mut merged = self.inner.list().await?;
801 merged.retain(|profile| !self.profiles.contains_key(profile.name.as_str()));
802 merged.extend(
803 self.profiles
804 .iter()
805 .map(|(name, profile)| self.stored(name, profile)),
806 );
807 merged.sort_by(|a, b| a.name.cmp(&b.name));
808 Ok(merged)
809 }
810
811 async fn update(
812 &self,
813 name: &str,
814 profile: &Profile,
815 expected_revision: u64,
816 ) -> Result<meerkat_mob::StoredRealmProfile, meerkat_mob::MobStoreError> {
817 if self.profiles.contains_key(name) {
818 return Err(meerkat_mob::MobStoreError::CasConflict(format!(
819 "realm profile '{name}' is provided by the mob definition"
820 )));
821 }
822 self.inner.update(name, profile, expected_revision).await
823 }
824
825 async fn delete(
826 &self,
827 name: &str,
828 expected_revision: u64,
829 ) -> Result<meerkat_mob::StoredRealmProfile, meerkat_mob::MobStoreError> {
830 if self.profiles.contains_key(name) {
831 return Err(meerkat_mob::MobStoreError::CasConflict(format!(
832 "realm profile '{name}' is provided by the mob definition"
833 )));
834 }
835 self.inner.delete(name, expected_revision).await
836 }
837}
838
839fn delegate_idle_retire_override_from_args(
840 tool_name: &str,
841 args: &mut Value,
842) -> Result<Option<DelegateIdleRetireOverride>, meerkat_core::ToolError> {
843 let Some(object) = args.as_object_mut() else {
844 return Ok(None);
845 };
846 let Some(value) = object.remove("idle_retire_secs") else {
847 return Ok(None);
848 };
849 if value.is_null() {
850 return Ok(Some(DelegateIdleRetireOverride::Disabled));
851 }
852 value
853 .as_u64()
854 .map(DelegateIdleRetireOverride::Seconds)
855 .map(Some)
856 .ok_or_else(|| {
857 meerkat_core::ToolError::invalid_arguments(
858 tool_name,
859 "idle_retire_secs must be a non-negative integer or null",
860 )
861 })
862}
863
864fn delegate_tool_def_with_idle_retire_secs(
865 tool: &meerkat_core::types::ToolDef,
866) -> meerkat_core::types::ToolDef {
867 let mut patched = tool.clone();
868 if !patched.description.contains("IDLE RETIREMENT:") {
869 patched.description.push_str(
870 "\n\nIDLE RETIREMENT:\n\
871 Omit idle_retire_secs to use the runtime default. Pass an integer \
872 number of seconds to override idle auto-retirement for this helper. \
873 Pass null to disable auto-retirement for this helper.",
874 );
875 }
876 if let Some(properties) = patched
877 .input_schema
878 .get_mut("properties")
879 .and_then(Value::as_object_mut)
880 {
881 properties
882 .entry("idle_retire_secs".to_string())
883 .or_insert_with(|| {
884 serde_json::json!({
885 "description": "Override idle auto-retirement for this helper. Omit to use the runtime default, use an integer number of seconds to override, or null to disable auto-retirement for this helper.",
886 "anyOf": [
887 {"type": "integer", "minimum": 0},
888 {"type": "null"}
889 ]
890 })
891 });
892 }
893 patched
894}
895
896fn mob_spawn_tool_def_with_idle_retire_secs(
897 tool: &meerkat_core::types::ToolDef,
898) -> meerkat_core::types::ToolDef {
899 let mut patched = tool.clone();
900 if !patched.description.contains("IDLE RETIREMENT:") {
901 patched.description.push_str(
902 "\n\nIDLE RETIREMENT:\n\
903 Omit idle_retire_secs to leave this spawned member out of auto-retirement. \
904 Pass an integer number of seconds to retire the member after it has been \
905 idle for that long. Pass null to explicitly disable auto-retirement.",
906 );
907 }
908 if let Some(properties) = patched
909 .input_schema
910 .get_mut("properties")
911 .and_then(Value::as_object_mut)
912 {
913 properties
914 .entry("idle_retire_secs".to_string())
915 .or_insert_with(|| {
916 serde_json::json!({
917 "description": "Opt this spawned member into idle auto-retirement. Omit to keep the member indefinitely, use an integer number of seconds to retire after that much idle time, or null to explicitly disable auto-retirement.",
918 "anyOf": [
919 {"type": "integer", "minimum": 0},
920 {"type": "null"}
921 ]
922 })
923 });
924 }
925 patched
926}
927
928fn install_agent_mob_tools(
929 definition: &MobDefinition,
930 slot: Arc<std::sync::RwLock<Option<Arc<dyn meerkat_core::service::MobToolsFactory>>>>,
931 session_service: Arc<dyn MobSessionService>,
932) -> (
933 Arc<meerkat_mob_mcp::MobMcpState>,
934 ImplicitDelegateRetirementOverrides,
935 SharedDefaultLlmClientSlot,
936 SharedConsoleSpawnSinkSlot,
937) {
938 let default_llm_client_slot = Arc::new(std::sync::RwLock::new(None::<Arc<dyn LlmClient>>));
939 let default_llm_client_provider_slot = Arc::clone(&default_llm_client_slot);
940 let mut state = meerkat_mob_mcp::MobMcpState::new(session_service);
941 if let Some(base_store) = state.realm_profile_store().cloned()
942 && let Some(store) = DefinitionSeededRealmProfileStore::new(definition, base_store)
943 {
944 state = state.with_realm_profile_store(Some(Arc::new(store)));
945 }
946 state = state
947 .with_realm_skill_sources(definition.skills.clone())
948 .with_default_llm_client_provider(Some(Arc::new(move || {
949 default_llm_client_provider_slot
950 .read()
951 .unwrap_or_else(std::sync::PoisonError::into_inner)
952 .clone()
953 })));
954 let state = Arc::new(state);
955 let implicit_delegate_retirement_overrides = ImplicitDelegateRetirementOverrides::default();
956 let console_spawn_sink = new_console_spawn_sink_slot();
957 let inner = Arc::new(meerkat_mob_mcp::AgentMobToolSurfaceFactory::new(
958 Arc::clone(&state),
959 ));
960 let factory = Arc::new(AutoWireParentMobToolsFactory {
961 inner,
962 implicit_delegate_retirement_overrides: implicit_delegate_retirement_overrides.clone(),
963 console_spawn_sink: Arc::clone(&console_spawn_sink),
964 });
965 *slot
966 .write()
967 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(factory);
968 (
969 state,
970 implicit_delegate_retirement_overrides,
971 default_llm_client_slot,
972 console_spawn_sink,
973 )
974}
975
976#[cfg(test)]
977#[allow(dead_code)]
978#[derive(Debug, Clone)]
979pub(crate) struct RuntimeTurnTrace {
980 pub(crate) session_id: String,
981 pub(crate) boundary: String,
982 pub(crate) contributing_input_count: usize,
983 pub(crate) outcome: String,
984}
985
986fn is_replay_unsafe_server_tool_content(name: &str, content: &Value) -> bool {
987 name == "web_search_annotations"
988 || content
989 .get("type")
990 .and_then(Value::as_str)
991 .is_some_and(|kind| kind.starts_with("response."))
992}
993
994fn sanitize_llm_request_for_stateless_replay(request: &LlmRequest) -> LlmRequest {
995 let mut sanitized = request.clone();
996 sanitized.messages = request
997 .messages
998 .iter()
999 .cloned()
1000 .map(sanitize_message_for_stateless_replay)
1001 .collect();
1002 sanitized
1003}
1004
1005fn sanitize_create_session_request_llm_override(req: &mut CreateSessionRequest) {
1006 let Some(build) = req.build.as_mut() else {
1007 return;
1008 };
1009 let Some(client) = build
1010 .llm_client_override
1011 .as_ref()
1012 .and_then(meerkat::decode_llm_client_override_from_service)
1013 else {
1014 return;
1015 };
1016 build.llm_client_override = Some(meerkat::encode_llm_client_override_for_service(
1017 ReplaySanitizingLlmClient::wrap(client),
1018 ));
1019}
1020
1021const SHELL_BUILTIN_TOOL_NAMES: [&str; 4] = [
1022 "shell",
1023 "shell_job_status",
1024 "shell_jobs",
1025 "shell_job_cancel",
1026];
1027const COMMS_TOOL_NAMES: [&str; 4] = ["peers", "send_message", "send_request", "send_response"];
1028
1029fn shell_and_comms_tool_filter() -> meerkat_core::ToolFilter {
1030 meerkat_core::ToolFilter::Allow(
1031 SHELL_BUILTIN_TOOL_NAMES
1032 .iter()
1033 .chain(COMMS_TOOL_NAMES.iter())
1034 .map(|name| (*name).to_string())
1035 .collect(),
1036 )
1037}
1038
1039pub fn ensure_shell_tooling_build_substrate(req: &mut CreateSessionRequest) {
1045 let Some(build) = req.build.as_mut() else {
1046 return;
1047 };
1048 if matches!(
1049 build.override_shell,
1050 meerkat_core::ToolCategoryOverride::Enable
1051 ) && matches!(
1052 build.override_builtins,
1053 meerkat_core::ToolCategoryOverride::Disable
1054 ) {
1055 build.override_builtins = meerkat_core::ToolCategoryOverride::Enable;
1056 if build.initial_tool_filter.is_none() {
1057 build.initial_tool_filter = Some(shell_and_comms_tool_filter());
1058 }
1059 }
1060}
1061
1062fn sanitize_message_for_stateless_replay(message: Message) -> Message {
1063 match message {
1064 Message::BlockAssistant(mut assistant) => {
1065 assistant.blocks = assistant
1066 .blocks
1067 .into_iter()
1068 .filter_map(|block| match block {
1069 AssistantBlock::ServerToolContent { kind, content, .. }
1072 if is_replay_unsafe_server_tool_content(kind.provider_name(), &content) =>
1073 {
1074 None
1075 }
1076 other => Some(other),
1077 })
1078 .collect();
1079 Message::BlockAssistant(assistant)
1080 }
1081 other => other,
1082 }
1083}
1084
1085fn build_persistent_runtime_store(store_path: &Path) -> Arc<dyn meerkat_runtime::RuntimeStore> {
1095 let runtime_db = store_path.join("runtime.sqlite");
1096 match meerkat_runtime::store::SqliteRuntimeStore::new(&runtime_db) {
1097 Ok(store) => Arc::new(store),
1098 Err(err) => {
1099 tracing::warn!(
1100 path = %runtime_db.display(),
1101 error = %err,
1102 "failed to open SqliteRuntimeStore; falling back to InMemoryRuntimeStore. \
1103 Sessions will not survive process restart and archive operations may fail.",
1104 );
1105 Arc::new(meerkat_runtime::InMemoryRuntimeStore::new())
1106 }
1107 }
1108}
1109
1110struct SessionStoreBackedRuntimeStore {
1118 inner: Arc<dyn meerkat_runtime::RuntimeStore>,
1119}
1120
1121impl SessionStoreBackedRuntimeStore {
1122 fn new(
1123 inner: Arc<dyn meerkat_runtime::RuntimeStore>,
1124 _session_store: Arc<dyn SessionStore>,
1125 ) -> Self {
1126 Self { inner }
1127 }
1128}
1129
1130#[async_trait]
1131impl meerkat_runtime::RuntimeStore for SessionStoreBackedRuntimeStore {
1132 fn auth_authority_key(&self) -> Option<String> {
1133 self.inner.auth_authority_key()
1134 }
1135
1136 fn persist_auth_oauth_flow_snapshot(
1137 &self,
1138 snapshot_json: &[u8],
1139 ) -> Result<(), meerkat_runtime::store::RuntimeStoreError> {
1140 self.inner.persist_auth_oauth_flow_snapshot(snapshot_json)
1141 }
1142
1143 fn load_auth_oauth_flow_snapshot(
1144 &self,
1145 ) -> Result<Option<Vec<u8>>, meerkat_runtime::store::RuntimeStoreError> {
1146 self.inner.load_auth_oauth_flow_snapshot()
1147 }
1148
1149 fn update_auth_oauth_flow_snapshot(
1150 &self,
1151 update: &mut meerkat_runtime::store::AuthOAuthFlowSnapshotUpdate<'_>,
1152 ) -> Result<(), meerkat_runtime::store::RuntimeStoreError> {
1153 self.inner.update_auth_oauth_flow_snapshot(update)
1154 }
1155
1156 async fn commit_session_snapshot(
1157 &self,
1158 runtime_id: &meerkat_runtime::LogicalRuntimeId,
1159 session_delta: meerkat_runtime::store::SessionDelta,
1160 ) -> Result<(), meerkat_runtime::store::RuntimeStoreError> {
1161 self.inner
1162 .commit_session_snapshot(runtime_id, session_delta)
1163 .await
1164 }
1165
1166 async fn commit_session_transcript_rewrite_snapshot(
1167 &self,
1168 runtime_id: &meerkat_runtime::LogicalRuntimeId,
1169 session_delta: meerkat_runtime::store::SessionDelta,
1170 commit: &meerkat_core::TranscriptRewriteCommit,
1171 ) -> Result<(), meerkat_runtime::store::RuntimeStoreError> {
1172 self.inner
1173 .commit_session_transcript_rewrite_snapshot(runtime_id, session_delta, commit)
1174 .await
1175 }
1176
1177 async fn atomic_apply(
1178 &self,
1179 runtime_id: &meerkat_runtime::LogicalRuntimeId,
1180 session_delta: Option<meerkat_runtime::store::SessionDelta>,
1181 receipt: meerkat_core::lifecycle::RunBoundaryReceipt,
1182 input_updates: Vec<InputStatePersistenceRecord>,
1183 session_store_key: Option<meerkat_core::types::SessionId>,
1184 ) -> Result<(), meerkat_runtime::store::RuntimeStoreError> {
1185 self.inner
1186 .atomic_apply(
1187 runtime_id,
1188 session_delta,
1189 receipt,
1190 input_updates,
1191 session_store_key,
1192 )
1193 .await
1194 }
1195
1196 async fn load_input_states(
1197 &self,
1198 runtime_id: &meerkat_runtime::LogicalRuntimeId,
1199 ) -> Result<Vec<StoredInputState>, meerkat_runtime::store::RuntimeStoreError> {
1200 self.inner.load_input_states(runtime_id).await
1201 }
1202
1203 async fn load_boundary_receipt(
1204 &self,
1205 runtime_id: &meerkat_runtime::LogicalRuntimeId,
1206 run_id: &meerkat_core::lifecycle::RunId,
1207 sequence: u64,
1208 ) -> Result<
1209 Option<meerkat_core::lifecycle::RunBoundaryReceipt>,
1210 meerkat_runtime::store::RuntimeStoreError,
1211 > {
1212 self.inner
1213 .load_boundary_receipt(runtime_id, run_id, sequence)
1214 .await
1215 }
1216
1217 async fn load_session_snapshot(
1218 &self,
1219 runtime_id: &meerkat_runtime::LogicalRuntimeId,
1220 ) -> Result<Option<Vec<u8>>, meerkat_runtime::store::RuntimeStoreError> {
1221 self.inner.load_session_snapshot(runtime_id).await
1222 }
1223
1224 async fn clear_session_snapshot(
1225 &self,
1226 runtime_id: &meerkat_runtime::LogicalRuntimeId,
1227 ) -> Result<(), meerkat_runtime::store::RuntimeStoreError> {
1228 self.inner.clear_session_snapshot(runtime_id).await
1229 }
1230
1231 async fn replace_session_snapshot_if_current(
1232 &self,
1233 runtime_id: &meerkat_runtime::LogicalRuntimeId,
1234 expected_current: &[u8],
1235 replacement: Vec<u8>,
1236 ) -> Result<bool, meerkat_runtime::store::RuntimeStoreError> {
1237 self.inner
1238 .replace_session_snapshot_if_current(runtime_id, expected_current, replacement)
1239 .await
1240 }
1241
1242 async fn clear_session_snapshot_if_current(
1243 &self,
1244 runtime_id: &meerkat_runtime::LogicalRuntimeId,
1245 expected_current: &[u8],
1246 ) -> Result<bool, meerkat_runtime::store::RuntimeStoreError> {
1247 self.inner
1248 .clear_session_snapshot_if_current(runtime_id, expected_current)
1249 .await
1250 }
1251
1252 async fn persist_input_state(
1253 &self,
1254 runtime_id: &meerkat_runtime::LogicalRuntimeId,
1255 state: &InputStatePersistenceRecord,
1256 ) -> Result<(), meerkat_runtime::store::RuntimeStoreError> {
1257 self.inner.persist_input_state(runtime_id, state).await
1258 }
1259
1260 async fn load_input_state(
1261 &self,
1262 runtime_id: &meerkat_runtime::LogicalRuntimeId,
1263 input_id: &meerkat_core::lifecycle::InputId,
1264 ) -> Result<Option<StoredInputState>, meerkat_runtime::store::RuntimeStoreError> {
1265 self.inner.load_input_state(runtime_id, input_id).await
1266 }
1267
1268 async fn load_machine_lifecycle_record(
1269 &self,
1270 runtime_id: &meerkat_runtime::LogicalRuntimeId,
1271 ) -> Result<Option<Vec<u8>>, meerkat_runtime::store::RuntimeStoreError> {
1272 self.inner.load_machine_lifecycle_record(runtime_id).await
1273 }
1274
1275 async fn commit_machine_lifecycle(
1276 &self,
1277 runtime_id: &meerkat_runtime::LogicalRuntimeId,
1278 commit: MachineLifecycleCommit,
1279 input_states: &[InputStatePersistenceRecord],
1280 ) -> Result<(), meerkat_runtime::store::RuntimeStoreError> {
1281 self.inner
1282 .commit_machine_lifecycle(runtime_id, commit, input_states)
1283 .await
1284 }
1285
1286 async fn persist_ops_lifecycle(
1287 &self,
1288 runtime_id: &meerkat_runtime::LogicalRuntimeId,
1289 snapshot: &meerkat_runtime::ops_lifecycle::PersistedOpsSnapshot,
1290 ) -> Result<(), meerkat_runtime::store::RuntimeStoreError> {
1291 self.inner.persist_ops_lifecycle(runtime_id, snapshot).await
1292 }
1293
1294 async fn load_ops_lifecycle(
1295 &self,
1296 runtime_id: &meerkat_runtime::LogicalRuntimeId,
1297 ) -> Result<
1298 Option<meerkat_runtime::ops_lifecycle::PersistedOpsSnapshot>,
1299 meerkat_runtime::store::RuntimeStoreError,
1300 > {
1301 self.inner.load_ops_lifecycle(runtime_id).await
1302 }
1303}
1304
1305#[cfg(test)]
1306static RUNTIME_TURN_TRACES: OnceLock<Mutex<Vec<RuntimeTurnTrace>>> = OnceLock::new();
1307
1308#[cfg(test)]
1309fn runtime_turn_traces() -> &'static Mutex<Vec<RuntimeTurnTrace>> {
1310 RUNTIME_TURN_TRACES.get_or_init(|| Mutex::new(Vec::new()))
1311}
1312
1313#[cfg(test)]
1314#[allow(clippy::expect_used)]
1315fn record_runtime_turn_trace(trace: RuntimeTurnTrace) {
1316 runtime_turn_traces()
1317 .lock()
1318 .expect("runtime turn traces mutex")
1319 .push(trace);
1320}
1321
1322#[cfg(test)]
1323#[allow(dead_code)]
1324#[allow(clippy::expect_used)]
1325pub(crate) fn take_runtime_turn_traces() -> Vec<RuntimeTurnTrace> {
1326 std::mem::take(
1327 &mut *runtime_turn_traces()
1328 .lock()
1329 .expect("runtime turn traces mutex"),
1330 )
1331}
1332
1333#[cfg(not(test))]
1334#[allow(dead_code)]
1335fn record_runtime_turn_trace(_trace: ()) {}
1336
1337fn runtime_turn_diagnostics_enabled() -> bool {
1338 std::env::var_os("MOBKIT_TRACE_RUNTIME_TURNS").is_some()
1339}
1340
1341fn summarize_runtime_prompt(prompt: &meerkat_core::ContentInput) -> String {
1342 match prompt {
1343 meerkat_core::ContentInput::Text(text) => {
1344 text.lines().take(6).collect::<Vec<_>>().join(" ")
1345 }
1346 meerkat_core::ContentInput::Blocks(blocks) => blocks
1347 .iter()
1348 .map(|block| block.text_projection().to_string())
1349 .collect::<Vec<_>>()
1350 .join(" ")
1351 .lines()
1352 .take(6)
1353 .collect::<Vec<_>>()
1354 .join(" "),
1355 }
1356}
1357
1358pub fn mob_definition_may_use_image_generation(definition: &MobDefinition) -> bool {
1364 definition.profiles.values().any(|binding| {
1365 binding
1366 .as_inline()
1367 .is_none_or(|profile| profile.tools.image_generation)
1368 })
1369}
1370
1371pub fn mob_definition_may_use_shell(definition: &MobDefinition) -> bool {
1374 definition.profiles.values().any(|binding| {
1375 binding
1376 .as_inline()
1377 .is_none_or(|profile| profile.tools.shell)
1378 })
1379}
1380
1381fn normalize_runtime_turn_request(
1382 mut req: meerkat_core::service::StartTurnRequest,
1383) -> meerkat_core::service::StartTurnRequest {
1384 req.runtime.handling_mode = meerkat_core::types::HandlingMode::Queue;
1390 if let Some(metadata) = req.runtime.turn_metadata.as_mut() {
1393 metadata.render_metadata = None;
1394 }
1395 req
1396}
1397
1398macro_rules! delegate_mob_session_service {
1401 ($wrapper:ty) => {
1402 #[async_trait]
1403 impl meerkat_core::service::SessionService for $wrapper {
1404 async fn create_session(
1405 &self,
1406 mut req: CreateSessionRequest,
1407 ) -> Result<meerkat_core::types::RunResult, SessionError> {
1408 (self.hook)(&mut req).await?;
1409 ensure_shell_tooling_build_substrate(&mut req);
1410 sanitize_create_session_request_llm_override(&mut req);
1411
1412 let ctx = SessionCreatedContext {
1414 model: req.model.clone(),
1415 labels: req.labels.clone().unwrap_or_default(),
1416 system_prompt: req
1417 .system_prompt
1418 .as_set_prompt()
1419 .map(ToString::to_string),
1420 };
1421 let result = self.inner.create_session(req).await?;
1422
1423 if let Some(ref after_hook) = self.after_create_hook {
1425 after_hook(result.session_id.clone(), ctx).await;
1426 }
1427
1428 Ok(result)
1429 }
1430 async fn start_turn(
1431 &self,
1432 id: &meerkat_core::types::SessionId,
1433 req: meerkat_core::service::StartTurnRequest,
1434 ) -> Result<meerkat_core::types::RunResult, SessionError> {
1435 self.inner.start_turn(id, req).await
1436 }
1437 async fn interrupt(
1438 &self,
1439 id: &meerkat_core::types::SessionId,
1440 ) -> Result<(), SessionError> {
1441 self.inner.interrupt(id).await
1442 }
1443 async fn cancel_after_boundary(
1444 &self,
1445 id: &meerkat_core::types::SessionId,
1446 ) -> Result<(), SessionError> {
1447 self.inner.cancel_after_boundary(id).await
1448 }
1449 async fn set_session_client(
1450 &self,
1451 id: &meerkat_core::types::SessionId,
1452 client: Arc<dyn meerkat_core::AgentLlmClient>,
1453 ) -> Result<(), SessionError> {
1454 self.inner
1455 .set_session_client(id, ReplaySanitizingAgentLlmClient::wrap(client))
1456 .await
1457 }
1458 async fn hot_swap_session_llm_identity(
1459 &self,
1460 id: &meerkat_core::types::SessionId,
1461 client: Arc<dyn meerkat_core::AgentLlmClient>,
1462 identity: meerkat_core::session::SessionLlmIdentity,
1463 request_policy: meerkat_core::SessionLlmRequestPolicy,
1464 ) -> Result<(), SessionError> {
1465 self.inner
1466 .hot_swap_session_llm_identity(
1467 id,
1468 ReplaySanitizingAgentLlmClient::wrap(client),
1469 identity,
1470 request_policy,
1471 )
1472 .await
1473 }
1474 async fn update_session_mob_authority_context(
1475 &self,
1476 id: &meerkat_core::types::SessionId,
1477 authority_context: Option<meerkat_core::service::MobToolAuthorityContext>,
1478 ) -> Result<(), SessionError> {
1479 self.inner
1480 .update_session_mob_authority_context(id, authority_context)
1481 .await
1482 }
1483 async fn has_live_session(
1484 &self,
1485 id: &meerkat_core::types::SessionId,
1486 ) -> Result<bool, SessionError> {
1487 self.inner.has_live_session(id).await
1488 }
1489 async fn set_session_tool_visibility_state(
1490 &self,
1491 id: &meerkat_core::types::SessionId,
1492 state: Option<meerkat_core::SessionToolVisibilityState>,
1493 ) -> Result<(), SessionError> {
1494 self.inner
1495 .set_session_tool_visibility_state(id, state)
1496 .await
1497 }
1498 async fn set_session_tool_filter(
1499 &self,
1500 id: &meerkat_core::types::SessionId,
1501 filter: meerkat_core::ToolFilter,
1502 ) -> Result<(), SessionError> {
1503 self.inner.set_session_tool_filter(id, filter).await
1504 }
1505 async fn read(
1506 &self,
1507 id: &meerkat_core::types::SessionId,
1508 ) -> Result<meerkat_core::service::SessionView, SessionError> {
1509 self.inner.read(id).await
1510 }
1511 async fn list(
1512 &self,
1513 query: meerkat_core::service::SessionQuery,
1514 ) -> Result<Vec<meerkat_core::service::SessionSummary>, SessionError> {
1515 self.inner.list(query).await
1516 }
1517 async fn archive(
1518 &self,
1519 id: &meerkat_core::types::SessionId,
1520 ) -> Result<(), SessionError> {
1521 self.inner.archive(id).await
1522 }
1523 async fn subscribe_session_events(
1524 &self,
1525 id: &meerkat_core::types::SessionId,
1526 ) -> Result<meerkat_core::comms::EventStream, meerkat_core::comms::StreamError> {
1527 meerkat_core::service::SessionService::subscribe_session_events(
1528 self.inner.as_ref(),
1529 id,
1530 )
1531 .await
1532 }
1533 }
1534
1535 #[async_trait]
1536 impl meerkat_core::service::SessionServiceCommsExt for $wrapper {
1537 async fn comms_runtime(
1538 &self,
1539 id: &meerkat_core::types::SessionId,
1540 ) -> Option<Arc<dyn meerkat_core::agent::CommsRuntime>> {
1541 self.inner.comms_runtime(id).await
1542 }
1543
1544 async fn event_injector(
1545 &self,
1546 id: &meerkat_core::types::SessionId,
1547 ) -> Option<Arc<dyn meerkat_core::EventInjector>> {
1548 self.inner.event_injector(id).await
1549 }
1550
1551 async fn interaction_event_injector(
1552 &self,
1553 id: &meerkat_core::types::SessionId,
1554 ) -> Option<Arc<dyn meerkat_core::event_injector::SubscribableInjector>> {
1555 self.inner.interaction_event_injector(id).await
1556 }
1557 }
1558
1559 #[async_trait]
1560 impl meerkat_core::service::SessionServiceControlExt for $wrapper {
1561 async fn append_system_context(
1562 &self,
1563 id: &meerkat_core::types::SessionId,
1564 req: meerkat_core::service::AppendSystemContextRequest,
1565 ) -> Result<
1566 meerkat_core::service::AppendSystemContextResult,
1567 meerkat_core::service::SessionControlError,
1568 > {
1569 self.inner.append_system_context(id, req).await
1570 }
1571 async fn stage_tool_results(
1572 &self,
1573 id: &meerkat_core::types::SessionId,
1574 req: meerkat_core::service::StageToolResultsRequest,
1575 ) -> Result<meerkat_core::service::StageToolResultsResult, SessionError> {
1576 self.inner.stage_tool_results(id, req).await
1577 }
1578 }
1579
1580 #[async_trait]
1581 impl meerkat_core::service::SessionServiceHistoryExt for $wrapper {
1582 async fn read_history(
1583 &self,
1584 id: &meerkat_core::types::SessionId,
1585 query: meerkat_core::service::SessionHistoryQuery,
1586 ) -> Result<meerkat_core::service::SessionHistoryPage, SessionError> {
1587 self.inner.read_history(id, query).await
1588 }
1589 }
1590
1591 #[async_trait]
1592 impl MobSessionService for $wrapper {
1593 fn supports_persistent_sessions(&self) -> bool {
1594 self.inner.supports_persistent_sessions()
1595 }
1596 fn runtime_adapter(&self) -> Option<Arc<meerkat_runtime::MeerkatMachine>> {
1597 self.runtime_adapter_override
1598 .clone()
1599 .or_else(|| self.inner.runtime_adapter())
1600 }
1601 async fn interrupt_with_machine_authority(
1602 &self,
1603 session_id: &meerkat_core::types::SessionId,
1604 authority: meerkat_runtime::MachineSessionControlAuthority,
1605 ) -> Result<(), SessionError> {
1606 self.inner
1607 .interrupt_with_machine_authority(session_id, authority)
1608 .await
1609 }
1610 async fn cancel_after_boundary_with_machine_authority(
1611 &self,
1612 session_id: &meerkat_core::types::SessionId,
1613 authority: meerkat_runtime::MachineSessionControlAuthority,
1614 ) -> Result<(), SessionError> {
1615 self.inner
1616 .cancel_after_boundary_with_machine_authority(session_id, authority)
1617 .await
1618 }
1619 async fn session_belongs_to_mob(
1620 &self,
1621 session_id: &meerkat_core::types::SessionId,
1622 mob_id: &meerkat_mob::MobId,
1623 ) -> bool {
1624 self.inner.session_belongs_to_mob(session_id, mob_id).await
1625 }
1626 async fn load_persisted_session(
1627 &self,
1628 session_id: &meerkat_core::types::SessionId,
1629 ) -> Result<Option<meerkat_core::session::Session>, SessionError> {
1630 self.inner.load_persisted_session(session_id).await
1631 }
1632 async fn subscribe_session_events(
1633 &self,
1634 session_id: &meerkat_core::types::SessionId,
1635 ) -> Result<meerkat_core::comms::EventStream, meerkat_core::comms::StreamError> {
1636 meerkat_mob::MobSessionService::subscribe_session_events(
1637 self.inner.as_ref(),
1638 session_id,
1639 )
1640 .await
1641 }
1642 async fn archive_with_mob_lifecycle_authority(
1643 &self,
1644 session_id: &meerkat_core::types::SessionId,
1645 ) -> Result<(), SessionError> {
1646 match self
1647 .inner
1648 .archive_with_mob_lifecycle_authority(session_id)
1649 .await
1650 {
1651 Err(err)
1652 if is_stopped_session_archive_retire_rejection(&err.to_string()) =>
1653 {
1654 tracing::warn!(
1665 session_id = %session_id,
1666 error = %err,
1667 "archive: tolerating Retire rejection for stopped idle session; archive document committed"
1668 );
1669 Ok(())
1670 }
1671 other => other,
1672 }
1673 }
1674 async fn execution_snapshot(
1675 &self,
1676 session_id: &meerkat_core::types::SessionId,
1677 ) -> Result<Option<meerkat_core::agent::AgentExecutionSnapshot>, SessionError> {
1678 self.inner.execution_snapshot(session_id).await
1679 }
1680 async fn tool_scope_snapshot(
1681 &self,
1682 session_id: &meerkat_core::types::SessionId,
1683 ) -> Result<Option<meerkat_core::ToolScopeSnapshot>, SessionError> {
1684 self.inner.tool_scope_snapshot(session_id).await
1685 }
1686 async fn external_tool_surface_snapshot(
1687 &self,
1688 session_id: &meerkat_core::types::SessionId,
1689 ) -> Result<Option<meerkat_core::ExternalToolSurfaceSnapshot>, SessionError> {
1690 self.inner.external_tool_surface_snapshot(session_id).await
1691 }
1692 async fn peer_ingress_runtime_snapshot(
1693 &self,
1694 session_id: &meerkat_core::types::SessionId,
1695 ) -> Result<Option<meerkat_core::PeerIngressRuntimeSnapshot>, SessionError> {
1696 self.inner.peer_ingress_runtime_snapshot(session_id).await
1697 }
1698 async fn apply_runtime_turn(
1699 &self,
1700 session_id: &meerkat_core::types::SessionId,
1701 run_id: meerkat_core::lifecycle::RunId,
1702 req: meerkat_core::service::StartTurnRequest,
1703 boundary: meerkat_core::lifecycle::run_primitive::RunApplyBoundary,
1704 contributing_input_ids: Vec<meerkat_core::lifecycle::InputId>,
1705 ) -> Result<meerkat_core::lifecycle::core_executor::CoreApplyOutput, SessionError> {
1706 #[cfg(test)]
1707 let boundary_name = format!("{boundary:?}");
1708 #[cfg(test)]
1709 let contributing_count = contributing_input_ids.len();
1710 let run_id_for_log = run_id.to_string();
1711 let prompt_summary = if runtime_turn_diagnostics_enabled() {
1712 Some(summarize_runtime_prompt(&req.prompt))
1713 } else {
1714 None
1715 };
1716 if let Some(summary) = prompt_summary.as_ref() {
1717 tracing::warn!(
1718 session_id = %session_id,
1719 run_id = %run_id_for_log,
1720 boundary = ?boundary,
1721 contributing_inputs = contributing_input_ids.len(),
1722 prompt = %summary,
1723 runtime = ?req.runtime,
1724 "mobkit runtime turn start"
1725 );
1726 }
1727 let result = self
1728 .inner
1729 .apply_runtime_turn(
1730 session_id,
1731 run_id,
1732 normalize_runtime_turn_request(req),
1733 boundary,
1734 contributing_input_ids,
1735 )
1736 .await;
1737 #[cfg(test)]
1738 record_runtime_turn_trace(RuntimeTurnTrace {
1739 session_id: session_id.to_string(),
1740 boundary: boundary_name,
1741 contributing_input_count: contributing_count,
1742 outcome: match &result {
1743 Ok(_) => "ok".to_string(),
1744 Err(error) => format!("err:{error}"),
1745 },
1746 });
1747 if runtime_turn_diagnostics_enabled() {
1748 match &result {
1749 Ok(_) => tracing::warn!(
1750 session_id = %session_id,
1751 run_id = %run_id_for_log,
1752 "mobkit runtime turn ok"
1753 ),
1754 Err(error) => tracing::error!(
1755 session_id = %session_id,
1756 run_id = %run_id_for_log,
1757 error = %error,
1758 error_debug = ?error,
1759 "mobkit runtime turn error"
1760 ),
1761 }
1762 }
1763 result
1764 }
1765 async fn apply_runtime_context_appends(
1766 &self,
1767 session_id: &meerkat_core::types::SessionId,
1768 run_id: meerkat_core::lifecycle::RunId,
1769 appends: Vec<meerkat_core::session::PendingSystemContextAppend>,
1770 contributing_input_ids: Vec<meerkat_core::lifecycle::InputId>,
1771 ) -> Result<meerkat_core::lifecycle::core_executor::CoreApplyOutput, SessionError> {
1772 self.inner
1773 .apply_runtime_context_appends(
1774 session_id,
1775 run_id,
1776 appends,
1777 contributing_input_ids,
1778 )
1779 .await
1780 }
1781 async fn apply_runtime_context_appends_with_boundary(
1782 &self,
1783 session_id: &meerkat_core::types::SessionId,
1784 run_id: meerkat_core::lifecycle::RunId,
1785 appends: Vec<meerkat_core::session::PendingSystemContextAppend>,
1786 boundary: meerkat_core::lifecycle::run_primitive::RunApplyBoundary,
1787 contributing_input_ids: Vec<meerkat_core::lifecycle::InputId>,
1788 ) -> Result<meerkat_core::lifecycle::core_executor::CoreApplyOutput, SessionError> {
1789 self.inner
1790 .apply_runtime_context_appends_with_boundary(
1791 session_id,
1792 run_id,
1793 appends,
1794 boundary,
1795 contributing_input_ids,
1796 )
1797 .await
1798 }
1799 async fn apply_runtime_system_context_for_turn(
1800 &self,
1801 session_id: &meerkat_core::types::SessionId,
1802 appends: Vec<meerkat_core::session::PendingSystemContextAppend>,
1803 ) -> Result<(), SessionError> {
1804 self.inner
1805 .apply_runtime_system_context_for_turn(session_id, appends)
1806 .await
1807 }
1808 async fn stage_runtime_system_context_for_active_turn(
1809 &self,
1810 session_id: &meerkat_core::types::SessionId,
1811 expected_run_id: &meerkat_core::lifecycle::RunId,
1812 appends: Vec<meerkat_core::session::PendingSystemContextAppend>,
1813 ) -> Result<Option<Vec<u8>>, SessionError> {
1814 self.inner
1815 .stage_runtime_system_context_for_active_turn(
1816 session_id,
1817 expected_run_id,
1818 appends,
1819 )
1820 .await
1821 }
1822 async fn discard_runtime_system_context_for_active_turn(
1823 &self,
1824 session_id: &meerkat_core::types::SessionId,
1825 expected_run_id: &meerkat_core::lifecycle::RunId,
1826 idempotency_keys: Vec<String>,
1827 ) -> Result<(), SessionError> {
1828 self.inner
1829 .discard_runtime_system_context_for_active_turn(
1830 session_id,
1831 expected_run_id,
1832 idempotency_keys,
1833 )
1834 .await
1835 }
1836 async fn active_turn_system_context_boundary_available(
1837 &self,
1838 session_id: &meerkat_core::types::SessionId,
1839 ) -> Result<Option<bool>, SessionError> {
1840 self.inner
1841 .active_turn_system_context_boundary_available(session_id)
1842 .await
1843 }
1844 async fn discard_live_session(
1845 &self,
1846 session_id: &meerkat_core::types::SessionId,
1847 ) -> Result<(), SessionError> {
1848 self.inner.discard_live_session(session_id).await
1849 }
1850 async fn checkpoint_committed_runtime_session_snapshot(
1851 &self,
1852 session_id: &meerkat_core::types::SessionId,
1853 session_snapshot: &[u8],
1854 ) -> Result<(), SessionError> {
1855 self.inner
1856 .checkpoint_committed_runtime_session_snapshot(session_id, session_snapshot)
1857 .await
1858 }
1859 async fn cancel_all_checkpointers(&self) {
1860 self.inner.cancel_all_checkpointers().await;
1861 }
1862 async fn rearm_all_checkpointers(&self) {
1863 self.inner.rearm_all_checkpointers().await;
1864 }
1865 }
1866 };
1867}
1868
1869delegate_mob_session_service!(PreBuildMobSessionService);
1870
1871struct AfterCreateMobSessionService {
1876 inner: Arc<dyn MobSessionService>,
1877 after_hook: AfterCreateHook,
1878}
1879
1880#[async_trait]
1881impl meerkat_core::service::SessionService for AfterCreateMobSessionService {
1882 async fn create_session(
1883 &self,
1884 mut req: CreateSessionRequest,
1885 ) -> Result<meerkat_core::types::RunResult, SessionError> {
1886 sanitize_create_session_request_llm_override(&mut req);
1887 ensure_shell_tooling_build_substrate(&mut req);
1888 let ctx = SessionCreatedContext {
1898 model: req.model.clone(),
1899 labels: req.labels.clone().unwrap_or_default(),
1900 system_prompt: req.system_prompt.as_set_prompt().map(ToString::to_string),
1901 };
1902 let result = self.inner.create_session(req).await?;
1903 (self.after_hook)(result.session_id.clone(), ctx).await;
1904 Ok(result)
1905 }
1906 async fn start_turn(
1907 &self,
1908 id: &meerkat_core::types::SessionId,
1909 req: meerkat_core::service::StartTurnRequest,
1910 ) -> Result<meerkat_core::types::RunResult, SessionError> {
1911 self.inner.start_turn(id, req).await
1912 }
1913 async fn interrupt(&self, id: &meerkat_core::types::SessionId) -> Result<(), SessionError> {
1914 self.inner.interrupt(id).await
1915 }
1916 async fn cancel_after_boundary(
1917 &self,
1918 id: &meerkat_core::types::SessionId,
1919 ) -> Result<(), SessionError> {
1920 self.inner.cancel_after_boundary(id).await
1921 }
1922 async fn set_session_client(
1923 &self,
1924 id: &meerkat_core::types::SessionId,
1925 client: Arc<dyn meerkat_core::AgentLlmClient>,
1926 ) -> Result<(), SessionError> {
1927 self.inner
1928 .set_session_client(id, ReplaySanitizingAgentLlmClient::wrap(client))
1929 .await
1930 }
1931 async fn hot_swap_session_llm_identity(
1932 &self,
1933 id: &meerkat_core::types::SessionId,
1934 client: Arc<dyn meerkat_core::AgentLlmClient>,
1935 identity: meerkat_core::session::SessionLlmIdentity,
1936 request_policy: meerkat_core::SessionLlmRequestPolicy,
1937 ) -> Result<(), SessionError> {
1938 self.inner
1939 .hot_swap_session_llm_identity(
1940 id,
1941 ReplaySanitizingAgentLlmClient::wrap(client),
1942 identity,
1943 request_policy,
1944 )
1945 .await
1946 }
1947 async fn update_session_mob_authority_context(
1948 &self,
1949 id: &meerkat_core::types::SessionId,
1950 authority_context: Option<meerkat_core::service::MobToolAuthorityContext>,
1951 ) -> Result<(), SessionError> {
1952 self.inner
1953 .update_session_mob_authority_context(id, authority_context)
1954 .await
1955 }
1956 async fn has_live_session(
1957 &self,
1958 id: &meerkat_core::types::SessionId,
1959 ) -> Result<bool, SessionError> {
1960 self.inner.has_live_session(id).await
1961 }
1962 async fn set_session_tool_visibility_state(
1963 &self,
1964 id: &meerkat_core::types::SessionId,
1965 state: Option<meerkat_core::SessionToolVisibilityState>,
1966 ) -> Result<(), SessionError> {
1967 self.inner
1968 .set_session_tool_visibility_state(id, state)
1969 .await
1970 }
1971 async fn set_session_tool_filter(
1972 &self,
1973 id: &meerkat_core::types::SessionId,
1974 filter: meerkat_core::ToolFilter,
1975 ) -> Result<(), SessionError> {
1976 self.inner.set_session_tool_filter(id, filter).await
1977 }
1978 async fn read(
1979 &self,
1980 id: &meerkat_core::types::SessionId,
1981 ) -> Result<meerkat_core::service::SessionView, SessionError> {
1982 self.inner.read(id).await
1983 }
1984 async fn list(
1985 &self,
1986 query: meerkat_core::service::SessionQuery,
1987 ) -> Result<Vec<meerkat_core::service::SessionSummary>, SessionError> {
1988 self.inner.list(query).await
1989 }
1990 async fn archive(&self, id: &meerkat_core::types::SessionId) -> Result<(), SessionError> {
1991 self.inner.archive(id).await
1992 }
1993 async fn subscribe_session_events(
1994 &self,
1995 id: &meerkat_core::types::SessionId,
1996 ) -> Result<meerkat_core::comms::EventStream, meerkat_core::comms::StreamError> {
1997 meerkat_core::service::SessionService::subscribe_session_events(self.inner.as_ref(), id)
1998 .await
1999 }
2000}
2001
2002#[async_trait]
2003impl meerkat_core::service::SessionServiceCommsExt for AfterCreateMobSessionService {
2004 async fn comms_runtime(
2005 &self,
2006 id: &meerkat_core::types::SessionId,
2007 ) -> Option<Arc<dyn meerkat_core::agent::CommsRuntime>> {
2008 self.inner.comms_runtime(id).await
2009 }
2010
2011 async fn event_injector(
2012 &self,
2013 id: &meerkat_core::types::SessionId,
2014 ) -> Option<Arc<dyn meerkat_core::EventInjector>> {
2015 self.inner.event_injector(id).await
2016 }
2017
2018 async fn interaction_event_injector(
2019 &self,
2020 id: &meerkat_core::types::SessionId,
2021 ) -> Option<Arc<dyn meerkat_core::event_injector::SubscribableInjector>> {
2022 self.inner.interaction_event_injector(id).await
2023 }
2024}
2025
2026#[async_trait]
2027impl meerkat_core::service::SessionServiceControlExt for AfterCreateMobSessionService {
2028 async fn append_system_context(
2029 &self,
2030 id: &meerkat_core::types::SessionId,
2031 req: meerkat_core::service::AppendSystemContextRequest,
2032 ) -> Result<
2033 meerkat_core::service::AppendSystemContextResult,
2034 meerkat_core::service::SessionControlError,
2035 > {
2036 self.inner.append_system_context(id, req).await
2037 }
2038 async fn stage_tool_results(
2039 &self,
2040 id: &meerkat_core::types::SessionId,
2041 req: meerkat_core::service::StageToolResultsRequest,
2042 ) -> Result<meerkat_core::service::StageToolResultsResult, SessionError> {
2043 self.inner.stage_tool_results(id, req).await
2044 }
2045}
2046
2047#[async_trait]
2048impl meerkat_core::service::SessionServiceHistoryExt for AfterCreateMobSessionService {
2049 async fn read_history(
2050 &self,
2051 id: &meerkat_core::types::SessionId,
2052 query: meerkat_core::service::SessionHistoryQuery,
2053 ) -> Result<meerkat_core::service::SessionHistoryPage, SessionError> {
2054 self.inner.read_history(id, query).await
2055 }
2056}
2057
2058#[async_trait]
2059impl MobSessionService for AfterCreateMobSessionService {
2060 fn supports_persistent_sessions(&self) -> bool {
2061 self.inner.supports_persistent_sessions()
2062 }
2063 fn runtime_adapter(&self) -> Option<Arc<meerkat_runtime::MeerkatMachine>> {
2064 self.inner.runtime_adapter()
2065 }
2066 async fn interrupt_with_machine_authority(
2067 &self,
2068 session_id: &meerkat_core::types::SessionId,
2069 authority: meerkat_runtime::MachineSessionControlAuthority,
2070 ) -> Result<(), SessionError> {
2071 self.inner
2072 .interrupt_with_machine_authority(session_id, authority)
2073 .await
2074 }
2075 async fn cancel_after_boundary_with_machine_authority(
2076 &self,
2077 session_id: &meerkat_core::types::SessionId,
2078 authority: meerkat_runtime::MachineSessionControlAuthority,
2079 ) -> Result<(), SessionError> {
2080 self.inner
2081 .cancel_after_boundary_with_machine_authority(session_id, authority)
2082 .await
2083 }
2084 async fn session_belongs_to_mob(
2085 &self,
2086 session_id: &meerkat_core::types::SessionId,
2087 mob_id: &meerkat_mob::MobId,
2088 ) -> bool {
2089 self.inner.session_belongs_to_mob(session_id, mob_id).await
2090 }
2091 async fn load_persisted_session(
2092 &self,
2093 session_id: &meerkat_core::types::SessionId,
2094 ) -> Result<Option<meerkat_core::session::Session>, SessionError> {
2095 self.inner.load_persisted_session(session_id).await
2096 }
2097 async fn subscribe_session_events(
2098 &self,
2099 session_id: &meerkat_core::types::SessionId,
2100 ) -> Result<meerkat_core::comms::EventStream, meerkat_core::comms::StreamError> {
2101 meerkat_mob::MobSessionService::subscribe_session_events(self.inner.as_ref(), session_id)
2102 .await
2103 }
2104 async fn archive_with_mob_lifecycle_authority(
2105 &self,
2106 session_id: &meerkat_core::types::SessionId,
2107 ) -> Result<(), SessionError> {
2108 self.inner
2109 .archive_with_mob_lifecycle_authority(session_id)
2110 .await
2111 }
2112 async fn execution_snapshot(
2113 &self,
2114 session_id: &meerkat_core::types::SessionId,
2115 ) -> Result<Option<meerkat_core::agent::AgentExecutionSnapshot>, SessionError> {
2116 self.inner.execution_snapshot(session_id).await
2117 }
2118 async fn tool_scope_snapshot(
2119 &self,
2120 session_id: &meerkat_core::types::SessionId,
2121 ) -> Result<Option<meerkat_core::ToolScopeSnapshot>, SessionError> {
2122 self.inner.tool_scope_snapshot(session_id).await
2123 }
2124 async fn external_tool_surface_snapshot(
2125 &self,
2126 session_id: &meerkat_core::types::SessionId,
2127 ) -> Result<Option<meerkat_core::ExternalToolSurfaceSnapshot>, SessionError> {
2128 self.inner.external_tool_surface_snapshot(session_id).await
2129 }
2130 async fn peer_ingress_runtime_snapshot(
2131 &self,
2132 session_id: &meerkat_core::types::SessionId,
2133 ) -> Result<Option<meerkat_core::PeerIngressRuntimeSnapshot>, SessionError> {
2134 self.inner.peer_ingress_runtime_snapshot(session_id).await
2135 }
2136 async fn apply_runtime_turn(
2137 &self,
2138 session_id: &meerkat_core::types::SessionId,
2139 run_id: meerkat_core::lifecycle::RunId,
2140 req: meerkat_core::service::StartTurnRequest,
2141 boundary: meerkat_core::lifecycle::run_primitive::RunApplyBoundary,
2142 contributing_input_ids: Vec<meerkat_core::lifecycle::InputId>,
2143 ) -> Result<meerkat_core::lifecycle::core_executor::CoreApplyOutput, SessionError> {
2144 self.inner
2145 .apply_runtime_turn(session_id, run_id, req, boundary, contributing_input_ids)
2146 .await
2147 }
2148 async fn apply_runtime_context_appends(
2149 &self,
2150 session_id: &meerkat_core::types::SessionId,
2151 run_id: meerkat_core::lifecycle::RunId,
2152 appends: Vec<meerkat_core::session::PendingSystemContextAppend>,
2153 contributing_input_ids: Vec<meerkat_core::lifecycle::InputId>,
2154 ) -> Result<meerkat_core::lifecycle::core_executor::CoreApplyOutput, SessionError> {
2155 self.inner
2156 .apply_runtime_context_appends(session_id, run_id, appends, contributing_input_ids)
2157 .await
2158 }
2159 async fn apply_runtime_context_appends_with_boundary(
2160 &self,
2161 session_id: &meerkat_core::types::SessionId,
2162 run_id: meerkat_core::lifecycle::RunId,
2163 appends: Vec<meerkat_core::session::PendingSystemContextAppend>,
2164 boundary: meerkat_core::lifecycle::run_primitive::RunApplyBoundary,
2165 contributing_input_ids: Vec<meerkat_core::lifecycle::InputId>,
2166 ) -> Result<meerkat_core::lifecycle::core_executor::CoreApplyOutput, SessionError> {
2167 self.inner
2168 .apply_runtime_context_appends_with_boundary(
2169 session_id,
2170 run_id,
2171 appends,
2172 boundary,
2173 contributing_input_ids,
2174 )
2175 .await
2176 }
2177 async fn apply_runtime_system_context_for_turn(
2178 &self,
2179 session_id: &meerkat_core::types::SessionId,
2180 appends: Vec<meerkat_core::session::PendingSystemContextAppend>,
2181 ) -> Result<(), SessionError> {
2182 self.inner
2183 .apply_runtime_system_context_for_turn(session_id, appends)
2184 .await
2185 }
2186 async fn stage_runtime_system_context_for_active_turn(
2187 &self,
2188 session_id: &meerkat_core::types::SessionId,
2189 expected_run_id: &meerkat_core::lifecycle::RunId,
2190 appends: Vec<meerkat_core::session::PendingSystemContextAppend>,
2191 ) -> Result<Option<Vec<u8>>, SessionError> {
2192 self.inner
2193 .stage_runtime_system_context_for_active_turn(session_id, expected_run_id, appends)
2194 .await
2195 }
2196 async fn discard_runtime_system_context_for_active_turn(
2197 &self,
2198 session_id: &meerkat_core::types::SessionId,
2199 expected_run_id: &meerkat_core::lifecycle::RunId,
2200 idempotency_keys: Vec<String>,
2201 ) -> Result<(), SessionError> {
2202 self.inner
2203 .discard_runtime_system_context_for_active_turn(
2204 session_id,
2205 expected_run_id,
2206 idempotency_keys,
2207 )
2208 .await
2209 }
2210 async fn active_turn_system_context_boundary_available(
2211 &self,
2212 session_id: &meerkat_core::types::SessionId,
2213 ) -> Result<Option<bool>, SessionError> {
2214 self.inner
2215 .active_turn_system_context_boundary_available(session_id)
2216 .await
2217 }
2218 async fn discard_live_session(
2219 &self,
2220 session_id: &meerkat_core::types::SessionId,
2221 ) -> Result<(), SessionError> {
2222 self.inner.discard_live_session(session_id).await
2223 }
2224 async fn checkpoint_committed_runtime_session_snapshot(
2225 &self,
2226 session_id: &meerkat_core::types::SessionId,
2227 session_snapshot: &[u8],
2228 ) -> Result<(), SessionError> {
2229 self.inner
2230 .checkpoint_committed_runtime_session_snapshot(session_id, session_snapshot)
2231 .await
2232 }
2233 async fn cancel_all_checkpointers(&self) {
2234 self.inner.cancel_all_checkpointers().await;
2235 }
2236 async fn rearm_all_checkpointers(&self) {
2237 self.inner.rearm_all_checkpointers().await;
2238 }
2239}
2240
2241pub struct MobBootstrapSpec {
2243 pub definition: MobDefinition,
2244 pub storage: MobStorage,
2245 pub session_service: Arc<dyn MobSessionService>,
2246 pub binary_blob_store: Option<Arc<dyn BinaryBlobStore>>,
2247 pub(crate) agent_mob_mcp_state: Option<Arc<meerkat_mob_mcp::MobMcpState>>,
2248 pub(crate) implicit_delegate_retirement_overrides: Option<ImplicitDelegateRetirementOverrides>,
2249 pub(crate) agent_mob_default_llm_client_slot: Option<SharedDefaultLlmClientSlot>,
2250 pub(crate) console_spawn_sink_slot: Option<SharedConsoleSpawnSinkSlot>,
2251 pub options: MobBootstrapOptions,
2252 pub runtime_adapter: Option<Arc<meerkat_runtime::MeerkatMachine>>,
2258 pub(crate) spawn_member_customizer: Option<Arc<dyn meerkat_mob::SpawnMemberCustomizer>>,
2262 pub(crate) _ephemeral_dir: Option<Arc<tempfile::TempDir>>,
2265}
2266
2267impl MobBootstrapSpec {
2268 pub fn new(
2269 definition: MobDefinition,
2270 storage: MobStorage,
2271 session_service: Arc<dyn MobSessionService>,
2272 ) -> Self {
2273 let session_service = Arc::new(PreBuildMobSessionService {
2274 inner: session_service,
2275 hook: no_op_pre_build_hook(),
2276 after_create_hook: None,
2277 runtime_adapter_override: None,
2278 }) as Arc<dyn MobSessionService>;
2279 Self {
2280 definition,
2281 storage,
2282 session_service,
2283 binary_blob_store: None,
2284 agent_mob_mcp_state: None,
2285 implicit_delegate_retirement_overrides: None,
2286 agent_mob_default_llm_client_slot: None,
2287 console_spawn_sink_slot: None,
2288 options: MobBootstrapOptions {
2289 allow_ephemeral_sessions: true,
2290 notify_orchestrator_on_resume: true,
2291 default_llm_client: None,
2292 },
2293 runtime_adapter: None,
2294 spawn_member_customizer: None,
2295 _ephemeral_dir: None,
2296 }
2297 }
2298
2299 pub fn with_options(mut self, options: MobBootstrapOptions) -> Self {
2300 self.options = options;
2301 self
2302 }
2303
2304 pub fn with_session_runtime_adapter(
2312 mut self,
2313 adapter: Arc<meerkat_runtime::MeerkatMachine>,
2314 ) -> Self {
2315 self.session_service = Arc::new(PreBuildMobSessionService {
2316 inner: self.session_service,
2317 hook: no_op_pre_build_hook(),
2318 after_create_hook: None,
2319 runtime_adapter_override: Some(adapter),
2320 });
2321 self
2322 }
2323
2324 pub fn with_after_create_hook(mut self, hook: AfterCreateHook) -> Self {
2330 self.session_service = Arc::new(AfterCreateMobSessionService {
2331 inner: self.session_service,
2332 after_hook: hook,
2333 });
2334 self
2335 }
2336
2337 pub fn ephemeral(
2342 definition: MobDefinition,
2343 storage: MobStorage,
2344 store_path: PathBuf,
2345 max_sessions: usize,
2346 session_store: Option<Arc<dyn AgentSessionStore>>,
2347 ) -> Self {
2348 Self::ephemeral_inner(
2349 definition,
2350 storage,
2351 store_path,
2352 max_sessions,
2353 session_store,
2354 None,
2355 CapabilityFlags::default(),
2356 None,
2357 None,
2358 )
2359 }
2360
2361 pub fn ephemeral_with_hook(
2365 definition: MobDefinition,
2366 storage: MobStorage,
2367 store_path: PathBuf,
2368 max_sessions: usize,
2369 session_store: Option<Arc<dyn AgentSessionStore>>,
2370 hook: impl Fn(
2371 &mut CreateSessionRequest,
2372 ) -> std::pin::Pin<
2373 Box<dyn std::future::Future<Output = Result<(), SessionError>> + Send + '_>,
2374 > + Send
2375 + Sync
2376 + 'static,
2377 ) -> Self {
2378 Self::ephemeral_inner(
2379 definition,
2380 storage,
2381 store_path,
2382 max_sessions,
2383 session_store,
2384 Some(Arc::new(hook)),
2385 CapabilityFlags::default(),
2386 None,
2387 None,
2388 )
2389 }
2390
2391 #[allow(clippy::too_many_arguments)]
2392 pub(crate) fn ephemeral_inner(
2393 definition: MobDefinition,
2394 storage: MobStorage,
2395 store_path: PathBuf,
2396 max_sessions: usize,
2397 session_store: Option<Arc<dyn AgentSessionStore>>,
2398 hook: Option<PreBuildHook>,
2399 mut caps: CapabilityFlags,
2400 after_create_hook: Option<AfterCreateHook>,
2401 agent_config: Option<Config>,
2402 ) -> Self {
2403 caps.image_generation |= mob_definition_may_use_image_generation(&definition);
2404 let binary_blob_store: Arc<dyn BinaryBlobStore> = Arc::new(ObjectStoreBlobStore::memory());
2405 let blob_store: Arc<dyn meerkat_core::BlobStore> =
2406 Arc::new(Base64BlobStoreAdapter::new(binary_blob_store.clone()));
2407 let runtime_adapter = if caps.image_generation {
2408 let runtime_store: Arc<dyn meerkat_runtime::RuntimeStore> =
2409 Arc::new(meerkat_runtime::InMemoryRuntimeStore::new());
2410 Some(Arc::new(meerkat_runtime::MeerkatMachine::persistent(
2411 runtime_store,
2412 Arc::clone(&blob_store),
2413 )))
2414 } else {
2415 None
2416 };
2417 let mut factory = AgentFactory::new(&store_path)
2418 .builtins(caps.builtins)
2419 .shell(caps.shell)
2420 .mob(caps.mob)
2421 .comms(caps.comms)
2422 .memory(caps.memory);
2423 if let Some(machine) = runtime_adapter.clone() {
2424 factory = factory.with_image_generation_machine(machine);
2425 }
2426 let config = agent_config.unwrap_or_default();
2427 let mut builder = FactoryAgentBuilder::new(factory, config);
2428 builder.default_blob_store = Some(blob_store);
2429 if let Some(store) = session_store {
2430 builder.default_session_store = Some(store);
2431 }
2432 let mob_tools_slot = Arc::clone(&builder.default_mob_tools);
2433 let session_service: Arc<dyn MobSessionService> = Arc::new(
2434 meerkat_session::EphemeralSessionService::new(builder, max_sessions),
2435 );
2436 let hook = hook.unwrap_or_else(no_op_pre_build_hook);
2437 let after_create_hook = if let Some(runtime_adapter) = runtime_adapter.clone() {
2438 let user_after_create_hook = after_create_hook.clone();
2439 Some(Arc::new(
2440 move |session_id: meerkat_core::types::SessionId, ctx: SessionCreatedContext| {
2441 let runtime_adapter = runtime_adapter.clone();
2442 let user_after_create_hook = user_after_create_hook.clone();
2443 Box::pin(async move {
2444 if let Err(error) =
2448 runtime_adapter.register_session(session_id.clone()).await
2449 {
2450 tracing::error!(
2451 session_id = %session_id,
2452 error = %error,
2453 "post-create session runtime registration failed"
2454 );
2455 }
2456 if let Some(user_after_create_hook) = user_after_create_hook {
2457 user_after_create_hook(session_id, ctx).await;
2458 }
2459 })
2460 as std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send>>
2461 },
2462 ) as AfterCreateHook)
2463 } else {
2464 after_create_hook
2465 };
2466 let session_service = Arc::new(PreBuildMobSessionService {
2467 inner: session_service,
2468 hook,
2469 after_create_hook,
2470 runtime_adapter_override: runtime_adapter.clone(),
2471 }) as Arc<dyn MobSessionService>;
2472 let (
2473 agent_mob_mcp_state,
2474 implicit_delegate_retirement_overrides,
2475 agent_mob_default_llm_client_slot,
2476 console_spawn_sink_slot,
2477 ) = install_agent_mob_tools(&definition, mob_tools_slot, Arc::clone(&session_service));
2478 let mut spec = Self::new(definition, storage, session_service);
2479 spec.agent_mob_mcp_state = Some(agent_mob_mcp_state);
2480 spec.implicit_delegate_retirement_overrides = Some(implicit_delegate_retirement_overrides);
2481 spec.agent_mob_default_llm_client_slot = Some(agent_mob_default_llm_client_slot);
2482 spec.console_spawn_sink_slot = Some(console_spawn_sink_slot);
2483 spec.runtime_adapter = runtime_adapter;
2484 spec.binary_blob_store = Some(binary_blob_store);
2485 spec
2486 }
2487
2488 pub fn persistent(
2495 definition: MobDefinition,
2496 storage: MobStorage,
2497 store_path: PathBuf,
2498 max_sessions: usize,
2499 session_store: Arc<dyn SessionStore>,
2500 ) -> Self {
2501 Self::persistent_inner(
2502 definition,
2503 storage,
2504 store_path,
2505 max_sessions,
2506 session_store,
2507 None,
2508 None,
2509 CapabilityFlags::default(),
2510 None,
2511 None,
2512 )
2513 }
2514
2515 pub fn persistent_with_hook(
2519 definition: MobDefinition,
2520 storage: MobStorage,
2521 store_path: PathBuf,
2522 max_sessions: usize,
2523 session_store: Arc<dyn SessionStore>,
2524 hook: impl Fn(
2525 &mut CreateSessionRequest,
2526 ) -> std::pin::Pin<
2527 Box<dyn std::future::Future<Output = Result<(), SessionError>> + Send + '_>,
2528 > + Send
2529 + Sync
2530 + 'static,
2531 ) -> Self {
2532 Self::persistent_inner(
2533 definition,
2534 storage,
2535 store_path,
2536 max_sessions,
2537 session_store,
2538 None,
2539 Some(Arc::new(hook)),
2540 CapabilityFlags::default(),
2541 None,
2542 None,
2543 )
2544 }
2545
2546 #[allow(clippy::too_many_arguments)]
2547 pub(crate) fn persistent_inner(
2548 definition: MobDefinition,
2549 storage: MobStorage,
2550 store_path: PathBuf,
2551 max_sessions: usize,
2552 session_store: Arc<dyn SessionStore>,
2553 custom_blob_store: Option<Arc<dyn meerkat_core::BlobStore>>,
2554 hook: Option<PreBuildHook>,
2555 mut caps: CapabilityFlags,
2556 after_create_hook: Option<AfterCreateHook>,
2557 agent_config: Option<Config>,
2558 ) -> Self {
2559 caps.image_generation |= mob_definition_may_use_image_generation(&definition);
2560 let (binary_blob_store, blob_store): (
2561 Arc<dyn BinaryBlobStore>,
2562 Arc<dyn meerkat_core::BlobStore>,
2563 ) = if let Some(blob_store) = custom_blob_store {
2564 (
2565 Arc::new(BinaryBlobStoreAdapter::new(blob_store.clone())),
2566 blob_store,
2567 )
2568 } else {
2569 let binary_blob_store: Arc<dyn BinaryBlobStore> = match ObjectStoreBlobStore::local(
2570 store_path.join("blobs"),
2571 ) {
2572 Ok(store) => Arc::new(store),
2573 Err(err) => {
2574 tracing::warn!(
2575 error = %err,
2576 "failed to initialize persistent binary blob store; falling back to in-memory blobs"
2577 );
2578 Arc::new(ObjectStoreBlobStore::memory())
2579 }
2580 };
2581 let blob_store: Arc<dyn meerkat_core::BlobStore> =
2582 Arc::new(Base64BlobStoreAdapter::new(binary_blob_store.clone()));
2583 (binary_blob_store, blob_store)
2584 };
2585 let runtime_store: Arc<dyn meerkat_runtime::RuntimeStore> =
2595 build_persistent_runtime_store(&store_path);
2596 let runtime_adapter = Arc::new(meerkat_runtime::MeerkatMachine::persistent(
2597 Arc::clone(&runtime_store),
2598 Arc::clone(&blob_store),
2599 ));
2600 let mut factory = AgentFactory::new(&store_path)
2601 .builtins(caps.builtins)
2602 .shell(caps.shell)
2603 .mob(caps.mob)
2604 .comms(caps.comms)
2605 .memory(caps.memory);
2606 if caps.image_generation {
2607 factory = factory.with_image_generation_machine(runtime_adapter.clone());
2608 }
2609 let config = agent_config.unwrap_or_default();
2610 let mut builder = FactoryAgentBuilder::new(factory, config);
2611 builder.default_session_store = Some(Arc::new(StoreAdapter::new(session_store.clone())));
2612 builder.default_blob_store = Some(blob_store.clone());
2613 let mob_tools_slot = Arc::clone(&builder.default_mob_tools);
2614 let session_service: Arc<dyn MobSessionService> =
2615 Arc::new(meerkat_session::PersistentSessionService::new(
2616 builder,
2617 max_sessions,
2618 session_store,
2619 Some(runtime_store),
2620 blob_store,
2621 ));
2622 let hook = hook.unwrap_or_else(no_op_pre_build_hook);
2623 let session_service = Arc::new(PreBuildMobSessionService {
2624 inner: session_service,
2625 hook,
2626 after_create_hook,
2627 runtime_adapter_override: None,
2628 }) as Arc<dyn MobSessionService>;
2629 let (
2630 agent_mob_mcp_state,
2631 implicit_delegate_retirement_overrides,
2632 agent_mob_default_llm_client_slot,
2633 console_spawn_sink_slot,
2634 ) = install_agent_mob_tools(&definition, mob_tools_slot, Arc::clone(&session_service));
2635 let mut spec = Self::new(definition, storage, session_service);
2636 spec.agent_mob_mcp_state = Some(agent_mob_mcp_state);
2637 spec.implicit_delegate_retirement_overrides = Some(implicit_delegate_retirement_overrides);
2638 spec.agent_mob_default_llm_client_slot = Some(agent_mob_default_llm_client_slot);
2639 spec.console_spawn_sink_slot = Some(console_spawn_sink_slot);
2640 spec.runtime_adapter = Some(runtime_adapter);
2641 spec.binary_blob_store = Some(binary_blob_store);
2642 spec
2643 }
2644
2645 #[allow(clippy::too_many_arguments)]
2646 pub(crate) fn ephemeral_runtime_backed_inner(
2647 definition: MobDefinition,
2648 storage: MobStorage,
2649 store_path: PathBuf,
2650 max_sessions: usize,
2651 custom_session_store: Option<Arc<dyn SessionStore>>,
2652 custom_blob_store: Option<Arc<dyn meerkat_core::BlobStore>>,
2653 hook: Option<PreBuildHook>,
2654 mut caps: CapabilityFlags,
2655 after_create_hook: Option<AfterCreateHook>,
2656 agent_config: Option<Config>,
2657 ) -> Self {
2658 caps.image_generation |= mob_definition_may_use_image_generation(&definition);
2659 let config = agent_config.unwrap_or_default();
2660 let session_store: Arc<dyn SessionStore> = custom_session_store
2661 .clone()
2662 .unwrap_or_else(|| Arc::new(meerkat_store::MemoryStore::new()));
2663 let (binary_blob_store, blob_store): (
2664 Arc<dyn BinaryBlobStore>,
2665 Arc<dyn meerkat_core::BlobStore>,
2666 ) = if let Some(blob_store) = custom_blob_store {
2667 (
2668 Arc::new(BinaryBlobStoreAdapter::new(blob_store.clone())),
2669 blob_store,
2670 )
2671 } else {
2672 let binary_blob_store: Arc<dyn BinaryBlobStore> =
2673 Arc::new(ObjectStoreBlobStore::memory());
2674 let blob_store: Arc<dyn meerkat_core::BlobStore> =
2675 Arc::new(Base64BlobStoreAdapter::new(binary_blob_store.clone()));
2676 (binary_blob_store, blob_store)
2677 };
2678 let base_runtime_store: Arc<dyn meerkat_runtime::RuntimeStore> =
2686 Arc::new(meerkat_runtime::InMemoryRuntimeStore::new());
2687 let runtime_store: Arc<dyn meerkat_runtime::RuntimeStore> =
2688 if let Some(custom_session_store) = custom_session_store.clone() {
2689 Arc::new(SessionStoreBackedRuntimeStore::new(
2690 Arc::clone(&base_runtime_store),
2691 custom_session_store,
2692 ))
2693 } else {
2694 Arc::clone(&base_runtime_store)
2695 };
2696 let runtime_adapter = Arc::new(meerkat_runtime::MeerkatMachine::persistent(
2697 Arc::clone(&runtime_store),
2698 Arc::clone(&blob_store),
2699 ));
2700 let mut factory = AgentFactory::new(&store_path)
2701 .builtins(caps.builtins)
2702 .shell(caps.shell)
2703 .mob(caps.mob)
2704 .comms(caps.comms)
2705 .memory(caps.memory);
2706 if caps.image_generation {
2707 factory = factory.with_image_generation_machine(runtime_adapter.clone());
2708 }
2709 let mut builder = FactoryAgentBuilder::new(factory, config);
2710 builder.default_session_store = Some(Arc::new(StoreAdapter::new(session_store.clone())));
2711 builder.default_blob_store = Some(blob_store.clone());
2712 let mob_tools_slot = Arc::clone(&builder.default_mob_tools);
2713 let session_service: Arc<dyn MobSessionService> =
2714 if let Some(custom_session_store) = custom_session_store {
2715 Arc::new(meerkat_session::PersistentSessionService::new(
2716 builder,
2717 max_sessions,
2718 custom_session_store,
2719 Some(runtime_store.clone()),
2720 blob_store,
2721 ))
2722 } else {
2723 Arc::new(meerkat_session::EphemeralSessionService::new(
2724 builder,
2725 max_sessions,
2726 ))
2727 };
2728 let hook = hook.unwrap_or_else(no_op_pre_build_hook);
2729 let runtime_adapter_for_after_create = runtime_adapter.clone();
2730 let combined_after_create_hook: AfterCreateHook = Arc::new(move |session_id, ctx| {
2731 let runtime_adapter = runtime_adapter_for_after_create.clone();
2732 let after_create_hook = after_create_hook.clone();
2733 Box::pin(async move {
2734 if let Err(error) = runtime_adapter.register_session(session_id.clone()).await {
2738 tracing::error!(
2739 session_id = %session_id,
2740 error = %error,
2741 "post-create session runtime registration failed"
2742 );
2743 }
2744 if let Some(after_create_hook) = after_create_hook {
2745 after_create_hook(session_id, ctx).await;
2746 }
2747 })
2748 });
2749 let session_service = Arc::new(PreBuildMobSessionService {
2750 inner: session_service,
2751 hook,
2752 after_create_hook: Some(combined_after_create_hook),
2753 runtime_adapter_override: Some(runtime_adapter.clone()),
2754 }) as Arc<dyn MobSessionService>;
2755 let (
2756 agent_mob_mcp_state,
2757 implicit_delegate_retirement_overrides,
2758 agent_mob_default_llm_client_slot,
2759 console_spawn_sink_slot,
2760 ) = install_agent_mob_tools(&definition, mob_tools_slot, Arc::clone(&session_service));
2761 let mut spec = Self::new(definition, storage, session_service);
2762 spec.agent_mob_mcp_state = Some(agent_mob_mcp_state);
2763 spec.implicit_delegate_retirement_overrides = Some(implicit_delegate_retirement_overrides);
2764 spec.agent_mob_default_llm_client_slot = Some(agent_mob_default_llm_client_slot);
2765 spec.console_spawn_sink_slot = Some(console_spawn_sink_slot);
2766 spec.runtime_adapter = Some(runtime_adapter);
2767 spec.binary_blob_store = Some(binary_blob_store);
2768 spec
2769 }
2770}
2771
2772#[derive(Debug)]
2774pub enum MobRuntimeError {
2775 Mob(MobError),
2776 InvalidInput(&'static str),
2777 InvalidConfig(String),
2778}
2779
2780impl std::fmt::Display for MobRuntimeError {
2781 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2782 match self {
2783 Self::Mob(err) => write!(f, "{err}"),
2784 Self::InvalidInput(message) => write!(f, "{message}"),
2785 Self::InvalidConfig(message) => write!(f, "{message}"),
2786 }
2787 }
2788}
2789
2790impl std::error::Error for MobRuntimeError {}
2791
2792impl From<MobError> for MobRuntimeError {
2793 fn from(value: MobError) -> Self {
2794 Self::Mob(value)
2795 }
2796}
2797
2798#[derive(Clone, Debug)]
2807pub struct SessionCreatedContext {
2808 pub model: String,
2809 pub labels: std::collections::BTreeMap<String, String>,
2810 pub system_prompt: Option<String>,
2811}
2812
2813#[async_trait]
2820pub trait SessionHook: Send + Sync {
2821 async fn before_create(&self, _req: &mut CreateSessionRequest) -> Result<(), SessionError> {
2825 Ok(())
2826 }
2827
2828 async fn after_create(
2831 &self,
2832 _session_id: &meerkat_core::types::SessionId,
2833 _ctx: &SessionCreatedContext,
2834 ) {
2835 }
2836}
2837
2838#[derive(Clone, Copy, Debug)]
2840pub struct CapabilityFlags {
2841 pub builtins: bool,
2842 pub shell: bool,
2843 pub mob: bool,
2844 pub comms: bool,
2845 pub memory: bool,
2846 pub image_generation: bool,
2847}
2848
2849impl Default for CapabilityFlags {
2850 fn default() -> Self {
2851 Self {
2852 builtins: true,
2853 shell: true,
2854 mob: true,
2855 comms: true,
2856 memory: true,
2857 image_generation: false,
2858 }
2859 }
2860}
2861
2862pub type RealMobRuntime = MobRuntime;
2864
2865#[derive(Clone)]
2867pub struct MobRuntime {
2868 handle: MobHandle,
2869 session_service: Option<Arc<dyn MobSessionService>>,
2870 agent_mob_mcp_state: Option<Arc<meerkat_mob_mcp::MobMcpState>>,
2871 implicit_delegate_retirement_overrides: Option<ImplicitDelegateRetirementOverrides>,
2872 binary_blob_store: Option<Arc<dyn BinaryBlobStore>>,
2873 baseline_member_specs: Arc<tokio::sync::RwLock<Vec<SpawnMemberSpec>>>,
2874 console_spawn_sink_slot: Option<SharedConsoleSpawnSinkSlot>,
2877 _ephemeral_dir: Option<Arc<tempfile::TempDir>>,
2880}
2881
2882impl MobRuntime {
2883 pub async fn bootstrap(spec: MobBootstrapSpec) -> Result<Self, MobRuntimeError> {
2884 let ephemeral_dir = spec._ephemeral_dir.clone();
2885 let session_service = spec.session_service.clone();
2886 let binary_blob_store = spec.binary_blob_store.clone();
2887 let mob_id = spec.definition.id.clone();
2888 let agent_mob_mcp_state = spec.agent_mob_mcp_state.clone();
2889 let implicit_delegate_retirement_overrides =
2890 spec.implicit_delegate_retirement_overrides.clone();
2891 let console_spawn_sink_slot = spec.console_spawn_sink_slot.clone();
2892 let default_llm_client = spec
2893 .options
2894 .default_llm_client
2895 .clone()
2896 .map(ReplaySanitizingLlmClient::wrap);
2897 if let Some(slot) = spec.agent_mob_default_llm_client_slot.as_ref() {
2898 *slot
2899 .write()
2900 .unwrap_or_else(std::sync::PoisonError::into_inner) = default_llm_client.clone();
2901 }
2902 let effective_runtime_adapter = spec
2903 .runtime_adapter
2904 .clone()
2905 .or_else(|| session_service.runtime_adapter());
2906
2907 let mut builder = MobBuilder::new(spec.definition, spec.storage);
2908
2909 if let Some(adapter) = effective_runtime_adapter {
2916 builder = builder.with_runtime_adapter(adapter);
2917 }
2918
2919 builder = builder
2920 .with_session_service(session_service.clone())
2921 .allow_ephemeral_sessions(spec.options.allow_ephemeral_sessions)
2922 .notify_orchestrator_on_resume(spec.options.notify_orchestrator_on_resume);
2923
2924 if let Some(customizer) = spec.spawn_member_customizer.clone() {
2925 builder = builder.with_spawn_member_customizer(customizer);
2926 }
2927
2928 if let Some(client) = default_llm_client {
2929 builder = builder.with_default_llm_client(client);
2930 }
2931
2932 let handle = builder.create().await?;
2933 if let Some(state) = agent_mob_mcp_state.as_ref() {
2934 state.mob_insert_handle(mob_id, handle.clone()).await;
2935 }
2936 Ok(Self {
2937 handle,
2938 session_service: Some(session_service),
2939 agent_mob_mcp_state,
2940 implicit_delegate_retirement_overrides,
2941 binary_blob_store,
2942 baseline_member_specs: Arc::new(tokio::sync::RwLock::new(Vec::new())),
2943 console_spawn_sink_slot,
2944 _ephemeral_dir: ephemeral_dir,
2945 })
2946 }
2947
2948 pub fn from_handle(handle: MobHandle) -> Self {
2949 Self {
2950 handle,
2951 session_service: None,
2952 agent_mob_mcp_state: None,
2953 implicit_delegate_retirement_overrides: None,
2954 binary_blob_store: None,
2955 baseline_member_specs: Arc::new(tokio::sync::RwLock::new(Vec::new())),
2956 console_spawn_sink_slot: None,
2957 _ephemeral_dir: None,
2958 }
2959 }
2960
2961 pub fn handle(&self) -> MobHandle {
2962 self.handle.clone()
2963 }
2964
2965 pub fn agent_mob_mcp_state(&self) -> Option<Arc<meerkat_mob_mcp::MobMcpState>> {
2966 self.agent_mob_mcp_state.clone()
2967 }
2968
2969 pub(crate) fn install_console_spawn_sink(&self, sink: ConsoleSpawnSink) {
2972 if let Some(slot) = self.console_spawn_sink_slot.as_ref() {
2973 *slot
2974 .write()
2975 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(sink);
2976 }
2977 }
2978
2979 pub(crate) fn console_spawn_sink(&self) -> Option<ConsoleSpawnSink> {
2981 self.console_spawn_sink_slot.as_ref().and_then(|slot| {
2982 slot.read()
2983 .unwrap_or_else(std::sync::PoisonError::into_inner)
2984 .clone()
2985 })
2986 }
2987
2988 pub(crate) async fn console_identity_labels(
2991 &self,
2992 ) -> BTreeMap<String, BTreeMap<String, String>> {
2993 match self.console_spawn_sink() {
2994 Some(sink) => sink.identity_labels_snapshot().await,
2995 None => BTreeMap::new(),
2996 }
2997 }
2998
2999 pub(crate) fn implicit_delegate_retirement_overrides(
3000 &self,
3001 ) -> Option<ImplicitDelegateRetirementOverrides> {
3002 self.implicit_delegate_retirement_overrides.clone()
3003 }
3004
3005 pub async fn set_baseline_member_specs(&self, specs: Vec<SpawnMemberSpec>) {
3006 *self.baseline_member_specs.write().await = specs;
3007 }
3008
3009 pub async fn baseline_member_specs(&self) -> Vec<SpawnMemberSpec> {
3010 self.baseline_member_specs.read().await.clone()
3011 }
3012
3013 pub async fn read_session_history(
3014 &self,
3015 session_id_str: &str,
3016 offset: usize,
3017 limit: Option<usize>,
3018 ) -> Result<SessionHistoryPage, MobRuntimeError> {
3019 if session_id_str.trim().is_empty() {
3020 return Err(MobRuntimeError::InvalidInput(
3021 "session_id must not be empty",
3022 ));
3023 }
3024 let Some(session_service) = self.session_service.as_ref() else {
3025 return Err(MobRuntimeError::InvalidInput(
3026 "session history unavailable for this runtime",
3027 ));
3028 };
3029 let session_id = meerkat_core::types::SessionId::parse(session_id_str)
3030 .map_err(|_| MobRuntimeError::InvalidInput("invalid session_id format"))?;
3031 SessionServiceHistoryExt::read_history(
3032 session_service.as_ref(),
3033 &session_id,
3034 SessionHistoryQuery { offset, limit },
3035 )
3036 .await
3037 .map_err(|err| MobRuntimeError::Mob(MobError::Internal(err.to_string())))
3038 }
3039
3040 #[allow(dead_code)]
3041 pub(crate) async fn runtime_state_for_session(
3042 &self,
3043 session_id_str: &str,
3044 ) -> Result<Option<meerkat_runtime::RuntimeState>, MobRuntimeError> {
3045 if session_id_str.trim().is_empty() {
3046 return Err(MobRuntimeError::InvalidInput(
3047 "session_id must not be empty",
3048 ));
3049 }
3050 let Some(session_service) = self.session_service.as_ref() else {
3051 return Ok(None);
3052 };
3053 let Some(runtime_adapter) = session_service.runtime_adapter() else {
3054 return Ok(None);
3055 };
3056 let session_id = meerkat_core::types::SessionId::parse(session_id_str)
3057 .map_err(|_| MobRuntimeError::InvalidInput("invalid session_id format"))?;
3058 let state = meerkat_runtime::service_ext::SessionServiceRuntimeExt::runtime_state(
3059 runtime_adapter.as_ref(),
3060 &session_id,
3061 )
3062 .await
3063 .map_err(|err| MobRuntimeError::Mob(MobError::Internal(err.to_string())))?;
3064 Ok(Some(state))
3065 }
3066
3067 #[allow(dead_code)]
3068 pub(crate) async fn comms_runtime_for_session(
3069 &self,
3070 session_id_str: &str,
3071 ) -> Result<Option<Arc<dyn CommsRuntime>>, MobRuntimeError> {
3072 if session_id_str.trim().is_empty() {
3073 return Err(MobRuntimeError::InvalidInput(
3074 "session_id must not be empty",
3075 ));
3076 }
3077 let Some(session_service) = self.session_service.as_ref() else {
3078 return Ok(None);
3079 };
3080 let session_id = meerkat_core::types::SessionId::parse(session_id_str)
3081 .map_err(|_| MobRuntimeError::InvalidInput("invalid session_id format"))?;
3082 Ok(
3083 meerkat_core::service::SessionServiceCommsExt::comms_runtime(
3084 session_service.as_ref(),
3085 &session_id,
3086 )
3087 .await,
3088 )
3089 }
3090
3091 #[allow(dead_code)]
3092 pub(crate) async fn active_input_ids_for_session(
3093 &self,
3094 session_id_str: &str,
3095 ) -> Result<Option<Vec<String>>, MobRuntimeError> {
3096 if session_id_str.trim().is_empty() {
3097 return Err(MobRuntimeError::InvalidInput(
3098 "session_id must not be empty",
3099 ));
3100 }
3101 let Some(session_service) = self.session_service.as_ref() else {
3102 return Ok(None);
3103 };
3104 let Some(runtime_adapter) = session_service.runtime_adapter() else {
3105 return Ok(None);
3106 };
3107 let session_id = meerkat_core::types::SessionId::parse(session_id_str)
3108 .map_err(|_| MobRuntimeError::InvalidInput("invalid session_id format"))?;
3109 let input_ids = meerkat_runtime::service_ext::SessionServiceRuntimeExt::list_active_inputs(
3110 runtime_adapter.as_ref(),
3111 &session_id,
3112 )
3113 .await
3114 .map_err(|err| MobRuntimeError::Mob(MobError::Internal(err.to_string())))?;
3115 Ok(Some(
3116 input_ids.into_iter().map(|id| id.to_string()).collect(),
3117 ))
3118 }
3119
3120 #[allow(dead_code)]
3121 pub(crate) async fn ensure_comms_drain_for_session(
3122 &self,
3123 session_id_str: &str,
3124 ) -> Result<Option<bool>, MobRuntimeError> {
3125 if session_id_str.trim().is_empty() {
3126 return Err(MobRuntimeError::InvalidInput(
3127 "session_id must not be empty",
3128 ));
3129 }
3130 let Some(session_service) = self.session_service.as_ref() else {
3131 return Ok(None);
3132 };
3133 let Some(runtime_adapter) = session_service.runtime_adapter() else {
3134 return Ok(None);
3135 };
3136 let session_id = meerkat_core::types::SessionId::parse(session_id_str)
3137 .map_err(|_| MobRuntimeError::InvalidInput("invalid session_id format"))?;
3138 let comms_runtime = meerkat_core::service::SessionServiceCommsExt::comms_runtime(
3139 session_service.as_ref(),
3140 &session_id,
3141 )
3142 .await;
3143 if let Some(comms) = comms_runtime {
3144 let _handle = meerkat_runtime::comms_drain::spawn_comms_drain(
3145 runtime_adapter.clone(),
3146 session_id,
3147 comms,
3148 None,
3149 );
3150 Ok(Some(true))
3151 } else {
3152 Ok(Some(false))
3153 }
3154 }
3155
3156 pub fn session_service(&self) -> Option<&Arc<dyn MobSessionService>> {
3162 self.session_service.as_ref()
3163 }
3164
3165 pub fn binary_blob_store(&self) -> Option<Arc<dyn BinaryBlobStore>> {
3166 self.binary_blob_store.clone()
3167 }
3168}
3169
3170pub fn member_entry_to_json(entry: &meerkat_mob::runtime::MobMemberListEntry) -> serde_json::Value {
3177 let mut value = serde_json::to_value(entry).unwrap_or(serde_json::Value::Null);
3178 if let Some(object) = value.as_object_mut() {
3182 if let Some(serde_json::Value::String(id)) = object.get_mut("agent_identity") {
3183 *id = crate::member_comms_id::runtime_alias_str(id).into_owned();
3184 }
3185 if let Some(serde_json::Value::Array(peers)) = object.get_mut("wired_to") {
3186 for peer in peers {
3187 if let serde_json::Value::String(peer_id) = peer {
3188 *peer_id = crate::member_comms_id::runtime_alias_str(peer_id).into_owned();
3189 }
3190 }
3191 }
3192 object.insert(
3198 "state".to_string(),
3199 serde_json::Value::String(member_status_state_string(entry.status)),
3200 );
3201 }
3202 value
3203}
3204
3205#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
3206pub struct ResolvedToolsSnapshot {
3207 pub identity: String,
3208 pub session_id: String,
3209 pub tools: Vec<String>,
3210}
3211
3212pub async fn resolved_tools_for_session(
3213 session_service: Option<&Arc<dyn MobSessionService>>,
3214 identity: &str,
3215 session_id: meerkat_core::types::SessionId,
3216) -> Result<ResolvedToolsSnapshot, MobRuntimeError> {
3217 let Some(session_service) = session_service else {
3218 return Err(MobRuntimeError::InvalidInput(
3219 "resolved tools unavailable for this runtime",
3220 ));
3221 };
3222 let scope = session_service
3223 .tool_scope_snapshot(&session_id)
3224 .await
3225 .map_err(|err| MobRuntimeError::Mob(MobError::Internal(err.to_string())))?
3226 .ok_or(MobRuntimeError::InvalidInput(
3227 "identity tool scope is unavailable",
3228 ))?;
3229 let mut tools = scope
3230 .visible_names
3231 .into_iter()
3232 .map(meerkat_core::types::ToolName::into_string)
3233 .collect::<Vec<_>>();
3234 tools.sort();
3235 Ok(ResolvedToolsSnapshot {
3236 identity: identity.to_string(),
3237 session_id: session_id.to_string(),
3238 tools,
3239 })
3240}
3241
3242pub async fn resolved_tools_for_member(
3243 handle: &MobHandle,
3244 session_service: Option<&Arc<dyn MobSessionService>>,
3245 member_id: &str,
3246) -> Result<ResolvedToolsSnapshot, MobRuntimeError> {
3247 if member_id.trim().is_empty() {
3248 return Err(MobRuntimeError::InvalidInput("identity must not be empty"));
3249 }
3250 let mid = crate::member_comms_id::mob_member_id(member_id);
3251 let status = handle.member_status(&mid).await?;
3252 let Some(session_id) = status.current_session_id else {
3253 return Err(MobRuntimeError::InvalidInput(
3254 "identity has no current session",
3255 ));
3256 };
3257 resolved_tools_for_session(session_service, member_id, session_id).await
3258}
3259
3260pub fn console_agent_event_payload(event: &meerkat_core::AgentEvent) -> Value {
3273 use meerkat_core::AgentEvent;
3274 use meerkat_core::event::agent_event_type;
3275
3276 let mut payload = serde_json::to_value(event).unwrap_or_else(|_| serde_json::json!({}));
3277 let record = match payload.as_object_mut() {
3278 Some(record) => record,
3279 None => return payload,
3280 };
3281 let is_tool_event = matches!(
3282 agent_event_type(event),
3283 "tool_call_requested"
3284 | "tool_result_received"
3285 | "tool_execution_started"
3286 | "tool_execution_completed"
3287 | "tool_execution_timed_out"
3288 );
3289 if is_tool_event
3290 && !record.contains_key("tool_call_id")
3291 && let Some(id) = record.get("id").cloned()
3292 {
3293 record.insert("tool_call_id".to_string(), id);
3294 }
3295 if let AgentEvent::ToolExecutionCompleted { content, .. } = event
3296 && !record.contains_key("result")
3297 {
3298 record.insert(
3299 "result".to_string(),
3300 Value::String(meerkat_core::types::text_content(content)),
3301 );
3302 }
3303 payload
3304}
3305
3306pub fn content_input_has_images(content: &meerkat_core::ContentInput) -> bool {
3307 match content {
3308 meerkat_core::ContentInput::Text(_) => false,
3309 meerkat_core::ContentInput::Blocks(blocks) => blocks
3310 .iter()
3311 .any(|block| matches!(block, meerkat_core::ContentBlock::Image { .. })),
3312 }
3313}
3314
3315pub fn model_capabilities_for_model(
3316 provider: Provider,
3317 model: &str,
3318) -> crate::runtime::ConsoleModelCapabilities {
3319 let image_input = meerkat_models::profile_for(provider, model)
3320 .map(|profile| profile.vision)
3321 .unwrap_or(false);
3322 crate::runtime::ConsoleModelCapabilities { image_input }
3323}
3324
3325pub fn model_capabilities_for_profile(
3326 profile: &Profile,
3327) -> crate::runtime::ConsoleModelCapabilities {
3328 let image_input = meerkat_models::infer_provider(&profile.model)
3329 .and_then(|provider| meerkat_models::profile_for(provider, &profile.model))
3330 .map(|profile| profile.vision)
3331 .unwrap_or(false);
3332 crate::runtime::ConsoleModelCapabilities { image_input }
3333}
3334
3335pub fn model_capabilities_for_role(
3336 definition: &MobDefinition,
3337 role: &str,
3338) -> crate::runtime::ConsoleModelCapabilities {
3339 let profile_name = ProfileName::from(role);
3340 definition
3341 .resolve_inline_profile(&profile_name)
3342 .map(model_capabilities_for_profile)
3343 .unwrap_or(crate::runtime::ConsoleModelCapabilities { image_input: false })
3344}
3345
3346pub fn model_capabilities_for_member_entry(
3347 definition: &MobDefinition,
3348 entry: &meerkat_mob::runtime::MobMemberListEntry,
3349) -> crate::runtime::ConsoleModelCapabilities {
3350 model_capabilities_for_role(definition, entry.role.as_str())
3351}
3352
3353pub async fn model_capabilities_for_member(
3354 handle: &MobHandle,
3355 session_service: Option<&Arc<dyn MobSessionService>>,
3356 member_id: &meerkat_mob::ids::AgentIdentity,
3357) -> crate::runtime::ConsoleModelCapabilities {
3358 if let Some(service) = session_service
3359 && let Some(session_id) = handle.resolve_bridge_session_id(member_id).await
3360 && let Ok(view) = service.read(&session_id).await
3361 {
3362 return model_capabilities_for_model(view.state.provider, &view.state.model);
3363 }
3364
3365 handle
3368 .get_member(member_id)
3369 .await
3370 .ok()
3371 .flatten()
3372 .map(|member| model_capabilities_for_role(handle.definition(), member.role.as_str()))
3373 .unwrap_or(crate::runtime::ConsoleModelCapabilities { image_input: false })
3374}
3375
3376pub async fn assert_member_accepts_images(
3377 handle: &MobHandle,
3378 session_service: Option<&Arc<dyn MobSessionService>>,
3379 member_id: &str,
3380 content: &meerkat_core::ContentInput,
3381) -> Result<(), MobRuntimeError> {
3382 if !content_input_has_images(content) {
3383 return Ok(());
3384 }
3385 let mid = crate::member_comms_id::mob_member_id(member_id);
3388 let Some(member) = handle
3389 .get_member(&mid)
3390 .await
3391 .map_err(|_| MobRuntimeError::InvalidInput("member lookup failed"))?
3392 else {
3393 return Err(MobRuntimeError::InvalidInput("member not found"));
3394 };
3395 let caps = model_capabilities_for_member(handle, session_service, &member.agent_identity).await;
3396 if caps.image_input {
3397 Ok(())
3398 } else {
3399 Err(MobRuntimeError::InvalidInput(
3400 "target member model cannot accept image input",
3401 ))
3402 }
3403}
3404
3405pub async fn send_message_on_mob(
3415 handle: &MobHandle,
3416 member_id: &str,
3417 content: impl Into<meerkat_core::ContentInput>,
3418) -> Result<String, MobRuntimeError> {
3419 send_message_on_mob_with_mode(
3420 handle,
3421 member_id,
3422 content,
3423 meerkat_core::types::HandlingMode::Queue,
3424 )
3425 .await
3426}
3427
3428pub async fn send_message_on_mob_with_mode(
3431 handle: &MobHandle,
3432 member_id: &str,
3433 content: impl Into<meerkat_core::ContentInput>,
3434 handling_mode: meerkat_core::types::HandlingMode,
3435) -> Result<String, MobRuntimeError> {
3436 if member_id.trim().is_empty() {
3437 return Err(MobRuntimeError::InvalidInput("member_id must not be empty"));
3438 }
3439 let content = content.into();
3440 let is_empty = match &content {
3441 meerkat_core::ContentInput::Text(s) => s.trim().is_empty(),
3442 meerkat_core::ContentInput::Blocks(blocks) => blocks.is_empty(),
3443 };
3444 if is_empty {
3445 return Err(MobRuntimeError::InvalidInput("content must not be empty"));
3446 }
3447 let mid = crate::member_comms_id::mob_member_id(member_id);
3450 let _receipt = handle
3451 .member(&mid)
3452 .await?
3453 .send(content, handling_mode)
3454 .await?;
3455 if let Some(session_id) = handle.resolve_bridge_session_id(&mid).await {
3456 return Ok(session_id.to_string());
3457 }
3458
3459 let status = handle.member_status(&mid).await?;
3460 if status.external_member.is_some() {
3461 return Ok(String::new());
3462 }
3463
3464 Err(MobRuntimeError::Mob(MobError::Internal(
3465 "member has no bridge session after send".to_string(),
3466 )))
3467}
3468
3469#[cfg(test)]
3470#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
3471mod tests {
3472 use super::*;
3473
3474 struct EmptyDispatcher;
3475
3476 #[async_trait::async_trait]
3477 impl meerkat_core::AgentToolDispatcher for EmptyDispatcher {
3478 fn tools(&self) -> Arc<[Arc<meerkat_core::types::ToolDef>]> {
3479 Vec::<Arc<meerkat_core::types::ToolDef>>::new().into()
3480 }
3481
3482 async fn dispatch(
3483 &self,
3484 call: meerkat_core::types::ToolCallView<'_>,
3485 ) -> Result<meerkat_core::ToolDispatchOutcome, meerkat_core::ToolError> {
3486 Err(meerkat_core::ToolError::not_found(call.name))
3487 }
3488
3489 fn capabilities(&self) -> meerkat_core::agent::DispatcherCapabilities {
3490 meerkat_core::agent::DispatcherCapabilities::default()
3491 }
3492 }
3493
3494 fn wrapper_with_overrides(
3495 overrides: ImplicitDelegateRetirementOverrides,
3496 ) -> AutoWireParentMobToolDispatcher {
3497 AutoWireParentMobToolDispatcher {
3498 inner: Arc::new(EmptyDispatcher),
3499 implicit_delegate_retirement_overrides: overrides,
3500 console_spawn_sink: new_console_spawn_sink_slot(),
3501 spawner_comms_name: None,
3502 }
3503 }
3504
3505 #[test]
3506 fn delegate_tool_schema_exposes_idle_retire_secs() {
3507 let tool = meerkat_core::types::ToolDef::new(
3508 "delegate",
3509 "Delegate work",
3510 serde_json::json!({
3511 "type": "object",
3512 "properties": {
3513 "task": {"type": "string"}
3514 },
3515 "required": ["task"]
3516 }),
3517 );
3518
3519 let patched = delegate_tool_def_with_idle_retire_secs(&tool);
3520 let idle_retire_secs = &patched.input_schema["properties"]["idle_retire_secs"];
3521
3522 assert!(patched.description.contains("IDLE RETIREMENT:"));
3523 assert_eq!(idle_retire_secs["anyOf"][0]["type"], "integer");
3524 assert_eq!(idle_retire_secs["anyOf"][0]["minimum"], 0);
3525 assert_eq!(idle_retire_secs["anyOf"][1]["type"], "null");
3526 }
3527
3528 #[test]
3529 fn mob_spawn_tool_schema_exposes_opt_in_idle_retire_secs() {
3530 let tool = meerkat_core::types::ToolDef::new(
3531 "mob_spawn_member",
3532 "Spawn member",
3533 serde_json::json!({
3534 "type": "object",
3535 "properties": {
3536 "profile": {"type": "string"},
3537 "member_id": {"type": "string"}
3538 },
3539 "required": ["profile", "member_id"]
3540 }),
3541 );
3542
3543 let patched = mob_spawn_tool_def_with_idle_retire_secs(&tool);
3544 let idle_retire_secs = &patched.input_schema["properties"]["idle_retire_secs"];
3545
3546 assert!(
3547 patched
3548 .description
3549 .contains("Omit idle_retire_secs to leave this spawned member out")
3550 );
3551 assert_eq!(idle_retire_secs["anyOf"][0]["type"], "integer");
3552 assert_eq!(idle_retire_secs["anyOf"][0]["minimum"], 0);
3553 assert_eq!(idle_retire_secs["anyOf"][1]["type"], "null");
3554 }
3555
3556 #[tokio::test]
3557 async fn auto_wire_wrapper_preserves_ops_lifecycle_binding() {
3558 use meerkat_core::AgentToolDispatcher;
3559 use std::sync::atomic::{AtomicBool, Ordering};
3560
3561 struct BindAwareDispatcher {
3562 bound: Arc<AtomicBool>,
3563 }
3564
3565 #[async_trait::async_trait]
3566 impl meerkat_core::AgentToolDispatcher for BindAwareDispatcher {
3567 fn tools(&self) -> Arc<[Arc<meerkat_core::types::ToolDef>]> {
3568 Vec::<Arc<meerkat_core::types::ToolDef>>::new().into()
3569 }
3570
3571 async fn dispatch(
3572 &self,
3573 call: meerkat_core::types::ToolCallView<'_>,
3574 ) -> Result<meerkat_core::ToolDispatchOutcome, meerkat_core::ToolError> {
3575 Err(meerkat_core::ToolError::not_found(call.name))
3576 }
3577
3578 fn capabilities(&self) -> meerkat_core::agent::DispatcherCapabilities {
3579 meerkat_core::agent::DispatcherCapabilities {
3580 ops_lifecycle: true,
3581 }
3582 }
3583
3584 fn bind_ops_lifecycle(
3585 self: Arc<Self>,
3586 _registry: Arc<dyn meerkat_core::ops_lifecycle::OpsLifecycleRegistry>,
3587 _owner_bridge_session_id: meerkat_core::types::SessionId,
3588 ) -> Result<meerkat_core::agent::BindOutcome, meerkat_core::agent::OpsLifecycleBindError>
3589 {
3590 self.bound.store(true, Ordering::SeqCst);
3591 Ok(meerkat_core::agent::BindOutcome::Bound(self))
3592 }
3593 }
3594
3595 let bound = Arc::new(AtomicBool::new(false));
3596 let dispatcher = Arc::new(AutoWireParentMobToolDispatcher {
3597 inner: Arc::new(BindAwareDispatcher {
3598 bound: Arc::clone(&bound),
3599 }),
3600 implicit_delegate_retirement_overrides: ImplicitDelegateRetirementOverrides::default(),
3601 console_spawn_sink: new_console_spawn_sink_slot(),
3602 spawner_comms_name: None,
3603 });
3604
3605 assert!(dispatcher.capabilities().ops_lifecycle);
3606 let outcome = dispatcher
3607 .bind_ops_lifecycle(
3608 Arc::new(meerkat_runtime::ops_lifecycle::RuntimeOpsLifecycleRegistry::new()),
3609 meerkat_core::types::SessionId::new(),
3610 )
3611 .expect("wrapper should delegate ops lifecycle binding");
3612
3613 assert!(outcome.was_bound());
3614 assert!(bound.load(Ordering::SeqCst));
3615 assert!(outcome.into_dispatcher().capabilities().ops_lifecycle);
3616 }
3617
3618 #[test]
3619 fn delegate_idle_retire_secs_arg_is_stripped_and_parsed() {
3620 let mut args = serde_json::json!({
3621 "task": "inspect",
3622 "idle_retire_secs": 42
3623 });
3624
3625 let parsed = delegate_idle_retire_override_from_args("delegate", &mut args)
3626 .expect("valid idle retire arg");
3627
3628 assert_eq!(parsed, Some(DelegateIdleRetireOverride::Seconds(42)));
3629 assert!(args.get("idle_retire_secs").is_none());
3630 }
3631
3632 #[test]
3633 fn delegate_idle_retire_secs_null_disables_member_retirement() {
3634 let mut args = serde_json::json!({
3635 "task": "inspect",
3636 "idle_retire_secs": null
3637 });
3638
3639 let parsed = delegate_idle_retire_override_from_args("delegate", &mut args)
3640 .expect("valid idle retire arg");
3641
3642 assert_eq!(parsed, Some(DelegateIdleRetireOverride::Disabled));
3643 assert!(args.get("idle_retire_secs").is_none());
3644 }
3645
3646 #[test]
3647 fn delegate_idle_retire_secs_omitted_inherits_runtime_default() {
3648 let mut args = serde_json::json!({"task": "inspect"});
3649
3650 let parsed = delegate_idle_retire_override_from_args("delegate", &mut args)
3651 .expect("omitted idle retire arg");
3652
3653 assert_eq!(parsed, None);
3654 assert_eq!(args, serde_json::json!({"task": "inspect"}));
3655 }
3656
3657 #[test]
3658 fn delegate_idle_retire_secs_rejects_negative_or_fractional_values() {
3659 let mut negative = serde_json::json!({"task": "inspect", "idle_retire_secs": -1});
3660 let mut fractional = serde_json::json!({"task": "inspect", "idle_retire_secs": 1.5});
3661
3662 assert!(delegate_idle_retire_override_from_args("delegate", &mut negative).is_err());
3663 assert!(delegate_idle_retire_override_from_args("delegate", &mut fractional).is_err());
3664 }
3665
3666 #[test]
3667 fn mob_spawn_idle_retire_targets_use_args_when_result_omits_mob_id() {
3668 let args = serde_json::json!({
3669 "mob_id": "ob3",
3670 "profile": "review-worker",
3671 "member_id": "review-worker-vibe-forward",
3672 });
3673 let fallback_targets = idle_retire_targets_from_spawn_args(&args);
3674
3675 assert_eq!(
3676 fallback_targets,
3677 vec![IdleRetireTarget {
3678 mob_id: "ob3".to_string(),
3679 member_id: "review-worker-vibe-forward".to_string(),
3680 }]
3681 );
3682 assert_eq!(
3683 idle_retire_targets_from_outcome_text(
3684 r#"{"agent_identity":"review-worker-vibe-forward","member_ref":"opaque"}"#,
3685 &fallback_targets,
3686 ),
3687 fallback_targets
3688 );
3689 }
3690
3691 #[test]
3692 fn mob_spawn_idle_retire_targets_support_canonical_specs_shape() {
3693 let args = serde_json::json!({
3694 "mob_id": "ob3",
3695 "specs": [
3696 {"profile": "person-worker", "agent_identity": "person-worker-a"},
3697 {"profile": "person-worker", "member_id": "person-worker-b", "mob_id": "other"}
3698 ]
3699 });
3700 let fallback_targets = idle_retire_targets_from_spawn_args(&args);
3701
3702 assert_eq!(
3703 fallback_targets,
3704 vec![
3705 IdleRetireTarget {
3706 mob_id: "ob3".to_string(),
3707 member_id: "person-worker-a".to_string(),
3708 },
3709 IdleRetireTarget {
3710 mob_id: "other".to_string(),
3711 member_id: "person-worker-b".to_string(),
3712 },
3713 ]
3714 );
3715 assert_eq!(
3716 idle_retire_targets_from_outcome_text(
3717 r#"{"members":[{"agent_identity":"person-worker-a"},{"agent_identity":"person-worker-b","mob_id":"other"}]}"#,
3718 &fallback_targets,
3719 ),
3720 fallback_targets
3721 );
3722 }
3723
3724 #[tokio::test]
3725 async fn implicit_delegate_retirement_overrides_round_trip_per_member() {
3726 let overrides = ImplicitDelegateRetirementOverrides::default();
3727
3728 overrides
3729 .set("mob-a", "worker-1", DelegateIdleRetireOverride::Seconds(12))
3730 .await;
3731 overrides
3732 .set("mob-a", "worker-2", DelegateIdleRetireOverride::Disabled)
3733 .await;
3734
3735 assert_eq!(
3736 overrides.get("mob-a", "worker-1").await,
3737 Some(DelegateIdleRetireOverride::Seconds(12))
3738 );
3739 assert_eq!(
3740 overrides.get("mob-a", "worker-2").await,
3741 Some(DelegateIdleRetireOverride::Disabled)
3742 );
3743 assert_eq!(overrides.get("mob-a", "worker-3").await, None);
3744 }
3745
3746 #[tokio::test]
3747 async fn mob_spawn_idle_retire_registration_uses_spawn_args_when_result_omits_mob_id() {
3748 let overrides = ImplicitDelegateRetirementOverrides::default();
3749 let dispatcher = wrapper_with_overrides(overrides.clone());
3750 let fallback_targets = idle_retire_targets_from_spawn_args(&serde_json::json!({
3751 "mob_id": "ob3",
3752 "member_id": "review-worker-vibe-forward",
3753 }));
3754 let outcome =
3755 meerkat_core::ToolDispatchOutcome::sync_result(meerkat_core::types::ToolResult::new(
3756 "spawn-1".to_string(),
3757 r#"{"agent_identity":"review-worker-vibe-forward","member_ref":"opaque"}"#
3758 .to_string(),
3759 false,
3760 ));
3761
3762 dispatcher
3763 .register_idle_retire_override_from_outcome(
3764 &outcome,
3765 Some(DelegateIdleRetireOverride::Seconds(900)),
3766 &fallback_targets,
3767 )
3768 .await;
3769
3770 assert_eq!(
3771 overrides.get("ob3", "review-worker-vibe-forward").await,
3772 Some(DelegateIdleRetireOverride::Seconds(900))
3773 );
3774 }
3775
3776 #[test]
3777 fn image_generation_substrate_defaults_off_for_inline_profiles() {
3778 let definition = meerkat_mob::MobDefinition::from_toml(
3779 r#"
3780[mob]
3781id = "test"
3782
3783[profiles.worker]
3784model = "gpt-5.5"
3785
3786[profiles.worker.tools]
3787builtins = true
3788"#,
3789 )
3790 .unwrap_or_else(|e| panic!("{e}"));
3791
3792 assert!(
3793 !mob_definition_may_use_image_generation(&definition),
3794 "inline profiles should not wire the image substrate unless a profile opts in"
3795 );
3796 }
3797
3798 #[test]
3799 fn image_generation_substrate_follows_profile_tool_config() {
3800 let definition = meerkat_mob::MobDefinition::from_toml(
3801 r#"
3802[mob]
3803id = "test"
3804
3805[profiles.commander]
3806model = "gpt-5.5"
3807
3808[profiles.commander.tools]
3809builtins = true
3810image_generation = true
3811
3812[profiles.investigator]
3813model = "gpt-5.5"
3814
3815[profiles.investigator.tools]
3816builtins = true
3817image_generation = false
3818"#,
3819 )
3820 .unwrap_or_else(|e| panic!("{e}"));
3821
3822 let commander = definition.profiles["commander"].as_inline().unwrap();
3823 let investigator = definition.profiles["investigator"].as_inline().unwrap();
3824 assert!(commander.tools.image_generation);
3825 assert!(!investigator.tools.image_generation);
3826 assert!(
3827 mob_definition_may_use_image_generation(&definition),
3828 "one opt-in profile is enough to wire substrate; Meerkat gates visibility per profile"
3829 );
3830 }
3831
3832 #[test]
3833 fn shell_substrate_defaults_off_for_inline_profiles() {
3834 let definition = meerkat_mob::MobDefinition::from_toml(
3835 r#"
3836[mob]
3837id = "test"
3838
3839[profiles.worker]
3840model = "gpt-5.5"
3841
3842[profiles.worker.tools]
3843builtins = true
3844"#,
3845 )
3846 .unwrap_or_else(|e| panic!("{e}"));
3847
3848 assert!(
3849 !mob_definition_may_use_shell(&definition),
3850 "inline profiles should not wire the shell substrate unless a profile opts in"
3851 );
3852 }
3853
3854 #[test]
3855 fn shell_substrate_follows_profile_tool_config() {
3856 let definition = meerkat_mob::MobDefinition::from_toml(
3857 r#"
3858[mob]
3859id = "test"
3860
3861[profiles.domain]
3862model = "gpt-5.5"
3863
3864[profiles.domain.tools]
3865builtins = true
3866shell = false
3867
3868[profiles.security]
3869model = "gpt-5.5"
3870
3871[profiles.security.tools]
3872builtins = true
3873shell = true
3874"#,
3875 )
3876 .unwrap_or_else(|e| panic!("{e}"));
3877
3878 let domain = definition.profiles["domain"].as_inline().unwrap();
3879 let security = definition.profiles["security"].as_inline().unwrap();
3880 assert!(!domain.tools.shell);
3881 assert!(security.tools.shell);
3882 assert!(
3883 mob_definition_may_use_shell(&definition),
3884 "one opt-in profile is enough to wire substrate; Meerkat gates visibility per profile"
3885 );
3886 }
3887
3888 #[test]
3889 fn shell_tooling_forces_builtin_substrate_without_exposing_broad_builtins() {
3890 let mut req = CreateSessionRequest {
3891 model: "gpt-5.5".to_string(),
3892 prompt: meerkat_core::ContentInput::Text("test".to_string()),
3893 system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
3894 max_tokens: None,
3895 event_tx: None,
3896 initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
3897 build: Some(meerkat_core::service::SessionBuildOptions {
3898 override_builtins: meerkat_core::ToolCategoryOverride::Disable,
3899 override_shell: meerkat_core::ToolCategoryOverride::Enable,
3900 ..Default::default()
3901 }),
3902 labels: None,
3903 deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::default(),
3904 injected_context: Vec::new(),
3905 };
3906
3907 ensure_shell_tooling_build_substrate(&mut req);
3908
3909 let build = req.build.expect("build options");
3910 assert_eq!(
3911 build.override_shell,
3912 meerkat_core::ToolCategoryOverride::Enable
3913 );
3914 assert_eq!(
3915 build.override_builtins,
3916 meerkat_core::ToolCategoryOverride::Enable,
3917 "shell-only profiles must still enable Meerkat's builtin substrate"
3918 );
3919 let allow = match build.initial_tool_filter.expect("shell visibility filter") {
3920 meerkat_core::ToolFilter::Allow(allow) => allow,
3921 other => panic!("expected shell/comms allow filter, got {other:?}"),
3922 };
3923 for tool in SHELL_BUILTIN_TOOL_NAMES
3924 .iter()
3925 .chain(COMMS_TOOL_NAMES.iter())
3926 {
3927 assert!(allow.contains(tool), "missing expected tool {tool}");
3928 }
3929 for broad_builtin in ["task_list", "task_create", "apply_patch", "browse_skills"] {
3930 assert!(
3931 !allow.contains(broad_builtin),
3932 "shell-only filter must not expose broad builtin {broad_builtin}",
3933 );
3934 }
3935 }
3936
3937 #[test]
3938 fn non_shell_profiles_keep_builtin_override_unchanged() {
3939 let mut req = CreateSessionRequest {
3940 model: "gpt-5.5".to_string(),
3941 prompt: meerkat_core::ContentInput::Text("test".to_string()),
3942 system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
3943 max_tokens: None,
3944 event_tx: None,
3945 initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
3946 build: Some(meerkat_core::service::SessionBuildOptions {
3947 override_builtins: meerkat_core::ToolCategoryOverride::Disable,
3948 override_shell: meerkat_core::ToolCategoryOverride::Disable,
3949 ..Default::default()
3950 }),
3951 labels: None,
3952 deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::default(),
3953 injected_context: Vec::new(),
3954 };
3955
3956 ensure_shell_tooling_build_substrate(&mut req);
3957
3958 let build = req.build.expect("build options");
3959 assert_eq!(
3960 build.override_builtins,
3961 meerkat_core::ToolCategoryOverride::Disable
3962 );
3963 assert_eq!(
3964 build.override_shell,
3965 meerkat_core::ToolCategoryOverride::Disable
3966 );
3967 }
3968
3969 #[test]
3970 fn image_generation_profiles_can_disable_builtins_with_meerkat_062() {
3971 let definition = meerkat_mob::MobDefinition::from_toml(
3972 r#"
3973[mob]
3974id = "test"
3975
3976[profiles.commander]
3977model = "gpt-5.5"
3978
3979[profiles.commander.tools]
3980builtins = false
3981image_generation = true
3982"#,
3983 )
3984 .unwrap_or_else(|e| panic!("{e}"));
3985
3986 let commander = definition.profiles["commander"].as_inline().unwrap();
3987 assert!(!commander.tools.builtins);
3988 assert!(commander.tools.image_generation);
3989 assert!(
3990 mob_definition_may_use_image_generation(&definition),
3991 "image generation now has its own Meerkat tool gate"
3992 );
3993 }
3994
3995 #[test]
3996 fn image_generation_substrate_is_conservative_for_realm_profile_refs() {
3997 let definition = meerkat_mob::MobDefinition::from_toml(
3998 r#"
3999[mob]
4000id = "test"
4001
4002[profiles.worker]
4003realm_profile = "worker-v2"
4004"#,
4005 )
4006 .unwrap_or_else(|e| panic!("{e}"));
4007
4008 assert!(
4009 mob_definition_may_use_image_generation(&definition),
4010 "realm profiles resolve at spawn time, so MobKit wires substrate and lets Meerkat enforce profile policy"
4011 );
4012 }
4013
4014 #[test]
4015 fn sanitize_llm_request_drops_replay_unsafe_server_tool_blocks() {
4016 let request = meerkat_client::LlmRequest::new(
4017 "gpt-5.5",
4018 vec![meerkat_core::Message::BlockAssistant(
4019 meerkat_core::BlockAssistantMessage::new(
4020 vec![
4021 meerkat_core::AssistantBlock::Text {
4022 text: "done".to_string(),
4023 meta: None,
4024 },
4025 meerkat_core::AssistantBlock::ServerToolContent {
4026 id: Some("ws-stream".to_string()),
4027 kind: meerkat_core::ServerToolKind::WebSearch,
4028 content: serde_json::json!({
4029 "type": "response.web_search_call.searching",
4030 "item_id": "ws_123"
4031 }),
4032 meta: None,
4033 },
4034 meerkat_core::AssistantBlock::ServerToolContent {
4035 id: Some("ws_123".to_string()),
4036 kind: meerkat_core::ServerToolKind::ProviderNative {
4037 name: "web_search_call".to_string(),
4038 },
4039 content: serde_json::json!({
4040 "type": "web_search_call",
4041 "id": "ws_123",
4042 "status": "completed"
4043 }),
4044 meta: None,
4045 },
4046 meerkat_core::AssistantBlock::ServerToolContent {
4047 id: None,
4048 kind: meerkat_core::ServerToolKind::ProviderNative {
4049 name: "web_search_annotations".to_string(),
4050 },
4051 content: serde_json::json!({
4052 "type": "message_annotations",
4053 "annotations": []
4054 }),
4055 meta: None,
4056 },
4057 ],
4058 meerkat_core::StopReason::EndTurn,
4059 ),
4060 )],
4061 );
4062
4063 let sanitized = sanitize_llm_request_for_stateless_replay(&request);
4064 let meerkat_core::Message::BlockAssistant(assistant) = &sanitized.messages[0] else {
4065 panic!("expected block assistant");
4066 };
4067
4068 assert_eq!(assistant.blocks.len(), 2);
4069 assert!(matches!(
4070 assistant.blocks[0],
4071 meerkat_core::AssistantBlock::Text { .. }
4072 ));
4073 assert!(matches!(
4074 assistant.blocks[1],
4075 meerkat_core::AssistantBlock::ServerToolContent { ref kind, .. }
4076 if kind.provider_name() == "web_search_call"
4077 ));
4078 }
4079
4080 #[test]
4081 fn sanitize_llm_request_preserves_generated_images_for_meerkat_062() {
4082 let request = meerkat_client::LlmRequest::new(
4083 "gpt-5.5",
4084 vec![meerkat_core::Message::BlockAssistant(
4085 meerkat_core::BlockAssistantMessage::new(
4086 vec![
4087 meerkat_core::AssistantBlock::Text {
4088 text: "visible".to_string(),
4089 meta: None,
4090 },
4091 generated_image_block_for_test(),
4092 ],
4093 meerkat_core::StopReason::EndTurn,
4094 ),
4095 )],
4096 );
4097
4098 let sanitized = sanitize_llm_request_for_stateless_replay(&request);
4099
4100 let meerkat_core::Message::BlockAssistant(original_assistant) = &request.messages[0] else {
4101 panic!("expected original block assistant");
4102 };
4103 assert!(
4104 original_assistant
4105 .blocks
4106 .iter()
4107 .any(|block| matches!(block, meerkat_core::AssistantBlock::Image { .. })),
4108 "request-view sanitization must not rewrite canonical caller-owned messages"
4109 );
4110
4111 let meerkat_core::Message::BlockAssistant(sanitized_assistant) = &sanitized.messages[0]
4112 else {
4113 panic!("expected sanitized block assistant");
4114 };
4115 assert!(
4116 sanitized_assistant
4117 .blocks
4118 .iter()
4119 .any(|block| matches!(block, meerkat_core::AssistantBlock::Image { .. })),
4120 "Meerkat 0.6.2 owns provider replay projection for generated images"
4121 );
4122 }
4123
4124 #[derive(Default)]
4125 struct CapturingLlmClient {
4126 projected_messages: std::sync::Mutex<Vec<meerkat_core::Message>>,
4127 }
4128
4129 #[async_trait]
4130 impl LlmClient for CapturingLlmClient {
4131 fn project_replay_messages(
4132 &self,
4133 messages: &[meerkat_core::Message],
4134 ) -> Result<Vec<meerkat_core::Message>, meerkat_client::LlmError> {
4135 *self
4136 .projected_messages
4137 .lock()
4138 .unwrap_or_else(std::sync::PoisonError::into_inner) = messages.to_vec();
4139 Ok(messages.to_vec())
4140 }
4141
4142 fn stream<'a>(&'a self, _request: &'a LlmRequest) -> LlmStream<'a> {
4143 Box::pin(futures::stream::iter([Ok(
4144 meerkat_client::LlmEvent::Done {
4145 outcome: meerkat_client::LlmDoneOutcome::Success {
4146 stop_reason: meerkat_core::StopReason::EndTurn,
4147 },
4148 },
4149 )]))
4150 }
4151
4152 fn provider(&self) -> meerkat_core::Provider {
4153 meerkat_core::Provider::OpenAI
4154 }
4155
4156 async fn health_check(&self) -> Result<(), meerkat_client::LlmError> {
4157 Ok(())
4158 }
4159 }
4160
4161 #[test]
4162 fn replay_sanitizing_llm_client_delegates_provider_projection() {
4163 let capture = Arc::new(CapturingLlmClient::default());
4164 let inner: Arc<dyn LlmClient> = capture.clone();
4165 let wrapped = ReplaySanitizingLlmClient::new(inner);
4166 let messages = vec![meerkat_core::Message::BlockAssistant(
4167 meerkat_core::BlockAssistantMessage::new(
4168 vec![
4169 meerkat_core::AssistantBlock::Text {
4170 text: "visible".to_string(),
4171 meta: None,
4172 },
4173 meerkat_core::AssistantBlock::ServerToolContent {
4174 id: Some("ws-stream".to_string()),
4175 kind: meerkat_core::ServerToolKind::WebSearch,
4176 content: serde_json::json!({
4177 "type": "response.web_search_call.searching",
4178 "item_id": "ws_123"
4179 }),
4180 meta: None,
4181 },
4182 ],
4183 meerkat_core::StopReason::EndTurn,
4184 ),
4185 )];
4186
4187 let projected = wrapped
4188 .project_replay_messages(&messages)
4189 .expect("wrapped client should delegate provider projection");
4190
4191 let seen = capture
4192 .projected_messages
4193 .lock()
4194 .unwrap_or_else(std::sync::PoisonError::into_inner)
4195 .clone();
4196 let meerkat_core::Message::BlockAssistant(assistant) = &seen[0] else {
4197 panic!("expected block assistant");
4198 };
4199 assert_eq!(
4200 assistant.blocks.len(),
4201 1,
4202 "MobKit sanitization must happen before Meerkat provider projection"
4203 );
4204 assert!(matches!(
4205 assistant.blocks[0],
4206 meerkat_core::AssistantBlock::Text { .. }
4207 ));
4208 assert_eq!(
4209 serde_json::to_value(&projected).expect("projected messages serialize"),
4210 serde_json::to_value(&seen).expect("seen messages serialize")
4211 );
4212 }
4213
4214 #[derive(Default)]
4215 struct CapturingAgentLlmClient {
4216 seen_messages: std::sync::Mutex<Vec<meerkat_core::Message>>,
4217 }
4218
4219 #[async_trait]
4220 impl meerkat_core::AgentLlmClient for CapturingAgentLlmClient {
4221 async fn stream_response(
4222 &self,
4223 messages: &[meerkat_core::Message],
4224 _tools: &[Arc<meerkat_core::ToolDef>],
4225 _max_tokens: u32,
4226 _temperature: Option<f32>,
4227 _provider_params: Option<
4228 &meerkat_core::lifecycle::run_primitive::ProviderParamsOverride,
4229 >,
4230 ) -> Result<meerkat_core::agent::LlmStreamResult, meerkat_core::AgentError> {
4231 *self
4232 .seen_messages
4233 .lock()
4234 .unwrap_or_else(std::sync::PoisonError::into_inner) = messages.to_vec();
4235 Ok(meerkat_core::agent::LlmStreamResult::new(
4236 Vec::new(),
4237 meerkat_core::StopReason::EndTurn,
4238 meerkat_core::Usage::default(),
4239 ))
4240 }
4241
4242 fn provider(&self) -> meerkat_core::Provider {
4243 meerkat_core::Provider::OpenAI
4244 }
4245
4246 fn model(&self) -> &'static str {
4247 "gpt-5.5"
4248 }
4249 }
4250
4251 #[tokio::test]
4252 async fn sanitize_agent_llm_client_drops_replay_unsafe_server_tool_blocks() {
4253 let capture = Arc::new(CapturingAgentLlmClient::default());
4254 let inner: Arc<dyn meerkat_core::AgentLlmClient> = capture.clone();
4255 let wrapped = ReplaySanitizingAgentLlmClient::wrap(inner);
4256 let messages = vec![meerkat_core::Message::BlockAssistant(
4257 meerkat_core::BlockAssistantMessage::new(
4258 vec![
4259 meerkat_core::AssistantBlock::Text {
4260 text: "visible".to_string(),
4261 meta: None,
4262 },
4263 meerkat_core::AssistantBlock::ServerToolContent {
4264 id: Some("ws-stream".to_string()),
4265 kind: meerkat_core::ServerToolKind::WebSearch,
4266 content: serde_json::json!({
4267 "type": "response.web_search_call.searching",
4268 "item_id": "ws_123"
4269 }),
4270 meta: None,
4271 },
4272 meerkat_core::AssistantBlock::ServerToolContent {
4273 id: Some("ok".to_string()),
4274 kind: meerkat_core::ServerToolKind::ProviderNative {
4275 name: "web_search_call".to_string(),
4276 },
4277 content: serde_json::json!({
4278 "type": "web_search_call",
4279 "id": "ws_123",
4280 "status": "completed"
4281 }),
4282 meta: None,
4283 },
4284 ],
4285 meerkat_core::StopReason::EndTurn,
4286 ),
4287 )];
4288 let tools: Vec<Arc<meerkat_core::ToolDef>> = Vec::new();
4289
4290 wrapped
4291 .stream_response(&messages, &tools, 512, None, None)
4292 .await
4293 .expect("wrapped client should delegate");
4294
4295 let seen = capture
4296 .seen_messages
4297 .lock()
4298 .unwrap_or_else(std::sync::PoisonError::into_inner)
4299 .clone();
4300 let meerkat_core::Message::BlockAssistant(assistant) = &seen[0] else {
4301 panic!("expected block assistant");
4302 };
4303 assert_eq!(assistant.blocks.len(), 2);
4304 assert!(matches!(
4305 assistant.blocks[0],
4306 meerkat_core::AssistantBlock::Text { .. }
4307 ));
4308 assert!(matches!(
4309 assistant.blocks[1],
4310 meerkat_core::AssistantBlock::ServerToolContent { ref kind, .. }
4311 if kind.provider_name() == "web_search_call"
4312 ));
4313 }
4314
4315 fn generated_image_block_for_test() -> meerkat_core::AssistantBlock {
4316 serde_json::from_value(serde_json::json!({
4317 "block_type": "image",
4318 "data": {
4319 "image_id": "00000000-0000-0000-0000-000000000051",
4320 "blob_ref": {
4321 "blob_id": "sha256:test-generated-image",
4322 "media_type": "image/png"
4323 },
4324 "media_type": "image/png",
4325 "width": 1024,
4326 "height": 1024,
4327 "revised_prompt": { "disposition": "not_requested" },
4328 "meta": { "provider": "not_emitted" }
4329 }
4330 }))
4331 .expect("test image block should deserialize")
4332 }
4333
4334 #[test]
4335 fn sanitize_message_preserves_assistant_image_blocks() {
4336 let message =
4337 meerkat_core::Message::BlockAssistant(meerkat_core::BlockAssistantMessage::new(
4338 vec![
4339 meerkat_core::AssistantBlock::Text {
4340 text: "Here is the image.".to_string(),
4341 meta: None,
4342 },
4343 generated_image_block_for_test(),
4344 ],
4345 meerkat_core::StopReason::EndTurn,
4346 ));
4347
4348 let sanitized = sanitize_message_for_stateless_replay(message);
4349 let meerkat_core::Message::BlockAssistant(assistant) = sanitized else {
4350 panic!("expected block assistant");
4351 };
4352
4353 assert_eq!(assistant.blocks.len(), 2);
4354 assert!(matches!(
4355 assistant.blocks[0],
4356 meerkat_core::AssistantBlock::Text { .. }
4357 ));
4358 assert!(
4359 matches!(
4360 assistant.blocks[1],
4361 meerkat_core::AssistantBlock::Image { .. }
4362 ),
4363 "generated image blocks should reach Meerkat's provider projection"
4364 );
4365 }
4366
4367 #[test]
4370 fn persistent_with_hook_wraps_session_service() {
4371 let dir = tempfile::tempdir().unwrap_or_else(|e| panic!("{e}"));
4372 let store_path = dir.path().to_path_buf();
4373 let Ok(sqlite) = meerkat_store::SqliteSessionStore::open(store_path.join("sessions.db"))
4374 else {
4375 panic!("failed to open sqlite session store");
4376 };
4377 let session_store: Arc<dyn SessionStore> = Arc::new(sqlite);
4378 let Ok(definition) = meerkat_mob::MobDefinition::from_toml("[mob]\nid = \"test\"\n") else {
4379 panic!("failed to parse minimal mob definition");
4380 };
4381
4382 let hook_called = Arc::new(std::sync::atomic::AtomicBool::new(false));
4383 let hook_called_clone = hook_called.clone();
4384
4385 let spec = MobBootstrapSpec::persistent_with_hook(
4386 definition,
4387 meerkat_mob::MobStorage::in_memory(),
4388 store_path.clone(),
4389 4,
4390 session_store,
4391 move |_req: &mut CreateSessionRequest| {
4392 hook_called_clone.store(true, std::sync::atomic::Ordering::Relaxed);
4393 Box::pin(async { Ok(()) })
4394 },
4395 );
4396
4397 assert!(
4405 spec.runtime_adapter.is_some(),
4406 "persistent_with_hook must provide a runtime adapter via spec.runtime_adapter"
4407 );
4408 assert!(
4409 spec.session_service.runtime_adapter().is_some(),
4410 "session service must own a runtime_store so archive/retire don't \
4411 hit the store-only-projection rejection in meerkat-session"
4412 );
4413 assert!(
4414 store_path.join("runtime.sqlite").exists(),
4415 "persistent_inner must open a SqliteRuntimeStore at <store_path>/runtime.sqlite"
4416 );
4417
4418 assert!(
4423 !hook_called.load(std::sync::atomic::Ordering::Relaxed),
4424 "hook must not be called before create_session"
4425 );
4426 }
4427
4428 #[test]
4430 fn ephemeral_with_hook_creates_spec() {
4431 let dir = tempfile::tempdir().unwrap_or_else(|e| panic!("{e}"));
4432 let store_path = dir.path().to_path_buf();
4433 let Ok(definition) = meerkat_mob::MobDefinition::from_toml("[mob]\nid = \"test\"\n") else {
4434 panic!("failed to parse minimal mob definition");
4435 };
4436
4437 let hook_called = Arc::new(std::sync::atomic::AtomicBool::new(false));
4438 let hook_called_clone = hook_called.clone();
4439
4440 let spec = MobBootstrapSpec::ephemeral_with_hook(
4441 definition,
4442 meerkat_mob::MobStorage::in_memory(),
4443 store_path,
4444 4,
4445 None,
4446 move |_req: &mut CreateSessionRequest| {
4447 hook_called_clone.store(true, std::sync::atomic::Ordering::Relaxed);
4448 Box::pin(async { Ok(()) })
4449 },
4450 );
4451
4452 assert!(spec.runtime_adapter.is_none());
4454
4455 assert!(
4457 !hook_called.load(std::sync::atomic::Ordering::Relaxed),
4458 "hook must not be called before create_session"
4459 );
4460 }
4461
4462 #[tokio::test]
4466 async fn pre_build_hook_mutates_create_session_request() {
4467 use std::sync::Mutex;
4468
4469 let captured = Arc::new(Mutex::new(None::<(String, Option<String>)>));
4470 let captured_clone = captured.clone();
4471
4472 let dir = tempfile::tempdir().unwrap_or_else(|e| panic!("{e}"));
4474 let factory = AgentFactory::new(dir.path()).builtins(true);
4475 let config = Config::default();
4476 let builder = FactoryAgentBuilder::new(factory, config);
4477 let inner: Arc<dyn MobSessionService> =
4478 Arc::new(meerkat_session::EphemeralSessionService::new(builder, 4));
4479
4480 let hook: PreBuildHook = Arc::new(move |req: &mut CreateSessionRequest| {
4482 req.model = "hooked-model".to_string();
4483 req.system_prompt =
4484 meerkat_core::config::SystemPromptOverride::Set("injected-prompt".to_string());
4485 let labels = req.labels.get_or_insert_with(Default::default);
4486 labels.insert("hook_label".to_string(), "hook_value".to_string());
4487 let mut lock = captured_clone
4489 .lock()
4490 .unwrap_or_else(std::sync::PoisonError::into_inner);
4491 *lock = Some((
4492 req.model.clone(),
4493 req.system_prompt.as_set_prompt().map(ToString::to_string),
4494 ));
4495 Box::pin(async { Ok(()) })
4496 });
4497 let wrapped = PreBuildMobSessionService {
4498 inner,
4499 hook,
4500 after_create_hook: None,
4501 runtime_adapter_override: None,
4502 };
4503
4504 let req = CreateSessionRequest {
4505 model: "original-model".to_string(),
4506 prompt: meerkat_core::ContentInput::Text("test".to_string()),
4507 system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
4508 max_tokens: None,
4509 event_tx: None,
4510 initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
4511 build: None,
4512 labels: None,
4513 deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::default(),
4514 injected_context: Vec::new(),
4515 };
4516
4517 let _ = meerkat_core::service::SessionService::create_session(&wrapped, req).await;
4519
4520 let (model, prompt) = captured
4521 .lock()
4522 .unwrap_or_else(std::sync::PoisonError::into_inner)
4523 .clone()
4524 .expect("hook must have been called");
4525 assert_eq!(model, "hooked-model", "hook must mutate the model");
4526 assert_eq!(
4527 prompt.as_deref(),
4528 Some("injected-prompt"),
4529 "hook must set the system prompt"
4530 );
4531 }
4532
4533 #[tokio::test]
4544 async fn config_retry_reaches_agent_effective_retry_policy() {
4545 use meerkat_session::SessionAgentBuilder as _;
4546
4547 let dir = tempfile::tempdir().unwrap_or_else(|e| panic!("{e}"));
4548 let factory = AgentFactory::new(dir.path()).builtins(true);
4549
4550 let mut config = Config::default();
4552 assert_ne!(
4553 config.retry.max_retries, 11,
4554 "test sentinel must differ from default"
4555 );
4556 config.retry.max_retries = 11;
4557 config.retry.initial_delay = std::time::Duration::from_millis(125);
4558 config.retry.max_delay = std::time::Duration::from_secs(7);
4559 config.retry.multiplier = 3.5;
4560
4561 let mut builder = FactoryAgentBuilder::new(factory, config);
4562 builder.default_llm_client = Some(Arc::new(CapturingLlmClient::default()));
4564
4565 let req = CreateSessionRequest {
4566 model: "mock-model".to_string(),
4567 prompt: meerkat_core::ContentInput::Text("noop".to_string()),
4568 system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
4569 max_tokens: None,
4570 event_tx: None,
4571 initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
4572 build: None,
4573 labels: None,
4574 deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::Discard,
4575 injected_context: Vec::new(),
4576 };
4577
4578 let (event_tx, _event_rx) = tokio::sync::mpsc::channel(8);
4579 let factory_agent = builder
4580 .build_agent(&req, event_tx)
4581 .await
4582 .unwrap_or_else(|e| panic!("build_agent should succeed offline: {e}"));
4583
4584 let effective = factory_agent.agent().retry_policy();
4585 assert_eq!(
4586 effective.max_retries, 11,
4587 "Config.retry.max_retries must be plumbed into the agent's effective \
4588 RetryPolicy, not the default 3"
4589 );
4590 assert_eq!(
4591 effective.initial_delay,
4592 std::time::Duration::from_millis(125),
4593 "Config.retry.initial_delay must be plumbed, not the 500ms default"
4594 );
4595 assert_eq!(
4596 effective.max_delay,
4597 std::time::Duration::from_secs(7),
4598 "Config.retry.max_delay must be plumbed, not the 30s default"
4599 );
4600 assert!(
4601 (effective.multiplier - 3.5).abs() < f64::EPSILON,
4602 "Config.retry.multiplier must be plumbed, not the 2.0 default"
4603 );
4604 }
4605
4606 #[derive(Default)]
4607 struct ForwardingProbe {
4608 calls: Mutex<Vec<&'static str>>,
4609 }
4610
4611 impl ForwardingProbe {
4612 fn record(&self, call: &'static str) {
4613 self.calls
4614 .lock()
4615 .unwrap_or_else(std::sync::PoisonError::into_inner)
4616 .push(call);
4617 }
4618
4619 fn calls(&self) -> Vec<&'static str> {
4620 self.calls
4621 .lock()
4622 .unwrap_or_else(std::sync::PoisonError::into_inner)
4623 .clone()
4624 }
4625 }
4626
4627 #[async_trait]
4628 impl meerkat_core::service::SessionService for ForwardingProbe {
4629 async fn create_session(
4630 &self,
4631 _req: CreateSessionRequest,
4632 ) -> Result<meerkat_core::types::RunResult, SessionError> {
4633 Err(SessionError::Unsupported("create_session".to_string()))
4634 }
4635
4636 async fn start_turn(
4637 &self,
4638 _id: &meerkat_core::types::SessionId,
4639 _req: meerkat_core::service::StartTurnRequest,
4640 ) -> Result<meerkat_core::types::RunResult, SessionError> {
4641 Err(SessionError::Unsupported("start_turn".to_string()))
4642 }
4643
4644 async fn interrupt(
4645 &self,
4646 _id: &meerkat_core::types::SessionId,
4647 ) -> Result<(), SessionError> {
4648 self.record("interrupt");
4649 Ok(())
4650 }
4651
4652 async fn read(
4653 &self,
4654 id: &meerkat_core::types::SessionId,
4655 ) -> Result<meerkat_core::service::SessionView, SessionError> {
4656 Err(SessionError::NotFound { id: id.clone() })
4657 }
4658
4659 async fn list(
4660 &self,
4661 _query: meerkat_core::service::SessionQuery,
4662 ) -> Result<Vec<meerkat_core::service::SessionSummary>, SessionError> {
4663 Ok(Vec::new())
4664 }
4665
4666 async fn archive(&self, _id: &meerkat_core::types::SessionId) -> Result<(), SessionError> {
4667 self.record("archive");
4668 Ok(())
4669 }
4670 }
4671
4672 #[async_trait]
4673 impl meerkat_core::service::SessionServiceCommsExt for ForwardingProbe {}
4674
4675 #[async_trait]
4676 impl meerkat_core::service::SessionServiceControlExt for ForwardingProbe {
4677 async fn append_system_context(
4678 &self,
4679 _id: &meerkat_core::types::SessionId,
4680 _req: meerkat_core::service::AppendSystemContextRequest,
4681 ) -> Result<
4682 meerkat_core::service::AppendSystemContextResult,
4683 meerkat_core::service::SessionControlError,
4684 > {
4685 self.record("append_system_context");
4686 Ok(meerkat_core::service::AppendSystemContextResult {
4687 status: meerkat_core::service::AppendSystemContextStatus::Applied,
4688 })
4689 }
4690
4691 async fn stage_tool_results(
4692 &self,
4693 _id: &meerkat_core::types::SessionId,
4694 _req: meerkat_core::service::StageToolResultsRequest,
4695 ) -> Result<meerkat_core::service::StageToolResultsResult, SessionError> {
4696 self.record("stage_tool_results");
4697 Ok(meerkat_core::service::StageToolResultsResult {
4698 accepted_result_count: 7,
4699 })
4700 }
4701 }
4702
4703 #[async_trait]
4704 impl meerkat_core::service::SessionServiceHistoryExt for ForwardingProbe {
4705 async fn read_history(
4706 &self,
4707 id: &meerkat_core::types::SessionId,
4708 _query: meerkat_core::service::SessionHistoryQuery,
4709 ) -> Result<meerkat_core::service::SessionHistoryPage, SessionError> {
4710 Err(SessionError::NotFound { id: id.clone() })
4711 }
4712 }
4713
4714 #[async_trait]
4715 impl MobSessionService for ForwardingProbe {
4716 fn supports_persistent_sessions(&self) -> bool {
4717 true
4718 }
4719
4720 fn runtime_adapter(&self) -> Option<Arc<meerkat_runtime::MeerkatMachine>> {
4721 Some(Arc::new(meerkat_runtime::MeerkatMachine::ephemeral()))
4722 }
4723
4724 async fn archive_with_mob_lifecycle_authority(
4725 &self,
4726 _session_id: &meerkat_core::types::SessionId,
4727 ) -> Result<(), SessionError> {
4728 self.record("archive_with_mob_lifecycle_authority");
4729 Ok(())
4730 }
4731
4732 async fn stage_runtime_system_context_for_active_turn(
4733 &self,
4734 _session_id: &meerkat_core::types::SessionId,
4735 _expected_run_id: &meerkat_core::lifecycle::RunId,
4736 _appends: Vec<meerkat_core::session::PendingSystemContextAppend>,
4737 ) -> Result<Option<Vec<u8>>, SessionError> {
4738 self.record("stage_runtime_system_context_for_active_turn");
4739 Ok(Some(b"snapshot".to_vec()))
4740 }
4741
4742 async fn discard_runtime_system_context_for_active_turn(
4743 &self,
4744 _session_id: &meerkat_core::types::SessionId,
4745 _expected_run_id: &meerkat_core::lifecycle::RunId,
4746 _idempotency_keys: Vec<String>,
4747 ) -> Result<(), SessionError> {
4748 self.record("discard_runtime_system_context_for_active_turn");
4749 Ok(())
4750 }
4751
4752 async fn active_turn_system_context_boundary_available(
4753 &self,
4754 _session_id: &meerkat_core::types::SessionId,
4755 ) -> Result<Option<bool>, SessionError> {
4756 self.record("active_turn_system_context_boundary_available");
4757 Ok(Some(true))
4758 }
4759 }
4760
4761 #[tokio::test]
4762 async fn pre_build_wrapper_forwards_mob_authority_and_control_extensions() {
4763 let probe = Arc::new(ForwardingProbe::default());
4764 let inner: Arc<dyn MobSessionService> = probe.clone();
4765 let wrapped = PreBuildMobSessionService {
4766 inner,
4767 hook: no_op_pre_build_hook(),
4768 after_create_hook: None,
4769 runtime_adapter_override: Some(Arc::new(meerkat_runtime::MeerkatMachine::ephemeral())),
4770 };
4771 let session_id = meerkat_core::types::SessionId::new();
4772
4773 MobSessionService::archive_with_mob_lifecycle_authority(&wrapped, &session_id)
4774 .await
4775 .expect("archive_with_mob_lifecycle_authority should forward to inner service");
4776 let staged = meerkat_core::service::SessionServiceControlExt::stage_tool_results(
4777 &wrapped,
4778 &session_id,
4779 meerkat_core::service::StageToolResultsRequest {
4780 results: Vec::new(),
4781 },
4782 )
4783 .await
4784 .expect("stage_tool_results should forward to inner service");
4785
4786 assert_eq!(staged.accepted_result_count, 7);
4787 let boundary_available = wrapped
4788 .active_turn_system_context_boundary_available(&session_id)
4789 .await
4790 .expect("active-turn boundary probe should forward");
4791 assert_eq!(boundary_available, Some(true));
4792 let snapshot = wrapped
4793 .stage_runtime_system_context_for_active_turn(
4794 &session_id,
4795 &meerkat_core::lifecycle::RunId::new(),
4796 vec![meerkat_core::session::PendingSystemContextAppend {
4797 content: meerkat_core::lifecycle::run_primitive::CoreRenderable::Text {
4798 text: "steer".to_string(),
4799 },
4800 source: Some("test".to_string()),
4801 idempotency_key: Some("test".to_string()),
4802 source_kind: meerkat_core::session::SystemContextSource::default(),
4803 peer_response_terminal: None,
4804 accepted_at: meerkat_core::time_compat::SystemTime::now(),
4805 }],
4806 )
4807 .await
4808 .expect("active-turn staging should forward");
4809 assert_eq!(snapshot.as_deref(), Some(&b"snapshot"[..]));
4810 wrapped
4811 .discard_runtime_system_context_for_active_turn(
4812 &session_id,
4813 &meerkat_core::lifecycle::RunId::new(),
4814 vec!["test".to_string()],
4815 )
4816 .await
4817 .expect("active-turn rollback should forward");
4818 assert_eq!(
4819 probe.calls(),
4820 vec![
4821 "archive_with_mob_lifecycle_authority",
4822 "stage_tool_results",
4823 "active_turn_system_context_boundary_available",
4824 "stage_runtime_system_context_for_active_turn",
4825 "discard_runtime_system_context_for_active_turn",
4826 ]
4827 );
4828 }
4829
4830 #[test]
4850 fn persistent_bootstrap_uses_sqlite_runtime_store() {
4851 let dir = tempfile::tempdir().unwrap_or_else(|e| panic!("{e}"));
4852 let store_path = dir.path().to_path_buf();
4853 let Ok(sqlite) = meerkat_store::SqliteSessionStore::open(store_path.join("sessions.db"))
4854 else {
4855 panic!("failed to open sqlite session store");
4856 };
4857 let session_store: Arc<dyn SessionStore> = Arc::new(sqlite);
4858 let Ok(definition) = meerkat_mob::MobDefinition::from_toml("[mob]\nid = \"test\"\n") else {
4859 panic!("failed to parse minimal mob definition");
4860 };
4861 let spec = MobBootstrapSpec::persistent(
4862 definition,
4863 meerkat_mob::MobStorage::in_memory(),
4864 store_path.clone(),
4865 4,
4866 session_store,
4867 );
4868 assert!(
4869 spec.runtime_adapter.is_some(),
4870 "persistent bootstrap must provide its own runtime adapter via spec.runtime_adapter"
4871 );
4872 assert!(
4873 spec.session_service.runtime_adapter().is_some(),
4874 "session service must own a runtime_store so archive/retire don't \
4875 hit the store-only-projection rejection"
4876 );
4877 assert!(
4878 store_path.join("runtime.sqlite").exists(),
4879 "persistent_inner must open a SqliteRuntimeStore at <store_path>/runtime.sqlite"
4880 );
4881 }
4882
4883 #[test]
4887 fn ephemeral_runtime_backed_uses_session_service_runtime_adapter() {
4888 let dir = tempfile::tempdir().unwrap_or_else(|e| panic!("{e}"));
4889 let store_path = dir.path().to_path_buf();
4890 let Ok(definition) = meerkat_mob::MobDefinition::from_toml("[mob]\nid = \"test\"\n") else {
4891 panic!("failed to parse minimal mob definition");
4892 };
4893 let spec = MobBootstrapSpec::ephemeral_runtime_backed_inner(
4894 definition,
4895 meerkat_mob::MobStorage::in_memory(),
4896 store_path,
4897 4,
4898 None,
4899 None,
4900 None,
4901 CapabilityFlags::default(),
4902 None,
4903 None,
4904 );
4905 assert!(
4906 spec.runtime_adapter.is_some(),
4907 "ephemeral_runtime_backed_inner must expose the shared runtime authority"
4908 );
4909 assert!(
4910 spec.session_service.runtime_adapter().is_some(),
4911 "session service must still expose a runtime adapter so autonomous-host comms can wire"
4912 );
4913 }
4914
4915 #[tokio::test]
4916 async fn agent_mob_tools_expose_definition_profiles_as_realm_profiles() {
4917 let dir = tempfile::tempdir().unwrap_or_else(|e| panic!("{e}"));
4918 let store_path = dir.path().to_path_buf();
4919 let Ok(definition) = meerkat_mob::MobDefinition::from_toml(
4920 "[mob]\nid = \"test\"\n\n[profiles.investigation-worker]\nmodel = \"gpt-5.5\"\n[profiles.investigation-worker.tools]\ncomms = true\nmob = true\n\n[profiles.person-worker]\nmodel = \"gpt-5.5\"\n[profiles.person-worker.tools]\ncomms = true\n",
4921 ) else {
4922 panic!("failed to parse mob definition with worker profiles");
4923 };
4924
4925 let spec = MobBootstrapSpec::ephemeral_runtime_backed_inner(
4926 definition,
4927 meerkat_mob::MobStorage::in_memory(),
4928 store_path,
4929 4,
4930 None,
4931 None,
4932 None,
4933 CapabilityFlags::default(),
4934 None,
4935 None,
4936 );
4937 let state = spec
4938 .agent_mob_mcp_state
4939 .expect("agent mob MCP state should be installed");
4940
4941 let profiles = state
4942 .realm_profile_list()
4943 .await
4944 .expect("definition profiles should list through agent mob tools");
4945 let names = profiles
4946 .iter()
4947 .map(|profile| profile.name.as_str())
4948 .collect::<Vec<_>>();
4949 assert!(
4950 names.contains(&"investigation-worker"),
4951 "definition profiles must be visible to mob_profile_list so agents can create mobs that reference them"
4952 );
4953 assert!(names.contains(&"person-worker"));
4954
4955 let worker = state
4956 .realm_profile_get("investigation-worker")
4957 .await
4958 .expect("definition profile lookup should succeed")
4959 .expect("definition profile should exist");
4960 assert_eq!(worker.profile.model, "gpt-5.5");
4961 assert_eq!(
4962 worker.revision, 0,
4963 "definition-backed profiles are immutable runtime seeds, not persisted realm revisions"
4964 );
4965 }
4966
4967 #[tokio::test]
4968 async fn agent_created_mobs_can_spawn_definition_seeded_realm_profiles() {
4969 let dir = tempfile::tempdir().unwrap_or_else(|e| panic!("{e}"));
4970 let store_path = dir.path().to_path_buf();
4971 let Ok(parent_definition) = meerkat_mob::MobDefinition::from_toml(
4972 "[mob]\nid = \"parent\"\n\n[profiles.investigation-worker]\nmodel = \"gpt-5.5\"\n[profiles.investigation-worker.tools]\ncomms = true\nmob = true\n",
4973 ) else {
4974 panic!("failed to parse parent mob definition");
4975 };
4976
4977 let mut spec = MobBootstrapSpec::ephemeral_runtime_backed_inner(
4978 parent_definition,
4979 meerkat_mob::MobStorage::in_memory(),
4980 store_path,
4981 4,
4982 None,
4983 None,
4984 None,
4985 CapabilityFlags::default(),
4986 None,
4987 None,
4988 );
4989 spec.options.default_llm_client = Some(Arc::new(meerkat_client::TestClient::default()));
4990 let runtime = MobRuntime::bootstrap(spec)
4991 .await
4992 .unwrap_or_else(|e| panic!("{e}"));
4993 let state = runtime
4994 .agent_mob_mcp_state
4995 .clone()
4996 .expect("agent mob MCP state should be installed");
4997 let Ok(child_definition) = meerkat_mob::MobDefinition::from_toml(
4998 "[mob]\nid = \"child\"\n\n[profiles.investigation-worker]\nrealm_profile = \"investigation-worker\"\n",
4999 ) else {
5000 panic!("failed to parse child mob definition");
5001 };
5002
5003 let mob_id = Box::pin(state.mob_create_definition(child_definition))
5004 .await
5005 .expect("child mob should be created");
5006 Box::pin(state.mob_spawn_spec(
5007 &mob_id,
5008 SpawnMemberSpec::new(
5009 ProfileName::from("investigation-worker"),
5010 meerkat_mob::AgentIdentity::from("investigation-worker-one"),
5013 ),
5014 ))
5015 .await
5016 .expect("created mob should resolve definition-seeded realm profile at spawn time");
5017 }
5018
5019 #[tokio::test]
5020 async fn ephemeral_runtime_backed_custom_session_store_persists_created_session() {
5021 let dir = tempfile::tempdir().unwrap_or_else(|e| panic!("{e}"));
5022 let store_path = dir.path().to_path_buf();
5023 let Ok(definition) = meerkat_mob::MobDefinition::from_toml(
5024 "[mob]\nid = \"test\"\n\n[profiles.worker]\nmodel = \"gpt-5.5\"\n[profiles.worker.tools]\ncomms = true\n",
5025 ) else {
5026 panic!("failed to parse minimal mob definition");
5027 };
5028 let custom_store: Arc<dyn SessionStore> = Arc::new(meerkat_store::MemoryStore::new());
5029 let mut spec = MobBootstrapSpec::ephemeral_runtime_backed_inner(
5030 definition,
5031 meerkat_mob::MobStorage::in_memory(),
5032 store_path,
5033 4,
5034 Some(custom_store.clone()),
5035 None,
5036 None,
5037 CapabilityFlags::default(),
5038 None,
5039 None,
5040 );
5041 spec.options.default_llm_client = Some(Arc::new(meerkat_client::TestClient::default()));
5042
5043 let runtime = MobRuntime::bootstrap(spec)
5044 .await
5045 .unwrap_or_else(|e| panic!("{e}"));
5046 let mid = meerkat_mob::ids::AgentIdentity::from("worker-one");
5049 Box::pin(runtime.handle.spawn_spec(SpawnMemberSpec::new(
5050 ProfileName::from("worker"),
5051 mid.clone(),
5052 )))
5053 .await
5054 .unwrap_or_else(|e| panic!("{e}"));
5055 let session_id = runtime
5056 .handle
5057 .resolve_bridge_session_id(&mid)
5058 .await
5059 .unwrap_or_else(|| panic!("spawned worker has no bridge session id"));
5060
5061 let stored = custom_store
5062 .load(&session_id)
5063 .await
5064 .unwrap_or_else(|e| panic!("{e}"));
5065 assert!(
5066 stored.is_some(),
5067 "ephemeral runtime-backed builds with a custom store must persist through that store"
5068 );
5069 }
5070
5071 #[tokio::test]
5072 async fn ephemeral_runtime_backed_custom_session_store_resumes_after_runtime_restart() {
5073 let dir = tempfile::tempdir().unwrap_or_else(|e| panic!("{e}"));
5074 let store_path = dir.path().to_path_buf();
5075 let definition_toml = "[mob]\nid = \"test\"\n\n[profiles.worker]\nmodel = \"gpt-5.5\"\n[profiles.worker.tools]\ncomms = true\n";
5076 let Ok(definition) = meerkat_mob::MobDefinition::from_toml(definition_toml) else {
5077 panic!("failed to parse minimal mob definition");
5078 };
5079 let custom_store: Arc<dyn SessionStore> = Arc::new(meerkat_store::MemoryStore::new());
5080 let mut spec = MobBootstrapSpec::ephemeral_runtime_backed_inner(
5081 definition,
5082 meerkat_mob::MobStorage::in_memory(),
5083 store_path.clone(),
5084 4,
5085 Some(custom_store.clone()),
5086 None,
5087 None,
5088 CapabilityFlags::default(),
5089 None,
5090 None,
5091 );
5092 spec.options.default_llm_client = Some(Arc::new(meerkat_client::TestClient::default()));
5093
5094 let runtime = MobRuntime::bootstrap(spec)
5095 .await
5096 .unwrap_or_else(|e| panic!("{e}"));
5097 let mid = meerkat_mob::ids::AgentIdentity::from("worker-one");
5100 Box::pin(runtime.handle.spawn_spec(SpawnMemberSpec::new(
5101 ProfileName::from("worker"),
5102 mid.clone(),
5103 )))
5104 .await
5105 .unwrap_or_else(|e| panic!("{e}"));
5106 let session_id = runtime
5107 .handle
5108 .resolve_bridge_session_id(&mid)
5109 .await
5110 .unwrap_or_else(|| panic!("spawned worker has no bridge session id"));
5111 drop(runtime);
5112
5113 let Ok(definition) = meerkat_mob::MobDefinition::from_toml(definition_toml) else {
5114 panic!("failed to parse minimal mob definition");
5115 };
5116 let mut restarted_spec = MobBootstrapSpec::ephemeral_runtime_backed_inner(
5117 definition,
5118 meerkat_mob::MobStorage::in_memory(),
5119 store_path,
5120 4,
5121 Some(custom_store),
5122 None,
5123 None,
5124 CapabilityFlags::default(),
5125 None,
5126 None,
5127 );
5128 restarted_spec.options.default_llm_client =
5129 Some(Arc::new(meerkat_client::TestClient::default()));
5130
5131 let restarted = MobRuntime::bootstrap(restarted_spec)
5132 .await
5133 .unwrap_or_else(|e| panic!("{e}"));
5134 let mut resume_spec = SpawnMemberSpec::new(ProfileName::from("worker"), mid.clone());
5135 resume_spec.launch_mode = meerkat_mob::MemberLaunchMode::Resume {
5136 bridge_session_id: session_id.clone(),
5137 };
5138 Box::pin(restarted.handle.spawn_spec(resume_spec))
5139 .await
5140 .unwrap_or_else(|e| panic!("resume should load the external session snapshot: {e}"));
5141
5142 let resumed_session_id = restarted
5143 .handle
5144 .resolve_bridge_session_id(&mid)
5145 .await
5146 .unwrap_or_else(|| panic!("resumed worker has no bridge session id"));
5147 assert_eq!(resumed_session_id, session_id);
5148 }
5149
5150 #[test]
5153 fn ephemeral_bootstrap_without_image_generation_stays_direct() {
5154 let dir = tempfile::tempdir().unwrap_or_else(|e| panic!("{e}"));
5155 let store_path = dir.path().to_path_buf();
5156 let Ok(definition) = meerkat_mob::MobDefinition::from_toml(
5157 "[mob]\nid = \"test\"\n\n[profiles.worker]\nmodel = \"gpt-5.5\"\nruntime_mode = \"autonomous_host\"\n[profiles.worker.tools]\ncomms = true\n",
5158 ) else {
5159 panic!("failed to parse definition");
5160 };
5161 let spec = MobBootstrapSpec::ephemeral_inner(
5162 definition,
5163 meerkat_mob::MobStorage::in_memory(),
5164 store_path,
5165 4,
5166 None,
5167 None,
5168 CapabilityFlags::default(),
5169 None,
5170 None,
5171 );
5172 assert!(
5173 spec.runtime_adapter.is_none(),
5174 "public ephemeral builds only need a runtime adapter when the definition may use image generation"
5175 );
5176 }
5177
5178 #[test]
5183 fn ephemeral_bootstrap_with_image_generation_shares_runtime_adapter() {
5184 let dir = tempfile::tempdir().unwrap_or_else(|e| panic!("{e}"));
5185 let store_path = dir.path().to_path_buf();
5186 let Ok(definition) = meerkat_mob::MobDefinition::from_toml(
5187 r#"
5188[mob]
5189id = "test"
5190
5191[profiles.commander]
5192model = "gpt-5.5"
5193
5194[profiles.commander.tools]
5195builtins = true
5196image_generation = true
5197"#,
5198 ) else {
5199 panic!("failed to parse image-generation definition");
5200 };
5201 let spec = MobBootstrapSpec::ephemeral(
5202 definition,
5203 meerkat_mob::MobStorage::in_memory(),
5204 store_path,
5205 4,
5206 None,
5207 );
5208 let spec_adapter = spec
5209 .runtime_adapter
5210 .as_ref()
5211 .expect("image-generation ephemeral builds must expose a runtime adapter");
5212 let service_adapter = spec
5213 .session_service
5214 .runtime_adapter()
5215 .expect("session service must expose the same runtime adapter");
5216 assert!(
5217 spec_adapter.shares_runtime_persistence_with(&service_adapter),
5218 "image-generation tool state and session state must share one runtime authority"
5219 );
5220 }
5221
5222 #[test]
5225 fn normalize_runtime_turn_request_strips_runtime_owned_semantics() {
5226 let req = meerkat_core::service::StartTurnRequest {
5227 prompt: meerkat_core::ContentInput::Text("checkpoint".to_string()),
5228 injected_context: Vec::new(),
5229 system_prompt: Some("system".to_string()),
5230 event_tx: None,
5231 runtime: meerkat_core::service::StartTurnRuntimeSemantics {
5232 handling_mode: meerkat_core::types::HandlingMode::Steer,
5233 flow_tool_overlay: None,
5234 pre_turn_context_appends: Vec::new(),
5235 typed_turn_appends: Vec::new(),
5236 turn_metadata: Some(
5239 meerkat_core::lifecycle::run_primitive::RuntimeTurnMetadata {
5240 render_metadata: Some(meerkat_core::types::RenderMetadata {
5241 class: meerkat_core::types::RenderClass::OpsProgress,
5242 salience: meerkat_core::types::RenderSalience::Urgent,
5243 }),
5244 ..Default::default()
5245 },
5246 ),
5247 },
5248 };
5249
5250 let expected_prompt = req.prompt.clone();
5251 let expected_system_prompt = req.system_prompt.clone();
5252
5253 let normalized = normalize_runtime_turn_request(req);
5254
5255 assert_eq!(
5256 normalized.runtime.handling_mode,
5257 meerkat_core::types::HandlingMode::Queue,
5258 "runtime-applied turns must downgrade Steer before reaching direct session services"
5259 );
5260 assert!(
5261 normalized
5262 .runtime
5263 .turn_metadata
5264 .as_ref()
5265 .is_none_or(|metadata| metadata.render_metadata.is_none()),
5266 "runtime-owned render metadata must not be forwarded through the direct agent path"
5267 );
5268 assert_eq!(normalized.prompt, expected_prompt);
5269 assert_eq!(normalized.system_prompt, expected_system_prompt);
5270 }
5271
5272 #[test]
5274 fn session_created_context_fields() {
5275 let ctx = SessionCreatedContext {
5276 model: "claude-sonnet-4-5".to_string(),
5277 labels: std::collections::BTreeMap::from([(
5278 "agent_type".to_string(),
5279 "lead".to_string(),
5280 )]),
5281 system_prompt: Some("You are a lead agent.".to_string()),
5282 };
5283 assert_eq!(ctx.model, "claude-sonnet-4-5");
5284 assert_eq!(ctx.labels["agent_type"], "lead");
5285 assert_eq!(ctx.system_prompt.as_deref(), Some("You are a lead agent."));
5286 }
5287
5288 #[tokio::test]
5290 async fn session_hook_default_impls_are_noop() {
5291 struct EmptyHook;
5292 #[async_trait]
5293 impl SessionHook for EmptyHook {}
5294
5295 let hook = EmptyHook;
5296 let mut req = CreateSessionRequest {
5297 model: "test".to_string(),
5298 prompt: meerkat_core::ContentInput::Text("test".to_string()),
5299 system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
5300 max_tokens: None,
5301 event_tx: None,
5302 initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
5303 build: None,
5304 labels: None,
5305 deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::default(),
5306 injected_context: Vec::new(),
5307 };
5308 hook.before_create(&mut req).await.unwrap();
5310 let ctx = SessionCreatedContext {
5312 model: "test".to_string(),
5313 labels: Default::default(),
5314 system_prompt: None,
5315 };
5316 hook.after_create(&meerkat_core::types::SessionId::new(), &ctx)
5317 .await;
5318 }
5319
5320 #[tokio::test]
5322 async fn session_hook_before_create_can_abort() {
5323 struct AbortHook;
5324 #[async_trait]
5325 impl SessionHook for AbortHook {
5326 async fn before_create(
5327 &self,
5328 _req: &mut CreateSessionRequest,
5329 ) -> Result<(), SessionError> {
5330 Err(SessionError::Unsupported("hook abort".into()))
5331 }
5332 }
5333
5334 let hook = AbortHook;
5335 let mut req = CreateSessionRequest {
5336 model: "test".to_string(),
5337 prompt: meerkat_core::ContentInput::Text("test".to_string()),
5338 system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
5339 max_tokens: None,
5340 event_tx: None,
5341 initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
5342 build: None,
5343 labels: None,
5344 deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::default(),
5345 injected_context: Vec::new(),
5346 };
5347 let result = hook.before_create(&mut req).await;
5348 assert!(result.is_err());
5349 }
5350
5351 #[tokio::test]
5353 async fn session_hook_before_create_mutates_request() {
5354 struct MutatingHook;
5355 #[async_trait]
5356 impl SessionHook for MutatingHook {
5357 async fn before_create(
5358 &self,
5359 req: &mut CreateSessionRequest,
5360 ) -> Result<(), SessionError> {
5361 req.model = "hook-overridden".to_string();
5362 req.system_prompt =
5363 meerkat_core::config::SystemPromptOverride::Set("injected by hook".to_string());
5364 Ok(())
5365 }
5366 }
5367
5368 let hook = MutatingHook;
5369 let mut req = CreateSessionRequest {
5370 model: "original".to_string(),
5371 prompt: meerkat_core::ContentInput::Text("test".to_string()),
5372 system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
5373 max_tokens: None,
5374 event_tx: None,
5375 initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
5376 build: None,
5377 labels: None,
5378 deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::default(),
5379 injected_context: Vec::new(),
5380 };
5381 hook.before_create(&mut req).await.unwrap();
5382 assert_eq!(req.model, "hook-overridden");
5383 assert_eq!(req.system_prompt.as_set_prompt(), Some("injected by hook"));
5384 }
5385
5386 #[test]
5387 fn recoverable_lifecycle_cleanup_accepts_ambiguous_member_cleanup() {
5388 let error = "previous member cleanup ambiguous for member rt:deep-investigator:singleton:0";
5389
5390 assert!(is_previous_member_cleanup_ambiguous_error(error));
5391 assert!(is_recoverable_lifecycle_cleanup_error(error));
5392 }
5393
5394 #[test]
5395 fn topology_restore_failed_peer_ids_extracts_tolerated_peers() {
5396 let identity = meerkat_mob::AgentIdentity::from("rt:review:singleton:0");
5397 let receipt = meerkat_mob::MemberRespawnReceipt::new(
5398 identity.clone(),
5399 meerkat_mob::AgentRuntimeId::new(identity, meerkat_mob::ids::Generation::INITIAL),
5400 meerkat_mob::FenceToken::new(1),
5401 meerkat_mob::FenceToken::new(2),
5402 );
5403 let err = meerkat_mob::MobRespawnError::TopologyRestoreFailed {
5404 receipt,
5405 failed_peer_ids: vec![meerkat_mob::RespawnTopologyPeerId::from(
5406 "initiative:broken",
5407 )],
5408 };
5409
5410 assert_eq!(
5411 topology_restore_failed_peer_ids(&err),
5412 Some(vec!["initiative:broken".to_string()])
5413 );
5414 assert_eq!(
5415 topology_restore_failed_peer_ids(&meerkat_mob::MobRespawnError::NoRuntimeControl {
5416 identity: meerkat_mob::AgentIdentity::from("rt:review:singleton:0"),
5417 }),
5418 None
5419 );
5420 }
5421
5422 #[test]
5423 fn topology_restore_warning_json_surfaces_isolated_peers() {
5424 let failed_peer_ids = vec!["initiative:broken".to_string(), "helper:cold".to_string()];
5425
5426 assert_eq!(
5427 topology_restore_warning_json(&failed_peer_ids),
5428 serde_json::json!({
5429 "kind": "topology_restore_degraded",
5430 "failed_peer_ids": ["initiative:broken", "helper:cold"],
5431 })
5432 );
5433 }
5434
5435 #[test]
5436 fn recoverable_lifecycle_cleanup_preserves_archive_cancel_race() {
5437 let error = "internal error: disposal completed but ArchiveSession failed: \
5438 session error: agent error: Internal error: runtime cancel-before-retire failed \
5439 for 019e3c52-0f1b-73d3-a5c7-4b21c2bbf131: Runtime not ready: running";
5440
5441 assert!(is_recoverable_lifecycle_cleanup_error(error));
5442 }
5443
5444 #[test]
5445 fn recoverable_lifecycle_cleanup_rejects_unrelated_errors() {
5446 assert!(!is_recoverable_lifecycle_cleanup_error(
5447 "actor task dropped"
5448 ));
5449 assert!(!is_recoverable_lifecycle_cleanup_error(
5450 "model provider returned rate limit"
5451 ));
5452 }
5453
5454 #[test]
5459 fn recoverable_lifecycle_cleanup_accepts_stopped_guard_archive_retire() {
5460 let error = "internal error: disposal completed but ArchiveSession failed: \
5461 session error: agent error: Internal error: machine archive retire failed \
5462 after registration: Internal error: DSL authority (Retire): guard rejected \
5463 transition from Stopped for input::Retire";
5464
5465 assert!(is_recoverable_lifecycle_cleanup_error(error));
5466 }
5467
5468 #[test]
5474 fn recoverable_lifecycle_cleanup_accepts_stale_continuity_save_on_retire() {
5475 let record_gone = "internal error: disposal completed but ArchiveSession failed: \
5476 session error: agent error: Internal error: continuity save: \
5477 continuity record not found for identity:luka";
5478 let stale_fence = "internal error: disposal completed but ArchiveSession failed: \
5479 session error: agent error: Internal error: continuity save: \
5480 stale fencing token for identity:luka: presented 1, current 6";
5481
5482 assert!(is_recoverable_lifecycle_cleanup_error(record_gone));
5483 assert!(is_recoverable_lifecycle_cleanup_error(stale_fence));
5484 }
5485
5486 #[test]
5493 fn stopped_session_archive_retire_rejection_matches_only_stopped_guard() {
5494 assert!(is_stopped_session_archive_retire_rejection(
5495 "session error: agent error: Internal error: machine archive retire failed \
5496 after registration: Internal error: DSL authority (Retire): guard rejected \
5497 transition from Stopped for input::Retire"
5498 ));
5499 assert!(is_stopped_session_archive_retire_rejection(
5501 "agent error: Internal error: machine archive retire failed: Internal error: \
5502 DSL authority (Retire): guard rejected transition from Stopped for input::Retire"
5503 ));
5504 assert!(!is_stopped_session_archive_retire_rejection(
5506 "machine archive retire failed after registration: Internal error: \
5507 DSL authority (Retire): guard rejected transition from Running for input::Retire"
5508 ));
5509 assert!(!is_stopped_session_archive_retire_rejection(
5510 "guard rejected transition from Stopped for input::Retire"
5511 ));
5512 assert!(!is_stopped_session_archive_retire_rejection(
5513 "machine archive retire failed: store unavailable"
5514 ));
5515 }
5516
5517 #[test]
5521 fn recoverable_lifecycle_cleanup_requires_completed_disposal() {
5522 assert!(!is_recoverable_lifecycle_cleanup_error(
5523 "continuity save: continuity record not found for identity:luka"
5524 ));
5525 assert!(!is_recoverable_lifecycle_cleanup_error(
5526 "DSL authority (Retire): guard rejected transition from Stopped for input::Retire"
5527 ));
5528 assert!(!is_recoverable_lifecycle_cleanup_error(
5529 "disposal aborted at ArchiveSession: continuity save: stale fencing token"
5530 ));
5531 }
5532
5533 #[test]
5540 fn session_owned_retire_cleanup_accepts_archive_not_found_for_registered_runtime_session() {
5541 let reset_error = "internal error: disposal completed but ArchiveSession failed: \
5545 session error: agent error: Internal error: mob archive authority returned \
5546 NotFound for registered runtime session 019ee136-33bc-7bc3-80f9-2aac38736291";
5547 let delete_error = "internal error: disposal completed but ArchiveSession failed: \
5548 session error: agent error: Internal error: mob archive authority returned \
5549 NotFound for registered runtime session 019ee136-340d-7203-a089-c3357835c824";
5550
5551 assert!(is_recoverable_session_owned_retire_cleanup_error(
5552 reset_error
5553 ));
5554 assert!(is_recoverable_session_owned_retire_cleanup_error(
5555 delete_error
5556 ));
5557 }
5558
5559 #[test]
5564 fn mob_member_lifecycle_gate_still_rejects_archive_not_found() {
5565 let error = "internal error: disposal completed but ArchiveSession failed: \
5566 session error: agent error: Internal error: mob archive authority returned \
5567 NotFound for registered runtime session 019ee136-33bc-7bc3-80f9-2aac38736291";
5568
5569 assert!(!is_recoverable_lifecycle_cleanup_error(error));
5570 assert!(is_session_owned_archive_absent_cleanup_error(error));
5572 }
5573
5574 #[test]
5576 fn session_owned_retire_cleanup_still_accepts_shared_recoverable_cases() {
5577 let cancel_race = "internal error: disposal completed but ArchiveSession failed: \
5578 session error: agent error: Internal error: runtime cancel-before-retire failed \
5579 for 019e3c52-0f1b-73d3-a5c7-4b21c2bbf131: Runtime not ready: running";
5580
5581 assert!(is_recoverable_session_owned_retire_cleanup_error(
5582 cancel_race
5583 ));
5584 }
5585
5586 #[test]
5589 fn session_owned_retire_cleanup_requires_completed_disposal() {
5590 assert!(!is_recoverable_session_owned_retire_cleanup_error(
5591 "disposal aborted at ArchiveSession: mob archive authority returned NotFound \
5592 for registered runtime session 019ee136-33bc-7bc3-80f9-2aac38736291"
5593 ));
5594 assert!(!is_recoverable_session_owned_retire_cleanup_error(
5595 "mob archive authority returned NotFound for registered runtime session \
5596 019ee136-33bc-7bc3-80f9-2aac38736291"
5597 ));
5598 assert!(!is_recoverable_session_owned_retire_cleanup_error(
5599 "model provider returned rate limit"
5600 ));
5601 }
5602
5603 mod console_spawn_projection {
5604 use super::*;
5605 use crate::console_spawn::new_console_spawn_sink_slot;
5606 use crate::unified_runtime::ConsoleEventStore;
5607 use meerkat_core::AgentToolDispatcher;
5608
5609 struct CannedDispatcher {
5612 payload: Value,
5613 is_error: bool,
5614 }
5615
5616 #[async_trait::async_trait]
5617 impl meerkat_core::AgentToolDispatcher for CannedDispatcher {
5618 fn tools(&self) -> Arc<[Arc<meerkat_core::types::ToolDef>]> {
5619 Vec::<Arc<meerkat_core::types::ToolDef>>::new().into()
5620 }
5621
5622 async fn dispatch(
5623 &self,
5624 call: meerkat_core::types::ToolCallView<'_>,
5625 ) -> Result<meerkat_core::ToolDispatchOutcome, meerkat_core::ToolError> {
5626 Ok(meerkat_core::ToolDispatchOutcome::sync_result(
5627 meerkat_core::types::ToolResult::new(
5628 call.id.to_string(),
5629 self.payload.to_string(),
5630 self.is_error,
5631 ),
5632 ))
5633 }
5634
5635 fn capabilities(&self) -> meerkat_core::agent::DispatcherCapabilities {
5636 meerkat_core::agent::DispatcherCapabilities::default()
5637 }
5638 }
5639
5640 fn console_wrapper(
5641 payload: Value,
5642 is_error: bool,
5643 store: Option<&ConsoleEventStore>,
5644 ) -> AutoWireParentMobToolDispatcher {
5645 let console_spawn_sink = new_console_spawn_sink_slot();
5646 if let Some(store) = store {
5647 *console_spawn_sink
5648 .write()
5649 .unwrap_or_else(std::sync::PoisonError::into_inner) =
5650 Some(ConsoleSpawnSink::new(store.clone()));
5651 }
5652 AutoWireParentMobToolDispatcher {
5653 inner: Arc::new(CannedDispatcher { payload, is_error }),
5654 implicit_delegate_retirement_overrides:
5655 ImplicitDelegateRetirementOverrides::default(),
5656 console_spawn_sink,
5657 spawner_comms_name: Some("ob3/orchestrator/ops-lead".to_string()),
5658 }
5659 }
5660
5661 async fn dispatch_tool(
5662 dispatcher: &AutoWireParentMobToolDispatcher,
5663 name: &str,
5664 args: Value,
5665 ) -> meerkat_core::ToolDispatchOutcome {
5666 let raw = serde_json::value::RawValue::from_string(args.to_string()).expect("raw args");
5667 dispatcher
5668 .dispatch(meerkat_core::types::ToolCallView {
5669 id: "call-1",
5670 name,
5671 args: &raw,
5672 })
5673 .await
5674 .expect("dispatch succeeds")
5675 }
5676
5677 async fn kickoff_events(
5678 store: &ConsoleEventStore,
5679 ) -> Vec<crate::console_contracts::ConsoleIdentityEventEnvelope> {
5680 store
5681 .replay_all(None)
5682 .await
5683 .expect("replay")
5684 .into_iter()
5685 .filter(|event| event.event_type == "user_input")
5686 .collect()
5687 }
5688
5689 #[tokio::test]
5690 async fn mob_spawn_member_projects_kickoff_into_console() {
5691 let store = ConsoleEventStore::new();
5692 let dispatcher = console_wrapper(
5693 serde_json::json!({
5694 "agent_identity": "worker-3",
5695 "member_ref": "opaque-ref"
5696 }),
5697 false,
5698 Some(&store),
5699 );
5700
5701 let outcome = dispatch_tool(
5702 &dispatcher,
5703 "mob_spawn_member",
5704 serde_json::json!({
5705 "mob_id": "ob3",
5706 "profile": "person-worker",
5707 "member_id": "worker-3",
5708 "initial_message": "Find the person",
5709 "labels": { "group": "workers" }
5710 }),
5711 )
5712 .await;
5713 assert!(!outcome.result.is_error);
5714
5715 let kickoffs = kickoff_events(&store).await;
5716 assert_eq!(kickoffs.len(), 1, "spawn must project one kickoff");
5717 let kickoff = &kickoffs[0];
5718 assert_eq!(kickoff.identity, "worker-3");
5719 assert!(kickoff.event_id.starts_with("spawn-kickoff:ob3:worker-3:"));
5720 assert_eq!(kickoff.data["content"][0]["text"], "Find the person");
5721 assert_eq!(kickoff.data["via_tool"], "mob_spawn_member");
5722 assert_eq!(kickoff.data["parent_identity"], "ops-lead");
5723
5724 let labels = store
5725 .identity_labels("worker-3")
5726 .await
5727 .expect("spawn registers console identity labels");
5728 assert_eq!(labels.get("group").map(String::as_str), Some("workers"));
5729 assert_eq!(
5730 labels.get("spawned_by").map(String::as_str),
5731 Some("ops-lead")
5732 );
5733 }
5734
5735 #[tokio::test]
5736 async fn repeated_spawn_dispatch_keeps_one_kickoff() {
5737 let store = ConsoleEventStore::new();
5738 let dispatcher = console_wrapper(
5739 serde_json::json!({ "agent_identity": "worker-3" }),
5740 false,
5741 Some(&store),
5742 );
5743 let args = serde_json::json!({
5744 "mob_id": "ob3",
5745 "profile": "person-worker",
5746 "member_id": "worker-3",
5747 "initial_message": "Find the person"
5748 });
5749
5750 dispatch_tool(&dispatcher, "mob_spawn_member", args.clone()).await;
5751 dispatch_tool(&dispatcher, "mob_spawn_member", args).await;
5752
5753 assert_eq!(
5754 kickoff_events(&store).await.len(),
5755 1,
5756 "retry/double-spawn must not duplicate the kickoff frame"
5757 );
5758 }
5759
5760 #[tokio::test]
5761 async fn delegate_projects_task_kickoff_for_generated_helper() {
5762 let store = ConsoleEventStore::new();
5763 let dispatcher = console_wrapper(
5764 serde_json::json!({
5765 "agent_identity": "helper-3f2a",
5766 "member_ref": "opaque",
5767 "mob_id": "implicit-1",
5768 "wired": true
5769 }),
5770 false,
5771 Some(&store),
5772 );
5773
5774 dispatch_tool(
5775 &dispatcher,
5776 "delegate",
5777 serde_json::json!({ "task": "Review the diff" }),
5778 )
5779 .await;
5780
5781 let kickoffs = kickoff_events(&store).await;
5782 assert_eq!(kickoffs.len(), 1);
5783 assert_eq!(kickoffs[0].identity, "helper-3f2a");
5784 assert_eq!(kickoffs[0].data["content"][0]["text"], "Review the diff");
5785 assert_eq!(kickoffs[0].data["via_tool"], "delegate");
5786 }
5787
5788 #[tokio::test]
5789 async fn spawn_without_console_sink_leaves_outcome_unchanged() {
5790 let dispatcher = console_wrapper(
5791 serde_json::json!({ "agent_identity": "worker-3" }),
5792 false,
5793 None,
5794 );
5795
5796 let outcome = dispatch_tool(
5797 &dispatcher,
5798 "mob_spawn_member",
5799 serde_json::json!({
5800 "mob_id": "ob3",
5801 "profile": "person-worker",
5802 "member_id": "worker-3",
5803 "initial_message": "Find the person"
5804 }),
5805 )
5806 .await;
5807
5808 assert!(!outcome.result.is_error);
5809 assert!(
5810 outcome
5811 .result
5812 .text_content()
5813 .contains("\"agent_identity\":\"worker-3\""),
5814 "no console store → spawn outcome passes through untouched"
5815 );
5816 }
5817
5818 #[tokio::test]
5819 async fn failed_spawn_projects_nothing() {
5820 let store = ConsoleEventStore::new();
5821 let dispatcher = console_wrapper(
5822 serde_json::json!({ "error": "spawn rejected" }),
5823 true,
5824 Some(&store),
5825 );
5826
5827 dispatch_tool(
5828 &dispatcher,
5829 "mob_spawn_member",
5830 serde_json::json!({
5831 "mob_id": "ob3",
5832 "profile": "person-worker",
5833 "member_id": "worker-3",
5834 "initial_message": "Find the person"
5835 }),
5836 )
5837 .await;
5838
5839 assert!(
5840 kickoff_events(&store).await.is_empty(),
5841 "failed spawns must not seed console chats"
5842 );
5843 assert!(store.identity_labels("worker-3").await.is_none());
5844 }
5845 }
5846}