#![allow(clippy::expect_used, clippy::unwrap_used, clippy::panic)]
use meerkat::session_runtime::errors::LiveOpenPrecheckError;
use meerkat::session_runtime::live_orchestration::{
LiveChannelCloseFailure, LiveChannelRefreshFailure, LiveConfigPropagationReport,
LiveHotSwapSkipReason, apply_precheck_gates, live_channel_requires_close_for_identity_change,
should_apply_global_model_hot_swap, should_fire_live_propagation,
};
use meerkat_core::types::SessionId;
use meerkat_core::{Provider, SessionLlmIdentity};
#[test]
fn precheck_b19_fires_before_b18_for_non_realtime_non_openai() {
let err = apply_precheck_gates(Provider::Anthropic, "claude-opus-4-8", false)
.expect_err("non-realtime should fail B19");
match err {
LiveOpenPrecheckError::ModelNotRealtime { model, provider } => {
assert_eq!(model, "claude-opus-4-8");
assert_eq!(provider, "anthropic");
}
other => panic!("expected ModelNotRealtime, got {other:?}"),
}
}
#[test]
fn precheck_allows_realtime_capable_non_openai_for_factory_resolution() {
apply_precheck_gates(Provider::Anthropic, "synthetic-rt-anthropic", true)
.expect("provider adapter support is checked by the live session factory");
}
#[test]
fn precheck_accepts_realtime_capable_openai() {
apply_precheck_gates(Provider::OpenAI, "gpt-realtime-2", true)
.expect("realtime-capable OpenAI must pass both gates");
}
#[test]
fn live_channel_close_on_model_swap() {
let prev = SessionLlmIdentity {
model: "gpt-realtime-2".into(),
provider: Provider::OpenAI,
self_hosted_server_id: None,
provider_params: None,
auth_binding: None,
};
let next = SessionLlmIdentity {
model: "gpt-realtime-3".into(),
..prev.clone()
};
assert!(live_channel_requires_close_for_identity_change(
&prev, &next
));
}
#[test]
fn live_channel_close_on_provider_swap() {
let prev = SessionLlmIdentity {
model: "shared".into(),
provider: Provider::OpenAI,
self_hosted_server_id: None,
provider_params: None,
auth_binding: None,
};
let next = SessionLlmIdentity {
provider: Provider::Anthropic,
..prev.clone()
};
assert!(live_channel_requires_close_for_identity_change(
&prev, &next
));
}
#[test]
fn live_channel_in_place_refresh_when_identity_unchanged() {
let identity = SessionLlmIdentity {
model: "gpt-realtime-2".into(),
provider: Provider::OpenAI,
self_hosted_server_id: None,
provider_params: None,
auth_binding: None,
};
assert!(!live_channel_requires_close_for_identity_change(
&identity, &identity
));
}
#[test]
fn hot_swap_propagates_when_prior_global_unknown_and_models_differ() {
assert!(should_apply_global_model_hot_swap(
"gpt-realtime-2",
"gpt-realtime-3"
));
}
#[test]
fn hot_swap_skips_when_session_already_matches_new_global() {
assert!(!should_apply_global_model_hot_swap(
"gpt-realtime-3",
"gpt-realtime-3"
));
assert!(!should_apply_global_model_hot_swap(
"gpt-realtime-3",
"gpt-realtime-3"
));
}
#[test]
fn hot_swap_propagates_when_session_was_tracking_prior_global() {
assert!(should_apply_global_model_hot_swap(
"gpt-realtime-2",
"gpt-realtime-3"
));
}
#[test]
fn hot_swap_propagates_when_session_diverged_from_prior_global() {
assert!(should_apply_global_model_hot_swap(
"gpt-realtime-prior-override",
"gpt-realtime-3"
));
}
#[test]
fn hot_swap_skips_when_session_already_matches_new_global_after_divergence() {
assert!(!should_apply_global_model_hot_swap(
"gpt-realtime-3",
"gpt-realtime-3"
));
}
#[test]
fn should_fire_live_propagation_returns_false_for_unrelated_field_change() {
let prior = meerkat_core::config::Config::default();
let mut new = prior.clone();
new.max_tokens = Some(prior.max_tokens.unwrap_or_default().saturating_add(1));
assert_ne!(prior.max_tokens, new.max_tokens);
assert_eq!(prior.agent.model, new.agent.model);
assert!(!should_fire_live_propagation(&prior, &new));
}
#[test]
fn should_fire_live_propagation_returns_true_when_agent_model_changes() {
let prior = meerkat_core::config::Config::default();
let mut new = prior.clone();
new.agent.model = "gpt-realtime-3".to_string();
assert_ne!(prior.agent.model, new.agent.model);
assert!(should_fire_live_propagation(&prior, &new));
}
#[test]
fn should_fire_live_propagation_returns_false_when_agent_model_unchanged() {
let prior = meerkat_core::config::Config::default();
let new = prior.clone();
assert!(!should_fire_live_propagation(&prior, &new));
}
#[test]
fn should_fire_live_propagation_only_consults_agent_model_today() {
let prior = meerkat_core::config::Config::default();
let mut new = prior.clone();
new.agent.system_prompt = Some("flipped".to_string());
assert_eq!(prior.agent.model, new.agent.model);
assert!(!should_fire_live_propagation(&prior, &new));
}
#[test]
fn propagation_report_default_is_clean_and_empty() {
let report = LiveConfigPropagationReport::default();
assert!(report.is_clean());
assert!(report.swapped.is_empty());
assert!(report.skipped.is_empty());
assert!(report.swap_failed.is_empty());
assert!(report.refreshed.is_empty());
assert!(report.closed.is_empty());
assert!(report.refresh_failed.is_empty());
}
#[test]
fn propagation_report_with_forced_channel_failure_is_not_clean_and_enumerates_it() {
let refreshed = SessionId::new();
let failed = SessionId::new();
let report = LiveConfigPropagationReport {
refreshed: vec![refreshed.clone()],
refresh_failed: vec![(
failed.clone(),
LiveChannelRefreshFailure::EnqueueFailed("live channel closed".to_string()),
)],
..Default::default()
};
assert!(
!report.is_clean(),
"a dropped per-channel refresh must not read as a clean propagation"
);
assert_eq!(report.refreshed, vec![refreshed]);
assert_eq!(report.refresh_failed.len(), 1);
let (session, failure) = &report.refresh_failed[0];
assert_eq!(session, &failed);
match failure {
LiveChannelRefreshFailure::EnqueueFailed(detail) => {
assert_eq!(detail, "live channel closed");
}
other => panic!("expected EnqueueFailed, got {other:?}"),
}
}
#[test]
fn propagation_report_with_swap_failure_is_not_clean() {
let failed = SessionId::new();
let report = LiveConfigPropagationReport {
swap_failed: vec![(failed.clone(), "reconfigure rejected".to_string())],
..Default::default()
};
assert!(!report.is_clean());
assert_eq!(report.swap_failed.len(), 1);
assert_eq!(report.swap_failed[0].0, failed);
}
#[test]
fn propagation_report_deliberate_skip_and_close_stay_clean() {
let skipped = SessionId::new();
let closed = SessionId::new();
let report = LiveConfigPropagationReport {
skipped: vec![(skipped, LiveHotSwapSkipReason::NoOpOrOverride)],
closed: vec![closed],
..Default::default()
};
assert!(report.is_clean());
}
#[test]
fn propagation_report_with_close_failure_is_not_clean() {
let failed = SessionId::new();
let report = LiveConfigPropagationReport {
close_failed: vec![(
failed.clone(),
LiveChannelCloseFailure::CommitHandoffMissing,
)],
..Default::default()
};
assert!(!report.is_clean());
assert!(report.closed.is_empty());
assert_eq!(report.close_failed.len(), 1);
assert_eq!(report.close_failed[0].0, failed);
}
#[cfg(all(
feature = "session-store",
feature = "live",
feature = "memory-store",
not(target_arch = "wasm32")
))]
mod orchestrator_e2e {
use std::sync::Arc;
use meerkat::session_runtime::admission::StagedCapacityAdmissions;
use meerkat::session_runtime::errors::LiveOpenPrecheckError;
use meerkat::session_runtime::live_orchestration::LiveOrchestrator;
use meerkat::session_runtime::runtime_state::ArchiveRuntimeCleanup;
use meerkat::surface::build_runtime_backed_service_with_capacities;
use meerkat::{
AgentBuildConfig, AgentFactory, Config, FactoryAgentBuilder, PersistenceBundle,
StagedSessionRegistry, StagedSlot,
};
use meerkat_core::SessionLlmIdentity;
use meerkat_core::types::SessionId;
use meerkat_runtime::MeerkatMachine;
use meerkat_store::MemoryBlobStore;
struct Fixture {
service: Arc<meerkat_session::PersistentSessionService<FactoryAgentBuilder>>,
staged_sessions: Arc<StagedSessionRegistry>,
staged_capacity_admissions: StagedCapacityAdmissions,
archive_runtime_cleanup: ArchiveRuntimeCleanup,
runtime_adapter: Arc<MeerkatMachine>,
_temp: tempfile::TempDir,
}
fn build_fixture() -> Fixture {
let session_store: Arc<dyn meerkat::SessionStore> = Arc::new(meerkat::MemoryStore::new());
let persistence = PersistenceBundle::new(
session_store,
Arc::new(meerkat_runtime::InMemoryRuntimeStore::new()),
Arc::new(MemoryBlobStore::new()),
);
let temp = tempfile::tempdir().expect("tempdir");
let factory = AgentFactory::new(temp.path().join("sessions")).builtins(false);
let builder = FactoryAgentBuilder::new(factory, Config::default());
let staged_sessions = Arc::new(StagedSessionRegistry::new());
let (service, runtime_adapter) =
build_runtime_backed_service_with_capacities(builder, 4, 16, persistence);
let service = Arc::new(service);
let staged_capacity_admissions: StagedCapacityAdmissions =
Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
let archive_runtime_cleanup = ArchiveRuntimeCleanup {
runtime_adapter: Arc::clone(&runtime_adapter),
pending_session_event_streams: None,
mcp_state: None,
mob_state: None,
};
Fixture {
service,
staged_sessions,
staged_capacity_admissions,
archive_runtime_cleanup,
runtime_adapter,
_temp: temp,
}
}
fn orchestrator(fx: &Fixture) -> LiveOrchestrator<'_> {
LiveOrchestrator {
service: &fx.service,
staged_sessions: &fx.staged_sessions,
staged_capacity_admissions: &fx.staged_capacity_admissions,
runtime_adapter: &fx.runtime_adapter,
host: None,
config_runtime: None,
default_llm_client: None,
agent_llm_client_decorator: None,
external_tools: None,
archive_runtime_cleanup: fx.archive_runtime_cleanup.clone(),
realm_id: None,
instance_id: None,
backend: None,
}
}
fn now_secs() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0)
}
fn staged_slot(model: &str, provider: meerkat_core::Provider) -> StagedSlot {
let mut build_config = AgentBuildConfig::new(model.to_string());
build_config.provider = Some(provider);
let now = now_secs();
StagedSlot::new_staged(
&SessionId::new(),
build_config,
SessionLlmIdentity {
model: model.to_string(),
provider,
self_hosted_server_id: None,
provider_params: None,
auth_binding: None,
},
None,
None,
Vec::new(),
now,
now,
false,
)
.expect("staged slot fixture must satisfy generated staged-session authority")
}
#[tokio::test]
async fn propagate_config_to_live_channels_no_host_is_noop() {
let fx = build_fixture();
let orch = orchestrator(&fx);
let report = orch.propagate_config_to_live_channels().await;
assert!(report.is_clean());
assert_eq!(report, Default::default());
}
#[tokio::test]
async fn precheck_live_open_rejects_non_realtime_staged_session() {
let fx = build_fixture();
let orch = orchestrator(&fx);
let session_id = SessionId::new();
let slot = staged_slot("claude-opus-4-8", meerkat_core::Provider::Anthropic);
fx.staged_sessions
.stage(session_id.clone(), slot)
.await
.expect("stage deferred session");
let err = orch
.precheck_live_open(&session_id)
.await
.expect_err("non-realtime staged session must be rejected");
match err {
LiveOpenPrecheckError::ModelNotRealtime { model, provider } => {
assert_eq!(model, "claude-opus-4-8");
assert_eq!(provider, "anthropic");
}
other => panic!("expected ModelNotRealtime, got {other:?}"),
}
}
}