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};
#[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_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"
);
}
}