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::{
AssistantBlock, HandlingMode, Message, RunResult, SessionId, StopReason, Usage,
};
use meerkat_session::EphemeralSessionService;
use meerkat_session::ephemeral::{
HeadCanonicalRuntimeBoundaryAcknowledgeOutcome, HeadCanonicalRuntimeBoundaryAuthority,
HeadCanonicalRuntimeBoundaryPrepareRequest, ObservedSessionTailKind,
PreparedHeadCanonicalRuntimeBoundary, SessionAgent, SessionAgentBuilder, SessionAgentTurnInput,
SessionSnapshot, SessionTranscriptAuthoritySnapshot,
};
use std::sync::Arc;
#[cfg(not(target_arch = "wasm32"))]
use tokio::sync::{mpsc, watch};
#[cfg(target_arch = "wasm32")]
use tokio_with_wasm::alias::sync::{mpsc, watch};
#[cfg(feature = "session-store")]
use crate::PersistenceBundle;
use crate::{AgentBuildConfig, AgentFactory, BuildAgentError, DynAgent};
use meerkat_client::LlmClient;
pub struct FactoryAgent {
agent: DynAgent,
session_context: Option<Arc<dyn meerkat_core::handles::SessionContextHandle>>,
pending_head_canonical_boundary: Option<PendingFactoryHeadCanonicalBoundary>,
acknowledged_head_canonical_boundary: Option<PendingFactoryHeadCanonicalBoundary>,
}
#[derive(Clone)]
struct PendingFactoryHeadCanonicalBoundary {
authority: HeadCanonicalRuntimeBoundaryAuthority,
observed_head_token: Option<String>,
request_projection_token: String,
prepared: PreparedHeadCanonicalRuntimeBoundary,
}
fn build_agent_error_to_session_error(
error: BuildAgentError,
provider: Option<meerkat_core::Provider>,
auth_binding: Option<&meerkat_core::AuthBindingRef>,
) -> SessionError {
match error {
#[cfg(feature = "comms")]
BuildAgentError::SessionIdentityInUse(session_id) => SessionError::Agent(
meerkat_core::error::AgentError::SessionIdentityInUse(session_id),
),
BuildAgentError::LlmClient(meerkat_client::FactoryError::ProviderAuth(
meerkat_llm_core::provider_runtime::ProviderAuthError::Auth(auth_error),
)) => match provider {
Some(provider) => {
let kind = auth_error.kind();
tracing::warn!(
provider = %provider.as_str(),
auth_kind = %kind.as_str(),
"provider authentication prevented session materialization"
);
SessionError::provider_auth_failure(
meerkat_core::service::SessionProviderAuthFailure {
kind,
provider,
realm_id: auth_binding.map(|binding| binding.realm.clone()),
binding_id: auth_binding.map(|binding| binding.binding.clone()),
},
)
}
None => SessionError::build_llm_identity_unresolvable(format!(
"provider authentication failed: {}",
auth_error.kind().as_str()
)),
},
BuildAgentError::LlmClient(error) => {
SessionError::build_llm_identity_unresolvable(error.to_string())
}
other => SessionError::Agent(meerkat_core::error::AgentError::BuildError(
other.to_string(),
)),
}
}
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> {
self.agent.send_comms(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 {
fn validate_live_bridge_member_eligibility(
&self,
) -> Result<(), meerkat_core::error::AgentError> {
self.agent.validate_noncommitting_live_bridge_eligibility()
}
fn validate_live_bridge_operation(
&self,
request: &meerkat_session::LiveBridgeSessionOperationRequest,
) -> Result<(), meerkat_core::error::AgentError> {
self.agent.validate_noncommitting_live_bridge_policy()?;
if request.operation_id.as_ref() != request.dispatch_admission.operation_id() {
return Err(meerkat_core::error::AgentError::ConfigError(
"live bridge operation id does not match its sealed dispatch admission".to_string(),
));
}
if request.snapshot.id() != self.agent.session().id() {
return Err(meerkat_core::error::AgentError::ConfigError(
"live bridge snapshot does not belong to this durable member session".to_string(),
));
}
let admitted_revision = request
.snapshot
.canonical_context_revision()
.map_err(|error| meerkat_core::error::AgentError::ConfigError(error.to_string()))?;
request.run_permit.validate_binding(
request.operation_id.as_ref(),
request.snapshot.id(),
&admitted_revision,
)?;
let current_revision = self
.agent
.session()
.canonical_context_revision()
.map_err(|error| meerkat_core::error::AgentError::ConfigError(error.to_string()))?;
if admitted_revision != current_revision {
return Err(meerkat_core::error::AgentError::ConfigError(
"live bridge snapshot revision is no longer the durable member head".to_string(),
));
}
Ok(())
}
fn prepare_live_bridge_operation(
&self,
request: meerkat_session::LiveBridgeSessionOperationRequest,
cancellation: watch::Receiver<bool>,
) -> Result<meerkat_session::LiveBridgePreparedSessionOperation, meerkat_core::error::AgentError>
{
self.agent.prepare_live_bridge_noncommitting(
request.snapshot,
request.semantic_request,
request.dispatch_admission,
request.run_permit,
cancellation,
)
}
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 reconcile_runtime_compaction_projections(
&mut self,
intents: &[meerkat_core::CompactionProjectionIntent],
) -> Result<(), meerkat_core::error::AgentError> {
self.agent
.reconcile_runtime_compaction_projections(intents)
.await
}
async fn settle_inflight_sticky_model_fallback(
&mut self,
) -> Result<(), meerkat_core::error::AgentError> {
self.agent.settle_inflight_sticky_model_fallback().await
}
async fn abort_uncommitted_compaction_projections(
&mut self,
) -> Result<(), meerkat_core::error::AgentError> {
self.agent.abort_uncommitted_compaction_projections().await
}
fn take_runtime_terminal_failure_witness(
&mut self,
) -> Result<Option<meerkat_core::TurnErrorMetadata>, meerkat_core::error::AgentError> {
self.agent.take_runtime_terminal_failure_witness()
}
async fn run_turn_with_events(
&mut self,
input: SessionAgentTurnInput,
event_tx: mpsc::Sender<AgentEvent>,
) -> Result<RunResult, meerkat_core::error::AgentError> {
if input.handling_mode != HandlingMode::Queue {
return Err(meerkat_core::error::AgentError::ConfigError(format!(
"handling_mode {:?} requires a runtime-backed surface; direct session-service path supports Queue only",
input.handling_mode,
)));
}
if input.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.set_runtime_execution_kind(input.execution_kind);
if input.typed_turn_appends.is_empty()
&& input.transcript_identity.is_none()
&& input.injected_context.is_empty()
{
self.agent.run_with_events(input.prompt, event_tx).await
} else {
self.agent
.run_with_events_and_typed_turn_appends(
input.prompt,
input.typed_turn_appends,
input.injected_context,
input.transcript_identity,
event_tx,
)
.await
}
}
async fn run_pending_with_events(
&mut self,
transcript_identity: Option<meerkat_core::types::TranscriptMessageIdentity>,
execution_kind: Option<meerkat_core::lifecycle::RuntimeExecutionKind>,
event_tx: mpsc::Sender<AgentEvent>,
) -> Result<RunResult, meerkat_core::error::AgentError> {
self.agent.set_runtime_execution_kind(execution_kind);
self.agent
.set_active_transcript_identity(transcript_identity);
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_turn_tool_overlay(
&mut self,
overlay: Option<TurnToolOverlay>,
) -> Result<(), meerkat_core::error::AgentError> {
self.agent
.set_turn_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 self
.agent
.apply_pending_callback_tool_results(results.clone())?
{
return Ok(());
}
if results.is_empty() {
return Ok(());
}
self.agent
.session_mut()
.push(Message::tool_results(results));
Ok(())
}
fn replace_client(
&mut self,
client: std::sync::Arc<dyn meerkat_core::AgentLlmClient>,
) -> Result<(), meerkat_core::error::AgentError> {
self.agent.replace_client(client)
}
fn hot_swap_llm_identity(
&mut self,
client: std::sync::Arc<dyn meerkat_core::AgentLlmClient>,
identity: meerkat_core::SessionLlmIdentity,
request_policy: meerkat_core::SessionLlmRequestPolicy,
) -> Result<(), meerkat_core::error::AgentError> {
self.agent
.hot_swap_llm_identity(client, identity, request_policy)
}
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 append_system_messages(
&mut self,
contents: Vec<String>,
) -> Result<(), meerkat_core::error::AgentError> {
for content in contents {
self.agent.session_mut().append_system_message(content);
}
Ok(())
}
fn append_system_message_control(
&mut self,
req: meerkat_core::AppendSystemContextRequest,
) -> Result<meerkat_core::service::AppendSystemContextStatus, meerkat_core::error::AgentError>
{
self.agent
.session_mut()
.append_system_message_idempotent(
req.content.render_text(),
req.source,
req.idempotency_key,
meerkat_core::types::message_timestamp_now(),
)
.map_err(|error| meerkat_core::error::AgentError::ConfigError(error.to_string()))
}
fn activate_instruction_control(
&mut self,
request: meerkat_core::InstructionActivationRequest,
) -> Result<meerkat_core::InstructionActivationMutation, meerkat_core::InstructionActivationError>
{
self.agent.session_mut().activate_instruction(request)
}
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 {
let authorized_state = self
.agent
.tool_scope()
.authorized_visibility_state()
.map_err(|err| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to authorize tool visibility state for persistence: {err}"
))
})?;
debug_assert_eq!(authorized_state.as_state(), &state);
self.agent
.session_mut()
.set_tool_visibility_state(authorized_state)
.map_err(|err| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to persist tool visibility state into session metadata: {err}"
))
})
} else {
let authorized_state = self
.agent
.tool_scope()
.authorized_visibility_state()
.map_err(|err| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to authorize default tool visibility state for persistence: {err}"
))
})?;
debug_assert_eq!(
authorized_state.as_state(),
&meerkat_core::SessionToolVisibilityState::default()
);
self.agent
.session_mut()
.set_tool_visibility_state(authorized_state)
.map_err(|err| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to persist default tool visibility state into session metadata: {err}"
))
})
}
}
async fn dispatch_external_tool_call(
&mut self,
call: meerkat_core::ToolCall,
) -> Result<meerkat_core::ops::ToolDispatchOutcome, meerkat_core::error::AgentError> {
self.agent.dispatch_external_tool_call(call).await
}
async fn dispatch_external_tool_call_with_timeout_policy(
&mut self,
call: meerkat_core::ToolCall,
timeout_policy: meerkat_core::ToolDispatchTimeoutPolicy,
) -> Result<meerkat_core::ops::ToolDispatchOutcome, meerkat_core::error::AgentError> {
self.agent
.dispatch_external_tool_call_with_timeout_policy(call, timeout_policy)
.await
}
#[cfg(all(feature = "session-store", not(target_arch = "wasm32")))]
fn sync_session_from_durable_snapshot(
&mut self,
mut session: Session,
) -> Result<(), meerkat_core::error::AgentError> {
if session.id() != self.agent.session().id() {
return Err(meerkat_core::error::AgentError::InternalError(format!(
"durable snapshot session id {} does not match live session {}",
session.id(),
self.agent.session().id()
)));
}
let deferred_turn_state = session.try_deferred_turn_state().map_err(|err| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to restore durable deferred-turn state during live session sync: {err}"
))
})?;
if let Some(state) = deferred_turn_state {
session.set_deferred_turn_state(state).map_err(|err| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to serialize restored durable deferred-turn state during live session sync: {err}"
))
})?;
}
let visibility_state = session
.tool_visibility_state()
.map_err(|err| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to decode durable tool visibility state during live session sync: {err}"
))
})?
.unwrap_or_default();
self.agent
.tool_scope()
.set_visibility_state(visibility_state)
.map_err(|err| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to synchronize live tool visibility state from durable session: {err}"
))
})?;
*self.agent.session_mut() = session;
Ok(())
}
fn cancel(&mut self) {
self.agent.cancel();
}
fn cancel_after_boundary_handle(&self) -> Option<meerkat_core::CancelAfterBoundarySender> {
Some(self.agent.cancel_after_boundary_handle())
}
fn turn_state_handle(&self) -> Option<Arc<dyn meerkat_core::TurnStateHandle>> {
self.agent.turn_state_handle()
}
fn session_context_handle(
&self,
) -> Option<Arc<dyn meerkat_core::handles::SessionContextHandle>> {
self.session_context.as_ref().map(Arc::clone)
}
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 execution_snapshot(
&self,
) -> Result<Option<meerkat_core::AgentExecutionSnapshot>, meerkat_core::SnapshotProjectionError>
{
self.agent.execution_snapshot()
}
fn tool_scope_snapshot(&self) -> Option<meerkat_core::ToolScopeSnapshot> {
self.agent.tool_scope_snapshot()
}
fn visible_tool_defs(&self) -> Vec<meerkat_core::ToolDef> {
self.agent.visible_tool_defs()
}
fn external_tool_surface_snapshot(&self) -> Option<meerkat_core::ExternalToolSurfaceSnapshot> {
self.agent.external_tool_surface_snapshot()
}
fn session_clone(&self) -> Result<Session, meerkat_core::error::AgentError> {
Ok(self.agent.session().clone())
}
fn session_transcript_authority(
&self,
) -> Result<SessionTranscriptAuthoritySnapshot, meerkat_core::error::AgentError> {
SessionTranscriptAuthoritySnapshot::from_session(self.agent.session())
}
async fn prepare_head_canonical_runtime_boundary(
&mut self,
request: HeadCanonicalRuntimeBoundaryPrepareRequest,
) -> Result<PreparedHeadCanonicalRuntimeBoundary, meerkat_core::error::AgentError> {
let (authority, observed_head, deferred_turn_state, request_projection_token, blob_store) =
request.into_parts();
let observed_head_token = observed_head
.as_ref()
.map(meerkat_core::session_store::session_head_cas_token)
.transpose()
.map_err(|error| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to bind observed head for head-canonical preparation: {error}"
))
})?;
if let Some(pending) = self.pending_head_canonical_boundary.as_ref() {
if pending.authority == authority
&& pending.observed_head_token == observed_head_token
&& pending.request_projection_token == request_projection_token
{
return Ok(pending.prepared.clone());
}
return Err(meerkat_core::error::AgentError::InternalError(
"head-canonical re-prepare does not match the unacknowledged exact request projection"
.to_string(),
));
}
let suffix_start_u64 = observed_head.as_ref().map_or(0, |head| head.message_count);
let suffix_start = usize::try_from(suffix_start_u64).map_err(|_| {
meerkat_core::error::AgentError::InternalError(
"head-canonical predecessor message count exceeds the host index range".to_string(),
)
})?;
self.agent
.session_mut()
.externalize_media(blob_store.as_ref(), suffix_start)
.await
.map_err(|error| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to externalize head-canonical session suffix: {error}"
))
})?;
self.agent
.session_mut()
.set_deferred_turn_state(deferred_turn_state)
.map_err(|error| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to project deferred-turn state into head-canonical successor: {error}"
))
})?;
let mutation = match &authority {
HeadCanonicalRuntimeBoundaryAuthority::Root => {
if observed_head.is_some() {
return Err(meerkat_core::error::AgentError::InternalError(
"head-canonical root preparation observed an existing physical head"
.to_string(),
));
}
meerkat_core::session_store::PreparedHeadCanonicalMutation::prepare_root(
self.agent.session(),
)
.map(
meerkat_core::lifecycle::core_executor::PreparedHeadCanonicalPhysicalMutation::from,
)
}
HeadCanonicalRuntimeBoundaryAuthority::Successor {
boundary_head, ..
} => {
let observed_head = observed_head.ok_or_else(|| {
meerkat_core::error::AgentError::InternalError(
"head-canonical successor preparation is missing its physical predecessor head"
.to_string(),
)
})?;
let rewrite_required =
meerkat_core::session_store::PreparedHeadCanonicalRewriteMutation::is_required(
self.agent.session(),
&observed_head,
)
.map_err(|error| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to classify head-canonical physical mutation: {error}"
))
})?;
if rewrite_required {
meerkat_core::session_store::PreparedHeadCanonicalRewriteMutation::prepare_successor(
self.agent.session(),
boundary_head,
observed_head,
)
.map(
meerkat_core::lifecycle::core_executor::PreparedHeadCanonicalPhysicalMutation::from,
)
} else {
meerkat_core::session_store::PreparedHeadCanonicalMutation::prepare(
self.agent.session(),
Some(observed_head),
)
.map(
meerkat_core::lifecycle::core_executor::PreparedHeadCanonicalPhysicalMutation::from,
)
}
}
}
.map_err(|error| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to prepare head-canonical boundary: {error}"
))
})?;
let committed =
meerkat_core::lifecycle::core_executor::BoundSessionCommit::head_canonical_physical_from_session(
self.agent.session(),
mutation,
)
.map_err(|error| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to seal head-canonical successor boundary: {error}"
))
})?;
let prepared = PreparedHeadCanonicalRuntimeBoundary::new(committed).map_err(|error| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to bind actor-prepared head-canonical boundary: {error}"
))
})?;
self.pending_head_canonical_boundary = Some(PendingFactoryHeadCanonicalBoundary {
authority,
observed_head_token,
request_projection_token,
prepared: prepared.clone(),
});
self.acknowledged_head_canonical_boundary = None;
Ok(prepared)
}
fn acknowledge_head_canonical_runtime_boundary(
&mut self,
successor_head_token: &str,
) -> Result<HeadCanonicalRuntimeBoundaryAcknowledgeOutcome, meerkat_core::error::AgentError>
{
let (carrier, outcome) = if let Some(pending) =
self.pending_head_canonical_boundary.as_ref()
{
(
pending,
HeadCanonicalRuntimeBoundaryAcknowledgeOutcome::Applied,
)
} else if let Some(acknowledged) = self.acknowledged_head_canonical_boundary.as_ref() {
(
acknowledged,
HeadCanonicalRuntimeBoundaryAcknowledgeOutcome::AlreadyAcknowledgedExact,
)
} else {
return Err(meerkat_core::error::AgentError::InternalError(
"factory agent has no prepared or acknowledged head-canonical boundary".to_string(),
));
};
let committed = carrier.prepared.committed().clone();
let boundary = committed.head_canonical().ok_or_else(|| {
meerkat_core::error::AgentError::InternalError(
"factory agent pending boundary is not head-canonical".to_string(),
)
})?;
if boundary.mutation().successor_head_token() != successor_head_token {
return Err(meerkat_core::error::AgentError::InternalError(
"head-canonical acknowledgement head token does not match the factory actor's successor"
.to_string(),
));
}
if outcome == HeadCanonicalRuntimeBoundaryAcknowledgeOutcome::AlreadyAcknowledgedExact {
return Ok(outcome);
}
boundary
.mutation()
.acknowledge_session(self.agent.session_mut(), successor_head_token)
.map_err(|error| {
meerkat_core::error::AgentError::InternalError(format!(
"failed to acknowledge committed head-canonical successor: {error}"
))
})?;
self.acknowledged_head_canonical_boundary = self.pending_head_canonical_boundary.take();
Ok(outcome)
}
fn durable_llm_identity(&self) -> Option<meerkat_core::SessionLlmIdentity> {
self.agent
.session()
.session_metadata()
.map(|metadata| metadata.llm_identity())
}
fn observed_session_tail(&self) -> ObservedSessionTailKind {
meerkat_core::pending_continuation::observe_session_tail(self.agent.session().messages())
}
fn append_external_user_content(
&mut self,
content: meerkat_core::types::ContentInput,
) -> Result<(), meerkat_core::error::AgentError> {
self.agent
.session_mut()
.append_external_user_content(content);
Ok(())
}
fn append_external_assistant_output(
&mut self,
blocks: Vec<AssistantBlock>,
stop_reason: StopReason,
usage: Usage,
) -> Result<(), meerkat_core::error::AgentError> {
let usage = meerkat_core::TurnUsage::try_from_usage(usage).map_err(|error| {
meerkat_core::error::AgentError::ConfigError(format!(
"external assistant usage requires normalized provider accounting: {error}"
))
})?;
self.agent.session_mut().append_external_assistant_blocks(
blocks,
stop_reason,
usage.clone(),
);
self.agent.budget().record_turn_usage(&usage);
Ok(())
}
fn append_realtime_transcript_event(
&mut self,
event: meerkat_core::RealtimeTranscriptEvent,
) -> Result<meerkat_core::RealtimeTranscriptApplyOutcome, meerkat_core::error::AgentError> {
let outcome = self
.agent
.session_mut()
.append_realtime_transcript_event(event);
for materialized in &outcome.materialized_messages {
if let meerkat_core::RealtimeTranscriptMaterializedMessage::Assistant {
usage: Some(usage),
..
} = materialized
{
self.agent.budget().record_turn_usage(usage);
}
}
Ok(outcome)
}
fn staged_realtime_assistant_segment_text(
&self,
response_id: &str,
item_id: &str,
content_index: u32,
) -> Option<String> {
self.agent.session().staged_realtime_assistant_segment_text(
response_id,
item_id,
content_index,
)
}
fn staged_realtime_assistant_segment_is_final(
&self,
response_id: &str,
item_id: &str,
content_index: u32,
) -> bool {
self.agent
.session()
.staged_realtime_assistant_segment_is_final(response_id, item_id, content_index)
}
fn admit_live_assistant_playback_target(
&mut self,
channel_id: &meerkat_core::LiveChannelId,
interaction_id: meerkat_core::InteractionId,
response_id: &str,
item_id: &str,
content_index: u32,
) -> Result<meerkat_core::LiveAssistantPlaybackTarget, meerkat_core::error::AgentError> {
self.agent
.session_mut()
.admit_live_assistant_playback_target(
channel_id,
interaction_id,
response_id,
item_id,
content_index,
)
}
fn live_assistant_playback_target(
&self,
channel_id: &meerkat_core::LiveChannelId,
item_id: &str,
content_index: u32,
) -> Option<meerkat_core::LiveAssistantPlaybackTarget> {
self.agent
.session()
.live_assistant_playback_target(channel_id, item_id, content_index)
}
fn live_assistant_playback_target_for_channel(
&self,
channel_id: &meerkat_core::LiveChannelId,
) -> Option<meerkat_core::LiveAssistantPlaybackTarget> {
self.agent
.session()
.live_assistant_playback_target_for_channel(channel_id)
}
fn resolve_live_assistant_playback_target(
&mut self,
channel_id: &meerkat_core::LiveChannelId,
interaction_id: meerkat_core::InteractionId,
response_id: &str,
item_id: &str,
content_index: u32,
) -> Result<(), meerkat_core::error::AgentError> {
self.agent
.session_mut()
.resolve_live_assistant_playback_target(
channel_id,
interaction_id,
response_id,
item_id,
content_index,
)
}
fn transient_turn_context_state(&self) -> meerkat_core::TransientTurnContextStateHandle {
self.agent.transient_turn_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()
}
fn observed_comms_sender(&self) -> Option<Arc<meerkat_core::ObservedCommsSender>> {
self.agent.observed_comms_sender()
}
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Clone)]
pub struct RealmInheritance {
source: Arc<dyn meerkat_core::RealmConfigSource>,
head: meerkat_core::connection::RealmId,
}
#[cfg(not(target_arch = "wasm32"))]
impl RealmInheritance {
pub fn new(
source: Arc<dyn meerkat_core::RealmConfigSource>,
head: meerkat_core::connection::RealmId,
) -> Self {
Self { source, head }
}
pub fn head(&self) -> &meerkat_core::connection::RealmId {
&self.head
}
pub async fn compose_over(
&self,
head_config: Config,
) -> Result<Config, meerkat_core::ConfigError> {
meerkat_core::EffectiveConfigReader::new(self.source.clone())
.effective_config_over_head(&self.head, head_config)
.await
}
}
#[derive(Clone)]
pub struct FactoryAgentBuilder {
factory: AgentFactory,
config_snapshot: Config,
#[cfg(not(target_arch = "wasm32"))]
config_store: Option<Arc<dyn ConfigStore>>,
#[cfg(not(target_arch = "wasm32"))]
pub realm_inheritance: Arc<std::sync::RwLock<Option<RealmInheritance>>>,
pub default_llm_client: Option<Arc<dyn LlmClient>>,
pub default_agent_llm_client_decorator:
Arc<std::sync::RwLock<Option<meerkat_core::AgentLlmClientDecorator>>>,
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>>>>,
pub default_workgraph_tools:
Arc<std::sync::RwLock<Option<Arc<dyn meerkat_core::AgentToolDispatcher>>>>,
pub default_workgraph_namespace_grant:
Arc<std::sync::RwLock<Option<meerkat_core::service::WorkGraphNamespaceGrant>>>,
pub default_blob_store: Option<Arc<dyn meerkat_core::BlobStore>>,
pub default_realm_id: Option<meerkat_core::RealmId>,
#[cfg(not(target_arch = "wasm32"))]
pub default_detached_job_store: Option<Arc<dyn meerkat_jobs::DetachedJobStore>>,
#[cfg(not(target_arch = "wasm32"))]
pub default_shell_job_delivery_projector: Option<crate::JobOutboxProjector>,
pub default_image_generation_executor:
Option<Arc<dyn meerkat_llm_core::ImageGenerationExecutor>>,
}
impl FactoryAgentBuilder {
pub fn new(factory: AgentFactory, config: Config) -> Self {
Self {
factory,
config_snapshot: config,
#[cfg(not(target_arch = "wasm32"))]
config_store: None,
#[cfg(not(target_arch = "wasm32"))]
realm_inheritance: Arc::new(std::sync::RwLock::new(None)),
default_llm_client: None,
default_agent_llm_client_decorator: Arc::new(std::sync::RwLock::new(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)),
default_workgraph_tools: Arc::new(std::sync::RwLock::new(None)),
default_workgraph_namespace_grant: Arc::new(std::sync::RwLock::new(None)),
default_blob_store: None,
default_realm_id: None,
#[cfg(not(target_arch = "wasm32"))]
default_detached_job_store: None,
#[cfg(not(target_arch = "wasm32"))]
default_shell_job_delivery_projector: None,
default_image_generation_executor: 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),
realm_inheritance: Arc::new(std::sync::RwLock::new(None)),
default_llm_client: None,
default_agent_llm_client_decorator: Arc::new(std::sync::RwLock::new(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)),
default_workgraph_tools: Arc::new(std::sync::RwLock::new(None)),
default_workgraph_namespace_grant: Arc::new(std::sync::RwLock::new(None)),
default_blob_store: None,
default_realm_id: None,
#[cfg(not(target_arch = "wasm32"))]
default_detached_job_store: None,
#[cfg(not(target_arch = "wasm32"))]
default_shell_job_delivery_projector: None,
default_image_generation_executor: None,
}
}
#[cfg(not(target_arch = "wasm32"))]
pub fn with_realm_inheritance(
self,
source: Arc<dyn meerkat_core::RealmConfigSource>,
head: meerkat_core::connection::RealmId,
) -> Self {
*self
.realm_inheritance
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) =
Some(RealmInheritance::new(source, head));
self
}
pub fn with_image_generation_machine(
mut self,
machine: Arc<dyn meerkat_tools::builtin::image_generation::ImageGenerationMachine>,
) -> Self {
self.factory = self.factory.with_image_generation_machine(machine);
self
}
async fn resolve_config(&self) -> Result<Config, SessionError> {
#[cfg(not(target_arch = "wasm32"))]
let head_config = if let Some(store) = &self.config_store {
store.get().await.map_err(|err| {
SessionError::Agent(meerkat_core::error::AgentError::ConfigError(format!(
"failed to read latest config from store: {err}"
)))
})?
} else {
self.config_snapshot.clone()
};
#[cfg(target_arch = "wasm32")]
let head_config = self.config_snapshot.clone();
#[cfg(not(target_arch = "wasm32"))]
{
let inheritance = self
.realm_inheritance
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
if let Some(inheritance) = inheritance {
let reader = meerkat_core::EffectiveConfigReader::new(inheritance.source.clone());
return reader
.effective_config_over_head(&inheritance.head, head_config)
.await
.map_err(|err| {
SessionError::Agent(meerkat_core::error::AgentError::ConfigError(format!(
"failed to compose effective config for realm '{}': {err}",
inheritance.head
)))
});
}
}
Ok(head_config)
}
pub fn factory(&self) -> &AgentFactory {
&self.factory
}
#[cfg(all(feature = "session-store", not(target_arch = "wasm32")))]
pub(crate) fn runtime_config_store(&self) -> Arc<dyn ConfigStore> {
self.config_store.clone().unwrap_or_else(|| {
Arc::new(meerkat_core::MemoryConfigStore::new(
self.config_snapshot.clone(),
meerkat_models::canonical(),
))
})
}
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 abort_absent_session_compaction_stages(
&self,
session_id: &meerkat_core::SessionId,
) -> Result<(), SessionError> {
#[cfg(all(feature = "memory-store-session", not(target_arch = "wasm32")))]
{
use meerkat_core::memory::{MemoryOwner, MemoryStore};
let memory_dir = self.factory.store_path.join("memory");
let memory_store_exists = memory_dir.try_exists().map_err(|error| {
SessionError::Agent(meerkat_core::error::AgentError::InternalError(format!(
"failed to inspect canonical memory store at {} before session materialization: {error}",
memory_dir.display()
)))
})?;
if !memory_store_exists {
return Ok(());
}
let store = meerkat_memory::HnswMemoryStore::open(&memory_dir).map_err(|error| {
SessionError::Agent(meerkat_core::error::AgentError::InternalError(format!(
"failed to open canonical memory store at {} before session materialization: {error}",
memory_dir.display()
)))
})?;
let receipt = store
.reconcile_compaction_stages(
&MemoryOwner::canonical_session(session_id.clone()),
&[],
)
.await
.map_err(|error| {
SessionError::Agent(meerkat_core::error::AgentError::InternalError(format!(
"failed to reconcile empty runtime compaction authority for absent session {session_id}: {error}"
)))
})?;
tracing::debug!(
%session_id,
aborted_orphans = receipt.aborted_orphans,
retained_committed = receipt.retained_committed,
"reconciled durable compaction stages before SessionTask materialization"
);
return Ok(());
}
#[cfg(not(all(feature = "memory-store-session", not(target_arch = "wasm32"))))]
{
let _ = session_id;
Ok(())
}
}
async fn model_supports_inline_video(
&self,
identity: &meerkat_core::SessionLlmIdentity,
) -> Option<bool> {
self.resolve_config()
.await
.ok()?
.model_registry(meerkat_models::canonical())
.ok()
.and_then(|registry| registry.profile_for_provider(identity.provider, &identity.model))
.map(|profile| profile.inline_video)
}
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.agent_llm_client_decorator.is_none()
&& let Some(decorator) = self
.default_agent_llm_client_decorator
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
{
build_config.agent_llm_client_decorator = Some(decorator);
}
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);
}
if build_config.workgraph_tools.is_none()
&& let Some(workgraph_dispatcher) = self
.default_workgraph_tools
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
{
build_config.workgraph_tools = Some(workgraph_dispatcher);
let default_grant = self
.default_workgraph_namespace_grant
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
if let (Some(requested), Some(default)) =
(&build_config.workgraph_namespace_grant, &default_grant)
&& requested != default
{
return Err(SessionError::Agent(
meerkat_core::error::AgentError::ConfigError(format!(
"requested WorkGraph namespace grant {}/{} does not match the injected dispatcher grant {}/{}",
requested.realm_id,
requested.namespace,
default.realm_id,
default.namespace
)),
));
}
build_config.workgraph_namespace_grant = default_grant;
}
if build_config.blob_store_override.is_none()
&& let Some(blob_store) = self.default_blob_store.clone()
{
build_config.blob_store_override = Some(blob_store);
}
if build_config.realm_id.is_none()
&& let Some(realm_id) = self.default_realm_id.clone()
{
build_config.realm_id = Some(realm_id);
}
#[cfg(not(target_arch = "wasm32"))]
if build_config.detached_job_store_override.is_none()
&& let Some(job_store) = self.default_detached_job_store.clone()
{
build_config.detached_job_store_override = Some(job_store);
}
#[cfg(not(target_arch = "wasm32"))]
if build_config.shell_job_delivery_projector_override.is_none()
&& let Some(projector) = self.default_shell_job_delivery_projector.clone()
{
let projector = match build_config.realm_id.as_ref() {
Some(realm_id) => projector.bound_to_realm(realm_id.to_string()),
None => projector,
};
build_config.shell_job_delivery_projector_override = Some(Arc::new(projector));
}
if build_config.image_generation_executor_override.is_none()
&& let Some(executor) = self.default_image_generation_executor.clone()
{
build_config.image_generation_executor_override = Some(executor);
}
let config = self.resolve_config().await?;
let session_context = match req.build.as_ref().map(|opts| &opts.runtime_build_mode) {
Some(meerkat_core::RuntimeBuildMode::SessionOwned(bindings)) => {
Some(Arc::clone(bindings.session_context()))
}
_ => None,
};
let provider_error_context = build_config.provider.or_else(|| {
config
.model_registry(meerkat_models::canonical())
.ok()
.and_then(|registry| {
registry
.entry(&build_config.model)
.map(|entry| entry.provider)
})
});
let auth_binding_error_context = provider_error_context.and_then(|provider| {
meerkat_core::resolve_auth_binding_candidates_for_provider(
&config,
provider,
build_config.auth_binding.as_ref(),
build_config.realm_id.as_ref(),
true,
)
.ok()
.and_then(|candidates| candidates.into_iter().next())
.filter(|target| !target.auth_binding.is_env_default())
.map(|target| target.auth_binding)
});
let factory = self.factory.clone();
let agent = meerkat_runtime::stack_relief::relieve_caller_stack(move || async move {
factory.build_agent(build_config, &config).await
})
.await
.map_err(|error| {
build_agent_error_to_session_error(
error,
provider_error_context,
auth_binding_error_context.as_ref(),
)
})?;
Ok(FactoryAgent {
agent,
session_context,
pending_head_canonical_boundary: None,
acknowledged_head_canonical_boundary: None,
})
}
}
pub fn build_ephemeral_service(
factory: AgentFactory,
config: Config,
max_sessions: usize,
) -> EphemeralSessionService<FactoryAgentBuilder> {
let builder = FactoryAgentBuilder::new(factory, config);
crate::surface::build_embedded_service_from_builder(builder, max_sessions)
}
#[cfg(feature = "session-store")]
fn set_default_workgraph_tools_from_persistence(
builder: &FactoryAgentBuilder,
persistence: &PersistenceBundle,
) {
#[cfg(not(target_arch = "wasm32"))]
if let Some(manifest) = persistence.manifest() {
let service = crate::WorkGraphService::with_scope(
persistence.workgraph_store(),
manifest.realm.as_str().to_owned(),
crate::WorkNamespace::default(),
);
crate::surface::set_default_workgraph_namespace_grant(
builder,
Some(service.namespace_grant().clone()),
);
crate::surface::set_default_workgraph_tools(
builder,
Some(Arc::new(crate::WorkGraphToolSurface::new(service))),
);
}
}
#[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::MeerkatMachine>,
) {
let builder = FactoryAgentBuilder::new(factory, config);
set_default_workgraph_tools_from_persistence(&builder, &persistence);
crate::surface::build_runtime_backed_service(builder, max_sessions, persistence)
}
#[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, SessionService};
use meerkat_core::{
Provider, ToolCallView, ToolDef, ToolDispatchOutcome, ToolError, ToolResult,
};
use meerkat_llm_core::{
ImageGenerationExecutor, ProviderImageGenerationOutput, ProviderImageGenerationRequest,
};
use meerkat_runtime::MeerkatMachine;
use meerkat_schedule::{MemoryScheduleStore, ScheduleService, ScheduleToolDispatcher};
use meerkat_session::ephemeral::SessionAgent;
use meerkat_store::MemoryBlobStore;
use std::pin::Pin;
use std::sync::Mutex;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use tempfile::TempDir;
#[cfg(feature = "experimental-gpt-live")]
struct AllowLiveBridgeModelOnly;
#[cfg(feature = "experimental-gpt-live")]
#[async_trait]
impl meerkat_core::ToolDispatchAdmission for AllowLiveBridgeModelOnly {
async fn await_dispatch_admission(
&self,
_call: ToolCallView<'_>,
_context: Option<&meerkat_core::ToolDispatchContext>,
_effect_kind: meerkat_core::LiveBridgeEffectKind,
) -> Result<(), ToolError> {
Ok(())
}
}
#[test]
fn llm_client_build_failure_keeps_typed_session_cause() {
let error = build_agent_error_to_session_error(
BuildAgentError::LlmClient(meerkat_client::FactoryError::ClientCreationFailed(
"backend diagnostic that callers must not classify by text".to_string(),
)),
Some(Provider::OpenAI),
None,
);
assert!(error.is_build_llm_identity_unresolvable());
}
#[test]
fn provider_auth_build_failure_keeps_typed_session_payload() {
let binding = meerkat_core::AuthBindingRef {
realm: meerkat_core::RealmId::parse("global").unwrap(),
binding: meerkat_core::BindingId::parse("openai").unwrap(),
profile: None,
origin: meerkat_core::BindingOrigin::Configured,
};
let error = build_agent_error_to_session_error(
BuildAgentError::LlmClient(meerkat_client::FactoryError::ProviderAuth(
meerkat_llm_core::provider_runtime::ProviderAuthError::Auth(
meerkat_core::AuthError::InteractiveLoginRequired,
),
)),
Some(Provider::OpenAI),
Some(&binding),
);
assert!(error.is_build_llm_identity_unresolvable());
assert_eq!(
error.provider_auth_failure_data(),
Some(meerkat_core::service::SessionProviderAuthFailure {
kind: meerkat_core::AuthErrorKind::InteractiveLoginRequired,
provider: Provider::OpenAI,
realm_id: Some(meerkat_core::RealmId::parse("global").unwrap()),
binding_id: Some(meerkat_core::BindingId::parse("openai").unwrap()),
})
);
let data = error.structured_data().expect("structured auth failure");
assert_eq!(data["cause"], "provider_auth");
assert_eq!(data["kind"], "interactive_login_required");
assert_eq!(data["provider"], "openai");
assert_eq!(data["realm_id"], "global");
assert_eq!(data["binding_id"], "openai");
}
#[test]
fn provider_auth_public_message_redacts_provider_diagnostics() {
let error = build_agent_error_to_session_error(
BuildAgentError::LlmClient(meerkat_client::FactoryError::ProviderAuth(
meerkat_llm_core::provider_runtime::ProviderAuthError::Auth(
meerkat_core::AuthError::RefreshFailed(
"sentinel-token-endpoint-body".to_string(),
),
),
)),
Some(Provider::OpenAI),
None,
);
assert!(!error.to_string().contains("sentinel-token-endpoint-body"));
assert_eq!(
error
.provider_auth_failure_data()
.map(|failure| failure.kind),
Some(meerkat_core::AuthErrorKind::RefreshFailed)
);
}
struct MockLlmClient {
delta: &'static str,
}
#[cfg(feature = "experimental-gpt-live")]
struct MidRunCancellationClient {
calls: AtomicUsize,
saw_member_tool: AtomicBool,
seen_tools: Mutex<Vec<String>>,
}
#[cfg(feature = "experimental-gpt-live")]
struct BlockingMemberToolDispatcher {
started: Arc<tokio::sync::Notify>,
dispatches: AtomicUsize,
tools: Arc<[Arc<ToolDef>]>,
}
#[cfg(feature = "experimental-gpt-live")]
struct CountingAgentSessionStore {
saves: AtomicUsize,
}
#[cfg(feature = "experimental-gpt-live")]
struct CountingSessionCheckpointer {
checkpoints: AtomicUsize,
}
#[cfg(feature = "experimental-gpt-live")]
struct CountingHookEngine {
calls: AtomicUsize,
}
#[cfg(feature = "experimental-gpt-live")]
struct PolicyProbeDispatcher {
tools: Arc<[Arc<ToolDef>]>,
catalog: Arc<[meerkat_core::ToolCatalogEntry]>,
exact_catalog: bool,
dispatches: AtomicUsize,
}
fn provider_for_successful_test_model(model: &str) -> Provider {
if model.starts_with("claude-") {
Provider::Anthropic
} else if model.starts_with("gpt-") || model.starts_with("o1-") {
Provider::OpenAI
} else if model.starts_with("gemini-") {
Provider::Gemini
} else {
Provider::Other
}
}
impl Default for MockLlmClient {
fn default() -> Self {
Self { delta: "ok" }
}
}
#[cfg(feature = "experimental-gpt-live")]
#[async_trait]
impl LlmClient for MidRunCancellationClient {
fn project_replay_messages(
&self,
messages: &[meerkat_core::Message],
) -> Result<Vec<meerkat_core::Message>, meerkat_client::LlmError> {
Ok(messages.to_vec())
}
fn stream<'a>(
&'a self,
request: &'a LlmRequest,
) -> Pin<
Box<dyn futures::Stream<Item = Result<LlmEvent, meerkat_client::LlmError>> + Send + 'a>,
> {
if self.calls.fetch_add(1, Ordering::SeqCst) == 1 {
*self.seen_tools.lock().expect("seen tools lock") = request
.tools
.iter()
.map(|tool| tool.name.to_string())
.collect();
self.saw_member_tool.store(
request
.tools
.iter()
.any(|tool| tool.name == "member_lookup"),
Ordering::SeqCst,
);
return Box::pin(stream::iter(vec![
Ok(LlmEvent::ToolCallComplete {
id: "live-bridge-member-call".to_string(),
name: "member_lookup".to_string(),
args: serde_json::json!({}),
meta: None,
}),
Ok(LlmEvent::UsageUpdate {
usage: meerkat_core::TurnUsage::host_declared(
provider_for_successful_test_model(&request.model),
&request.model,
meerkat_core::Usage {
input_tokens: 7,
..meerkat_core::Usage::default()
},
),
}),
Ok(LlmEvent::Done {
outcome: LlmDoneOutcome::Success {
stop_reason: meerkat_core::StopReason::ToolUse,
},
}),
]));
}
Box::pin(stream::iter(vec![
Ok(LlmEvent::TextDelta {
delta: "ok".to_string(),
meta: None,
}),
Ok(LlmEvent::UsageUpdate {
usage: meerkat_core::TurnUsage::host_declared(
provider_for_successful_test_model(&request.model),
&request.model,
meerkat_core::Usage::default(),
),
}),
Ok(LlmEvent::Done {
outcome: LlmDoneOutcome::Success {
stop_reason: meerkat_core::StopReason::EndTurn,
},
}),
]))
}
fn provider(&self) -> meerkat_core::Provider {
meerkat_core::Provider::Anthropic
}
async fn health_check(&self) -> Result<(), meerkat_client::LlmError> {
Ok(())
}
}
#[cfg(feature = "experimental-gpt-live")]
#[async_trait]
impl meerkat_core::AgentToolDispatcher for BlockingMemberToolDispatcher {
fn tools(&self) -> Arc<[Arc<ToolDef>]> {
Arc::clone(&self.tools)
}
async fn dispatch(
&self,
_call: ToolCallView<'_>,
) -> Result<ToolDispatchOutcome, ToolError> {
self.dispatches.fetch_add(1, Ordering::SeqCst);
self.started.notify_one();
futures::future::pending().await
}
fn live_bridge_effect_kind(&self, _tool_name: &str) -> meerkat_core::LiveBridgeEffectKind {
meerkat_core::LiveBridgeEffectKind::ToolDispatch
}
}
#[cfg(feature = "experimental-gpt-live")]
#[async_trait]
impl meerkat_core::HookEngine for CountingHookEngine {
async fn execute(
&self,
_invocation: meerkat_core::HookInvocation,
_overrides: Option<&meerkat_core::config::HookRunOverrides>,
) -> Result<meerkat_core::HookExecutionReport, meerkat_core::HookEngineError> {
self.calls.fetch_add(1, Ordering::SeqCst);
Ok(meerkat_core::HookExecutionReport::empty())
}
}
#[cfg(feature = "experimental-gpt-live")]
#[async_trait]
impl meerkat_core::AgentToolDispatcher for PolicyProbeDispatcher {
fn tools(&self) -> Arc<[Arc<ToolDef>]> {
Arc::clone(&self.tools)
}
fn tool_catalog_capabilities(&self) -> meerkat_core::ToolCatalogCapabilities {
meerkat_core::ToolCatalogCapabilities {
exact_catalog: self.exact_catalog,
may_require_catalog_control_plane: false,
}
}
fn tool_catalog(&self) -> Arc<[meerkat_core::ToolCatalogEntry]> {
Arc::clone(&self.catalog)
}
async fn dispatch(
&self,
_call: ToolCallView<'_>,
) -> Result<ToolDispatchOutcome, ToolError> {
self.dispatches.fetch_add(1, Ordering::SeqCst);
Err(ToolError::execution_failed(
"policy probe dispatcher must never run",
))
}
}
#[cfg(feature = "experimental-gpt-live")]
#[async_trait]
impl meerkat_core::AgentSessionStore for CountingAgentSessionStore {
async fn save(
&self,
_session: &meerkat_core::Session,
) -> Result<(), meerkat_core::AgentError> {
self.saves.fetch_add(1, Ordering::SeqCst);
Ok(())
}
async fn load(
&self,
_id: &str,
) -> Result<Option<meerkat_core::Session>, meerkat_core::AgentError> {
Ok(None)
}
}
#[cfg(feature = "experimental-gpt-live")]
#[async_trait]
impl meerkat_core::SessionCheckpointer for CountingSessionCheckpointer {
async fn checkpoint_run(
&self,
_session: &mut meerkat_core::Session,
_run_id: &meerkat_core::RunId,
_previous: Option<&meerkat_core::RunCheckpointReceipt>,
) -> Result<Option<meerkat_core::RunCheckpointReceipt>, meerkat_core::AgentError> {
self.checkpoints.fetch_add(1, Ordering::SeqCst);
Ok(None)
}
}
struct CaptureToolClient {
inner: meerkat_client::TestClient,
seen_tools: Mutex<Vec<String>>,
}
impl Default for CaptureToolClient {
fn default() -> Self {
Self {
inner: meerkat_client::TestClient::for_provider(Provider::Anthropic),
seen_tools: Mutex::new(Vec::new()),
}
}
}
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 project_replay_messages(
&self,
messages: &[meerkat_core::Message],
) -> Result<Vec<meerkat_core::Message>, meerkat_client::LlmError> {
Ok(messages.to_vec())
}
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.to_string())
.collect();
self.inner.stream(request)
}
fn provider(&self) -> meerkat_core::Provider {
self.inner.provider()
}
async fn health_check(&self) -> Result<(), meerkat_client::LlmError> {
self.inner.health_check().await
}
}
struct FakeImageGenerationExecutor;
#[cfg(not(target_arch = "wasm32"))]
fn self_hosted_inline_video_config(inline_video: bool) -> Config {
let mut config = Config::default();
config.self_hosted.servers.insert(
"local".to_string(),
meerkat_core::SelfHostedServerConfig {
transport: meerkat_core::SelfHostedTransport::OpenAiCompatible,
base_url: "http://127.0.0.1:11434".to_string(),
api_style: meerkat_core::SelfHostedApiStyle::Responses,
},
);
config.self_hosted.models.insert(
"video-alias".to_string(),
meerkat_core::SelfHostedModelConfig {
server: "local".to_string(),
remote_model: "video-model".to_string(),
display_name: "Video Alias".to_string(),
family: "video-family".to_string(),
tier: meerkat_core::model_profile::catalog::ModelTier::Supported,
context_window: None,
max_output_tokens: None,
vision: true,
image_tool_results: true,
inline_video,
supports_temperature: true,
supports_thinking: false,
supports_reasoning: false,
supports_web_search: false,
call_timeout_secs: None,
},
);
config
}
#[cfg(not(target_arch = "wasm32"))]
#[tokio::test]
async fn inline_video_capability_reads_current_config_store() {
let initial_config = self_hosted_inline_video_config(false);
let current_config = self_hosted_inline_video_config(true);
let store: Arc<dyn meerkat_core::ConfigStore> = Arc::new(
meerkat_core::MemoryConfigStore::new(current_config, meerkat_models::canonical()),
);
let builder = FactoryAgentBuilder::new_with_config_store(
AgentFactory::minimal(),
initial_config,
store,
);
let identity = meerkat_core::SessionLlmIdentity {
model: "video-alias".to_string(),
provider: Provider::SelfHosted,
self_hosted_server_id: Some("local".to_string()),
provider_params: None,
auth_binding: None,
};
assert_eq!(
builder.model_supports_inline_video(&identity).await,
Some(true)
);
}
#[cfg(not(target_arch = "wasm32"))]
struct FailingConfigStore;
#[cfg(not(target_arch = "wasm32"))]
#[async_trait]
impl meerkat_core::ConfigStore for FailingConfigStore {
async fn get(&self) -> Result<Config, meerkat_core::config::ConfigError> {
Err(meerkat_core::config::ConfigError::InternalError(
"store unavailable".to_string(),
))
}
async fn set(&self, _config: Config) -> Result<(), meerkat_core::config::ConfigError> {
Ok(())
}
async fn patch(
&self,
_delta: meerkat_core::config::ConfigDelta,
) -> Result<Config, meerkat_core::config::ConfigError> {
Err(meerkat_core::config::ConfigError::InternalError(
"store unavailable".to_string(),
))
}
}
#[cfg(not(target_arch = "wasm32"))]
#[tokio::test]
async fn build_agent_fails_closed_on_config_store_read_error() {
let initial_config = self_hosted_inline_video_config(false);
let store: Arc<dyn meerkat_core::ConfigStore> = Arc::new(FailingConfigStore);
let mut builder = FactoryAgentBuilder::new_with_config_store(
AgentFactory::minimal(),
initial_config,
store,
);
builder.default_llm_client = Some(Arc::new(MockLlmClient::default()));
let (event_tx, _event_rx) = mpsc::channel(8);
let req = CreateSessionRequest {
injected_context: Vec::new(),
model: "video-alias".to_string(),
prompt: "fail closed on stale config".to_string().into(),
system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
max_tokens: None,
event_tx: None,
initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::Discard,
build: Some(SessionBuildOptions::default()),
labels: None,
};
let result = builder.build_agent(&req, event_tx).await;
let Err(err) = result else {
panic!("build must fail closed on config-store read error");
};
assert!(
err.to_string()
.contains("failed to read latest config from store"),
"expected a typed config-store error, got: {err}"
);
}
#[cfg(not(target_arch = "wasm32"))]
#[tokio::test]
async fn session_service_live_identity_matches_factory_resolved_metadata() -> 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, self_hosted_inline_video_config(true));
builder.default_llm_client = Some(Arc::new(MockLlmClient::default()));
let service = EphemeralSessionService::new(builder, 10);
let result = service
.create_session(CreateSessionRequest {
injected_context: Vec::new(),
model: "video-alias".to_string(),
prompt: "defer identity parity".to_string().into(),
system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
max_tokens: None,
event_tx: None,
initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::Discard,
build: Some(SessionBuildOptions::default()),
labels: None,
})
.await
.map_err(|err| err.to_string())?;
let live_identity = service
.live_session_llm_identity(&result.session_id)
.await
.map_err(|err| err.to_string())?;
assert_eq!(live_identity.model, "video-alias");
assert_eq!(live_identity.provider, Provider::SelfHosted);
assert_eq!(
live_identity.self_hosted_server_id.as_deref(),
Some("local")
);
let exported = service
.export_session(&result.session_id)
.await
.map_err(|err| err.to_string())?;
let metadata = exported
.session_metadata()
.ok_or_else(|| "missing session metadata".to_string())?;
assert_eq!(metadata.llm_identity(), live_identity);
Ok(())
}
#[async_trait]
impl ImageGenerationExecutor for FakeImageGenerationExecutor {
async fn execute_image_generation(
&self,
request: ProviderImageGenerationRequest,
) -> Result<ProviderImageGenerationOutput, meerkat_llm_core::LlmError> {
Ok(ProviderImageGenerationOutput {
operation_id: request.operation_id,
terminal_observation:
meerkat_core::ImageProviderTerminalObservation::ExecutionFailed,
images: Vec::new(),
provider_text: None,
revised_prompt: meerkat_core::RevisedPromptDisposition::NotRequested,
native_metadata: meerkat_core::ProviderImageMetadata::NotEmitted,
warnings: Vec::new(),
})
}
}
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
impl LlmClient for MockLlmClient {
fn project_replay_messages(
&self,
messages: &[meerkat_core::Message],
) -> Result<Vec<meerkat_core::Message>, meerkat_client::LlmError> {
Ok(messages.to_vec())
}
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::UsageUpdate {
usage: meerkat_core::TurnUsage::host_declared(
provider_for_successful_test_model(&request.model),
&request.model,
meerkat_core::Usage::default(),
),
}),
Ok(LlmEvent::Done {
outcome: LlmDoneOutcome::Success {
stop_reason: meerkat_core::StopReason::EndTurn,
},
}),
]))
}
fn provider(&self) -> meerkat_core::Provider {
meerkat_core::Provider::Other
}
async fn health_check(&self) -> Result<(), meerkat_client::LlmError> {
Ok(())
}
}
struct CountingAgentLlmClient {
inner: Arc<dyn meerkat_core::AgentLlmClient>,
stream_calls: Arc<AtomicUsize>,
}
struct CountingAgentLlmRequestAttempt {
inner: Arc<dyn meerkat_core::AgentLlmRequestAttempt>,
stream_calls: Arc<AtomicUsize>,
}
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
impl meerkat_core::AgentLlmRequestAttempt for CountingAgentLlmRequestAttempt {
fn request_pressure(
&self,
) -> Result<Option<meerkat_core::ProviderRequestPressure>, meerkat_core::AgentError>
{
self.inner.request_pressure()
}
async fn stream_response(
&self,
) -> Result<meerkat_core::LlmStreamResult, meerkat_core::AgentError> {
self.stream_calls.fetch_add(1, Ordering::SeqCst);
self.inner.stream_response().await
}
}
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
impl meerkat_core::AgentLlmClient for CountingAgentLlmClient {
fn prepare_request_attempt(
self: Arc<Self>,
messages: Arc<Vec<meerkat_core::Message>>,
tools: Arc<[Arc<meerkat_core::ToolDef>]>,
max_tokens: u32,
temperature: Option<f32>,
provider_params: Option<meerkat_core::lifecycle::run_primitive::ProviderParamsOverride>,
) -> Result<Arc<dyn meerkat_core::AgentLlmRequestAttempt>, meerkat_core::AgentError>
{
let attempt = Arc::clone(&self.inner).prepare_request_attempt(
messages,
tools,
max_tokens,
temperature,
provider_params,
)?;
Ok(Arc::new(CountingAgentLlmRequestAttempt {
inner: attempt,
stream_calls: Arc::clone(&self.stream_calls),
}))
}
fn request_attempt_authority(&self) -> meerkat_core::RequestAttemptAuthority {
self.inner.request_attempt_authority()
}
async fn stream_response(
&self,
messages: &[meerkat_core::Message],
tools: &[Arc<meerkat_core::ToolDef>],
max_tokens: u32,
temperature: Option<f32>,
provider_params: Option<
&meerkat_core::lifecycle::run_primitive::ProviderParamsOverride,
>,
) -> Result<meerkat_core::LlmStreamResult, meerkat_core::AgentError> {
self.stream_calls.fetch_add(1, Ordering::SeqCst);
self.inner
.stream_response(messages, tools, max_tokens, temperature, provider_params)
.await
}
fn provider(&self) -> meerkat_core::Provider {
self.inner.provider()
}
fn model(&self) -> &str {
self.inner.model()
}
fn compile_schema(
&self,
output_schema: &meerkat_core::OutputSchema,
) -> Result<meerkat_core::CompiledSchema, meerkat_core::SchemaError> {
self.inner.compile_schema(output_schema)
}
}
fn counting_agent_llm_client_decorator(
constructions: Arc<AtomicUsize>,
stream_calls: Arc<AtomicUsize>,
) -> meerkat_core::AgentLlmClientDecorator {
Arc::new(move |client| {
constructions.fetch_add(1, Ordering::SeqCst);
Arc::new(CountingAgentLlmClient {
inner: client,
stream_calls: Arc::clone(&stream_calls),
})
})
}
#[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_bridge_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_bridge_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,
session_context: None,
pending_head_canonical_boundary: None,
acknowledged_head_canonical_boundary: None,
})
}
#[cfg(feature = "experimental-gpt-live")]
async fn build_policy_probe_factory_agent(
temp: &TempDir,
client: Arc<MidRunCancellationClient>,
dispatcher: Arc<PolicyProbeDispatcher>,
store: Arc<CountingAgentSessionStore>,
checkpointer: Arc<CountingSessionCheckpointer>,
hook_engine: Option<Arc<dyn meerkat_core::HookEngine>>,
) -> Result<FactoryAgent, String> {
let factory = AgentFactory::new(temp.path().join("sessions"));
let runtime = MeerkatMachine::ephemeral();
let durable_session = Session::new();
let bindings = runtime
.prepare_bindings(durable_session.id().clone())
.await
.map_err(|error| error.to_string())?;
let agent = factory
.build_agent(
AgentBuildConfig {
llm_client_override: Some(client),
tool_dispatcher_override: Some(dispatcher),
session_store_override: Some(store),
checkpointer: Some(checkpointer),
hook_engine_override: hook_engine,
resume_session: Some(durable_session),
runtime_build_mode: meerkat_core::RuntimeBuildMode::SessionOwned(bindings),
..AgentBuildConfig::new("claude-sonnet-4-5")
},
&Config::default(),
)
.await
.map_err(|error| error.to_string())?;
Ok(FactoryAgent {
agent,
session_context: None,
pending_head_canonical_boundary: None,
acknowledged_head_canonical_boundary: None,
})
}
#[cfg(feature = "experimental-gpt-live")]
async fn warm_policy_probe_member(agent: &mut FactoryAgent) -> Result<(), String> {
let (event_tx, _event_rx) = mpsc::channel(8);
SessionAgent::run_turn_with_events(
agent,
meerkat_session::ephemeral::SessionAgentTurnInput {
prompt: "initialize policy probe tool visibility".to_string().into(),
injected_context: Vec::new(),
handling_mode: meerkat_core::HandlingMode::Queue,
render_metadata: None,
typed_turn_appends: Vec::new(),
transcript_identity: None,
execution_kind: Some(meerkat_core::lifecycle::RuntimeExecutionKind::ContentTurn),
},
event_tx,
)
.await
.map(|_| ())
.map_err(|error| error.to_string())
}
#[cfg(feature = "experimental-gpt-live")]
fn policy_probe_bridge_request(
agent: &FactoryAgent,
operation_id: &str,
) -> Result<meerkat_session::LiveBridgeSessionOperationRequest, String> {
let snapshot = agent.session().clone();
let revision = snapshot
.canonical_context_revision()
.map_err(|error| error.to_string())?;
Ok(meerkat_session::LiveBridgeSessionOperationRequest {
operation_id: Arc::from(operation_id),
snapshot: snapshot.clone(),
semantic_request: "policy probe live request".to_string().into(),
dispatch_admission: meerkat_core::LiveBridgeToolDispatchAdmission::__test_new(
operation_id,
Arc::new(AllowLiveBridgeModelOnly),
),
run_permit: meerkat_core::LiveBridgeNoncommittingRunPermit::__test_new(
operation_id,
snapshot.id().clone(),
revision,
),
})
}
#[cfg(feature = "experimental-gpt-live")]
async fn assert_live_bridge_policy_rejected_without_effects(
agent: FactoryAgent,
request: meerkat_session::LiveBridgeSessionOperationRequest,
expected_reason: &str,
client: &MidRunCancellationClient,
dispatcher: &PolicyProbeDispatcher,
store: &CountingAgentSessionStore,
checkpointer: &CountingSessionCheckpointer,
hook_engine: Option<&CountingHookEngine>,
) -> Result<(), String> {
let canonical = serde_json::to_value(agent.session()).map_err(|error| error.to_string())?;
let revision = agent
.session()
.canonical_context_revision()
.map_err(|error| error.to_string())?;
let client_calls = client.calls.load(Ordering::SeqCst);
let dispatches = dispatcher.dispatches.load(Ordering::SeqCst);
let saves = store.saves.load(Ordering::SeqCst);
let checkpoints = checkpointer.checkpoints.load(Ordering::SeqCst);
let hook_calls = hook_engine.map_or(0, |engine| engine.calls.load(Ordering::SeqCst));
let eligibility = SessionAgent::validate_live_bridge_member_eligibility(&agent)
.expect_err("unsupported member policy must fail bridge eligibility preflight");
assert!(eligibility.to_string().contains(expected_reason));
let validation = SessionAgent::validate_live_bridge_operation(&agent, &request)
.expect_err("unsupported member policy must fail before actor acceptance");
assert!(validation.to_string().contains(expected_reason));
let (_cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
let execution =
match SessionAgent::prepare_live_bridge_operation(&agent, request, cancel_rx) {
Ok(_) => panic!("unsupported member policy must fail before provider execution"),
Err(error) => error,
};
assert!(execution.to_string().contains(expected_reason));
assert_eq!(client.calls.load(Ordering::SeqCst), client_calls);
assert_eq!(dispatcher.dispatches.load(Ordering::SeqCst), dispatches);
assert_eq!(store.saves.load(Ordering::SeqCst), saves);
assert_eq!(checkpointer.checkpoints.load(Ordering::SeqCst), checkpoints);
if let Some(engine) = hook_engine {
assert_eq!(engine.calls.load(Ordering::SeqCst), hook_calls);
}
assert_eq!(
serde_json::to_value(agent.session()).map_err(|error| error.to_string())?,
canonical,
"policy rejection must preserve the exact durable Session"
);
assert_eq!(
agent
.session()
.canonical_context_revision()
.map_err(|error| error.to_string())?,
revision,
"policy rejection must preserve the exact canonical revision"
);
Ok(())
}
#[cfg(feature = "experimental-gpt-live")]
#[tokio::test]
async fn factory_agent_live_bridge_restores_canonical_session_after_cancellation()
-> Result<(), String> {
let temp = tempfile::tempdir().map_err(|error| format!("tempdir: {error}"))?;
let started = Arc::new(tokio::sync::Notify::new());
let client = Arc::new(MidRunCancellationClient {
calls: AtomicUsize::new(0),
saw_member_tool: AtomicBool::new(false),
seen_tools: Mutex::new(Vec::new()),
});
let dispatcher = Arc::new(BlockingMemberToolDispatcher {
started: Arc::clone(&started),
dispatches: AtomicUsize::new(0),
tools: Arc::from([Arc::new(ToolDef::new(
"member_lookup",
"member-local read-only tool",
serde_json::json!({ "type": "object" }),
))]),
});
let store = Arc::new(CountingAgentSessionStore {
saves: AtomicUsize::new(0),
});
let checkpointer = Arc::new(CountingSessionCheckpointer {
checkpoints: AtomicUsize::new(0),
});
let factory = AgentFactory::new(temp.path().join("sessions"));
let runtime = MeerkatMachine::ephemeral();
let durable_session = Session::new();
let bindings = runtime
.prepare_bindings(durable_session.id().clone())
.await
.map_err(|error| error.to_string())?;
let agent = factory
.build_agent(
AgentBuildConfig {
llm_client_override: Some(client.clone()),
tool_dispatcher_override: Some(dispatcher.clone()),
session_store_override: Some(store.clone()),
checkpointer: Some(checkpointer.clone()),
resume_session: Some(durable_session),
runtime_build_mode: meerkat_core::RuntimeBuildMode::SessionOwned(bindings),
budget_limits: Some(
meerkat_core::BudgetLimits::unlimited()
.with_max_tokens(1_000)
.with_max_tool_calls(1_000),
),
..AgentBuildConfig::new("claude-sonnet-4-5")
},
&Config::default(),
)
.await
.map_err(|error| error.to_string())?;
let mut agent = FactoryAgent {
agent,
session_context: None,
pending_head_canonical_boundary: None,
acknowledged_head_canonical_boundary: None,
};
let (warmup_tx, _warmup_rx) = mpsc::channel(8);
SessionAgent::run_turn_with_events(
&mut agent,
meerkat_session::ephemeral::SessionAgentTurnInput {
prompt: "initialize durable member tool visibility"
.to_string()
.into(),
injected_context: Vec::new(),
handling_mode: meerkat_core::HandlingMode::Queue,
render_metadata: None,
typed_turn_appends: Vec::new(),
transcript_identity: None,
execution_kind: Some(meerkat_core::lifecycle::RuntimeExecutionKind::ContentTurn),
},
warmup_tx,
)
.await
.map_err(|error| format!("durable member warmup turn: {error}"))?;
let (tap_tx, mut tap_rx) = mpsc::channel(32);
*agent.agent().event_tap().lock() = Some(meerkat_core::EventTapState {
tx: tap_tx,
truncated: AtomicBool::new(false),
});
let tool_scope = SessionAgent::tool_scope_snapshot(&agent)
.ok_or_else(|| "factory agent exposes no tool scope".to_string())?;
assert!(
tool_scope
.visible_names
.iter()
.any(|name| name.as_str() == "member_lookup"),
"member-local tool must be visible on the durable FactoryAgent"
);
let visible_defs = SessionAgent::visible_tool_defs(&agent);
assert!(
visible_defs.iter().any(|tool| tool.name == "member_lookup"),
"member-local tool definition must be available before bridge execution: {:?}",
visible_defs
.iter()
.map(|tool| tool.name.to_string())
.collect::<Vec<_>>()
);
let canonical = agent.session().clone();
let ordinary_turn_state = agent
.agent()
.turn_state_handle()
.ok_or_else(|| "SessionOwned agent has no ordinary turn-state authority".to_string())?;
let ordinary_turn_state_before_bridge = ordinary_turn_state.snapshot();
let ordinary_transient_context_before_bridge =
format!("{:?}", agent.agent().transient_turn_context_state());
let ordinary_boundary_cancel_before_bridge = agent.agent().cancel_after_boundary_handle();
ordinary_boundary_cancel_before_bridge
.send(meerkat_core::agent::CancelAfterBoundaryCommand::for_run(
meerkat_core::RunId::new(),
))
.map_err(|error| format!("queue ordinary boundary-cancel command: {error}"))?;
assert_eq!(agent.agent().pending_cancel_after_boundary_commands(), 1);
let budget_tokens_before_bridge = agent
.agent()
.budget()
.token_usage()
.ok_or_else(|| "test member has no token budget".to_string())?
.0;
let canonical_revision = canonical
.canonical_context_revision()
.map_err(|error| error.to_string())?;
let request = meerkat_session::LiveBridgeSessionOperationRequest {
operation_id: Arc::from("test-live-bridge-operation"),
snapshot: canonical.clone(),
semantic_request: "temporary live request".to_string().into(),
dispatch_admission: meerkat_core::LiveBridgeToolDispatchAdmission::__test_new(
"test-live-bridge-operation",
Arc::new(AllowLiveBridgeModelOnly),
),
run_permit: meerkat_core::LiveBridgeNoncommittingRunPermit::__test_new(
"test-live-bridge-operation",
canonical.id().clone(),
canonical_revision.clone(),
),
};
SessionAgent::validate_live_bridge_operation(&agent, &request)
.map_err(|error| error.to_string())?;
let saves_before_bridge = store.saves.load(Ordering::SeqCst);
let checkpoints_before_bridge = checkpointer.checkpoints.load(Ordering::SeqCst);
let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
let execution = SessionAgent::prepare_live_bridge_operation(&agent, request, cancel_rx)
.map_err(|error| error.to_string())?;
let operation = tokio::spawn(execution);
if tokio::time::timeout(std::time::Duration::from_secs(3), started.notified())
.await
.is_err()
{
let terminal = if operation.is_finished() {
let result = operation.await.map_err(|error| error.to_string())?;
format!("{result:?}")
} else {
"still running".to_string()
};
return Err(format!(
"member tool dispatch did not start: visible={}, dispatches={}, client_calls={}, request_tools={:?}, terminal={terminal}",
client.saw_member_tool.load(Ordering::SeqCst),
dispatcher.dispatches.load(Ordering::SeqCst),
client.calls.load(Ordering::SeqCst),
client.seen_tools.lock().expect("seen tools lock").clone(),
));
}
assert_eq!(
agent.agent().pending_cancel_after_boundary_commands(),
1,
"in-flight bridge execution must not consume ordinary boundary commands"
);
assert_eq!(
ordinary_turn_state.snapshot(),
ordinary_turn_state_before_bridge,
"in-flight bridge execution must not mutate the ordinary member machine"
);
assert_eq!(store.saves.load(Ordering::SeqCst), saves_before_bridge);
assert_eq!(
checkpointer.checkpoints.load(Ordering::SeqCst),
checkpoints_before_bridge
);
let (overlap_event_tx, _overlap_event_rx) = mpsc::channel(8);
let overlap_result = tokio::time::timeout(
std::time::Duration::from_secs(1),
SessionAgent::run_turn_with_events(
&mut agent,
meerkat_session::ephemeral::SessionAgentTurnInput {
prompt: "ordinary turn overlapping live bridge".to_string().into(),
injected_context: Vec::new(),
handling_mode: meerkat_core::HandlingMode::Queue,
render_metadata: None,
typed_turn_appends: Vec::new(),
transcript_identity: None,
execution_kind: Some(
meerkat_core::lifecycle::RuntimeExecutionKind::ContentTurn,
),
},
overlap_event_tx,
),
)
.await
.map_err(|_| "ordinary member turn was serialized behind live bridge".to_string())?
.map_err(|error| error.to_string())?;
assert_eq!(overlap_result.text, "ok");
assert_eq!(overlap_result.session_id, canonical.id().clone());
assert_eq!(
client.calls.load(Ordering::SeqCst),
3,
"warmup, bridge, and overlapping ordinary turn must use the same raw member client"
);
assert!(
!operation.is_finished(),
"bridge tool execution must remain in flight when ordinary turn completes"
);
assert_eq!(
agent.agent().pending_cancel_after_boundary_commands(),
0,
"overlapping ordinary turn must drain its own queued boundary command"
);
let canonical_after_overlap = agent.session().clone();
let canonical_revision_after_overlap = canonical_after_overlap
.canonical_context_revision()
.map_err(|error| error.to_string())?;
let ordinary_turn_state_after_overlap = ordinary_turn_state.snapshot();
let saves_after_overlap = store.saves.load(Ordering::SeqCst);
let checkpoints_after_overlap = checkpointer.checkpoints.load(Ordering::SeqCst);
while tap_rx.try_recv().is_ok() {}
cancel_tx
.send(true)
.map_err(|error| format!("cancel signal: {error}"))?;
let cancelled = operation.await.map_err(|error| error.to_string())?;
let cancelled = cancelled.expect_err("in-flight bridge operation must cancel locally");
assert!(matches!(cancelled, meerkat_core::AgentError::Cancelled));
assert!(client.saw_member_tool.load(Ordering::SeqCst));
assert_eq!(dispatcher.dispatches.load(Ordering::SeqCst), 1);
assert_eq!(
ordinary_turn_state.snapshot(),
ordinary_turn_state_after_overlap,
"bridge cancellation must not mutate the ordinary member machine after the overlapping turn"
);
assert_eq!(
format!("{:?}", agent.agent().transient_turn_context_state()),
ordinary_transient_context_before_bridge,
"bridge cancellation must restore the exact ordinary transient-context coordinator"
);
assert!(
ordinary_boundary_cancel_before_bridge
.same_channel(&agent.agent().cancel_after_boundary_handle()),
"bridge cancellation must restore the exact ordinary boundary-cancel channel"
);
assert_eq!(
agent.agent().pending_cancel_after_boundary_commands(),
0,
"bridge cancellation must not recreate or consume ordinary boundary commands"
);
assert_eq!(
agent
.agent()
.budget()
.token_usage()
.ok_or_else(|| "test member lost its token budget".to_string())?
.0,
budget_tokens_before_bridge + 7,
"bridge provider usage must count against the durable member budget"
);
assert_eq!(
store.saves.load(Ordering::SeqCst),
saves_after_overlap,
"noncommitting bridge execution must not write the member session store"
);
assert_eq!(
checkpointer.checkpoints.load(Ordering::SeqCst),
checkpoints_after_overlap,
"noncommitting bridge execution must not checkpoint the member session"
);
assert!(
matches!(
tap_rx.try_recv(),
Err(tokio::sync::mpsc::error::TryRecvError::Empty)
),
"bridge execution must not leak events to the ordinary member tap"
);
assert_eq!(
serde_json::to_value(agent.session()).map_err(|error| error.to_string())?,
serde_json::to_value(&canonical_after_overlap).map_err(|error| error.to_string())?,
"bridge cancellation must preserve the overlapping ordinary turn's canonical document"
);
assert_eq!(
agent
.session()
.canonical_context_revision()
.map_err(|error| error.to_string())?,
canonical_revision_after_overlap
);
let (event_tx, _event_rx) = mpsc::channel(8);
let result = SessionAgent::run_turn_with_events(
&mut agent,
meerkat_session::ephemeral::SessionAgentTurnInput {
prompt: "ordinary durable turn".to_string().into(),
injected_context: Vec::new(),
handling_mode: meerkat_core::HandlingMode::Queue,
render_metadata: None,
typed_turn_appends: Vec::new(),
transcript_identity: None,
execution_kind: Some(meerkat_core::lifecycle::RuntimeExecutionKind::ContentTurn),
},
event_tx,
)
.await
.map_err(|error| error.to_string())?;
assert_eq!(result.text, "ok");
assert_eq!(
agent.agent().pending_cancel_after_boundary_commands(),
0,
"the successor ordinary turn must drain the stale exact-run command without cancelling itself"
);
tokio::time::timeout(std::time::Duration::from_secs(1), tap_rx.recv())
.await
.map_err(|_| "ordinary member tap received no event after restoration".to_string())?
.ok_or_else(|| "ordinary member tap closed after bridge restoration".to_string())?;
assert_ne!(
agent
.session()
.canonical_context_revision()
.map_err(|error| error.to_string())?,
canonical_revision,
"ordinary member work continues from the restored canonical head"
);
assert!(
checkpointer.checkpoints.load(Ordering::SeqCst) > checkpoints_before_bridge,
"the ordinary durable turn still uses the restored checkpointer"
);
Ok(())
}
#[cfg(feature = "experimental-gpt-live")]
#[tokio::test]
async fn session_actor_runs_ordinary_turn_while_live_bridge_tool_is_in_flight()
-> Result<(), String> {
let temp = tempfile::tempdir().map_err(|error| error.to_string())?;
let started = Arc::new(tokio::sync::Notify::new());
let client = Arc::new(MidRunCancellationClient {
calls: AtomicUsize::new(0),
saw_member_tool: AtomicBool::new(false),
seen_tools: Mutex::new(Vec::new()),
});
let dispatcher = Arc::new(BlockingMemberToolDispatcher {
started: Arc::clone(&started),
dispatches: AtomicUsize::new(0),
tools: Arc::from([Arc::new(ToolDef::new(
"member_lookup",
"member-local read-only tool",
serde_json::json!({ "type": "object" }),
))]),
});
let runtime = MeerkatMachine::ephemeral();
let durable_session = Session::new();
let session_id = durable_session.id().clone();
let bindings = runtime
.prepare_bindings(session_id.clone())
.await
.map_err(|error| error.to_string())?;
let mut builder = FactoryAgentBuilder::new(
AgentFactory::new(temp.path().join("sessions")),
Config::default(),
);
builder.default_llm_client = Some(client.clone());
builder.default_tool_dispatcher = Some(dispatcher.clone());
let service = EphemeralSessionService::new(builder, 1);
let created = service
.create_session(CreateSessionRequest {
injected_context: Vec::new(),
model: "claude-sonnet-4-5".to_string(),
prompt: "initialize durable member".to_string().into(),
system_prompt: meerkat_core::SystemPromptOverride::Inherit,
max_tokens: None,
event_tx: None,
initial_turn: meerkat_core::service::InitialTurnPolicy::RunImmediately,
deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::Discard,
build: Some(SessionBuildOptions {
resume_session: Some(durable_session),
runtime_build_mode: meerkat_core::RuntimeBuildMode::SessionOwned(bindings),
initial_turn_metadata: Some(
meerkat_core::lifecycle::run_primitive::RuntimeTurnMetadata {
execution_kind: Some(
meerkat_core::lifecycle::RuntimeExecutionKind::ContentTurn,
),
..Default::default()
},
),
..SessionBuildOptions::default()
}),
labels: None,
})
.await
.map_err(|error| error.to_string())?;
assert_eq!(created.session_id, session_id);
let actor_before = service
.live_session_actor_witness(&session_id)
.await
.ok_or_else(|| "missing live actor witness".to_string())?;
let snapshot = service
.export_session(&session_id)
.await
.map_err(|error| error.to_string())?;
let revision = snapshot
.canonical_context_revision()
.map_err(|error| error.to_string())?;
let operation_id = "actor-overlap-live-bridge";
let request = meerkat_session::LiveBridgeSessionOperationRequest {
operation_id: Arc::from(operation_id),
snapshot: snapshot.clone(),
semantic_request: "temporary live request".to_string().into(),
dispatch_admission: meerkat_core::LiveBridgeToolDispatchAdmission::__test_new(
operation_id,
Arc::new(AllowLiveBridgeModelOnly),
),
run_permit: meerkat_core::LiveBridgeNoncommittingRunPermit::__test_new(
operation_id,
session_id.clone(),
revision,
),
};
let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
let terminal = service
.start_live_bridge_operation(&session_id, request, cancel_rx)
.await
.map_err(|error| error.to_string())?;
tokio::time::timeout(std::time::Duration::from_secs(3), started.notified())
.await
.map_err(|_| "bridge tool dispatch did not enter".to_string())?;
let ordinary = tokio::time::timeout(
std::time::Duration::from_secs(1),
service.start_turn(
&session_id,
meerkat_core::service::StartTurnRequest {
prompt: "ordinary actor turn during bridge".to_string().into(),
injected_context: Vec::new(),
system_prompt: None,
event_tx: None,
runtime: meerkat_core::service::StartTurnRuntimeSemantics {
turn_metadata: Some(
meerkat_core::lifecycle::run_primitive::RuntimeTurnMetadata {
execution_kind: Some(
meerkat_core::lifecycle::RuntimeExecutionKind::ContentTurn,
),
..Default::default()
},
),
..Default::default()
},
},
),
)
.await
.map_err(|_| "ordinary actor turn was serialized behind bridge execution".to_string())?
.map_err(|error| error.to_string())?;
assert_eq!(ordinary.text, "ok");
assert_eq!(ordinary.session_id, session_id);
assert_eq!(client.calls.load(Ordering::SeqCst), 3);
assert_eq!(dispatcher.dispatches.load(Ordering::SeqCst), 1);
assert_eq!(
service
.live_session_actor_witness(&session_id)
.await
.ok_or_else(|| "actor disappeared during overlap".to_string())?,
actor_before,
"ordinary and bridge execution must use one exact actor incarnation"
);
cancel_tx.send(true).map_err(|error| error.to_string())?;
let terminal = terminal.await.map_err(|error| error.to_string())?;
assert!(matches!(terminal, Err(meerkat_core::AgentError::Cancelled)));
let after = service
.export_session(&session_id)
.await
.map_err(|error| error.to_string())?;
assert!(after.messages().iter().any(|message| matches!(
message,
Message::User(user)
if user.text_content() == "ordinary actor turn during bridge"
)));
Ok(())
}
#[cfg(feature = "experimental-gpt-live")]
#[tokio::test]
async fn unsupported_llm_decorator_rejects_bridge_before_acceptance_without_effects()
-> Result<(), String> {
let temp = tempfile::tempdir().map_err(|error| error.to_string())?;
let constructions = Arc::new(AtomicUsize::new(0));
let stream_calls = Arc::new(AtomicUsize::new(0));
let dispatcher = Arc::new(PolicyProbeDispatcher {
tools: Arc::from([]),
catalog: Arc::from([]),
exact_catalog: true,
dispatches: AtomicUsize::new(0),
});
let store = Arc::new(CountingAgentSessionStore {
saves: AtomicUsize::new(0),
});
let runtime = MeerkatMachine::ephemeral();
let durable_session = Session::new();
let session_id = durable_session.id().clone();
let bindings = runtime
.prepare_bindings(session_id.clone())
.await
.map_err(|error| error.to_string())?;
let mut builder = FactoryAgentBuilder::new(
AgentFactory::new(temp.path().join("sessions")),
Config::default(),
);
builder.default_llm_client = Some(Arc::new(MockLlmClient::default()));
*builder
.default_agent_llm_client_decorator
.write()
.map_err(|_| "decorator lock poisoned".to_string())? =
Some(counting_agent_llm_client_decorator(
Arc::clone(&constructions),
Arc::clone(&stream_calls),
));
builder.default_tool_dispatcher = Some(dispatcher.clone());
builder.default_session_store = Some(store.clone());
let service = EphemeralSessionService::new(builder, 1);
service
.create_session(CreateSessionRequest {
injected_context: Vec::new(),
model: "claude-sonnet-4-5".to_string(),
prompt: "initialize decorated member".to_string().into(),
system_prompt: meerkat_core::SystemPromptOverride::Inherit,
max_tokens: None,
event_tx: None,
initial_turn: meerkat_core::service::InitialTurnPolicy::RunImmediately,
deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::Discard,
build: Some(SessionBuildOptions {
resume_session: Some(durable_session),
runtime_build_mode: meerkat_core::RuntimeBuildMode::SessionOwned(bindings),
initial_turn_metadata: Some(
meerkat_core::lifecycle::run_primitive::RuntimeTurnMetadata {
execution_kind: Some(
meerkat_core::lifecycle::RuntimeExecutionKind::ContentTurn,
),
..Default::default()
},
),
..SessionBuildOptions::default()
}),
labels: None,
})
.await
.map_err(|error| error.to_string())?;
assert!(constructions.load(Ordering::SeqCst) > 0);
let snapshot = service
.export_session(&session_id)
.await
.map_err(|error| error.to_string())?;
let transcript_authority_before = service
.observe_session_transcript_authority(&session_id)
.await
.map_err(|error| error.to_string())?;
let mut canonical_before =
serde_json::to_value(&snapshot).map_err(|error| error.to_string())?;
canonical_before
.as_object_mut()
.expect("serialized Session must be an object")
.remove("updated_at");
let revision = snapshot
.canonical_context_revision()
.map_err(|error| error.to_string())?;
let operation_id = "unsupported-decorator-live-bridge";
let request = meerkat_session::LiveBridgeSessionOperationRequest {
operation_id: Arc::from(operation_id),
snapshot,
semantic_request: "must not execute".to_string().into(),
dispatch_admission: meerkat_core::LiveBridgeToolDispatchAdmission::__test_new(
operation_id,
Arc::new(AllowLiveBridgeModelOnly),
),
run_permit: meerkat_core::LiveBridgeNoncommittingRunPermit::__test_new(
operation_id,
session_id.clone(),
revision,
),
};
let stream_calls_before = stream_calls.load(Ordering::SeqCst);
let dispatches_before = dispatcher.dispatches.load(Ordering::SeqCst);
let saves_before = store.saves.load(Ordering::SeqCst);
let preflight_error = service
.validate_live_bridge_member_eligibility(&session_id)
.await
.expect_err("unsupported decorator must fail before live open");
assert!(preflight_error.to_string().contains("event-isolated"));
assert_eq!(stream_calls.load(Ordering::SeqCst), stream_calls_before);
assert_eq!(
dispatcher.dispatches.load(Ordering::SeqCst),
dispatches_before
);
assert_eq!(store.saves.load(Ordering::SeqCst), saves_before);
let (_cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
let error = match service
.start_live_bridge_operation(&session_id, request, cancel_rx)
.await
{
Ok(_) => return Err("unsupported decorator bridge was accepted".to_string()),
Err(error) => error,
};
assert!(error.to_string().contains("event-isolated"));
assert_eq!(stream_calls.load(Ordering::SeqCst), stream_calls_before);
assert_eq!(
dispatcher.dispatches.load(Ordering::SeqCst),
dispatches_before
);
assert_eq!(store.saves.load(Ordering::SeqCst), saves_before);
let transcript_authority_after = service
.observe_session_transcript_authority(&session_id)
.await
.map_err(|error| error.to_string())?;
assert!(
transcript_authority_after == transcript_authority_before,
"pre-acceptance rejection must preserve exact actor transcript authority"
);
let mut canonical_after = serde_json::to_value(
service
.export_session(&session_id)
.await
.map_err(|error| error.to_string())?,
)
.map_err(|error| error.to_string())?;
canonical_after
.as_object_mut()
.expect("serialized Session must be an object")
.remove("updated_at");
assert_eq!(
canonical_after, canonical_before,
"pre-acceptance rejection must preserve the full Session document; updated_at is excluded because export_session stamps only its returned deferred-state projection clone"
);
Ok(())
}
#[cfg(feature = "experimental-gpt-live")]
#[tokio::test]
async fn hooked_session_owned_member_rejects_live_bridge_before_all_effects()
-> Result<(), String> {
let temp = tempfile::tempdir().map_err(|error| error.to_string())?;
let client = Arc::new(MidRunCancellationClient {
calls: AtomicUsize::new(0),
saw_member_tool: AtomicBool::new(false),
seen_tools: Mutex::new(Vec::new()),
});
let dispatcher = Arc::new(PolicyProbeDispatcher {
tools: Arc::from(Vec::<Arc<ToolDef>>::new()),
catalog: Arc::from(Vec::<meerkat_core::ToolCatalogEntry>::new()),
exact_catalog: false,
dispatches: AtomicUsize::new(0),
});
let store = Arc::new(CountingAgentSessionStore {
saves: AtomicUsize::new(0),
});
let checkpointer = Arc::new(CountingSessionCheckpointer {
checkpoints: AtomicUsize::new(0),
});
let hooks = Arc::new(CountingHookEngine {
calls: AtomicUsize::new(0),
});
let agent = build_policy_probe_factory_agent(
&temp,
Arc::clone(&client),
Arc::clone(&dispatcher),
Arc::clone(&store),
Arc::clone(&checkpointer),
Some(Arc::clone(&hooks) as Arc<dyn meerkat_core::HookEngine>),
)
.await?;
let request = policy_probe_bridge_request(&agent, "hooked-live-bridge")?;
assert_live_bridge_policy_rejected_without_effects(
agent,
request,
"ordinary-turn hooks",
client.as_ref(),
dispatcher.as_ref(),
store.as_ref(),
checkpointer.as_ref(),
Some(hooks.as_ref()),
)
.await
}
#[cfg(feature = "experimental-gpt-live")]
#[tokio::test]
async fn callback_provenance_member_rejects_live_bridge_before_provider_or_dispatch()
-> Result<(), String> {
let temp = tempfile::tempdir().map_err(|error| error.to_string())?;
let client = Arc::new(MidRunCancellationClient {
calls: AtomicUsize::new(0),
saw_member_tool: AtomicBool::new(false),
seen_tools: Mutex::new(Vec::new()),
});
let tool = Arc::new(
ToolDef::new(
"callback_tool",
"callback-only member tool",
serde_json::json!({ "type": "object" }),
)
.with_provenance(meerkat_core::ToolProvenance {
kind: meerkat_core::ToolSourceKind::Callback,
source_id: "callback-owner".into(),
}),
);
let dispatcher = Arc::new(PolicyProbeDispatcher {
tools: Arc::from([Arc::clone(&tool)]),
catalog: Arc::from([meerkat_core::ToolCatalogEntry::session_inline(tool, true)]),
exact_catalog: false,
dispatches: AtomicUsize::new(0),
});
let store = Arc::new(CountingAgentSessionStore {
saves: AtomicUsize::new(0),
});
let checkpointer = Arc::new(CountingSessionCheckpointer {
checkpoints: AtomicUsize::new(0),
});
let mut agent = build_policy_probe_factory_agent(
&temp,
Arc::clone(&client),
Arc::clone(&dispatcher),
Arc::clone(&store),
Arc::clone(&checkpointer),
None,
)
.await?;
warm_policy_probe_member(&mut agent).await?;
assert!(
SessionAgent::visible_tool_defs(&agent)
.iter()
.any(|tool| tool.name == "callback_tool")
);
let request = policy_probe_bridge_request(&agent, "callback-live-bridge")?;
assert_live_bridge_policy_rejected_without_effects(
agent,
request,
"callback-provenance",
client.as_ref(),
dispatcher.as_ref(),
store.as_ref(),
checkpointer.as_ref(),
None,
)
.await
}
#[cfg(feature = "experimental-gpt-live")]
#[tokio::test]
async fn exact_non_fast_catalog_member_rejects_live_bridge_before_provider_or_dispatch()
-> Result<(), String> {
let temp = tempfile::tempdir().map_err(|error| error.to_string())?;
let client = Arc::new(MidRunCancellationClient {
calls: AtomicUsize::new(0),
saw_member_tool: AtomicBool::new(false),
seen_tools: Mutex::new(Vec::new()),
});
let tool = Arc::new(ToolDef::new(
"streaming_tool",
"streaming-only member tool",
serde_json::json!({ "type": "object" }),
));
let streaming = meerkat_core::StreamingToolExecutionPolicy::new(
std::time::Duration::from_secs(1),
std::time::Duration::from_secs(5),
)
.map_err(|error| error.to_string())?;
let contract = meerkat_core::ToolExecutionContract::new(
std::collections::BTreeSet::from([meerkat_core::ToolExecutionMode::Streaming]),
meerkat_core::ToolExecutionMode::Streaming,
Some(streaming),
None,
)
.map_err(|error| error.to_string())?;
let dispatcher = Arc::new(PolicyProbeDispatcher {
tools: Arc::from([Arc::clone(&tool)]),
catalog: Arc::from([meerkat_core::ToolCatalogEntry::session_inline(tool, true)
.with_execution_contract(contract)]),
exact_catalog: true,
dispatches: AtomicUsize::new(0),
});
let store = Arc::new(CountingAgentSessionStore {
saves: AtomicUsize::new(0),
});
let checkpointer = Arc::new(CountingSessionCheckpointer {
checkpoints: AtomicUsize::new(0),
});
let mut agent = build_policy_probe_factory_agent(
&temp,
Arc::clone(&client),
Arc::clone(&dispatcher),
Arc::clone(&store),
Arc::clone(&checkpointer),
None,
)
.await?;
warm_policy_probe_member(&mut agent).await?;
assert!(
SessionAgent::visible_tool_defs(&agent)
.iter()
.any(|tool| tool.name == "streaming_tool")
);
let request = policy_probe_bridge_request(&agent, "streaming-live-bridge")?;
assert_live_bridge_policy_rejected_without_effects(
agent,
request,
"Fast-only",
client.as_ref(),
dispatcher.as_ref(),
store.as_ref(),
checkpointer.as_ref(),
None,
)
.await
}
fn session_with_raw_metadata(
session: Session,
key: &'static str,
value: serde_json::Value,
) -> Session {
let mut raw = serde_json::to_value(session).expect("session should serialize");
raw.get_mut("metadata")
.and_then(serde_json::Value::as_object_mut)
.expect("session metadata should be an object")
.insert(key.to_string(), value);
serde_json::from_value(raw).expect("session should deserialize with raw metadata")
}
#[cfg(all(feature = "session-store", not(target_arch = "wasm32")))]
#[tokio::test]
async fn factory_agent_sync_rejects_malformed_deferred_turn_state() -> Result<(), String> {
let temp = tempfile::tempdir().map_err(|err| format!("tempdir: {err}"))?;
let mut agent = build_factory_agent_with_mock(
&temp,
AgentBuildConfig {
..AgentBuildConfig::new("claude-sonnet-4-5")
},
)
.await?;
let durable = session_with_raw_metadata(
agent.session().clone(),
meerkat_core::SESSION_DEFERRED_TURN_STATE_KEY,
serde_json::json!("not-a-deferred-turn-state"),
);
let err = SessionAgent::sync_session_from_durable_snapshot(&mut agent, durable)
.expect_err("malformed deferred-turn state must fail closed");
assert!(
err.to_string().contains("deferred-turn state"),
"unexpected error: {err}"
);
assert!(
agent.session().try_deferred_turn_state().unwrap().is_none(),
"failed sync must not install raw deferred-turn metadata into the live session"
);
Ok(())
}
#[cfg(all(feature = "session-store", not(target_arch = "wasm32")))]
#[tokio::test]
async fn factory_head_canonical_prepare_is_exactly_retryable_after_dropped_reply()
-> Result<(), String> {
let temp = tempfile::tempdir().map_err(|err| format!("tempdir: {err}"))?;
let mut agent = build_factory_agent_with_mock(
&temp,
AgentBuildConfig {
..AgentBuildConfig::new("claude-sonnet-4-5")
},
)
.await?;
let blob_store: Arc<dyn meerkat_core::BlobStore> = Arc::new(MemoryBlobStore::default());
let request = || {
HeadCanonicalRuntimeBoundaryPrepareRequest::try_new(
HeadCanonicalRuntimeBoundaryAuthority::root(),
None,
meerkat_core::SessionDeferredTurnState::default(),
meerkat_session::ephemeral::HeadCanonicalDeferredProjectionSource::LiveActorState,
None,
"factory-root-retry-test",
Arc::clone(&blob_store),
)
};
let first = SessionAgent::prepare_head_canonical_runtime_boundary(
&mut agent,
request().map_err(|error| error.to_string())?,
)
.await
.map_err(|error| error.to_string())?;
let expected_token = first.successor_head_token().to_string();
drop(first);
let retried = SessionAgent::prepare_head_canonical_runtime_boundary(
&mut agent,
request().map_err(|error| error.to_string())?,
)
.await
.map_err(|error| error.to_string())?;
assert_eq!(retried.successor_head_token(), expected_token.as_str());
let mismatched = HeadCanonicalRuntimeBoundaryPrepareRequest::try_new(
HeadCanonicalRuntimeBoundaryAuthority::root(),
None,
meerkat_core::SessionDeferredTurnState::default(),
meerkat_session::ephemeral::HeadCanonicalDeferredProjectionSource::LiveActorState,
None,
"factory-root-retry-test-changed-role",
blob_store,
)
.map_err(|error| error.to_string())?;
let error =
match SessionAgent::prepare_head_canonical_runtime_boundary(&mut agent, mismatched)
.await
{
Ok(_) => {
return Err(
"a changed request projection replaced the pending head-canonical boundary"
.to_string(),
);
}
Err(error) => error,
};
assert!(
error.to_string().contains("does not match"),
"unexpected mismatch error: {error}"
);
let applied =
SessionAgent::acknowledge_head_canonical_runtime_boundary(&mut agent, &expected_token)
.map_err(|error| error.to_string())?;
assert_eq!(
applied,
HeadCanonicalRuntimeBoundaryAcknowledgeOutcome::Applied
);
let retried_ack =
SessionAgent::acknowledge_head_canonical_runtime_boundary(&mut agent, &expected_token)
.map_err(|error| error.to_string())?;
assert_eq!(
retried_ack,
HeadCanonicalRuntimeBoundaryAcknowledgeOutcome::AlreadyAcknowledgedExact
);
Ok(())
}
#[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 = MeerkatMachine::ephemeral();
let session = Session::new();
let session_id = session.id().clone();
let bindings = runtime_adapter
.prepare_bindings(session_id.clone())
.await
.map_err(|err| format!("prepare bindings: {err}"))?;
let expected_registry = Arc::clone(bindings.ops_lifecycle());
let req = CreateSessionRequest {
injected_context: Vec::new(),
model: "claude-sonnet-4-5".to_string(),
prompt: "hello".to_string().into(),
system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
max_tokens: None,
event_tx: 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,
stream: meerkat_core::comms::InputStreamMode::None,
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 {
injected_context: Vec::new(),
model: "claude-sonnet-4-5".to_string(),
prompt: "ignored".to_string().into(),
system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
max_tokens: None,
event_tx: 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_applies_persistence_realm_when_request_has_none() -> 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" }));
builder.default_realm_id =
Some(meerkat_core::RealmId::parse("durable-shell-realm").map_err(|e| e.to_string())?);
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
.realm_id
.as_ref()
.map(meerkat_core::RealmId::as_str),
Some("durable-shell-realm")
);
Ok(())
}
#[tokio::test]
async fn factory_builder_default_agent_llm_client_decorator_wraps_default_llm_client()
-> 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 constructions = Arc::new(AtomicUsize::new(0));
let stream_calls = Arc::new(AtomicUsize::new(0));
*builder
.default_agent_llm_client_decorator
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) =
Some(counting_agent_llm_client_decorator(
Arc::clone(&constructions),
Arc::clone(&stream_calls),
));
let (event_tx, _event_rx) = mpsc::channel(8);
let mut agent = builder
.build_agent(&make_session_request("mock-model"), event_tx)
.await
.map_err(|err| format!("{err}"))?;
let (run_tx, _run_rx) = mpsc::channel(8);
let result = SessionAgent::run_with_events(&mut agent, "hello".to_string().into(), run_tx)
.await
.map_err(|err| format!("{err}"))?;
assert_eq!(result.text, "default");
assert_eq!(constructions.load(Ordering::SeqCst), 1);
assert_eq!(stream_calls.load(Ordering::SeqCst), 1);
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 = MeerkatMachine::ephemeral();
let session = Session::new();
let session_id = session.id().clone();
let bindings = runtime_adapter
.prepare_bindings(session_id)
.await
.map_err(|err| format!("prepare bindings: {err}"))?;
let req = CreateSessionRequest {
injected_context: Vec::new(),
model: "claude-sonnet-4-5".to_string(),
prompt: "hello".to_string().into(),
system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
max_tokens: None,
event_tx: 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_turn_with_events(
&mut agent,
SessionAgentTurnInput {
prompt: "inspect".to_string().into(),
injected_context: Vec::new(),
handling_mode: HandlingMode::Queue,
render_metadata: None,
typed_turn_appends: Vec::new(),
transcript_identity: None,
execution_kind: Some(meerkat_core::lifecycle::RuntimeExecutionKind::ContentTurn),
},
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(())
}
#[tokio::test]
async fn factory_builder_wires_generate_image_on_runtime_backed_path() -> Result<(), String> {
let temp = tempfile::tempdir().map_err(|err| format!("tempdir: {err}"))?;
let runtime_adapter = Arc::new(MeerkatMachine::ephemeral());
let factory = AgentFactory::new(temp.path().join("sessions"))
.builtins(true)
.with_image_generation_machine(runtime_adapter.clone());
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_blob_store = Some(Arc::new(MemoryBlobStore::default()));
builder.default_image_generation_executor = Some(Arc::new(FakeImageGenerationExecutor));
let session = Session::new();
let session_id = session.id().clone();
let bindings = runtime_adapter
.prepare_bindings(session_id)
.await
.map_err(|err| format!("prepare bindings: {err}"))?;
let req = CreateSessionRequest {
injected_context: Vec::new(),
model: "claude-sonnet-4-5".to_string(),
prompt: "hello".to_string().into(),
system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
max_tokens: None,
event_tx: 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_turn_with_events(
&mut agent,
SessionAgentTurnInput {
prompt: "inspect".to_string().into(),
injected_context: Vec::new(),
handling_mode: HandlingMode::Queue,
render_metadata: None,
typed_turn_appends: Vec::new(),
transcript_identity: None,
execution_kind: Some(meerkat_core::lifecycle::RuntimeExecutionKind::ContentTurn),
},
run_tx,
)
.await
.map_err(|err| format!("{err}"))?;
let tool_names = capture.tool_names();
assert!(tool_names.iter().any(|name| name == "generate_image"));
Ok(())
}
#[tokio::test]
async fn factory_agent_execution_snapshot_forwards_core_state() -> 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 snapshot = SessionAgent::execution_snapshot(&agent)
.map_err(|err| format!("snapshot should project: {err}"))?
.ok_or_else(|| "factory agent should expose execution snapshot".to_string())?;
assert_eq!(
snapshot.loop_state,
meerkat_core::state::LoopState::CallingLlm
);
assert_eq!(
snapshot.turn_phase,
meerkat_core::turn_execution_authority::TurnPhase::Ready
);
assert_eq!(snapshot.active_run_id, None);
assert_eq!(snapshot.applied_cursor, 0);
Ok(())
}
#[tokio::test]
async fn factory_agent_tool_scope_snapshot_forwards_core_state() -> 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 snapshot = SessionAgent::tool_scope_snapshot(&agent)
.ok_or_else(|| "factory agent should expose tool-scope snapshot".to_string())?;
assert_eq!(snapshot.base_filter, meerkat_core::ToolFilter::All);
assert_eq!(
snapshot.active_external_filter,
meerkat_core::ToolFilter::All
);
assert_eq!(
snapshot.staged_external_filter,
meerkat_core::ToolFilter::All
);
assert_eq!(snapshot.active_revision, meerkat_core::ToolScopeRevision(0));
assert_eq!(snapshot.staged_revision, meerkat_core::ToolScopeRevision(0));
assert_eq!(snapshot.known_base_names, snapshot.visible_names);
Ok(())
}
#[cfg(feature = "mcp")]
#[tokio::test]
async fn factory_agent_external_tool_surface_snapshot_forwards_core_state() -> Result<(), String>
{
let temp = tempfile::tempdir().map_err(|err| format!("tempdir: {err}"))?;
let runtime_adapter = MeerkatMachine::ephemeral();
let session = Session::new();
let bindings = runtime_adapter
.prepare_bindings(session.id().clone())
.await
.map_err(|err| format!("prepare bindings: {err}"))?;
let mut router = meerkat_mcp::McpRouter::new_with_surface_handle(Arc::clone(
bindings.external_tool_surface(),
));
router
.stage_add(meerkat_core::McpServerConfig::stdio(
"planner",
"/bin/echo",
Vec::<String>::new(),
std::collections::HashMap::new(),
))
.map_err(|error| error.to_string())?;
let dispatcher = Arc::new(meerkat_mcp::McpRouterAdapter::new(router))
as Arc<dyn meerkat_core::AgentToolDispatcher>;
let agent = build_factory_agent_with_mock(
&temp,
AgentBuildConfig {
tool_dispatcher_override: Some(dispatcher),
resume_session: Some(session),
runtime_build_mode: meerkat_core::RuntimeBuildMode::SessionOwned(bindings),
..AgentBuildConfig::new("claude-sonnet-4-5")
},
)
.await?;
let snapshot = SessionAgent::external_tool_surface_snapshot(&agent).ok_or_else(|| {
"factory agent should expose external-tool surface snapshot".to_string()
})?;
assert_eq!(
snapshot.phase,
meerkat_core::ExternalToolSurfaceGlobalPhase::Operating
);
assert_eq!(snapshot.snapshot_epoch, 0);
assert_eq!(snapshot.snapshot_aligned_epoch, 0);
assert_eq!(snapshot.entries.len(), 1);
let entry = &snapshot.entries[0];
assert_eq!(entry.surface_id, "planner");
assert!(!entry.visible);
assert_eq!(
entry.base_state,
meerkat_core::ExternalToolSurfaceBaseState::Absent
);
assert_eq!(
entry.pending_op,
meerkat_core::ExternalToolSurfacePendingOp::None
);
assert_eq!(
entry.staged_op,
meerkat_core::ExternalToolSurfaceStagedOp::Add
);
assert_eq!(entry.staged_intent_sequence, 1);
assert_eq!(entry.pending_task_sequence, 0);
assert_eq!(entry.pending_lineage_sequence, 0);
assert_eq!(entry.inflight_call_count, 0);
assert_eq!(
entry.last_delta_operation,
meerkat_core::ExternalToolSurfaceDeltaOperation::None
);
assert_eq!(
entry.last_delta_phase,
meerkat_core::ExternalToolSurfaceDeltaPhase::None
);
Ok(())
}
#[cfg(feature = "mcp")]
#[tokio::test]
async fn factory_builds_declarative_mcp_servers_into_the_dispatcher() -> Result<(), String> {
let temp = tempfile::tempdir().map_err(|err| format!("tempdir: {err}"))?;
let agent = build_factory_agent_with_mock(
&temp,
AgentBuildConfig {
mcp_servers: vec![meerkat_core::McpServerConfig::stdio(
"planner",
"/bin/echo",
Vec::<String>::new(),
std::collections::HashMap::new(),
)],
..AgentBuildConfig::new("claude-sonnet-4-5")
},
)
.await?;
let snapshot = SessionAgent::external_tool_surface_snapshot(&agent).ok_or_else(|| {
"declarative mcp_servers must surface through the composed dispatcher".to_string()
})?;
assert!(
snapshot
.entries
.iter()
.any(|entry| entry.surface_id == "planner"),
"the declared MCP server must be staged on the session surface: {snapshot:?}"
);
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 {
injected_context: Vec::new(),
model: model.to_string(),
prompt: "test".to_string().into(),
system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
max_tokens: None,
event_tx: None,
initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::Discard,
build: None,
labels: None,
}
}
fn make_session_request_with_connection(model: &str, binding: &str) -> CreateSessionRequest {
let mut req = make_session_request(model);
let build = meerkat_core::service::SessionBuildOptions {
auth_binding: Some(meerkat_core::AuthBindingRef {
realm: meerkat_core::RealmId::parse("default").expect("valid test realm"),
binding: meerkat_core::BindingId::parse(binding).expect("valid test binding"),
profile: None,
origin: meerkat_core::BindingOrigin::Configured,
}),
..Default::default()
};
req.build = Some(build);
req
}
#[tokio::test]
async fn system_message_append_is_only_an_ordered_transcript_event() -> Result<(), String> {
let mut builder = FactoryAgentBuilder::new(AgentFactory::minimal(), Config::default());
builder.default_llm_client = Some(Arc::new(MockLlmClient::default()));
let mut request = make_session_request("claude-sonnet-4-5");
request.system_prompt =
meerkat_core::SystemPromptOverride::Set("create-time prompt".to_string());
let (event_tx, _event_rx) = mpsc::channel(8);
let mut agent = builder
.build_agent(&request, event_tx)
.await
.map_err(|error| format!("build failed: {error}"))?;
SessionAgent::append_system_messages(&mut agent, vec!["turn instruction".to_string()])
.map_err(|error| format!("System append failed: {error}"))?;
let systems = agent
.session()
.messages()
.iter()
.filter_map(|message| match message {
Message::System(system) => Some(system.content.as_str()),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(systems, ["create-time prompt", "turn instruction"]);
Ok(())
}
#[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();
let mut section = meerkat_core::RealmConfigSection::from_inline_api_keys(&[
("anthropic", "test-anthropic-key"),
("openai", "test-openai-key"),
("gemini", "test-gemini-key"),
]);
let _ = &mut section;
config.realm.insert("default".to_string(), section);
let builder = FactoryAgentBuilder::new(factory, config);
let (tx1, _rx1) = mpsc::channel(8);
builder
.build_agent(
&make_session_request_with_connection("claude-sonnet-4-5", "default_anthropic"),
tx1,
)
.await
.map_err(|e| format!("anthropic model should build: {e}"))?;
let (tx2, _rx2) = mpsc::channel(8);
builder
.build_agent(
&make_session_request_with_connection("gpt-5.4", "default_openai"),
tx2,
)
.await
.map_err(|e| format!("openai model should build: {e}"))?;
let (tx3, _rx3) = mpsc::channel(8);
builder
.build_agent(
&make_session_request_with_connection("gemini-3.5-flash", "default_gemini"),
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();
let section =
meerkat_core::RealmConfigSection::from_inline_api_keys(&[("anthropic", "test-key")]);
config.realm.insert("default".to_string(), section);
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(())
}
}