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