use async_trait::async_trait;
use meerkat_core::Config;
#[cfg(not(target_arch = "wasm32"))]
use meerkat_core::ConfigStore;
use meerkat_core::Session;
use meerkat_core::comms::{CommsCommand, PeerDirectoryEntry, SendError, SendReceipt};
use meerkat_core::event::AgentEvent;
use meerkat_core::service::{CreateSessionRequest, SessionError, TurnToolOverlay};
use meerkat_core::types::{HandlingMode, Message, RenderMetadata, RunResult, SessionId};
use meerkat_session::EphemeralSessionService;
use meerkat_session::ephemeral::{SessionAgent, SessionAgentBuilder, SessionSnapshot};
use std::sync::Arc;
#[cfg(not(target_arch = "wasm32"))]
use tokio::sync::mpsc;
#[cfg(target_arch = "wasm32")]
use tokio_with_wasm::alias::sync::mpsc;
#[cfg(feature = "session-store")]
use crate::PersistenceBundle;
use crate::{AgentBuildConfig, AgentFactory, DynAgent};
use meerkat_client::LlmClient;
pub struct FactoryAgent {
agent: DynAgent,
}
impl FactoryAgent {
pub fn agent(&self) -> &DynAgent {
&self.agent
}
pub fn agent_mut(&mut self) -> &mut DynAgent {
&mut self.agent
}
pub fn session(&self) -> &Session {
self.agent.session()
}
pub async fn send(&self, cmd: CommsCommand) -> Result<SendReceipt, SendError> {
let runtime = self
.agent
.comms()
.ok_or_else(|| SendError::Unsupported("comms runtime is not configured".to_string()))?;
runtime.send(cmd).await
}
pub async fn peers(&self) -> Vec<PeerDirectoryEntry> {
match self.agent.comms() {
Some(runtime) => runtime.peers().await,
None => Vec::new(),
}
}
}
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
impl SessionAgent for FactoryAgent {
async fn run_with_events(
&mut self,
prompt: meerkat_core::types::ContentInput,
event_tx: mpsc::Sender<AgentEvent>,
) -> Result<RunResult, meerkat_core::error::AgentError> {
self.agent.run_with_events(prompt, event_tx).await
}
async fn run_turn_with_events(
&mut self,
prompt: meerkat_core::types::ContentInput,
handling_mode: HandlingMode,
render_metadata: Option<RenderMetadata>,
event_tx: mpsc::Sender<AgentEvent>,
) -> Result<RunResult, meerkat_core::error::AgentError> {
if handling_mode != HandlingMode::Queue {
return Err(meerkat_core::error::AgentError::ConfigError(format!(
"handling_mode {handling_mode:?} requires a runtime-backed surface; direct session-service path supports Queue only",
)));
}
if render_metadata.is_some() {
return Err(meerkat_core::error::AgentError::ConfigError(
"render_metadata requires a runtime-backed surface; direct session-service path does not support it".to_string(),
));
}
self.agent.run_with_events(prompt, event_tx).await
}
async fn run_pending_with_events(
&mut self,
event_tx: mpsc::Sender<AgentEvent>,
) -> Result<RunResult, meerkat_core::error::AgentError> {
self.agent.run_pending_with_events(event_tx).await
}
fn set_skill_references(&mut self, refs: Option<Vec<meerkat_core::skills::SkillKey>>) {
self.agent.pending_skill_references = refs;
}
fn set_flow_tool_overlay(
&mut self,
overlay: Option<TurnToolOverlay>,
) -> Result<(), meerkat_core::error::AgentError> {
self.agent
.set_flow_tool_overlay(overlay)
.map_err(|error| meerkat_core::error::AgentError::ConfigError(error.to_string()))
}
fn apply_pending_tool_results(
&mut self,
results: Vec<meerkat_core::ToolResult>,
) -> Result<(), meerkat_core::error::AgentError> {
if results.is_empty() {
return Ok(());
}
self.agent
.session_mut()
.push(Message::ToolResults { results });
Ok(())
}
fn replace_client(&mut self, client: std::sync::Arc<dyn meerkat_core::AgentLlmClient>) {
self.agent.replace_client(client);
}
fn update_keep_alive(&mut self, keep_alive: bool) {
if let Some(mut metadata) = self.agent.session().session_metadata() {
metadata.keep_alive = keep_alive;
if let Err(e) = self.agent.session_mut().set_session_metadata(metadata) {
tracing::warn!(error = %e, "failed to update keep_alive in session metadata");
}
}
}
fn update_mob_tool_authority_context(
&mut self,
authority_context: Option<meerkat_core::service::MobToolAuthorityContext>,
) -> Result<(), meerkat_core::error::AgentError> {
self.agent
.session_mut()
.set_mob_tool_authority_context(authority_context)
.map_err(|err| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to update mob tool authority context in session metadata: {err}"
))
})
}
fn update_system_prompt(
&mut self,
system_prompt: String,
) -> Result<(), meerkat_core::error::AgentError> {
self.agent.session_mut().set_system_prompt(system_prompt);
Ok(())
}
fn hot_swap_llm_identity(
&mut self,
client: std::sync::Arc<dyn meerkat_core::AgentLlmClient>,
identity: meerkat_core::SessionLlmIdentity,
) -> Result<(), meerkat_core::error::AgentError> {
let Some(mut metadata) = self.agent.session().session_metadata() else {
return Err(meerkat_core::error::AgentError::InternalError(
"session metadata missing during llm identity hot-swap".to_string(),
));
};
metadata.apply_llm_identity(&identity);
self.agent
.session_mut()
.set_session_metadata(metadata)
.map_err(|err| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to update session metadata during llm identity hot-swap: {err}"
))
})?;
self.agent.replace_client(client);
Ok(())
}
fn stage_external_tool_filter(
&mut self,
filter: meerkat_core::ToolFilter,
) -> Result<(), meerkat_core::error::AgentError> {
self.agent
.stage_external_tool_filter(filter)
.map(|_| ())
.map_err(|error| meerkat_core::error::AgentError::ConfigError(error.to_string()))
}
fn set_tool_visibility_state(
&mut self,
state: Option<meerkat_core::SessionToolVisibilityState>,
) -> Result<(), meerkat_core::error::AgentError> {
let visibility_state = state.clone().unwrap_or_default();
self.agent
.tool_scope()
.set_visibility_state(visibility_state)
.map_err(|err| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to replace tool visibility state on live scope: {err}"
))
})?;
if let Some(state) = state {
self.agent
.session_mut()
.set_tool_visibility_state(state)
.map_err(|err| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to persist tool visibility state into session metadata: {err}"
))
})
} else {
self.agent
.session_mut()
.remove_metadata(meerkat_core::SESSION_TOOL_VISIBILITY_STATE_KEY);
Ok(())
}
}
fn sync_system_context_state(&mut self) {
self.agent.sync_system_context_state_to_session();
}
fn cancel(&mut self) {
self.agent.cancel();
}
fn session_id(&self) -> SessionId {
self.agent.session().id().clone()
}
fn snapshot(&self) -> SessionSnapshot {
let s = self.agent.session();
SessionSnapshot {
created_at: s.created_at(),
updated_at: s.updated_at(),
message_count: s.messages().len(),
total_tokens: s.total_tokens(),
usage: s.total_usage(),
last_assistant_text: s.last_assistant_text(),
}
}
fn session_clone(&self) -> Session {
self.agent.session_with_system_context_state()
}
fn has_pending_boundary(&self) -> bool {
self.agent.session().has_pending_boundary()
}
fn apply_runtime_system_context(
&mut self,
appends: &[meerkat_core::PendingSystemContextAppend],
) {
self.agent
.session_mut()
.append_system_context_blocks(appends);
}
fn system_context_state(
&self,
) -> Arc<std::sync::Mutex<meerkat_core::SessionSystemContextState>> {
self.agent.system_context_state()
}
fn event_injector(&self) -> Option<Arc<dyn meerkat_core::EventInjector>> {
self.agent.comms_arc()?.event_injector()
}
#[doc(hidden)]
fn interaction_event_injector(
&self,
) -> Option<Arc<dyn meerkat_core::event_injector::SubscribableInjector>> {
self.agent.comms_arc()?.interaction_event_injector()
}
fn comms_runtime(&self) -> Option<Arc<dyn meerkat_core::agent::CommsRuntime>> {
self.agent.comms_arc()
}
}
pub struct FactoryAgentBuilder {
factory: AgentFactory,
config_snapshot: Config,
#[cfg(not(target_arch = "wasm32"))]
config_store: Option<Arc<dyn ConfigStore>>,
pub default_llm_client: Option<Arc<dyn LlmClient>>,
pub default_tool_dispatcher: Option<Arc<dyn meerkat_core::AgentToolDispatcher>>,
pub default_session_store: Option<Arc<dyn meerkat_core::AgentSessionStore>>,
pub default_mob_tools:
Arc<std::sync::RwLock<Option<Arc<dyn meerkat_core::service::MobToolsFactory>>>>,
pub default_schedule_tools:
Arc<std::sync::RwLock<Option<Arc<dyn meerkat_core::AgentToolDispatcher>>>>,
}
impl FactoryAgentBuilder {
pub fn new(factory: AgentFactory, config: Config) -> Self {
Self {
factory,
config_snapshot: config,
#[cfg(not(target_arch = "wasm32"))]
config_store: None,
default_llm_client: None,
default_tool_dispatcher: None,
default_session_store: None,
default_mob_tools: Arc::new(std::sync::RwLock::new(None)),
default_schedule_tools: Arc::new(std::sync::RwLock::new(None)),
}
}
#[cfg(not(target_arch = "wasm32"))]
pub fn new_with_config_store(
factory: AgentFactory,
initial_config: Config,
config_store: Arc<dyn ConfigStore>,
) -> Self {
Self {
factory,
config_snapshot: initial_config,
config_store: Some(config_store),
default_llm_client: None,
default_tool_dispatcher: None,
default_session_store: None,
default_mob_tools: Arc::new(std::sync::RwLock::new(None)),
default_schedule_tools: Arc::new(std::sync::RwLock::new(None)),
}
}
async fn resolve_config(&self) -> Config {
#[cfg(not(target_arch = "wasm32"))]
if let Some(store) = &self.config_store {
match store.get().await {
Ok(config) => return config,
Err(err) => {
tracing::warn!("Failed to read latest config from store: {err}");
}
}
}
self.config_snapshot.clone()
}
pub fn factory(&self) -> &AgentFactory {
&self.factory
}
pub fn config(&self) -> &Config {
&self.config_snapshot
}
}
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
impl SessionAgentBuilder for FactoryAgentBuilder {
type Agent = FactoryAgent;
async fn build_agent(
&self,
req: &CreateSessionRequest,
event_tx: mpsc::Sender<AgentEvent>,
) -> Result<FactoryAgent, SessionError> {
let mut build_config = AgentBuildConfig::from_create_session_request(req, event_tx);
if build_config.llm_client_override.is_none()
&& let Some(ref client) = self.default_llm_client
{
build_config.llm_client_override = Some(client.clone());
}
if build_config.tool_dispatcher_override.is_none()
&& let Some(ref dispatcher) = self.default_tool_dispatcher
{
build_config.tool_dispatcher_override = Some(dispatcher.clone());
}
if build_config.session_store_override.is_none()
&& let Some(ref store) = self.default_session_store
{
build_config.session_store_override = Some(store.clone());
}
if build_config.mob_tools.is_none()
&& let Some(mob_factory) = self
.default_mob_tools
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
{
build_config.mob_tools = Some(mob_factory);
}
if build_config.schedule_tools.is_none()
&& let Some(schedule_dispatcher) = self
.default_schedule_tools
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
{
build_config.schedule_tools = Some(schedule_dispatcher);
}
let config = self.resolve_config().await;
let agent = self
.factory
.build_agent(build_config, &config)
.await
.map_err(|e| {
SessionError::Agent(meerkat_core::error::AgentError::BuildError(e.to_string()))
})?;
Ok(FactoryAgent { agent })
}
}
pub fn build_ephemeral_service(
factory: AgentFactory,
config: Config,
max_sessions: usize,
) -> EphemeralSessionService<FactoryAgentBuilder> {
let builder = FactoryAgentBuilder::new(factory, config);
EphemeralSessionService::new(builder, max_sessions)
}
#[cfg(feature = "session-store")]
pub fn build_persistent_service_with_runtime_adapter(
factory: AgentFactory,
config: Config,
max_sessions: usize,
persistence: PersistenceBundle,
) -> (
meerkat_session::PersistentSessionService<FactoryAgentBuilder>,
Arc<meerkat_runtime::RuntimeSessionAdapter>,
) {
let runtime_adapter = persistence.runtime_adapter();
let mut builder = FactoryAgentBuilder::new(factory, config);
let (store, runtime_store, blob_store) = persistence.into_parts();
builder.default_session_store = Some(Arc::new(meerkat_store::StoreAdapter::new(Arc::clone(
&store,
))));
(
meerkat_session::PersistentSessionService::new(
builder,
max_sessions,
store,
runtime_store,
blob_store,
),
runtime_adapter,
)
}
#[cfg(feature = "session-store")]
pub fn build_persistent_service(
factory: AgentFactory,
config: Config,
max_sessions: usize,
persistence: PersistenceBundle,
) -> meerkat_session::PersistentSessionService<FactoryAgentBuilder> {
build_persistent_service_with_runtime_adapter(factory, config, max_sessions, persistence).0
}
#[cfg(test)]
#[allow(clippy::expect_used)]
mod tests {
use super::*;
use async_trait::async_trait;
use futures::stream;
use meerkat_client::{LlmClient, LlmDoneOutcome, LlmEvent, LlmRequest};
use meerkat_core::Config;
use meerkat_core::comms::InputSource;
use meerkat_core::ops_lifecycle::OpsLifecycleRegistry;
use meerkat_core::service::SessionBuildOptions;
use meerkat_core::{
Provider, ToolCallView, ToolDef, ToolDispatchOutcome, ToolError, ToolResult,
};
use meerkat_runtime::RuntimeSessionAdapter;
use meerkat_schedule::{MemoryScheduleStore, ScheduleService, ScheduleToolDispatcher};
use meerkat_session::ephemeral::SessionAgent;
use std::pin::Pin;
use std::sync::Mutex;
use std::sync::atomic::{AtomicBool, Ordering};
use tempfile::TempDir;
struct MockLlmClient {
delta: &'static str,
}
impl Default for MockLlmClient {
fn default() -> Self {
Self { delta: "ok" }
}
}
#[derive(Default)]
struct CaptureToolClient {
inner: meerkat_client::TestClient,
seen_tools: Mutex<Vec<String>>,
}
impl CaptureToolClient {
fn tool_names(&self) -> Vec<String> {
self.seen_tools.lock().expect("capture lock").clone()
}
}
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
impl LlmClient for CaptureToolClient {
fn stream<'a>(
&'a self,
request: &'a LlmRequest,
) -> Pin<
Box<dyn futures::Stream<Item = Result<LlmEvent, meerkat_client::LlmError>> + Send + 'a>,
> {
*self.seen_tools.lock().expect("capture lock") =
request.tools.iter().map(|tool| tool.name.clone()).collect();
self.inner.stream(request)
}
fn provider(&self) -> &'static str {
self.inner.provider()
}
async fn health_check(&self) -> Result<(), meerkat_client::LlmError> {
self.inner.health_check().await
}
}
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
impl LlmClient for MockLlmClient {
fn stream<'a>(
&'a self,
_request: &'a LlmRequest,
) -> Pin<
Box<dyn futures::Stream<Item = Result<LlmEvent, meerkat_client::LlmError>> + Send + 'a>,
> {
Box::pin(stream::iter(vec![
Ok(LlmEvent::TextDelta {
delta: self.delta.to_string(),
meta: None,
}),
Ok(LlmEvent::Done {
outcome: LlmDoneOutcome::Success {
stop_reason: meerkat_core::StopReason::EndTurn,
},
}),
]))
}
fn provider(&self) -> &'static str {
"mock"
}
async fn health_check(&self) -> Result<(), meerkat_client::LlmError> {
Ok(())
}
}
#[derive(Default)]
struct RegistryBindingProbe {
bound: AtomicBool,
seen_registry: Mutex<Option<Arc<dyn OpsLifecycleRegistry>>>,
seen_session_id: Mutex<Option<SessionId>>,
}
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
impl meerkat_core::AgentToolDispatcher for RegistryBindingProbe {
fn tools(&self) -> Arc<[Arc<ToolDef>]> {
Arc::from([])
}
async fn dispatch(&self, call: ToolCallView<'_>) -> Result<ToolDispatchOutcome, ToolError> {
Ok(ToolResult::new(call.id.to_string(), "noop".to_string(), false).into())
}
fn capabilities(&self) -> meerkat_core::agent::DispatcherCapabilities {
meerkat_core::agent::DispatcherCapabilities {
ops_lifecycle: true,
}
}
fn bind_ops_lifecycle(
self: Arc<Self>,
registry: Arc<dyn OpsLifecycleRegistry>,
owner_session_id: SessionId,
) -> Result<meerkat_core::agent::BindOutcome, meerkat_core::agent::OpsLifecycleBindError>
{
self.bound.store(true, Ordering::SeqCst);
*self.seen_registry.lock().expect("probe lock") = Some(registry);
*self.seen_session_id.lock().expect("probe lock") = Some(owner_session_id);
Ok(meerkat_core::agent::BindOutcome::Bound(self))
}
}
async fn build_factory_agent_with_mock(
temp: &TempDir,
mut build_config: AgentBuildConfig,
) -> Result<FactoryAgent, String> {
let factory = AgentFactory::new(temp.path().join("sessions"));
build_config.llm_client_override = Some(Arc::new(MockLlmClient::default()));
let agent = factory
.build_agent(build_config, &Config::default())
.await
.map_err(|err| format!("{err}"))?;
Ok(FactoryAgent { agent })
}
#[tokio::test]
async fn factory_builder_uses_runtime_session_registry_override() -> Result<(), String> {
let temp = tempfile::tempdir().map_err(|err| format!("tempdir: {err}"))?;
let factory = AgentFactory::new(temp.path().join("sessions"));
let mut builder = FactoryAgentBuilder::new(factory, Config::default());
builder.default_llm_client = Some(Arc::new(MockLlmClient { delta: "ok" }));
let probe = Arc::new(RegistryBindingProbe::default());
let probe_dispatcher: Arc<dyn meerkat_core::AgentToolDispatcher> = probe.clone();
builder.default_tool_dispatcher = Some(probe_dispatcher);
let runtime_adapter = RuntimeSessionAdapter::ephemeral();
let session = Session::new();
let session_id = session.id().clone();
runtime_adapter.register_session(session_id.clone()).await;
let expected_registry = runtime_adapter
.ops_lifecycle_registry(&session_id)
.await
.ok_or_else(|| "missing runtime registry".to_string())?
as Arc<dyn OpsLifecycleRegistry>;
let bindings = meerkat_core::SessionRuntimeBindings {
session_id: session_id.clone(),
epoch_id: meerkat_core::runtime_epoch::RuntimeEpochId::new(),
ops_lifecycle: expected_registry.clone(),
cursor_state: Arc::new(meerkat_core::EpochCursorState::new()),
};
let req = CreateSessionRequest {
model: "claude-sonnet-4-5".to_string(),
prompt: "hello".to_string().into(),
render_metadata: None,
system_prompt: None,
max_tokens: None,
event_tx: None,
skill_references: None,
initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::Discard,
build: Some(SessionBuildOptions {
resume_session: Some(session),
runtime_build_mode: meerkat_core::RuntimeBuildMode::SessionOwned(bindings),
..SessionBuildOptions::default()
}),
labels: None,
};
let (event_tx, _event_rx) = mpsc::channel(8);
let agent = builder
.build_agent(&req, event_tx)
.await
.map_err(|err| err.to_string())?;
drop(agent);
assert!(
probe.bound.load(Ordering::SeqCst),
"dispatcher should receive ops lifecycle binding"
);
let seen_registry = probe
.seen_registry
.lock()
.expect("probe lock")
.clone()
.ok_or_else(|| "dispatcher did not record registry".to_string())?;
let seen_session_id = probe
.seen_session_id
.lock()
.expect("probe lock")
.clone()
.ok_or_else(|| "dispatcher did not record session id".to_string())?;
assert!(
Arc::ptr_eq(&seen_registry, &expected_registry),
"factory should use runtime adapter's canonical registry, not a fresh fallback"
);
assert_eq!(seen_session_id, session_id);
Ok(())
}
fn mock_input_cmd(session_id: &SessionId) -> CommsCommand {
CommsCommand::Input {
session_id: session_id.clone(),
body: "hello".to_string(),
blocks: None,
source: InputSource::Rpc,
handling_mode: meerkat_core::types::HandlingMode::Queue,
allow_self_session: true,
}
}
#[tokio::test]
async fn test_factory_agent_send_without_comms_runtime_is_unsupported() -> Result<(), String> {
let temp = tempfile::tempdir().map_err(|err| format!("tempdir: {err}"))?;
let agent = build_factory_agent_with_mock(
&temp,
AgentBuildConfig {
..AgentBuildConfig::new("claude-sonnet-4-5")
},
)
.await?;
let session_id = agent.session().id().clone();
let result = agent.send(mock_input_cmd(&session_id)).await;
assert!(matches!(result, Err(SendError::Unsupported(_))));
Ok(())
}
#[tokio::test]
async fn test_session_llm_override_is_applied_end_to_end() -> Result<(), String> {
let temp = tempfile::tempdir().map_err(|err| format!("tempdir: {err}"))?;
let factory = AgentFactory::new(temp.path().join("sessions"));
let mut builder = FactoryAgentBuilder::new(factory, Config::default());
builder.default_llm_client = Some(Arc::new(MockLlmClient { delta: "default" }));
let build = SessionBuildOptions {
llm_client_override: Some(crate::encode_llm_client_override_for_service(Arc::new(
MockLlmClient { delta: "override" },
))),
..SessionBuildOptions::default()
};
let req = CreateSessionRequest {
model: "claude-sonnet-4-5".to_string(),
prompt: "ignored".to_string().into(),
render_metadata: None,
system_prompt: None,
max_tokens: None,
event_tx: None,
skill_references: None,
initial_turn: meerkat_core::service::InitialTurnPolicy::RunImmediately,
deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::Discard,
build: Some(build),
labels: None,
};
let (build_event_tx, _build_event_rx) = mpsc::channel(8);
let mut agent = builder
.build_agent(&req, build_event_tx)
.await
.map_err(|err| format!("{err}"))?;
let (run_event_tx, _run_event_rx) = mpsc::channel(8);
let result =
SessionAgent::run_with_events(&mut agent, "hello".to_string().into(), run_event_tx)
.await
.map_err(|err| format!("{err}"))?;
assert_eq!(result.text, "override");
Ok(())
}
#[tokio::test]
async fn test_default_llm_override_allows_unknown_model_names() -> Result<(), String> {
let temp = tempfile::tempdir().map_err(|err| format!("tempdir: {err}"))?;
let factory = AgentFactory::new(temp.path().join("sessions"));
let mut builder = FactoryAgentBuilder::new(factory, Config::default());
builder.default_llm_client = Some(Arc::new(MockLlmClient { delta: "override" }));
let (event_tx, _event_rx) = mpsc::channel(8);
let agent = builder
.build_agent(&make_session_request("mock-model"), event_tx)
.await
.map_err(|err| format!("{err}"))?;
let metadata = agent
.session()
.session_metadata()
.ok_or_else(|| "missing session metadata".to_string())?;
assert_eq!(metadata.provider, Provider::Other);
Ok(())
}
#[tokio::test]
async fn factory_builder_uses_default_schedule_tools_on_runtime_backed_resume()
-> Result<(), String> {
let temp = tempfile::tempdir().map_err(|err| format!("tempdir: {err}"))?;
let factory = AgentFactory::new(temp.path().join("sessions")).schedule(true);
let mut builder = FactoryAgentBuilder::new(factory, Config::default());
let capture: Arc<CaptureToolClient> = Arc::new(CaptureToolClient::default());
builder.default_llm_client = Some(capture.clone());
*builder
.default_schedule_tools
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) =
Some(Arc::new(ScheduleToolDispatcher::new(ScheduleService::new(
Arc::new(MemoryScheduleStore::default()),
))));
let runtime_adapter = RuntimeSessionAdapter::ephemeral();
let session = Session::new();
let session_id = session.id().clone();
runtime_adapter.register_session(session_id.clone()).await;
let bindings = meerkat_core::SessionRuntimeBindings {
session_id,
epoch_id: meerkat_core::runtime_epoch::RuntimeEpochId::new(),
ops_lifecycle: runtime_adapter
.ops_lifecycle_registry(session.id())
.await
.ok_or_else(|| "missing runtime registry".to_string())?,
cursor_state: Arc::new(meerkat_core::EpochCursorState::new()),
};
let req = CreateSessionRequest {
model: "claude-sonnet-4-5".to_string(),
prompt: "hello".to_string().into(),
render_metadata: None,
system_prompt: None,
max_tokens: None,
event_tx: None,
skill_references: None,
initial_turn: meerkat_core::service::InitialTurnPolicy::RunImmediately,
deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::Discard,
build: Some(SessionBuildOptions {
resume_session: Some(session),
runtime_build_mode: meerkat_core::RuntimeBuildMode::SessionOwned(bindings),
..SessionBuildOptions::default()
}),
labels: None,
};
let (event_tx, _event_rx) = mpsc::channel(8);
let mut agent = builder
.build_agent(&req, event_tx)
.await
.map_err(|err| format!("{err}"))?;
let (run_tx, _run_rx) = mpsc::channel(8);
SessionAgent::run_with_events(&mut agent, "inspect".to_string().into(), run_tx)
.await
.map_err(|err| format!("{err}"))?;
let tool_names = capture.tool_names();
assert!(
tool_names
.iter()
.any(|name| name == "meerkat_schedule_create")
);
assert!(
tool_names
.iter()
.any(|name| name == "meerkat_schedule_list")
);
Ok(())
}
#[test]
fn test_session_build_options_preserve_keep_alive_flag() {
let mut build = AgentBuildConfig::new("claude-sonnet-4-5");
build.apply_session_build_options(&SessionBuildOptions {
keep_alive: true,
..SessionBuildOptions::default()
});
assert!(build.keep_alive);
}
fn make_session_request(model: &str) -> CreateSessionRequest {
CreateSessionRequest {
model: model.to_string(),
prompt: "test".to_string().into(),
render_metadata: None,
system_prompt: None,
max_tokens: None,
event_tx: None,
skill_references: None,
initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::Discard,
build: None,
labels: None,
}
}
#[tokio::test]
async fn test_config_api_keys_resolve_different_providers_per_model() -> Result<(), String> {
let temp = tempfile::tempdir().map_err(|e| format!("tempdir: {e}"))?;
let factory = AgentFactory::new(temp.path().join("sessions"));
let mut config = Config::default();
config.providers.api_keys = Some(std::collections::HashMap::from([
("anthropic".into(), "test-anthropic-key".into()),
("openai".into(), "test-openai-key".into()),
("gemini".into(), "test-gemini-key".into()),
]));
let builder = FactoryAgentBuilder::new(factory, config);
let (tx1, _rx1) = mpsc::channel(8);
builder
.build_agent(&make_session_request("claude-sonnet-4-5"), tx1)
.await
.map_err(|e| format!("anthropic model should build: {e}"))?;
let (tx2, _rx2) = mpsc::channel(8);
builder
.build_agent(&make_session_request("gpt-5.2"), tx2)
.await
.map_err(|e| format!("openai model should build: {e}"))?;
let (tx3, _rx3) = mpsc::channel(8);
builder
.build_agent(&make_session_request("gemini-3-flash-preview"), tx3)
.await
.map_err(|e| format!("gemini model should build: {e}"))?;
Ok(())
}
#[tokio::test]
async fn test_default_llm_client_takes_precedence_over_config_keys() -> Result<(), String> {
let temp = tempfile::tempdir().map_err(|e| format!("tempdir: {e}"))?;
let factory = AgentFactory::new(temp.path().join("sessions"));
let config = Config::default();
let mut builder = FactoryAgentBuilder::new(factory, config);
builder.default_llm_client = Some(Arc::new(MockLlmClient {
delta: "from-default",
}));
let (tx, _rx) = mpsc::channel(8);
let mut agent = builder
.build_agent(&make_session_request("claude-sonnet-4-5"), tx)
.await
.map_err(|e| format!("build failed: {e}"))?;
let (run_tx, _run_rx) = mpsc::channel(8);
let result = SessionAgent::run_with_events(&mut agent, "hello".into(), run_tx)
.await
.map_err(|e| format!("run failed: {e}"))?;
assert_eq!(result.text, "from-default");
Ok(())
}
#[tokio::test]
async fn test_unknown_model_prefix_fails_even_with_config_keys() -> Result<(), String> {
let temp = tempfile::tempdir().map_err(|e| format!("tempdir: {e}"))?;
let factory = AgentFactory::new(temp.path().join("sessions"));
let mut config = Config::default();
config.providers.api_keys = Some(std::collections::HashMap::from([(
"anthropic".into(),
"test-key".into(),
)]));
let builder = FactoryAgentBuilder::new(factory, config);
let (tx, _rx) = mpsc::channel(8);
let result = builder
.build_agent(&make_session_request("llama-3.1-70b"), tx)
.await;
assert!(result.is_err(), "unknown model prefix should fail");
Ok(())
}
}