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