#[cfg(not(feature = "memory-store"))]
use async_trait::async_trait;
#[cfg(not(feature = "memory-store"))]
use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::Arc;
use meerkat_client::{
DefaultClientFactory, DefaultFactoryConfig, FactoryError, LlmClient, LlmClientAdapter,
LlmClientFactory, LlmProvider, ProviderResolver,
};
use meerkat_core::ops_lifecycle::OpsLifecycleRegistry;
use meerkat_core::service::{CreateSessionRequest, SessionBuildOptions};
#[cfg(target_arch = "wasm32")]
const DEFAULT_WASM_SYSTEM_PROMPT: &str = r"You are an autonomous agent. Your task is to accomplish the user's goal by systematically using the tools available to you.
# Core Behavior
- Break complex tasks into steps and execute them one by one.
- Use tools to gather information, take actions, and verify results.
- When multiple tool calls are independent, execute them in parallel.
- If a tool call fails, analyze the error and try alternative approaches.
- Continue working until the task is complete or you determine it cannot be completed.
# Decision Making
- Act on the information you have. Make reasonable assumptions when necessary.
- If critical information is missing and no tool can provide it, state what you need and why.
- Prioritize correctness over speed. Verify your work when possible.
# Output
- When the task is complete, provide a clear summary of what was accomplished.
- If the task cannot be completed, explain what blocked progress and what was attempted.";
use meerkat_core::{
Agent, AgentBuilder, AgentEvent, AgentLlmClient, AgentSessionStore, AgentToolDispatcher,
BlobStore, BudgetLimits, Config, HookRunOverrides, OutputSchema, Provider, Session,
SessionMetadata, SessionTooling, ToolCategoryOverride,
};
#[cfg(not(feature = "memory-store"))]
use meerkat_core::{SessionId, SessionMeta};
use meerkat_runtime::RuntimeOpsLifecycleRegistry;
#[cfg(feature = "jsonl-store")]
use meerkat_store::JsonlStore;
#[cfg(all(feature = "memory-store", not(feature = "jsonl-store")))]
use meerkat_store::MemoryStore;
#[cfg(not(feature = "memory-store"))]
use meerkat_store::SessionFilter;
use meerkat_store::{SessionStore, StoreAdapter};
#[cfg(not(target_arch = "wasm32"))]
use meerkat_tools::BuiltinDispatcherConfig;
use meerkat_tools::CompositeDispatcherError;
use meerkat_tools::EmptyToolDispatcher;
#[cfg(all(not(feature = "session-store"), not(target_arch = "wasm32")))]
use meerkat_tools::builtin::FileTaskStore;
#[cfg(all(not(feature = "session-store"), not(target_arch = "wasm32")))]
use meerkat_tools::builtin::MemoryTaskStore;
#[cfg(all(feature = "session-store", not(target_arch = "wasm32")))]
use meerkat_tools::builtin::SqliteTaskStore;
#[cfg(not(target_arch = "wasm32"))]
use meerkat_tools::builtin::shell::ShellConfig;
#[cfg(not(target_arch = "wasm32"))]
use meerkat_tools::builtin::{BuiltinToolConfig, CompositeDispatcher, TaskStore, ToolPolicyLayer};
#[cfg(all(not(feature = "memory-store"), not(target_arch = "wasm32")))]
use tokio::sync::RwLock;
#[cfg(not(target_arch = "wasm32"))]
use tokio::sync::mpsc;
#[cfg(all(not(feature = "memory-store"), target_arch = "wasm32"))]
use tokio_with_wasm::alias::sync::RwLock;
#[cfg(target_arch = "wasm32")]
use tokio_with_wasm::alias::sync::mpsc;
#[cfg(feature = "comms")]
use crate::compose_tools_with_comms;
#[cfg(not(target_arch = "wasm32"))]
use crate::{create_default_hook_engine, resolve_layered_hooks_config};
#[cfg(not(feature = "memory-store"))]
#[derive(Default)]
#[allow(dead_code)]
struct EphemeralSessionStore {
sessions: RwLock<HashMap<SessionId, Session>>,
}
#[cfg(not(feature = "memory-store"))]
impl EphemeralSessionStore {
fn new() -> Self {
Self {
sessions: RwLock::new(HashMap::new()),
}
}
}
#[cfg(not(feature = "memory-store"))]
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
impl SessionStore for EphemeralSessionStore {
async fn save(&self, session: &Session) -> Result<(), meerkat_store::SessionStoreError> {
self.sessions
.write()
.await
.insert(session.id().clone(), session.clone());
Ok(())
}
async fn load(
&self,
id: &SessionId,
) -> Result<Option<Session>, meerkat_store::SessionStoreError> {
Ok(self.sessions.read().await.get(id).cloned())
}
async fn list(
&self,
filter: SessionFilter,
) -> Result<Vec<SessionMeta>, meerkat_store::SessionStoreError> {
let mut metas: Vec<SessionMeta> = self
.sessions
.read()
.await
.values()
.map(SessionMeta::from)
.collect();
metas.sort_by(|a, b| b.updated_at.cmp(&a.updated_at));
if let Some(created_after) = filter.created_after {
metas.retain(|m| m.created_at >= created_after);
}
if let Some(updated_after) = filter.updated_after {
metas.retain(|m| m.updated_at >= updated_after);
}
if let Some(offset) = filter.offset {
metas = metas.into_iter().skip(offset).collect();
}
if let Some(limit) = filter.limit {
metas.truncate(limit);
}
Ok(metas)
}
async fn delete(&self, id: &SessionId) -> Result<(), meerkat_store::SessionStoreError> {
self.sessions.write().await.remove(id);
Ok(())
}
}
pub type DynAgent = Agent<dyn AgentLlmClient, dyn AgentToolDispatcher, dyn AgentSessionStore>;
#[derive(Clone)]
struct ErasedLlmClientOverride(Arc<dyn LlmClient>);
pub fn encode_llm_client_override_for_service(
client: Arc<dyn LlmClient>,
) -> Arc<dyn std::any::Any + Send + Sync> {
Arc::new(ErasedLlmClientOverride(client))
}
pub fn decode_llm_client_override_from_service(
value: &Arc<dyn std::any::Any + Send + Sync>,
) -> Option<Arc<dyn LlmClient>> {
if let Some(typed) = value.as_ref().downcast_ref::<ErasedLlmClientOverride>() {
return Some(typed.0.clone());
}
value.as_ref().downcast_ref::<Arc<dyn LlmClient>>().cloned()
}
pub struct AgentBuildConfig {
pub model: String,
pub provider: Option<Provider>,
pub max_tokens: Option<u32>,
pub system_prompt: Option<String>,
pub output_schema: Option<OutputSchema>,
pub structured_output_retries: u32,
pub hooks_override: HookRunOverrides,
pub keep_alive: bool,
pub comms_name: Option<String>,
pub peer_meta: Option<meerkat_core::PeerMeta>,
pub resume_session: Option<Session>,
pub budget_limits: Option<BudgetLimits>,
pub event_tx: Option<mpsc::Sender<AgentEvent>>,
pub llm_client_override: Option<Arc<dyn LlmClient>>,
pub provider_params: Option<serde_json::Value>,
pub external_tools: Option<Arc<dyn AgentToolDispatcher>>,
pub recoverable_tool_defs: Option<Vec<meerkat_core::ToolDef>>,
pub blob_store_override: Option<Arc<dyn BlobStore>>,
pub override_builtins: ToolCategoryOverride,
pub override_shell: ToolCategoryOverride,
pub override_memory: ToolCategoryOverride,
pub override_mob: ToolCategoryOverride,
pub mob_tool_authority_context: Option<meerkat_core::service::MobToolAuthorityContext>,
pub mob_tools: Option<Arc<dyn meerkat_core::service::MobToolsFactory>>,
pub preload_skills: Option<Vec<meerkat_core::skills::SkillId>>,
pub realm_id: Option<String>,
pub instance_id: Option<String>,
pub backend: Option<String>,
pub config_generation: Option<u64>,
pub silent_comms_intents: Vec<String>,
pub max_inline_peer_notifications: Option<i32>,
pub tool_dispatcher_override: Option<Arc<dyn AgentToolDispatcher>>,
pub session_store_override: Option<Arc<dyn AgentSessionStore>>,
pub hook_engine_override: Option<Arc<dyn meerkat_core::HookEngine>>,
pub skill_engine_override: Option<Arc<meerkat_core::skills::SkillRuntime>>,
pub app_context: Option<serde_json::Value>,
pub additional_instructions: Option<Vec<String>>,
pub wait_for_mcp: bool,
pub shell_env: Option<std::collections::HashMap<String, String>>,
pub checkpointer: Option<Arc<dyn meerkat_core::checkpoint::SessionCheckpointer>>,
pub call_timeout_override: meerkat_core::CallTimeoutOverride,
pub resume_override_mask: meerkat_core::service::ResumeOverrideMask,
pub runtime_build_mode: meerkat_core::RuntimeBuildMode,
}
impl std::fmt::Debug for AgentBuildConfig {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("AgentBuildConfig")
.field("model", &self.model)
.field("provider", &self.provider)
.field("max_tokens", &self.max_tokens)
.field(
"system_prompt",
&self
.system_prompt
.as_deref()
.map(|s| if s.len() > 64 { &s[..64] } else { s }),
)
.field("output_schema", &self.output_schema.is_some())
.field("structured_output_retries", &self.structured_output_retries)
.field("keep_alive", &self.keep_alive)
.field("resume_override_mask", &self.resume_override_mask)
.field("comms_name", &self.comms_name)
.field("peer_meta", &self.peer_meta)
.field("resume_session", &self.resume_session.is_some())
.field("budget_limits", &self.budget_limits)
.field("event_tx", &self.event_tx.is_some())
.field("llm_client_override", &self.llm_client_override.is_some())
.field("provider_params", &self.provider_params.is_some())
.field("external_tools", &self.external_tools.is_some())
.field("recoverable_tool_defs", &self.recoverable_tool_defs)
.field("blob_store_override", &self.blob_store_override.is_some())
.field("override_builtins", &self.override_builtins)
.field("override_shell", &self.override_shell)
.field("override_memory", &self.override_memory)
.field("override_mob", &self.override_mob)
.field(
"mob_tool_authority_context",
&self.mob_tool_authority_context.is_some(),
)
.field("mob_tools", &self.mob_tools.is_some())
.field("realm_id", &self.realm_id)
.field("instance_id", &self.instance_id)
.field("backend", &self.backend)
.field("config_generation", &self.config_generation)
.field(
"max_inline_peer_notifications",
&self.max_inline_peer_notifications,
)
.field(
"tool_dispatcher_override",
&self.tool_dispatcher_override.is_some(),
)
.field(
"session_store_override",
&self.session_store_override.is_some(),
)
.field("hook_engine_override", &self.hook_engine_override.is_some())
.field(
"skill_engine_override",
&self.skill_engine_override.is_some(),
)
.field("app_context", &self.app_context.is_some())
.field("additional_instructions", &self.additional_instructions)
.field("wait_for_mcp", &self.wait_for_mcp)
.field("runtime_build_mode", &self.runtime_build_mode)
.finish()
}
}
impl AgentBuildConfig {
pub fn new(model: impl Into<String>) -> Self {
Self {
model: model.into(),
provider: None,
max_tokens: None,
system_prompt: None,
output_schema: None,
structured_output_retries: 2,
hooks_override: HookRunOverrides::default(),
keep_alive: false,
comms_name: None,
peer_meta: None,
resume_session: None,
budget_limits: None,
event_tx: None,
llm_client_override: None,
provider_params: None,
external_tools: None,
recoverable_tool_defs: None,
blob_store_override: None,
override_builtins: ToolCategoryOverride::Inherit,
override_shell: ToolCategoryOverride::Inherit,
override_memory: ToolCategoryOverride::Inherit,
override_mob: ToolCategoryOverride::Inherit,
mob_tool_authority_context: None,
mob_tools: None,
preload_skills: None,
realm_id: None,
instance_id: None,
backend: None,
config_generation: None,
silent_comms_intents: Vec::new(),
max_inline_peer_notifications: None,
tool_dispatcher_override: None,
session_store_override: None,
hook_engine_override: None,
skill_engine_override: None,
app_context: None,
additional_instructions: None,
wait_for_mcp: false,
shell_env: None,
checkpointer: None,
call_timeout_override: meerkat_core::CallTimeoutOverride::default(),
resume_override_mask: meerkat_core::service::ResumeOverrideMask::default(),
runtime_build_mode: meerkat_core::RuntimeBuildMode::StandaloneEphemeral,
}
}
pub fn from_create_session_request(
req: &CreateSessionRequest,
event_tx: mpsc::Sender<AgentEvent>,
) -> Self {
let mut build = Self::new(req.model.clone());
build.system_prompt = req.system_prompt.clone();
build.max_tokens = req.max_tokens;
if let Some(options) = &req.build {
build.apply_session_build_options(options);
}
build.event_tx = Some(event_tx);
build
}
pub fn apply_persisted_mob_operator_access(
&mut self,
enable_mob: ToolCategoryOverride,
persisted_authority_context: Option<meerkat_core::service::MobToolAuthorityContext>,
) {
let (override_mob, authority_context) = meerkat_core::service::resolve_mob_operator_access(
enable_mob,
persisted_authority_context,
);
self.override_mob = override_mob;
self.mob_tool_authority_context = authority_context;
}
pub fn apply_generated_create_only_mob_operator_access(
&mut self,
enable_mob: ToolCategoryOverride,
) {
self.apply_persisted_mob_operator_access(enable_mob, None);
}
pub fn apply_session_build_options(&mut self, build: &SessionBuildOptions) {
self.provider = build.provider;
self.output_schema = build.output_schema.clone();
self.structured_output_retries = build.structured_output_retries;
self.hooks_override = build.hooks_override.clone();
self.comms_name = build.comms_name.clone();
self.peer_meta = build.peer_meta.clone();
self.resume_session = build.resume_session.clone();
self.budget_limits = build.budget_limits.clone();
self.provider_params = build.provider_params.clone();
self.external_tools = build.external_tools.clone();
self.recoverable_tool_defs = build.recoverable_tool_defs.clone();
self.blob_store_override = build.blob_store_override.clone();
self.llm_client_override = build
.llm_client_override
.as_ref()
.and_then(decode_llm_client_override_from_service);
self.override_builtins = build.override_builtins;
self.override_shell = build.override_shell;
self.override_memory = build.override_memory;
self.override_mob = build.override_mob;
self.mob_tool_authority_context = build.mob_tool_authority_context.clone();
self.mob_tools = build.mob_tools.clone();
self.preload_skills = build.preload_skills.clone();
self.realm_id = build.realm_id.clone();
self.instance_id = build.instance_id.clone();
self.backend = build.backend.clone();
self.config_generation = build.config_generation;
self.keep_alive = build.keep_alive;
self.silent_comms_intents
.clone_from(&build.silent_comms_intents);
self.max_inline_peer_notifications = build.max_inline_peer_notifications;
self.app_context = build.app_context.clone();
self.additional_instructions = build.additional_instructions.clone();
self.shell_env = build.shell_env.clone();
self.checkpointer = build.checkpointer.clone();
self.call_timeout_override = build.call_timeout_override.clone();
self.resume_override_mask = build.resume_override_mask;
self.runtime_build_mode = build.runtime_build_mode.clone();
}
pub fn to_session_build_options(&self) -> SessionBuildOptions {
SessionBuildOptions {
provider: self.provider,
output_schema: self.output_schema.clone(),
structured_output_retries: self.structured_output_retries,
hooks_override: self.hooks_override.clone(),
comms_name: self.comms_name.clone(),
peer_meta: self.peer_meta.clone(),
resume_session: self.resume_session.clone(),
budget_limits: self.budget_limits.clone(),
provider_params: self.provider_params.clone(),
external_tools: self.external_tools.clone(),
recoverable_tool_defs: self.recoverable_tool_defs.clone(),
blob_store_override: self.blob_store_override.clone(),
llm_client_override: self
.llm_client_override
.clone()
.map(encode_llm_client_override_for_service),
override_builtins: self.override_builtins,
override_shell: self.override_shell,
override_memory: self.override_memory,
override_mob: self.override_mob,
mob_tool_authority_context: self.mob_tool_authority_context.clone(),
mob_tools: self.mob_tools.clone(),
preload_skills: self.preload_skills.clone(),
realm_id: self.realm_id.clone(),
instance_id: self.instance_id.clone(),
backend: self.backend.clone(),
config_generation: self.config_generation,
keep_alive: self.keep_alive,
silent_comms_intents: self.silent_comms_intents.clone(),
max_inline_peer_notifications: self.max_inline_peer_notifications,
app_context: self.app_context.clone(),
additional_instructions: self.additional_instructions.clone(),
shell_env: self.shell_env.clone(),
checkpointer: self.checkpointer.clone(),
call_timeout_override: self.call_timeout_override.clone(),
resume_override_mask: self.resume_override_mask,
runtime_build_mode: self.runtime_build_mode.clone(),
}
}
}
#[derive(Debug, thiserror::Error)]
pub enum BuildAgentError {
#[error("Cannot infer provider from model '{model}'")]
UnknownProvider { model: String },
#[error("API key not set for provider '{provider}'")]
MissingApiKey { provider: String },
#[error("LLM client creation failed: {0}")]
LlmClient(#[from] FactoryError),
#[error("Tool dispatcher creation failed: {0}")]
ToolDispatcher(#[from] CompositeDispatcherError),
#[error("Comms runtime failed: {0}")]
#[cfg(feature = "comms")]
Comms(String),
#[error("Config error: {0}")]
Config(String),
#[error("keep_alive requires comms_name to be set")]
#[cfg(feature = "comms")]
KeepAliveRequiresCommsName,
}
struct ProfileBasedDefaultsResolver;
impl meerkat_core::ModelOperationalDefaultsResolver for ProfileBasedDefaultsResolver {
fn call_timeout_for(&self, provider: &str, model: &str) -> Option<std::time::Duration> {
meerkat_models::profile::profile_for(provider, model)
.and_then(|p| p.call_timeout_secs)
.map(std::time::Duration::from_secs)
}
}
pub fn provider_key(provider: Provider) -> &'static str {
provider.as_str()
}
#[derive(Clone)]
pub struct AgentFactory {
pub store_path: PathBuf,
pub runtime_root: Option<PathBuf>,
pub project_root: Option<PathBuf>,
pub context_root: Option<PathBuf>,
pub user_config_root: Option<PathBuf>,
pub enable_builtins: bool,
pub enable_shell: bool,
#[cfg(feature = "comms")]
pub enable_comms: bool,
pub enable_memory: bool,
pub enable_mob: bool,
#[cfg(feature = "skills")]
pub skill_source: Option<Arc<meerkat_skills::CompositeSkillSource>>,
custom_store: Option<Arc<dyn SessionStore>>,
pub mob_tools: Option<Arc<dyn meerkat_core::service::MobToolsFactory>>,
#[cfg(feature = "comms")]
pub comms_runtime: Option<Arc<meerkat_comms::CommsRuntime>>,
}
impl std::fmt::Debug for AgentFactory {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let mut d = f.debug_struct("AgentFactory");
d.field("store_path", &self.store_path)
.field("runtime_root", &self.runtime_root)
.field("project_root", &self.project_root)
.field("context_root", &self.context_root)
.field("user_config_root", &self.user_config_root)
.field("enable_builtins", &self.enable_builtins)
.field("enable_shell", &self.enable_shell)
.field("enable_memory", &self.enable_memory)
.field("enable_mob", &self.enable_mob);
#[cfg(feature = "comms")]
d.field("enable_comms", &self.enable_comms);
#[cfg(feature = "skills")]
d.field("skill_source", &self.skill_source.as_ref().map(|_| ".."));
d.field("custom_store", &self.custom_store.as_ref().map(|_| ".."));
d.field("mob_tools", &self.mob_tools.is_some());
#[cfg(feature = "comms")]
d.field("comms_runtime", &self.comms_runtime.is_some());
d.finish()
}
}
impl AgentFactory {
pub fn minimal() -> Self {
Self {
store_path: PathBuf::new(),
runtime_root: None,
project_root: None,
context_root: None,
user_config_root: None,
enable_builtins: false,
enable_shell: false,
#[cfg(feature = "comms")]
enable_comms: false,
enable_memory: false,
enable_mob: false,
#[cfg(feature = "skills")]
skill_source: None,
custom_store: None,
mob_tools: None,
#[cfg(feature = "comms")]
comms_runtime: None,
}
}
pub fn new(store_path: impl Into<PathBuf>) -> Self {
Self {
store_path: store_path.into(),
runtime_root: None,
project_root: None,
context_root: None,
user_config_root: None,
enable_builtins: false,
enable_shell: false,
#[cfg(feature = "comms")]
enable_comms: false,
enable_memory: false,
enable_mob: false,
#[cfg(feature = "skills")]
skill_source: None,
custom_store: None,
mob_tools: None,
#[cfg(feature = "comms")]
comms_runtime: None,
}
}
#[cfg(feature = "skills")]
pub fn skill_source(mut self, source: Arc<meerkat_skills::CompositeSkillSource>) -> Self {
self.skill_source = Some(source);
self
}
pub fn project_root(mut self, path: impl Into<PathBuf>) -> Self {
self.project_root = Some(path.into());
self
}
pub fn context_root(mut self, path: impl Into<PathBuf>) -> Self {
self.context_root = Some(path.into());
self
}
pub fn user_config_root(mut self, path: impl Into<PathBuf>) -> Self {
self.user_config_root = Some(path.into());
self
}
pub fn runtime_root(mut self, path: impl Into<PathBuf>) -> Self {
self.runtime_root = Some(path.into());
self
}
pub fn builtins(mut self, enabled: bool) -> Self {
self.enable_builtins = enabled;
self
}
pub fn shell(mut self, enabled: bool) -> Self {
self.enable_shell = enabled;
self
}
pub fn memory(mut self, enabled: bool) -> Self {
self.enable_memory = enabled;
self
}
pub fn mob(mut self, enabled: bool) -> Self {
self.enable_mob = enabled;
self
}
pub fn mob_tools_factory(
mut self,
factory: Arc<dyn meerkat_core::service::MobToolsFactory>,
) -> Self {
self.mob_tools = Some(factory);
self
}
#[cfg(feature = "comms")]
pub fn comms(mut self, enabled: bool) -> Self {
self.enable_comms = enabled;
self
}
#[cfg(feature = "comms")]
pub fn with_comms_runtime(mut self, runtime: Arc<meerkat_comms::CommsRuntime>) -> Self {
self.comms_runtime = Some(runtime);
self
}
#[cfg(feature = "skills")]
pub async fn build_skill_runtime(
&self,
config: &Config,
) -> Option<Arc<meerkat_core::skills::SkillRuntime>> {
let skill_source: Option<Arc<meerkat_skills::CompositeSkillSource>> =
if self.skill_source.is_some() {
self.skill_source.clone()
} else if !config.skills.enabled {
None
} else {
#[cfg(not(target_arch = "wasm32"))]
{
let conventions_context_root = self
.context_root
.as_deref()
.or(self.project_root.as_deref());
let conventions_user_root = self.user_config_root.as_deref();
let runtime_root = self
.runtime_root
.clone()
.or_else(|| self.project_root.clone())
.unwrap_or_else(|| self.store_path.clone());
match meerkat_skills::resolve_repositories_with_roots(
&config.skills,
conventions_context_root,
conventions_user_root,
Some(runtime_root.as_path()),
)
.await
{
Ok(source) => source.map(Arc::new),
Err(e) => {
tracing::warn!("Failed to resolve skill repositories: {e}");
None
}
}
}
#[cfg(target_arch = "wasm32")]
None
};
skill_source.map(|source| {
let available_caps: Vec<String> = meerkat_contracts::build_capabilities()
.into_iter()
.map(|c| c.id.to_string())
.collect();
let engine = Arc::new(
meerkat_skills::DefaultSkillEngine::new(source, available_caps)
.with_inventory_threshold(config.skills.inventory_threshold)
.with_max_injection_bytes(config.skills.max_injection_bytes),
);
Arc::new(meerkat_core::skills::SkillRuntime::new(engine))
})
}
pub fn session_store(mut self, store: Arc<dyn SessionStore>) -> Self {
self.custom_store = Some(store);
self
}
fn realm_scope_root(&self, _build_config: &AgentBuildConfig) -> PathBuf {
self.runtime_root
.clone()
.or_else(|| self.project_root.clone())
.unwrap_or_else(|| self.store_path.clone())
}
fn apply_resumed_session_metadata(
build_config: &mut AgentBuildConfig,
) -> Option<SessionMetadata> {
let metadata = build_config
.resume_session
.as_ref()
.and_then(Session::session_metadata)?;
let mask = build_config.resume_override_mask;
if !mask.model {
build_config.model = metadata.model.clone();
}
if !mask.max_tokens {
build_config.max_tokens = Some(metadata.max_tokens);
}
if !mask.structured_output_retries {
build_config.structured_output_retries = metadata.structured_output_retries;
}
if !mask.provider {
build_config.provider = Some(metadata.provider);
}
if !mask.provider_params {
build_config.provider_params = metadata.provider_params.clone();
}
if !mask.override_builtins {
build_config.override_builtins = metadata.tooling.builtins;
}
if !mask.override_shell {
build_config.override_shell = metadata.tooling.shell;
}
if !mask.override_memory {
build_config.override_memory = metadata.tooling.memory;
}
if !mask.override_mob {
build_config.override_mob = metadata.tooling.mob;
build_config.mob_tool_authority_context = build_config
.resume_session
.as_ref()
.and_then(Session::mob_tool_authority_context);
}
if !mask.preload_skills {
build_config.preload_skills = metadata.tooling.active_skills.clone();
}
if !mask.keep_alive {
build_config.keep_alive = metadata.keep_alive;
}
if !mask.comms_name {
build_config.comms_name = metadata.comms_name.clone();
}
if !mask.peer_meta {
build_config.peer_meta = metadata.peer_meta.clone();
}
Some(metadata)
}
pub async fn build_llm_adapter(
&self,
client: Arc<dyn LlmClient>,
model: impl Into<String>,
) -> LlmClientAdapter {
LlmClientAdapter::new(client, model.into())
}
pub async fn build_llm_adapter_with_events(
&self,
client: Arc<dyn LlmClient>,
model: impl Into<String>,
event_tx: Option<mpsc::Sender<AgentEvent>>,
) -> LlmClientAdapter {
match event_tx {
Some(tx) => LlmClientAdapter::with_event_channel(client, model.into(), tx),
None => LlmClientAdapter::new(client, model.into()),
}
}
pub async fn build_llm_client(
&self,
provider: Provider,
api_key: Option<String>,
base_url: Option<String>,
) -> Result<Arc<dyn LlmClient>, FactoryError> {
let mapped = match provider {
Provider::Anthropic => LlmProvider::Anthropic,
Provider::OpenAI => LlmProvider::OpenAi,
Provider::Gemini => LlmProvider::Gemini,
Provider::Other => return Err(FactoryError::UnsupportedProvider("other".to_string())),
};
let mut config = DefaultFactoryConfig::default();
if let Some(url) = base_url {
match mapped {
LlmProvider::Anthropic => config = config.with_anthropic_base_url(url),
LlmProvider::OpenAi => config = config.with_openai_base_url(url),
LlmProvider::Gemini => config = config.with_gemini_base_url(url),
}
}
let factory = DefaultClientFactory::with_config(config);
factory.create_client(mapped, api_key)
}
fn resolve_provider_credentials(
&self,
provider: Provider,
config: &Config,
) -> (Option<String>, Option<String>) {
let mut base_url = config
.providers
.base_urls
.as_ref()
.and_then(|map| map.get(provider_key(provider)).cloned());
let mut api_key = config
.providers
.api_keys
.as_ref()
.and_then(|map| map.get(provider_key(provider)).cloned());
match (&config.provider, provider) {
(
meerkat_core::ProviderConfig::Anthropic {
api_key: cfg_key,
base_url: cfg_url,
},
Provider::Anthropic,
) => {
if cfg_key.is_some() {
api_key = cfg_key.clone();
}
if cfg_url.is_some() {
base_url = cfg_url.clone();
}
}
(
meerkat_core::ProviderConfig::OpenAI {
api_key: cfg_key,
base_url: cfg_url,
},
Provider::OpenAI,
) => {
if cfg_key.is_some() {
api_key = cfg_key.clone();
}
if cfg_url.is_some() {
base_url = cfg_url.clone();
}
}
(meerkat_core::ProviderConfig::Gemini { api_key: cfg_key }, Provider::Gemini) => {
if cfg_key.is_some() {
api_key = cfg_key.clone();
}
}
_ => {}
}
if api_key.is_none() {
api_key = ProviderResolver::api_key_for(provider);
}
(api_key, base_url)
}
pub async fn build_store_adapter<S: SessionStore + 'static>(
&self,
store: Arc<S>,
) -> StoreAdapter<S> {
StoreAdapter::new(store)
}
#[cfg(not(target_arch = "wasm32"))]
#[allow(clippy::too_many_arguments)]
pub async fn build_composite_dispatcher(
&self,
store: Arc<dyn TaskStore>,
config: &BuiltinToolConfig,
project_root: Option<PathBuf>,
shell_config: Option<ShellConfig>,
external: Option<Arc<dyn AgentToolDispatcher>>,
session_id: Option<String>,
ops_lifecycle: Option<Arc<dyn OpsLifecycleRegistry>>,
image_tool_results: bool,
) -> Result<CompositeDispatcher, CompositeDispatcherError> {
CompositeDispatcher::new_with_ops_lifecycle(
store,
config,
project_root,
shell_config,
external,
session_id,
ops_lifecycle,
image_tool_results,
)
}
#[cfg(not(target_arch = "wasm32"))]
#[allow(clippy::too_many_arguments)]
pub async fn build_builtin_dispatcher(
&self,
store: Arc<dyn TaskStore>,
config: BuiltinToolConfig,
project_root: Option<PathBuf>,
shell_config: Option<ShellConfig>,
external: Option<Arc<dyn AgentToolDispatcher>>,
session_id: Option<String>,
ops_lifecycle: Option<Arc<dyn OpsLifecycleRegistry>>,
) -> Result<Arc<dyn AgentToolDispatcher>, CompositeDispatcherError> {
self.build_builtin_dispatcher_with_skills(
store,
config,
project_root,
shell_config,
external,
session_id,
ops_lifecycle,
None,
)
.await
}
#[cfg(not(target_arch = "wasm32"))]
#[allow(clippy::too_many_arguments)]
pub async fn build_builtin_dispatcher_with_skills(
&self,
store: Arc<dyn TaskStore>,
config: BuiltinToolConfig,
project_root: Option<PathBuf>,
shell_config: Option<ShellConfig>,
external: Option<Arc<dyn AgentToolDispatcher>>,
session_id: Option<String>,
ops_lifecycle: Option<Arc<dyn OpsLifecycleRegistry>>,
#[cfg_attr(not(feature = "skills"), allow(unused_variables))] skill_engine: Option<
Arc<meerkat_core::skills::SkillRuntime>,
>,
) -> Result<Arc<dyn AgentToolDispatcher>, CompositeDispatcherError> {
self.build_builtin_dispatcher_with_skills_internal(
store,
config,
project_root,
shell_config,
external,
session_id,
ops_lifecycle,
skill_engine,
true,
)
.await
}
#[cfg(not(target_arch = "wasm32"))]
#[allow(clippy::too_many_arguments)]
async fn build_builtin_dispatcher_with_skills_internal(
&self,
store: Arc<dyn TaskStore>,
config: BuiltinToolConfig,
project_root: Option<PathBuf>,
shell_config: Option<ShellConfig>,
external: Option<Arc<dyn AgentToolDispatcher>>,
session_id: Option<String>,
ops_lifecycle: Option<Arc<dyn OpsLifecycleRegistry>>,
#[cfg_attr(not(feature = "skills"), allow(unused_variables))] skill_engine: Option<
Arc<meerkat_core::skills::SkillRuntime>,
>,
image_tool_results: bool,
) -> Result<Arc<dyn AgentToolDispatcher>, CompositeDispatcherError> {
let BuiltinDispatcherConfig {
store,
config,
project_root,
shell_config,
external,
session_id,
ops_lifecycle,
image_tool_results,
} = BuiltinDispatcherConfig {
store,
config,
project_root,
shell_config,
external,
session_id,
ops_lifecycle,
image_tool_results,
};
#[cfg_attr(not(feature = "skills"), allow(unused_mut))]
let mut composite = self
.build_composite_dispatcher(
store,
&config,
project_root,
shell_config,
external,
session_id,
ops_lifecycle,
image_tool_results,
)
.await?;
#[cfg(feature = "skills")]
if let Some(engine) = skill_engine {
composite
.register_skill_tools(meerkat_tools::builtin::skills::SkillToolSet::new(engine));
}
Ok(Arc::new(composite))
}
pub async fn build_agent(
&self,
mut build_config: AgentBuildConfig,
config: &Config,
) -> Result<DynAgent, BuildAgentError> {
build_config.resume_override_mask.override_builtins |= !matches!(
build_config.override_builtins,
ToolCategoryOverride::Inherit
);
build_config.resume_override_mask.override_shell |=
!matches!(build_config.override_shell, ToolCategoryOverride::Inherit);
build_config.resume_override_mask.override_memory |=
!matches!(build_config.override_memory, ToolCategoryOverride::Inherit);
build_config.resume_override_mask.override_mob |=
!matches!(build_config.override_mob, ToolCategoryOverride::Inherit);
let explicit_mob_override =
!matches!(build_config.override_mob, ToolCategoryOverride::Inherit);
let resumed_session_metadata = Self::apply_resumed_session_metadata(&mut build_config);
if build_config.mob_tool_authority_context.is_none()
&& matches!(build_config.override_mob, ToolCategoryOverride::Enable)
&& (build_config.resume_session.is_none() || explicit_mob_override)
{
build_config
.apply_generated_create_only_mob_operator_access(ToolCategoryOverride::Enable);
}
if let Some(value) = build_config.max_inline_peer_notifications
&& value < -1
{
return Err(BuildAgentError::Config(format!(
"max_inline_peer_notifications={value} is invalid (allowed: -1, 0, or >0)"
)));
}
#[cfg(feature = "comms")]
if build_config.keep_alive && build_config.comms_name.is_none() {
return Err(BuildAgentError::KeepAliveRequiresCommsName);
}
let provider = match build_config.provider {
Some(p) => p,
None => {
let inferred = ProviderResolver::infer_from_model(&build_config.model);
if inferred != Provider::Other {
inferred
} else if let Some(client) = build_config.llm_client_override.as_ref() {
Provider::from_name(client.provider())
} else {
return Err(BuildAgentError::UnknownProvider {
model: build_config.model.clone(),
});
}
}
};
let llm_client: Arc<dyn LlmClient> = match build_config.llm_client_override.as_ref() {
Some(client) => Arc::clone(client),
None => {
let (api_key, base_url) = self.resolve_provider_credentials(provider, config);
if api_key.is_none() {
return Err(BuildAgentError::MissingApiKey {
provider: provider_key(provider).to_string(),
});
}
self.build_llm_client(provider, api_key, base_url)
.await
.map_err(BuildAgentError::LlmClient)?
}
};
let model = build_config.model.clone();
let event_tap = meerkat_core::new_event_tap();
let mut llm_adapter_inner = match build_config.event_tx.clone() {
Some(tx) => LlmClientAdapter::with_event_channel(llm_client, model.clone(), tx),
None => LlmClientAdapter::new(llm_client, model.clone()),
};
llm_adapter_inner = llm_adapter_inner.with_event_tap(event_tap.clone());
if let Some(params) = build_config.provider_params.clone() {
llm_adapter_inner = llm_adapter_inner.with_provider_params(Some(params));
}
let llm_adapter: Arc<dyn AgentLlmClient> = Arc::new(llm_adapter_inner);
let max_tokens = build_config.max_tokens.unwrap_or(config.max_tokens);
let _realm_scope_root = self.realm_scope_root(&build_config);
let _conventions_context_root = self
.context_root
.as_deref()
.or(self.project_root.as_deref());
let _conventions_user_root = self.user_config_root.as_deref();
#[cfg(feature = "skills")]
let skill_engine: Option<Arc<meerkat_core::skills::SkillRuntime>> =
if let Some(engine) = build_config.skill_engine_override.take() {
Some(engine)
} else {
let skill_source: Option<Arc<meerkat_skills::CompositeSkillSource>> =
if self.skill_source.is_some() {
self.skill_source.clone()
} else if !config.skills.enabled {
None
} else {
#[cfg(not(target_arch = "wasm32"))]
{
match meerkat_skills::resolve_repositories_with_roots(
&config.skills,
_conventions_context_root,
_conventions_user_root,
Some(_realm_scope_root.as_path()),
)
.await
{
Ok(source) => source.map(Arc::new),
Err(e) => {
tracing::warn!("Failed to resolve skill repositories: {e}");
None
}
}
}
#[cfg(target_arch = "wasm32")]
None
};
skill_source.map(|source| {
let available_caps: Vec<String> = meerkat_contracts::build_capabilities()
.into_iter()
.map(|c| c.id.to_string())
.collect();
let engine = Arc::new(
meerkat_skills::DefaultSkillEngine::new(source, available_caps)
.with_inventory_threshold(config.skills.inventory_threshold)
.with_max_injection_bytes(config.skills.max_injection_bytes),
);
Arc::new(meerkat_core::skills::SkillRuntime::new(engine))
})
}; #[cfg(not(feature = "skills"))]
let skill_engine: Option<Arc<meerkat_core::skills::SkillRuntime>> = None;
let persisted_system_prompt = build_config.system_prompt.clone();
let per_request_prompt = build_config.system_prompt.take();
let effective_builtins = build_config.override_builtins.resolve(self.enable_builtins);
#[allow(unused_variables)] let effective_shell = build_config.override_shell.resolve(self.enable_shell);
let session = build_config.resume_session.clone().unwrap_or_default();
let _session_id = session.id().to_string();
#[cfg(all(feature = "comms", not(target_arch = "wasm32")))]
let comms_runtime = if let Some(ref shared) = self.comms_runtime
&& build_config.comms_name.is_none()
{
Some(Arc::clone(shared))
} else if build_config.keep_alive || build_config.comms_name.is_some() {
let comms_name = build_config
.comms_name
.as_ref()
.ok_or(BuildAgentError::KeepAliveRequiresCommsName)?;
let silent_intents = Arc::new(
build_config
.silent_comms_intents
.iter()
.cloned()
.collect::<std::collections::HashSet<String>>(),
);
let mut runtime =
crate::build_session_scoped_comms_runtime_from_config_scoped_with_silent_intents(
config,
_realm_scope_root.as_path(),
self.user_config_root.as_deref(),
comms_name,
build_config.peer_meta.clone(),
build_config.realm_id.clone(),
session.id(),
silent_intents,
)
.await
.map_err(BuildAgentError::Comms)?;
if let Some(blob_store) = build_config.blob_store_override.clone() {
runtime.set_blob_store(blob_store);
}
Some(Arc::new(runtime))
} else {
None
};
#[cfg(all(feature = "comms", target_arch = "wasm32"))]
let comms_runtime = if let Some(ref shared) = self.comms_runtime
&& build_config.comms_name.is_none()
{
Some(Arc::clone(shared))
} else if build_config.keep_alive || build_config.comms_name.is_some() {
let comms_name = build_config
.comms_name
.as_ref()
.ok_or(BuildAgentError::KeepAliveRequiresCommsName)?;
let silent_intents = Arc::new(
build_config
.silent_comms_intents
.iter()
.cloned()
.collect::<std::collections::HashSet<String>>(),
);
let mut runtime = meerkat_comms::CommsRuntime::inproc_only_with_silent_intents(
comms_name,
build_config.realm_id.clone(),
silent_intents,
)
.map_err(|e| BuildAgentError::Comms(e.to_string()))?;
if let Some(ref meta) = build_config.peer_meta {
runtime.set_peer_meta(meta.clone());
}
if let Some(blob_store) = build_config.blob_store_override.clone() {
runtime.set_blob_store(blob_store);
}
Some(Arc::new(runtime))
} else {
None
};
#[cfg(not(feature = "comms"))]
#[allow(clippy::no_effect_underscore_binding)]
let _comms_runtime: Option<()> = None;
let image_tool_results = meerkat_models::profile::profile_for(provider.as_str(), &model)
.is_none_or(|p| p.image_tool_results);
use meerkat_core::runtime_epoch::RuntimeBuildMode;
let resolved_mode = &build_config.runtime_build_mode;
#[allow(unused_variables)]
let (ops_lifecycle, concrete_ops_lifecycle): (
Arc<dyn OpsLifecycleRegistry>,
Option<Arc<RuntimeOpsLifecycleRegistry>>,
) = match &resolved_mode {
RuntimeBuildMode::SessionOwned(bindings) => {
if bindings.session_id != *session.id() {
return Err(BuildAgentError::Config(format!(
"SessionRuntimeBindings.session_id ({}) does not match session ({}); \
bindings may have been prepared for a different session",
bindings.session_id,
session.id(),
)));
}
(Arc::clone(&bindings.ops_lifecycle), None)
}
RuntimeBuildMode::StandaloneEphemeral => {
let concrete = Arc::new(RuntimeOpsLifecycleRegistry::new());
(
Arc::clone(&concrete) as Arc<dyn OpsLifecycleRegistry>,
Some(concrete),
)
}
};
let completion_feed = ops_lifecycle.completion_feed();
let interrupt_baseline: Option<Arc<std::sync::atomic::AtomicU64>> =
if completion_feed.is_some() {
Some(Arc::new(std::sync::atomic::AtomicU64::new(0)))
} else {
None
};
#[allow(unused_mut)]
let (mut tools, mut tool_usage_instructions) =
if let Some(dispatcher) = build_config.tool_dispatcher_override.take() {
let usage = render_tool_usage_instructions(dispatcher.tools().as_ref());
(dispatcher, usage)
} else {
#[cfg(not(target_arch = "wasm32"))]
{
self.build_tool_dispatcher_for_agent_with_overrides(
config,
build_config.external_tools,
effective_builtins,
effective_shell,
skill_engine.clone(),
build_config.shell_env.take(),
_session_id.clone(),
Arc::clone(&ops_lifecycle),
image_tool_results,
)
.await?
}
#[cfg(target_arch = "wasm32")]
{
let usage = String::new();
(
Arc::new(EmptyToolDispatcher) as Arc<dyn AgentToolDispatcher>,
usage,
)
}
};
tracing::debug!(
base_tool_count = tools.tools().len(),
effective_builtins,
effective_shell,
"tool composition: base dispatcher built"
);
let store_adapter: Arc<dyn AgentSessionStore> =
if let Some(store) = build_config.session_store_override.take() {
store
} else if let Some(store) = &self.custom_store {
Arc::new(StoreAdapter::new(Arc::clone(store)))
} else {
#[cfg(feature = "jsonl-store")]
{
let store = JsonlStore::new(self.store_path.clone());
store
.init()
.await
.map_err(|e| BuildAgentError::Config(format!("Store init failed: {e}")))?;
Arc::new(StoreAdapter::new(Arc::new(store)))
}
#[cfg(all(not(feature = "jsonl-store"), feature = "memory-store"))]
{
Arc::new(self.build_store_adapter(Arc::new(MemoryStore::new())).await)
}
#[cfg(all(not(feature = "jsonl-store"), not(feature = "memory-store")))]
{
Arc::new(
self.build_store_adapter(Arc::new(EphemeralSessionStore::new()))
.await,
)
}
};
#[cfg(feature = "comms")]
if let Some(ref runtime) = comms_runtime {
let composed =
compose_tools_with_comms(tools, tool_usage_instructions, runtime.tool_material())
.map_err(|e| {
BuildAgentError::Config(format!("Failed to compose comms tools: {e}"))
})?;
tools = composed.0;
tool_usage_instructions = composed.1;
}
tracing::debug!(
tool_count_after_comms = tools.tools().len(),
"tool composition: after comms gateway"
);
let effective_mob = build_config.override_mob.resolve(self.enable_mob)
|| build_config.mob_tool_authority_context.is_some();
let mob_factory = build_config
.mob_tools
.take()
.or_else(|| self.mob_tools.clone());
if effective_mob && let Some(mob_factory) = mob_factory {
#[cfg(feature = "comms")]
let mob_comms: Option<Arc<dyn meerkat_core::agent::CommsRuntime>> = comms_runtime
.as_ref()
.map(|r| Arc::clone(r) as Arc<dyn meerkat_core::agent::CommsRuntime>);
#[cfg(not(feature = "comms"))]
let mob_comms: Option<Arc<dyn meerkat_core::agent::CommsRuntime>> = None;
let mob_args = meerkat_core::service::MobToolsBuildArgs {
session_id: session.id().clone(),
model: model.clone(),
authority_context: build_config.mob_tool_authority_context.clone(),
comms_name: build_config.comms_name.clone(),
comms_runtime: mob_comms,
};
let mob_dispatcher = mob_factory
.build_mob_tools(mob_args)
.await
.map_err(|e| BuildAgentError::Config(format!("Mob tool factory: {e}")))?;
let mob_usage = render_tool_usage_instructions(mob_dispatcher.tools().as_ref());
tools = Arc::new(meerkat_core::DynamicToolComposite::new(vec![
tools,
mob_dispatcher,
]));
if !mob_usage.is_empty() {
if !tool_usage_instructions.is_empty() {
tool_usage_instructions.push_str("\n\n");
}
tool_usage_instructions.push_str(&mob_usage);
}
}
#[cfg(feature = "comms")]
let mut bind_succeeded_wait = false;
#[cfg(feature = "comms")]
if let Some(ref runtime) = comms_runtime {
use meerkat_core::agent::CommsRuntime as CoreCommsRuntimeTrait;
if effective_builtins || tools.capabilities().wait_interrupt {
let notify = CoreCommsRuntimeTrait::actionable_input_notify(runtime.as_ref())
.ok()
.or_else(|| Some(runtime.inbox_notify()));
if let Some(actionable_notify) = notify {
#[cfg(not(target_arch = "wasm32"))]
let (tx, rx) = tokio::sync::watch::channel(
None::<meerkat_core::wait_interrupt::WaitInterrupt>,
);
#[cfg(target_arch = "wasm32")]
let (tx, rx) = tokio_with_wasm::alias::sync::watch::channel(
None::<meerkat_core::wait_interrupt::WaitInterrupt>,
);
bind_succeeded_wait = if !tools.capabilities().wait_interrupt {
tracing::debug!("Dispatcher does not support wait interrupt binding");
false
} else if Arc::strong_count(&tools) == 1 {
let outcome = tools.bind_wait_interrupt(rx).map_err(|e| {
BuildAgentError::Config(format!("Wait interrupt binding failed: {e}"))
})?;
let bound = outcome.was_bound();
tools = outcome.into_dispatcher();
bound
} else {
tracing::debug!(
"Shared dispatcher (refcount={}) — wait interrupt not bound",
Arc::strong_count(&tools)
);
false
};
if bind_succeeded_wait {
#[cfg(not(target_arch = "wasm32"))]
tokio::spawn(async move {
loop {
actionable_notify.notified().await;
if tx
.send(Some(meerkat_core::wait_interrupt::WaitInterrupt {
reason: "Incoming actionable peer message".to_string(),
}))
.is_err()
{
break;
}
}
});
#[cfg(target_arch = "wasm32")]
tokio_with_wasm::alias::task::spawn(async move {
loop {
actionable_notify.notified().await;
if tx
.send(Some(meerkat_core::wait_interrupt::WaitInterrupt {
reason: "Incoming actionable peer message".to_string(),
}))
.is_err()
{
break;
}
}
});
}
} else {
tracing::debug!(
"Comms runtime lacks actionable_input_notify — wait interrupt not bound"
);
}
} else {
tracing::debug!("Builtins disabled — skipping wait interrupt binding");
}
}
{
#[cfg(feature = "comms")]
let comms_wired_feed = bind_succeeded_wait;
#[cfg(not(feature = "comms"))]
let comms_wired_feed = false;
if !comms_wired_feed
&& tools.capabilities().completion_feed
&& Arc::strong_count(&tools) == 1
&& let (Some(feed), Some(baseline)) =
(completion_feed.clone(), interrupt_baseline.clone())
{
let outcome = tools.bind_completion_feed(feed, baseline).map_err(|e| {
BuildAgentError::Config(format!("Completion feed binding failed: {e}"))
})?;
tools = outcome.into_dispatcher();
}
}
if tools.capabilities().ops_lifecycle {
let outcome = tools
.bind_ops_lifecycle(Arc::clone(&ops_lifecycle), session.id().clone())
.map_err(|e| {
BuildAgentError::Config(format!("Ops lifecycle binding failed: {e}"))
})?;
tools = outcome.into_dispatcher();
}
tracing::debug!(
final_tool_count = tools.tools().len(),
tool_names = %tools.tools().iter().map(|t| t.name.as_str()).collect::<Vec<_>>().join(", "),
"tool composition: final dispatcher"
);
#[allow(
clippy::manual_map,
clippy::unnecessary_literal_unwrap,
clippy::needless_match
)]
let hook_engine = match build_config.hook_engine_override.take() {
Some(engine) => Some(engine),
None => {
#[cfg(not(target_arch = "wasm32"))]
{
let layered_hooks = resolve_layered_hooks_config(
_conventions_context_root,
_conventions_user_root,
config,
)
.await;
create_default_hook_engine(layered_hooks)
}
#[cfg(target_arch = "wasm32")]
{
None
}
}
};
#[cfg(feature = "skills")]
let skill_inventory_section = {
if let Some(ref engine) = skill_engine {
let inventory = match engine.inventory_section().await {
Ok(s) => s,
Err(e) => {
tracing::warn!("Failed to generate skill inventory section: {e}");
String::new()
}
};
let mut preload = build_config
.preload_skills
.take()
.and_then(|ids| if ids.is_empty() { None } else { Some(ids) });
if build_config.resume_session.is_some()
&& let Some(ids) = preload.as_mut()
{
let available: std::collections::HashSet<_> = engine
.list_skills(&meerkat_core::skills::SkillFilter::default())
.await
.map(|descs| descs.into_iter().map(|desc| desc.id).collect())
.unwrap_or_default();
let mut dropped = Vec::new();
ids.retain(|id| {
let keep = available.contains(id);
if !keep {
dropped.push(id.0.clone());
}
keep
});
if !dropped.is_empty() {
tracing::warn!(
dropped_skills = ?dropped,
"dropping persisted active skills that are unavailable on the current surface"
);
}
if ids.is_empty() {
preload = None;
}
}
let mut preloaded_sections = Vec::new();
if let Some(ref ids) = preload {
match engine.resolve_and_render(ids).await {
Ok(resolved) => {
for skill in &resolved {
preloaded_sections.push(skill.rendered_body.clone());
}
}
Err(e) => {
return Err(BuildAgentError::Config(format!(
"Failed to preload skill: {e}"
)));
}
}
}
let skill_ids = preload.clone();
(inventory, preloaded_sections, skill_ids)
} else {
(String::new(), Vec::new(), None)
}
};
#[cfg(not(feature = "skills"))]
let skill_inventory_section: (
String,
Vec<String>,
Option<Vec<meerkat_core::skills::SkillId>>,
) = (String::new(), Vec::new(), None);
let (inventory_section, preloaded_skill_sections, active_skill_ids) =
skill_inventory_section;
let mut extra_sections: Vec<&str> = Vec::new();
if !inventory_section.is_empty() && effective_builtins {
extra_sections.push(inventory_section.as_str());
}
for section in &preloaded_skill_sections {
extra_sections.push(section.as_str());
}
let additional_instruction_storage: Vec<String> = build_config
.additional_instructions
.take()
.unwrap_or_default();
for instruction in &additional_instruction_storage {
if !instruction.is_empty() {
extra_sections.push(instruction.as_str());
}
}
let should_apply_system_prompt =
build_config.resume_session.is_none() || per_request_prompt.is_some();
#[cfg(not(target_arch = "wasm32"))]
let system_prompt = if should_apply_system_prompt {
Some(
crate::assemble_system_prompt(
config,
per_request_prompt.as_deref(),
_conventions_context_root,
&extra_sections,
&tool_usage_instructions,
)
.await,
)
} else {
None
};
#[cfg(target_arch = "wasm32")]
let system_prompt = if should_apply_system_prompt {
Some({
let base = per_request_prompt
.or_else(|| config.agent.system_prompt.clone())
.unwrap_or_else(|| DEFAULT_WASM_SYSTEM_PROMPT.to_string());
let mut prompt = base;
for section in &extra_sections {
if !section.is_empty() {
prompt.push_str("\n\n");
prompt.push_str(section);
}
}
if let Some(ref config_tools) = config.agent.tool_instructions
&& !config_tools.is_empty()
{
prompt.push_str("\n\n");
prompt.push_str(config_tools);
}
if !tool_usage_instructions.is_empty() {
prompt.push_str("\n\n");
prompt.push_str(&tool_usage_instructions);
}
prompt
})
} else {
None
};
if build_config.wait_for_mcp {
let timeout = std::time::Duration::from_secs(60);
let started = meerkat_core::time_compat::Instant::now();
loop {
let update = tools.poll_external_updates().await;
if update.pending.is_empty() {
break;
}
if started.elapsed() >= timeout {
tracing::warn!(
"wait_for_mcp timed out after {}s with {} server(s) still pending",
timeout.as_secs(),
update.pending.len()
);
break;
}
#[cfg(not(target_arch = "wasm32"))]
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
#[cfg(target_arch = "wasm32")]
tokio_with_wasm::alias::time::sleep(std::time::Duration::from_millis(500)).await;
}
}
let persisted_build_state = meerkat_core::SessionBuildState {
system_prompt: persisted_system_prompt,
output_schema: build_config.output_schema.clone(),
hooks_override: build_config.hooks_override.clone(),
budget_limits: build_config.budget_limits.clone(),
recoverable_tool_defs: build_config
.recoverable_tool_defs
.clone()
.unwrap_or_default(),
silent_comms_intents: build_config.silent_comms_intents.clone(),
max_inline_peer_notifications: build_config.max_inline_peer_notifications,
app_context: build_config.app_context.clone(),
additional_instructions: build_config.additional_instructions.clone(),
shell_env: build_config.shell_env.clone(),
mob_tool_authority_context: build_config.mob_tool_authority_context.clone(),
call_timeout_override: build_config.call_timeout_override.clone(),
};
let budget_limits = build_config
.budget_limits
.unwrap_or_else(|| config.budget_limits());
let effective_call_timeout_override = {
let build_override = build_config.call_timeout_override;
if build_override.is_inherit() {
config.retry.call_timeout_override.clone()
} else {
build_override
}
};
let mut builder = AgentBuilder::new()
.model(model.clone())
.max_tokens_per_turn(max_tokens)
.budget(budget_limits)
.structured_output_retries(build_config.structured_output_retries)
.with_hook_run_overrides(build_config.hooks_override)
.with_model_defaults_resolver(Arc::new(ProfileBasedDefaultsResolver))
.with_call_timeout_override(effective_call_timeout_override);
if let Some(system_prompt) = system_prompt {
builder = builder.system_prompt(system_prompt);
}
if let Some(schema) = build_config.output_schema {
builder = builder.output_schema(schema);
}
let _is_resumed = build_config.resume_session.is_some();
builder = builder.resume_session(session);
#[cfg(feature = "comms")]
let _comms_enabled = comms_runtime.is_some();
#[cfg(not(feature = "comms"))]
let _comms_enabled = false;
#[cfg(feature = "comms")]
if let Some(runtime) = comms_runtime {
builder =
builder.with_comms_runtime(runtime as Arc<dyn meerkat_core::agent::CommsRuntime>);
}
if let Some(engine) = hook_engine {
builder = builder.with_hook_engine(engine);
}
#[allow(unused_variables)]
let effective_memory = build_config.override_memory.resolve(self.enable_memory);
#[cfg(feature = "memory-store-session")]
if effective_memory {
let memory_dir = self.store_path.join("memory");
match meerkat_memory::HnswMemoryStore::open(&memory_dir) {
Ok(store) => {
let store = Arc::new(store) as Arc<dyn meerkat_core::memory::MemoryStore>;
builder = builder.memory_store(Arc::clone(&store));
let memory_dispatcher =
meerkat_memory::MemorySearchDispatcher::new(Arc::clone(&store));
let gateway = meerkat_core::ToolGatewayBuilder::new()
.add_dispatcher(tools)
.add_dispatcher(Arc::new(memory_dispatcher))
.build()
.map_err(|e| {
BuildAgentError::Config(format!("Failed to compose memory tools: {e}"))
})?;
tools = Arc::new(gateway);
}
Err(e) => {
tracing::warn!(
"Failed to open HnswMemoryStore at {}: {e}",
memory_dir.display()
);
}
}
}
#[cfg(feature = "session-compaction")]
{
let compactor = Arc::new(meerkat_session::DefaultCompactor::new(
config.compaction.clone().into(),
));
builder = builder.compactor(compactor);
}
if let Some(engine) = skill_engine {
builder = builder.with_skill_engine(engine);
}
builder = builder.with_event_tap(event_tap);
if let Some(tx) = build_config.event_tx {
builder = builder.with_default_event_tx(tx);
}
if !build_config.silent_comms_intents.is_empty() {
builder = builder.with_silent_comms_intents(build_config.silent_comms_intents);
}
builder =
builder.with_max_inline_peer_notifications(build_config.max_inline_peer_notifications);
if let Some(cp) = build_config.checkpointer {
builder = builder.with_checkpointer(cp);
}
if let Some(blob_store) = build_config.blob_store_override {
builder = builder.with_blob_store(blob_store);
}
builder = builder.with_ops_lifecycle(Arc::clone(&ops_lifecycle));
if let RuntimeBuildMode::SessionOwned(bindings) = resolved_mode {
builder = builder.with_epoch_cursor_state(Arc::clone(&bindings.cursor_state));
}
if let Some(feed) = completion_feed {
builder = builder.with_completion_feed(feed);
}
if let Some(baseline) = interrupt_baseline {
builder = builder.with_interrupt_baseline(baseline);
}
if let Some(enrichment) = tools.completion_enrichment() {
builder = builder.with_completion_enrichment(enrichment);
}
let mut agent = builder.build(llm_adapter, tools, store_adapter).await;
if !image_tool_results {
let deny = std::collections::HashSet::from(["view_image".to_string()]);
if let Err(err) = agent.stage_external_tool_filter(meerkat_core::ToolFilter::Deny(deny))
{
tracing::warn!(error = %err, "failed to stage initial view_image deny filter");
}
}
let metadata = if let Some(mut metadata) = resumed_session_metadata {
metadata.model = model;
metadata.max_tokens = max_tokens;
metadata.structured_output_retries = build_config.structured_output_retries;
metadata.provider = provider;
metadata.provider_params = build_config.provider_params;
metadata.tooling.builtins = build_config.override_builtins;
metadata.tooling.shell = build_config.override_shell;
metadata.tooling.mob = build_config.override_mob;
metadata.tooling.memory = build_config.override_memory;
if build_config.resume_override_mask.preload_skills {
metadata.tooling.active_skills = active_skill_ids;
}
metadata.keep_alive = build_config.keep_alive;
metadata.comms_name = build_config.comms_name;
metadata.peer_meta = build_config.peer_meta;
metadata.realm_id = build_config.realm_id;
metadata.instance_id = build_config.instance_id;
metadata.backend = build_config.backend;
metadata.config_generation = build_config.config_generation;
metadata
} else {
SessionMetadata {
model,
max_tokens,
structured_output_retries: build_config.structured_output_retries,
provider,
provider_params: build_config.provider_params,
tooling: SessionTooling {
builtins: build_config.override_builtins,
shell: build_config.override_shell,
comms: ToolCategoryOverride::Inherit,
mob: build_config.override_mob,
memory: build_config.override_memory,
active_skills: active_skill_ids,
},
keep_alive: build_config.keep_alive,
comms_name: build_config.comms_name,
peer_meta: build_config.peer_meta,
realm_id: build_config.realm_id,
instance_id: build_config.instance_id,
backend: build_config.backend,
config_generation: build_config.config_generation,
}
};
if let Err(err) = agent.session_mut().set_session_metadata(metadata) {
tracing::warn!("Failed to store session metadata: {}", err);
}
if let Err(err) = agent.session_mut().set_build_state(persisted_build_state) {
tracing::warn!("Failed to store session build state: {}", err);
}
Ok(agent)
}
}
impl AgentFactory {
#[cfg(not(target_arch = "wasm32"))]
#[allow(clippy::too_many_arguments)]
async fn build_tool_dispatcher_for_agent_with_overrides(
&self,
_config: &Config,
external: Option<Arc<dyn AgentToolDispatcher>>,
effective_builtins: bool,
effective_shell: bool,
skill_engine: Option<Arc<meerkat_core::skills::SkillRuntime>>,
shell_env: Option<std::collections::HashMap<String, String>>,
session_id: String,
ops_lifecycle: Arc<dyn OpsLifecycleRegistry>,
image_tool_results: bool,
) -> Result<(Arc<dyn AgentToolDispatcher>, String), BuildAgentError> {
if !effective_builtins {
return match external {
Some(ext) => {
let usage = render_tool_usage_instructions(ext.tools().as_ref());
Ok((ext, usage))
}
None => Ok((Arc::new(EmptyToolDispatcher), String::new())),
};
}
#[cfg(feature = "session-store")]
let task_store: Arc<dyn TaskStore> = Arc::new(SqliteTaskStore::for_session(
self.store_path.join("tasks.db"),
&session_id,
));
#[cfg(not(feature = "session-store"))]
let task_store: Arc<dyn TaskStore> = match self.project_root.as_ref() {
Some(root) => Arc::new(FileTaskStore::in_project(root)),
None => Arc::new(MemoryTaskStore::new()),
};
let shell_config = if effective_shell {
let project_root = self
.project_root
.clone()
.unwrap_or_else(|| self.store_path.clone());
let mut config = ShellConfig::with_project_root(project_root);
if let Some(env) = shell_env {
config.env_vars = env;
}
Some(config)
} else {
None
};
let builtin_config = if effective_shell {
BuiltinToolConfig {
policy: ToolPolicyLayer::new()
.enable_tool("shell")
.enable_tool("shell_job_status")
.enable_tool("shell_jobs")
.enable_tool("shell_job_cancel"),
..Default::default()
}
} else {
BuiltinToolConfig::default()
};
let dispatcher = self
.build_builtin_dispatcher_with_skills_internal(
task_store,
builtin_config,
self.project_root.clone(),
shell_config,
external,
Some(session_id),
Some(ops_lifecycle),
skill_engine,
image_tool_results,
)
.await?;
let usage = render_tool_usage_instructions(dispatcher.tools().as_ref());
Ok((dispatcher, usage))
}
}
fn render_tool_usage_instructions(tools: &[Arc<meerkat_core::ToolDef>]) -> String {
if tools.is_empty() {
return String::new();
}
let mut out = String::from("# Available Tools\n\n");
for tool in tools {
out.push_str(&format!("## {}\n{}\n\n", tool.name, tool.description));
}
out
}