use meerkat_core::error::AgentError;
use meerkat_core::service::SessionError;
use meerkat_core::types::{Message, SessionId, SystemMessage};
use meerkat_core::{
PendingSystemContextAppend, Session, SessionLlmIdentity, SessionToolVisibilityState,
};
use meerkat_llm_core::realtime_session::RealtimeSessionOpenConfig;
use crate::session_runtime::errors::LiveOpenPrecheckError;
pub fn precheck_identity(identity: &SessionLlmIdentity) -> Result<(), LiveOpenPrecheckError> {
let realtime_capable = meerkat_models::capabilities_for(identity.provider, &identity.model)
.map(|caps| caps.realtime)
.unwrap_or(false);
apply_precheck_gates(identity.provider, &identity.model, realtime_capable)
}
pub fn apply_precheck_gates(
provider: meerkat_core::Provider,
model: &str,
realtime_capable: bool,
) -> Result<(), LiveOpenPrecheckError> {
if !realtime_capable {
return Err(LiveOpenPrecheckError::ModelNotRealtime {
model: model.to_string(),
provider: provider.as_str(),
});
}
Ok(())
}
#[must_use]
pub fn build_live_projection_snapshot_for_runtime(
session_id: &SessionId,
open_config: &RealtimeSessionOpenConfig,
) -> meerkat_core::live_adapter::LiveProjectionSnapshot {
meerkat_core::live_adapter::LiveProjectionSnapshot {
session_id: session_id.clone(),
snapshot_version: 0,
seed_messages: open_config.seed_messages.clone(),
visible_tools: open_config.visible_tools.clone(),
system_prompt: open_config.system_prompt.clone(),
model_id: open_config.llm_identity.model.clone(),
provider_id: open_config.llm_identity.provider,
audio_config: None,
runtime_system_context: open_config.runtime_system_context.clone(),
}
}
#[must_use]
pub fn live_channel_requires_close_for_identity_change(
bound_identity: &SessionLlmIdentity,
new_identity: &SessionLlmIdentity,
) -> bool {
bound_identity.model != new_identity.model || bound_identity.provider != new_identity.provider
}
#[must_use]
pub fn should_apply_global_model_hot_swap(
current_session_model: &str,
new_global_model: &str,
) -> bool {
current_session_model != new_global_model
}
#[must_use]
pub fn should_fire_live_propagation(
prior: &meerkat_core::config::Config,
new: &meerkat_core::config::Config,
) -> bool {
prior.agent.model != new.agent.model
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum LiveHotSwapSkipReason {
NoOpOrOverride,
IdentityLookupFailed(String),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum LiveChannelRefreshFailure {
OpenConfigBuildFailed(String),
SnapshotVersionFailed(String),
EnqueueFailed(String),
QueueAcceptanceRejected(String),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum LiveChannelCloseFailure {
SignalFailed(String),
CloseAuthorityRejected(String),
CommitHandoffMissing,
HostCommitFailed(String),
}
#[derive(Debug, Default, Clone, PartialEq, Eq)]
#[must_use]
pub struct LiveConfigPropagationReport {
pub swapped: Vec<SessionId>,
pub skipped: Vec<(SessionId, LiveHotSwapSkipReason)>,
pub swap_failed: Vec<(SessionId, String)>,
pub refreshed: Vec<SessionId>,
pub closed: Vec<SessionId>,
pub refresh_failed: Vec<(SessionId, LiveChannelRefreshFailure)>,
pub close_failed: Vec<(SessionId, LiveChannelCloseFailure)>,
}
impl LiveConfigPropagationReport {
#[must_use]
pub fn is_clean(&self) -> bool {
self.swap_failed.is_empty()
&& self.refresh_failed.is_empty()
&& self.close_failed.is_empty()
}
}
pub fn realtime_projection_root_system_message(
session: &Session,
) -> Result<Option<Message>, SessionError> {
let build_state = session.build_state().ok_or_else(|| {
SessionError::Agent(AgentError::InternalError(format!(
"session {} is missing session build state",
session.id()
)))
})?;
let mut content = match build_state.system_prompt.clone() {
meerkat_core::SystemPromptOverride::Set(prompt) => prompt,
meerkat_core::SystemPromptOverride::Disable => String::new(),
meerkat_core::SystemPromptOverride::Inherit => session
.messages()
.first()
.and_then(|message| match message {
Message::System(system) => Some(system.content.clone()),
Message::SystemNotice(notice) => Some(notice.model_projection_text()),
_ => None,
})
.unwrap_or_default(),
};
if let Some(additional_instructions) = &build_state.additional_instructions
&& !additional_instructions.is_empty()
{
if !content.trim().is_empty() {
content.push_str("\n\n");
}
content.push_str("[Session Build Instructions]");
for instruction in additional_instructions {
let instruction = instruction.trim();
if instruction.is_empty() {
continue;
}
content.push_str("\n- ");
content.push_str(instruction);
}
}
if content.trim().is_empty() {
Ok(None)
} else {
Ok(Some(Message::System(SystemMessage::new(content))))
}
}
pub fn realtime_projection_messages(session: &Session) -> Result<Vec<Message>, SessionError> {
let mut projected = session.messages().to_vec();
if let Some(root_system) = realtime_projection_root_system_message(session)? {
match projected.first() {
Some(Message::System(_) | Message::SystemNotice(_)) => projected[0] = root_system,
_ => projected.insert(0, root_system),
}
}
Ok(projected)
}
pub fn realtime_projection_runtime_system_context(
session: &Session,
) -> Result<Vec<PendingSystemContextAppend>, SessionError> {
let state = session
.try_system_context_state()
.map_err(|err| {
SessionError::Agent(meerkat_core::error::AgentError::InternalError(format!(
"generated system-context authority rejected realtime projection restore: {err}"
)))
})?
.unwrap_or_default();
Ok(state.realtime_projection_appends())
}
#[allow(clippy::expect_used)]
pub fn exported_tool_visibility_state(session: &Session) -> SessionToolVisibilityState {
session
.tool_visibility_state()
.expect("exported visibility state should decode")
.unwrap_or_default()
}
#[must_use]
pub fn builtin_tool_visibility_witness() -> meerkat_core::ToolVisibilityWitness {
let provenance = meerkat_core::ToolProvenance {
kind: meerkat_core::ToolSourceKind::Builtin,
source_id: "builtin".into(),
};
meerkat_core::ToolVisibilityWitness {
last_seen_provenance: Some(provenance),
}
}
#[cfg(all(
feature = "session-store",
feature = "live",
not(target_arch = "wasm32")
))]
pub use orchestrator::LiveOrchestrator;
#[cfg(all(
feature = "session-store",
feature = "live",
not(target_arch = "wasm32")
))]
mod orchestrator {
use std::sync::Arc;
use meerkat_core::service::{
CreateSessionRequest, InitialTurnPolicy, SessionError, SessionService,
};
use meerkat_core::types::{ContentInput, Message, SessionId};
use meerkat_core::{DeferredPromptPolicy, SessionLlmIdentity, SurfaceSessionRecoveryOverrides};
use meerkat_live::LiveAdapterHost;
use meerkat_llm_core::realtime_session::RealtimeSessionOpenConfig;
use meerkat_runtime::{MeerkatMachine, SessionLlmReconfigureRequest, SessionServiceRuntimeExt};
use meerkat_session::PersistentSessionService;
use super::{
LiveChannelCloseFailure, LiveChannelRefreshFailure, LiveConfigPropagationReport,
LiveHotSwapSkipReason, build_live_projection_snapshot_for_runtime,
live_channel_requires_close_for_identity_change, precheck_identity,
realtime_projection_messages, realtime_projection_root_system_message,
realtime_projection_runtime_system_context, should_apply_global_model_hot_swap,
};
use crate::service_factory::FactoryAgentBuilder;
use crate::session_runtime::admission::{
StagedCapacityAdmissions, take_staged_capacity_admission,
};
use crate::session_runtime::errors::LiveOpenPrecheckError;
use crate::session_runtime::recovery::{RecoveryContext, RecoveryRuntimeBindingMode};
use crate::session_runtime::runtime_state::ArchiveRuntimeCleanup;
use crate::session_runtime::staged_promotion::PendingPromotionCleanup;
use crate::{StagedLifecycleError, StagedSessionRegistry};
use meerkat_core::error::AgentError;
pub struct LiveOrchestrator<'a> {
pub service: &'a Arc<PersistentSessionService<FactoryAgentBuilder>>,
pub staged_sessions: &'a Arc<StagedSessionRegistry>,
pub staged_capacity_admissions: &'a StagedCapacityAdmissions,
pub runtime_adapter: &'a Arc<MeerkatMachine>,
pub host: Option<Arc<LiveAdapterHost>>,
pub config_runtime: Option<Arc<meerkat_core::ConfigRuntime>>,
pub default_llm_client: Option<Arc<dyn crate::LlmClient>>,
pub agent_llm_client_decorator: Option<meerkat_core::AgentLlmClientDecorator>,
pub external_tools: Option<Arc<dyn meerkat_core::AgentToolDispatcher>>,
pub archive_runtime_cleanup: ArchiveRuntimeCleanup,
pub realm_id: Option<&'a meerkat_core::connection::RealmId>,
pub instance_id: Option<&'a str>,
pub backend: Option<&'a str>,
}
impl LiveOrchestrator<'_> {
fn recovery_context(&self) -> RecoveryContext<'_> {
RecoveryContext {
service: self.service,
runtime_adapter: self.runtime_adapter,
realm_id: self.realm_id,
instance_id: self.instance_id,
backend: self.backend,
default_llm_client: self.default_llm_client.clone(),
agent_llm_client_decorator: self.agent_llm_client_decorator.clone(),
external_tools: self.external_tools.clone(),
config_runtime: self.config_runtime.clone(),
}
}
async fn cleanup_recovered_runtime_if_new(
&self,
session_id: &SessionId,
runtime_was_registered: bool,
) {
if !runtime_was_registered {
let _ = self.archive_runtime_cleanup.run(session_id).await;
}
}
pub async fn materialize_staged_session_for_realtime_open(
&self,
session_id: &SessionId,
) -> Result<(), SessionError> {
let pending_session = match self.staged_sessions.begin_promotion(session_id).await {
Ok(slot) => slot,
Err(StagedLifecycleError::AlreadyPromoting(_)) => {
return Err(SessionError::Busy {
id: session_id.clone(),
});
}
Err(e) => {
return Err(SessionError::Agent(
meerkat_core::error::AgentError::InternalError(format!(
"staged session lifecycle error for {session_id}: {e}"
)),
));
}
};
let Some(slot) = pending_session else {
return Ok(());
};
let staged_capacity_admission =
take_staged_capacity_admission(self.staged_capacity_admissions, session_id);
let mut promotion_cleanup = PendingPromotionCleanup::new(
Arc::clone(self.staged_sessions),
Arc::clone(self.staged_capacity_admissions),
session_id,
&slot,
staged_capacity_admission,
);
let crate::PromotingSlot {
build_config,
labels,
deferred_prompt,
deferred_injected_context,
..
} = slot;
if !deferred_injected_context.is_empty() {
return Err(SessionError::Unsupported(
"a deferred session created with injected_context cannot be promoted by \
realtime open; promote it with turn/start"
.to_string(),
));
}
let mut build_config = *build_config;
if build_config.llm_client_override.is_none()
&& let Some(client) = self.default_llm_client.as_ref()
{
build_config.llm_client_override = Some(Arc::clone(client));
promotion_cleanup.update_build_config(&build_config);
}
let runtime_generation = if build_config.config_generation.is_none() {
if let Some(runtime) = self.config_runtime.as_ref() {
runtime.get().await.ok().map(|snapshot| snapshot.generation)
} else {
None
}
} else {
None
};
let mut build = build_config.to_session_build_options();
build.realm_id = build.realm_id.or_else(|| self.realm_id.cloned());
build.instance_id = build
.instance_id
.or_else(|| self.instance_id.map(ToString::to_string));
build.backend = build.backend.or_else(|| {
self.backend
.and_then(meerkat_core::RecoveryBackendKind::parse)
});
build.config_generation = build.config_generation.or(runtime_generation);
let (prompt, deferred_prompt_policy) = match deferred_prompt {
Some(prompt) => (prompt, DeferredPromptPolicy::Stage),
None => (
ContentInput::Text(String::new()),
DeferredPromptPolicy::Discard,
),
};
let create_req = CreateSessionRequest {
injected_context: Vec::new(),
model: build_config.model.clone(),
prompt,
system_prompt: build_config.system_prompt.clone(),
max_tokens: build_config.max_tokens,
event_tx: None,
initial_turn: InitialTurnPolicy::Defer,
deferred_prompt_policy,
build: Some(build),
labels,
};
let admission = match promotion_cleanup.take_staged_capacity_admission() {
Some(adm) => adm,
None => self.service.reserve_create_session_admission().await?,
};
match self
.service
.create_session_with_reserved_admission(create_req, admission)
.await
{
Ok(_) => {
promotion_cleanup.mark_materialized();
let _ = promotion_cleanup.finish_now().await;
promotion_cleanup.disarm();
Ok(())
}
Err(error) => {
promotion_cleanup.restore_now().await;
Err(error)
}
}
}
pub async fn recover_live_session_for_realtime_open(
&self,
session_id: &SessionId,
) -> Result<(), SessionError> {
if self.service.has_live_session(session_id).await? {
return Ok(());
}
if self.staged_sessions.contains(session_id).await {
Box::pin(self.materialize_staged_session_for_realtime_open(session_id)).await?;
return Ok(());
}
let recovery_ctx = self.recovery_context();
let session = recovery_ctx
.load_persisted_session(session_id)
.await?
.ok_or_else(|| SessionError::NotFound {
id: session_id.clone(),
})?;
let keep_alive = session
.session_metadata()
.ok_or_else(|| {
SessionError::Agent(meerkat_core::error::AgentError::InternalError(format!(
"session {session_id} is missing session metadata"
)))
})?
.keep_alive;
let recovery_overrides = SurfaceSessionRecoveryOverrides {
keep_alive: Some(keep_alive),
..Default::default()
};
let recovered = recovery_ctx
.recovered_create_request_with_runtime_binding_mode(
session_id,
session,
recovery_overrides,
RecoveryRuntimeBindingMode::LocalResources,
)
.await
.map_err(recovery_error_to_session_error)?;
let runtime_was_registered = recovered.runtime_was_registered;
let admission = self.service.reserve_create_session_admission().await?;
if let Err(error) = self
.service
.create_session_with_reserved_admission(recovered.request, admission)
.await
{
self.cleanup_recovered_runtime_if_new(session_id, runtime_was_registered)
.await;
return Err(error);
}
Ok(())
}
pub async fn realtime_session_open_config(
&self,
session_id: &SessionId,
turning_mode: meerkat_contracts::RealtimeTurningMode,
) -> Result<RealtimeSessionOpenConfig, SessionError> {
Box::pin(self.recover_live_session_for_realtime_open(session_id)).await?;
let session = match self
.service
.export_realtime_open_session_snapshot(session_id)
.await
{
Ok(session) => session,
Err(SessionError::NotFound { .. }) => {
Box::pin(self.recover_live_session_for_realtime_open(session_id)).await?;
self.service
.export_realtime_open_session_snapshot(session_id)
.await?
}
Err(error) => return Err(error),
};
let llm_identity = self.service.live_session_llm_identity(session_id).await?;
let visible_tools = self.service.live_visible_tool_defs(session_id).await?;
Ok(RealtimeSessionOpenConfig::new(
turning_mode,
llm_identity,
visible_tools,
realtime_projection_messages(&session)?,
)
.with_runtime_system_context(realtime_projection_runtime_system_context(&session)?)
.with_system_prompt(
match realtime_projection_root_system_message(&session)? {
Some(Message::System(system)) => Some(system.content),
_ => None,
},
))
}
pub async fn live_open_config_for_session(
&self,
session_id: &SessionId,
turning_mode: meerkat_contracts::RealtimeTurningMode,
) -> Result<RealtimeSessionOpenConfig, SessionError> {
self.realtime_session_open_config(session_id, turning_mode)
.await
}
pub async fn live_llm_identity_for_session(
&self,
session_id: &SessionId,
) -> Result<SessionLlmIdentity, SessionError> {
if let Some(info) = self
.staged_sessions
.try_info(session_id)
.await
.map_err(|err| SessionError::Agent(AgentError::InternalError(err.to_string())))?
{
return Ok(info.effective_llm_identity);
}
Box::pin(self.recover_live_session_for_realtime_open(session_id)).await?;
match self.service.live_session_llm_identity(session_id).await {
Ok(identity) => Ok(identity),
Err(SessionError::NotFound { .. }) => {
Box::pin(self.recover_live_session_for_realtime_open(session_id)).await?;
self.service.live_session_llm_identity(session_id).await
}
Err(error) => Err(error),
}
}
pub async fn precheck_live_open(
&self,
session_id: &SessionId,
) -> Result<(), LiveOpenPrecheckError> {
let map_lookup_err = |err: SessionError| LiveOpenPrecheckError::SessionLookup {
session_id: session_id.clone(),
source: err,
};
if let Some(info) = self
.staged_sessions
.try_info(session_id)
.await
.map_err(|err| {
map_lookup_err(SessionError::Agent(AgentError::InternalError(
err.to_string(),
)))
})?
{
return precheck_identity(&info.effective_llm_identity);
}
Box::pin(self.recover_live_session_for_realtime_open(session_id))
.await
.map_err(map_lookup_err)?;
let identity = match self.service.live_session_llm_identity(session_id).await {
Ok(identity) => identity,
Err(SessionError::NotFound { .. }) => {
Box::pin(self.recover_live_session_for_realtime_open(session_id))
.await
.map_err(map_lookup_err)?;
self.service
.live_session_llm_identity(session_id)
.await
.map_err(map_lookup_err)?
}
Err(other) => return Err(map_lookup_err(other)),
};
precheck_identity(&identity)
}
async fn close_live_channel_for_config_rejection(
&self,
host: &meerkat_live::LiveAdapterHost,
session_id: &SessionId,
channel_id: &meerkat_live::LiveChannelId,
reason: meerkat_core::live_adapter::LiveConfigRejectionReason,
context: &'static str,
) -> Result<(), LiveChannelCloseFailure> {
let observation = host
.signal_terminal_error_observed(
channel_id,
meerkat_core::live_adapter::LiveAdapterErrorCode::ConfigRejected { reason },
)
.await
.map_err(|err| {
tracing::warn!(
target: "meerkat::session_runtime::live_orchestration",
?channel_id,
?session_id,
?err,
context,
"failed to signal terminal error on live channel after config rejection"
);
LiveChannelCloseFailure::SignalFailed(err.to_string())
})?;
let authority = self
.runtime_adapter
.resolve_live_close_result(session_id, &observation)
.await
.map_err(|err| {
tracing::warn!(
target: "meerkat::session_runtime::live_orchestration",
?channel_id,
?session_id,
?err,
context,
"live close authority rejected config-rejection terminal cleanup"
);
LiveChannelCloseFailure::CloseAuthorityRejected(err.to_string())
})?;
let Some(close_commit_authority) = authority.channel_close_commit_authority() else {
tracing::warn!(
target: "meerkat::session_runtime::live_orchestration",
?channel_id,
?session_id,
context,
"live close authority omitted config-rejection host commit handoff"
);
return Err(LiveChannelCloseFailure::CommitHandoffMissing);
};
host.commit_channel_close_observation(&observation, close_commit_authority)
.await
.map_err(|err| {
tracing::warn!(
target: "meerkat::session_runtime::live_orchestration",
?channel_id,
?session_id,
?err,
context,
"host close commit failed after config-rejection generated terminal cleanup"
);
LiveChannelCloseFailure::HostCommitFailed(err.to_string())
})
}
pub async fn propagate_config_to_live_channels(&self) -> LiveConfigPropagationReport {
let mut report = LiveConfigPropagationReport::default();
let Some(host) = self.host.as_ref() else {
return report;
};
let channels = host.active_channels().await;
let mut unique_sessions: Vec<SessionId> = Vec::new();
for channel_id in &channels {
if let Some(session_id) = self
.runtime_adapter
.live_session_for_active_channel(channel_id)
.await
&& !unique_sessions.iter().any(|sid| sid == &session_id)
{
unique_sessions.push(session_id);
}
}
if !unique_sessions.is_empty()
&& let Some(runtime) = self.config_runtime.as_ref()
&& let Ok(snapshot) = runtime.get().await
{
let new_global_model = snapshot.config.agent.model.clone();
for session_id in &unique_sessions {
let current_model =
match self.service.live_session_llm_identity(session_id).await {
Ok(identity) => identity.model,
Err(err) => {
report.skipped.push((
session_id.clone(),
LiveHotSwapSkipReason::IdentityLookupFailed(err.to_string()),
));
continue;
}
};
if !should_apply_global_model_hot_swap(¤t_model, &new_global_model) {
report
.skipped
.push((session_id.clone(), LiveHotSwapSkipReason::NoOpOrOverride));
continue;
}
let request = SessionLlmReconfigureRequest {
model: Some(new_global_model.clone()),
provider: None,
provider_params: None,
auth_binding: None,
};
if let Err(err) = self
.runtime_adapter
.reconfigure_session_llm_identity(session_id, request)
.await
{
report
.swap_failed
.push((session_id.clone(), err.to_string()));
} else {
report.swapped.push(session_id.clone());
}
}
}
for channel_id in channels {
let session_id = match self
.runtime_adapter
.live_session_for_active_channel(&channel_id)
.await
{
Some(id) => id,
None => {
tracing::debug!(
target: "meerkat::session_runtime::live_orchestration",
?channel_id,
"skipping live channel absent from generated active-channel authority"
);
continue;
}
};
if let Err(precheck_err) = Box::pin(self.precheck_live_open(&session_id)).await {
tracing::info!(
target: "meerkat::session_runtime::live_orchestration",
?channel_id,
?session_id,
?precheck_err,
"closing live channel: new resolution not realtime-capable"
);
let reason = meerkat_core::live_adapter::LiveConfigRejectionReason::NonRealtimeResolution {
detail: format!("{precheck_err:?}"),
};
match self
.close_live_channel_for_config_rejection(
host,
&session_id,
&channel_id,
reason,
"non_realtime",
)
.await
{
Ok(()) => report.closed.push(session_id.clone()),
Err(failure) => report.close_failed.push((session_id.clone(), failure)),
}
continue;
}
let open_config = match Box::pin(self.live_open_config_for_session(
&session_id,
meerkat_contracts::RealtimeTurningMode::ProviderManaged,
))
.await
{
Ok(config) => config,
Err(err) => {
tracing::warn!(
target: "meerkat::session_runtime::live_orchestration",
?channel_id,
?session_id,
?err,
"failed to build refreshed open_config for live channel"
);
report.refresh_failed.push((
session_id.clone(),
LiveChannelRefreshFailure::OpenConfigBuildFailed(err.to_string()),
));
continue;
}
};
let bound_identity = match self
.runtime_adapter
.live_channel_bound_llm_identity(&session_id, &channel_id)
.await
{
Ok(Some(identity)) => identity,
Ok(None) => {
tracing::warn!(
target: "meerkat::session_runtime::live_orchestration",
?channel_id,
?session_id,
"closing live channel: generated bound LLM identity authority is absent"
);
match self
.close_live_channel_for_config_rejection(
host,
&session_id,
&channel_id,
meerkat_core::live_adapter::LiveConfigRejectionReason::Other {
detail:
"missing generated live-channel bound identity authority"
.to_string(),
},
"missing_generated_identity",
)
.await
{
Ok(()) => report.closed.push(session_id.clone()),
Err(failure) => {
report.close_failed.push((session_id.clone(), failure));
}
}
continue;
}
Err(err) => {
tracing::warn!(
target: "meerkat::session_runtime::live_orchestration",
?channel_id,
?session_id,
?err,
"closing live channel: generated bound LLM identity authority lookup failed"
);
match self
.close_live_channel_for_config_rejection(
host,
&session_id,
&channel_id,
meerkat_core::live_adapter::LiveConfigRejectionReason::Other {
detail: format!(
"generated live-channel bound identity authority lookup failed: {err}"
),
},
"generated_identity_lookup_failed",
)
.await
{
Ok(()) => report.closed.push(session_id.clone()),
Err(failure) => {
report.close_failed.push((session_id.clone(), failure));
}
}
continue;
}
};
if live_channel_requires_close_for_identity_change(
&bound_identity,
&open_config.llm_identity,
) {
tracing::info!(
target: "meerkat::session_runtime::live_orchestration",
%channel_id,
%session_id,
old_model_id = %bound_identity.model,
new_model_id = %open_config.llm_identity.model,
old_provider_id = ?bound_identity.provider,
new_provider_id = ?open_config.llm_identity.provider,
reason = "model_swap",
"closing live channel: config patch swapped \
model/provider; SDK must reopen against new identity"
);
let reason =
meerkat_core::live_adapter::LiveConfigRejectionReason::ChannelIdentitySwap {
from_model: bound_identity.model.clone(),
from_provider: bound_identity.provider,
to_model: open_config.llm_identity.model.clone(),
to_provider: open_config.llm_identity.provider,
};
match self
.close_live_channel_for_config_rejection(
host,
&session_id,
&channel_id,
reason,
"model_swap",
)
.await
{
Ok(()) => report.closed.push(session_id.clone()),
Err(failure) => report.close_failed.push((session_id.clone(), failure)),
}
continue;
}
let mut snapshot =
build_live_projection_snapshot_for_runtime(&session_id, &open_config);
match host.next_snapshot_version(&channel_id).await {
Ok(v) => snapshot.snapshot_version = v,
Err(err) => {
tracing::debug!(
target: "meerkat::session_runtime::live_orchestration",
?channel_id,
?session_id,
?err,
"skipping live channel: snapshot version stamp failed"
);
report.refresh_failed.push((
session_id.clone(),
LiveChannelRefreshFailure::SnapshotVersionFailed(err.to_string()),
));
continue;
}
}
match host.enqueue_refresh(&channel_id, snapshot).await {
Ok(acceptance) => {
if let Err(err) = self
.runtime_adapter
.resolve_live_refresh_queued_result(&session_id, &acceptance)
.await
{
tracing::warn!(
target: "meerkat::session_runtime::live_orchestration",
?channel_id,
?session_id,
?err,
"live refresh queue acceptance was rejected by generated authority"
);
report.refresh_failed.push((
session_id.clone(),
LiveChannelRefreshFailure::QueueAcceptanceRejected(err.to_string()),
));
} else {
report.refreshed.push(session_id.clone());
}
}
Err(err) => {
tracing::warn!(
target: "meerkat::session_runtime::live_orchestration",
?channel_id,
?session_id,
?err,
"failed to enqueue Refresh command to live channel"
);
report.refresh_failed.push((
session_id.clone(),
LiveChannelRefreshFailure::EnqueueFailed(err.to_string()),
));
}
}
}
report
}
}
fn recovery_error_to_session_error(
error: crate::session_runtime::errors::RecoveryError,
) -> SessionError {
use crate::session_runtime::errors::RecoveryError;
match error {
RecoveryError::Recovery(error) => SessionError::Agent(
meerkat_core::error::AgentError::InternalError(error.to_string()),
),
RecoveryError::BindingPreparation { .. } => SessionError::Agent(
meerkat_core::error::AgentError::InternalError(error.to_string()),
),
RecoveryError::Session(session_error) => session_error,
}
}
}
#[cfg(test)]
mod prompt_truth_tests {
use super::build_live_projection_snapshot_for_runtime;
use meerkat_core::types::{Message, SessionId, SystemMessage, UserMessage};
use meerkat_core::{Provider, SessionLlmIdentity};
use meerkat_llm_core::realtime_session::RealtimeSessionOpenConfig;
fn test_identity() -> SessionLlmIdentity {
SessionLlmIdentity {
model: "gpt-realtime-2".to_string(),
provider: Provider::OpenAI,
provider_params: None,
self_hosted_server_id: None,
auth_binding: None,
}
}
#[test]
fn runtime_snapshot_reads_typed_system_prompt_field_not_seed_messages() {
let resolved_prompt = "you are a helpful meerkat".to_string();
let open_config = RealtimeSessionOpenConfig::new(
meerkat_contracts::RealtimeTurningMode::ProviderManaged,
test_identity(),
Vec::new(),
vec![Message::User(UserMessage::text("hi".to_string()))],
)
.with_system_prompt(Some(resolved_prompt.clone()));
let snapshot = build_live_projection_snapshot_for_runtime(&SessionId::new(), &open_config);
assert_eq!(
snapshot.system_prompt,
Some(resolved_prompt),
"snapshot.system_prompt must read the typed open_config field, not infer from seed_messages[0]"
);
}
#[test]
fn runtime_snapshot_system_prompt_none_when_typed_field_absent() {
let open_config = RealtimeSessionOpenConfig::new(
meerkat_contracts::RealtimeTurningMode::ProviderManaged,
test_identity(),
Vec::new(),
vec![Message::System(SystemMessage::new(
"stray seed system message",
))],
);
let snapshot = build_live_projection_snapshot_for_runtime(&SessionId::new(), &open_config);
assert_eq!(
snapshot.system_prompt, None,
"snapshot.system_prompt must mirror the absent typed field, not infer from seed_messages[0]"
);
}
}