use std::sync::Arc;
use async_trait::async_trait;
use meerkat_client::{FactoryError, LlmError};
use meerkat_contracts::RealtimeCapabilities;
use meerkat_core::{Config, ConfigError, ConfigStore, Provider, SessionLlmIdentity};
use meerkat_llm_core::realtime_session::{
RealtimeExternalSessionTarget, RealtimeSession, RealtimeSessionFactory,
RealtimeSessionOpenConfig,
};
use crate::{AgentFactory, RealmInheritance};
pub fn build_per_open_realtime_session_factory(
factory: &AgentFactory,
config_store: Arc<dyn ConfigStore>,
realm_config_source: Arc<dyn meerkat_core::RealmConfigSource>,
realm: meerkat_core::connection::RealmId,
) -> Arc<dyn RealtimeSessionFactory> {
let realtime_config_source = Arc::new(StoreBackedRealtimeConfigSource::new(
config_store,
Some(RealmInheritance::new(realm_config_source, realm)),
));
factory.build_openai_realtime_session_factory(realtime_config_source)
}
#[async_trait]
pub trait RealtimeCurrentConfigSource: Send + Sync {
async fn current_config(&self) -> Result<Config, ConfigError>;
}
pub struct StoreBackedRealtimeConfigSource {
store: Arc<dyn ConfigStore>,
inheritance: Option<RealmInheritance>,
}
impl StoreBackedRealtimeConfigSource {
pub fn new(store: Arc<dyn ConfigStore>, inheritance: Option<RealmInheritance>) -> Self {
Self { store, inheritance }
}
}
#[async_trait]
impl RealtimeCurrentConfigSource for StoreBackedRealtimeConfigSource {
async fn current_config(&self) -> Result<Config, ConfigError> {
let head = self.store.get().await?;
match &self.inheritance {
Some(inheritance) => inheritance.compose_over(head).await,
None => Ok(head),
}
}
}
fn realtime_open_resolution_error(error: FactoryError) -> LlmError {
match error {
FactoryError::ProviderAuth(error) => LlmError::AuthenticationFailed {
message: format!("realtime credential resolution failed: {error}"),
},
FactoryError::TokenStore(error) => LlmError::AuthenticationFailed {
message: format!("realtime credential store unavailable: {error}"),
},
other => LlmError::InvalidConfig {
message: format!("realtime provider resolution failed: {other}"),
},
}
}
pub struct PerOpenCredentialRealtimeSessionFactory {
factory: AgentFactory,
config_source: Arc<dyn RealtimeCurrentConfigSource>,
}
impl PerOpenCredentialRealtimeSessionFactory {
pub fn new(factory: AgentFactory, config_source: Arc<dyn RealtimeCurrentConfigSource>) -> Self {
Self {
factory,
config_source,
}
}
async fn resolve_provider_factory(
&self,
identity: &SessionLlmIdentity,
) -> Result<Arc<dyn RealtimeSessionFactory>, LlmError> {
let config =
self.config_source
.current_config()
.await
.map_err(|error| LlmError::InvalidConfig {
message: format!(
"realtime credential resolution could not read the current config: {error}"
),
})?;
self.factory
.resolve_realtime_session_factory_for_identity(&config, identity)
.await
.map_err(realtime_open_resolution_error)
}
}
#[async_trait]
impl RealtimeSessionFactory for PerOpenCredentialRealtimeSessionFactory {
fn capabilities(&self) -> RealtimeCapabilities {
meerkat_openai::live::openai_realtime_capabilities_default()
}
fn supports_provider(&self, provider: Provider) -> bool {
provider == Provider::OpenAI
}
async fn open_session(
&self,
open_config: &RealtimeSessionOpenConfig,
) -> Result<Box<dyn RealtimeSession>, LlmError> {
let provider_factory = self
.resolve_provider_factory(&open_config.llm_identity)
.await?;
provider_factory.open_session(open_config).await
}
async fn attach_external_session(
&self,
target: &RealtimeExternalSessionTarget,
open_config: &RealtimeSessionOpenConfig,
) -> Result<Box<dyn RealtimeSession>, LlmError> {
let provider_factory = self
.resolve_provider_factory(&open_config.llm_identity)
.await?;
provider_factory
.attach_external_session(target, open_config)
.await
}
async fn open_live_adapter(
&self,
open_config: &RealtimeSessionOpenConfig,
) -> Result<Arc<dyn meerkat_core::live_adapter::LiveAdapter>, LlmError> {
let provider_factory = self
.resolve_provider_factory(&open_config.llm_identity)
.await?;
provider_factory.open_live_adapter(open_config).await
}
}
#[cfg(test)]
#[allow(clippy::expect_used, clippy::unwrap_used, clippy::panic)]
mod tests {
use super::*;
use meerkat_contracts::RealtimeTurningMode;
use meerkat_core::MemoryConfigStore;
use meerkat_core::connection::{RealmConfigSection, RealmId};
struct MapRealmSource {
docs: std::collections::BTreeMap<String, Config>,
}
#[async_trait]
impl meerkat_core::RealmConfigSource for MapRealmSource {
async fn config_for_realm(&self, realm: &RealmId) -> Result<Option<Config>, ConfigError> {
Ok(self.docs.get(realm.as_str()).cloned())
}
}
#[tokio::test]
async fn store_backed_source_serves_current_store_config_not_a_snapshot() {
let store = Arc::new(MemoryConfigStore::new(
Config::default(),
meerkat_models::canonical(),
));
let source =
StoreBackedRealtimeConfigSource::new(Arc::clone(&store) as Arc<dyn ConfigStore>, None);
let initial = source.current_config().await.expect("initial config");
assert_eq!(initial.agent.model, Config::default().agent.model);
let mut updated = Config::default();
updated.agent.model = "gpt-5.5".to_string();
store
.set(updated)
.await
.expect("store update should persist");
let current = source.current_config().await.expect("current config");
assert_eq!(
current.agent.model, "gpt-5.5",
"per-open config source must read the live store, not a startup clone"
);
}
#[tokio::test]
async fn store_backed_source_composes_realm_inheritance_over_head() {
let mut global = Config::default();
global.models.openai = "g-openai".to_string();
global
.realm
.insert("global".to_string(), RealmConfigSection::default());
let mut docs = std::collections::BTreeMap::new();
docs.insert("global".to_string(), global);
let mut head = Config::default();
head.realm.insert(
"child".to_string(),
RealmConfigSection {
parent: Some(RealmId::global()),
..Default::default()
},
);
let store = Arc::new(MemoryConfigStore::new(head, meerkat_models::canonical()));
let source = StoreBackedRealtimeConfigSource::new(
store as Arc<dyn ConfigStore>,
Some(RealmInheritance::new(
Arc::new(MapRealmSource { docs }),
RealmId::parse("child").expect("valid realm"),
)),
);
let composed = source.current_config().await.expect("composed config");
assert_eq!(
composed.models.openai, "g-openai",
"per-open config source must compose the realm parent chain over the head document"
);
}
#[tokio::test]
async fn facade_per_open_factory_resolves_from_live_config_and_fails_closed() {
use meerkat_core::provider_matrix::openai::{OpenAiAuthMethod, OpenAiBackendKind};
use meerkat_core::{
AuthBindingRef, AuthProfileConfig, BackendProfileConfig, BindingId,
CredentialSourceSpec, ProviderBindingConfig,
};
const TEST_REALM: &str = "rt-facade-credentials";
const TEST_BINDING: &str = "openai-main";
const ABSENT_ENV_VAR: &str = "MEERKAT_FACADE_RT_CREDENTIALS_TEST_ABSENT_KEY";
fn openai_env_binding_section() -> RealmConfigSection {
let mut section = RealmConfigSection::default();
section.backend.insert(
"openai-backend".to_string(),
BackendProfileConfig {
provider: Provider::OpenAI.as_str().to_string(),
backend_kind: OpenAiBackendKind::OpenAiApi.as_str().to_string(),
base_url: None,
options: serde_json::Value::Null,
},
);
section.auth.insert(
"openai-auth".to_string(),
AuthProfileConfig {
provider: Provider::OpenAI.as_str().to_string(),
auth_method: OpenAiAuthMethod::ApiKey.as_str().to_string(),
source: CredentialSourceSpec::Env {
env: ABSENT_ENV_VAR.to_string(),
fallback: Vec::new(),
},
constraints: Default::default(),
metadata_defaults: Default::default(),
},
);
section.binding.insert(
TEST_BINDING.to_string(),
ProviderBindingConfig {
backend_profile: "openai-backend".to_string(),
auth_profile: "openai-auth".to_string(),
default_model: None,
policy: Default::default(),
provider_default: false,
},
);
section
}
let identity = SessionLlmIdentity {
model: "gpt-realtime".to_string(),
provider: Provider::OpenAI,
self_hosted_server_id: None,
provider_params: None,
auth_binding: Some(AuthBindingRef {
realm: RealmId::parse(TEST_REALM).expect("valid realm slug"),
binding: BindingId::parse(TEST_BINDING).expect("valid binding id"),
profile: None,
origin: Default::default(),
}),
};
let open_config = RealtimeSessionOpenConfig::new(
RealtimeTurningMode::ProviderManaged,
identity,
Vec::new(),
Vec::new(),
)
.expect("empty seed must be representable");
let temp = tempfile::TempDir::new().expect("tempdir");
let factory =
AgentFactory::new(temp.path().join("sessions")).without_provider_auth_persistence();
let store = Arc::new(MemoryConfigStore::new(
Config::default(),
meerkat_models::canonical(),
));
let realm_source = Arc::new(meerkat_store::FilesystemRealmConfigSource::new(
temp.path().join("state"),
temp.path().join("state").join("global-config.toml"),
meerkat_models::canonical(),
));
let wired = build_per_open_realtime_session_factory(
&factory,
Arc::clone(&store) as Arc<dyn ConfigStore>,
realm_source,
RealmId::parse(TEST_REALM).expect("valid realm slug"),
);
let err = wired
.open_session(&open_config)
.await
.err()
.expect("live open without a resolvable binding must fail closed");
assert!(
matches!(err, LlmError::InvalidConfig { .. }),
"unknown realm/binding must surface as the typed InvalidConfig \
per-open resolution failure, got: {err:?}"
);
let mut config = Config::default();
config
.realm
.insert(TEST_REALM.to_string(), openai_env_binding_section());
store.set(config).await.expect("config store update");
let err = wired
.open_session(&open_config)
.await
.err()
.expect("live open without resolvable credential material must fail closed");
assert!(
matches!(err, LlmError::AuthenticationFailed { .. }),
"missing credential material must surface as the typed \
AuthenticationFailed per-open resolution failure, got: {err:?}"
);
}
}