use std::sync::Arc;
use async_trait::async_trait;
use meerkat_contracts::wire::supervisor_bridge::{BridgeLiveControlOutcome, BridgeLiveControlVerb};
use meerkat_contracts::{LiveCloseStatus, LiveOpenResult, LiveOpenTransport, RealtimeTurningMode};
use meerkat_core::connection::RealmId;
use meerkat_core::types::SessionId;
use meerkat_live::{LiveAdapterHost, LiveChannelId, LiveWsState};
use meerkat_llm_core::realtime_session::RealtimeSessionFactory;
use meerkat_runtime::MeerkatMachine;
use meerkat_runtime::member_live::{
MEMBER_LIVE_OPEN_CEILING, MemberLiveError, MemberLiveHost, MemberLiveStatus,
};
use crate::service_factory::FactoryAgentBuilder;
use crate::session_runtime::admission::StagedCapacityAdmissions;
use crate::session_runtime::errors::{LiveChannelVerbError, LiveIngressError, LiveOpenError};
use crate::session_runtime::live_orchestration::{
LiveOrchestrator, LiveSessionIngressReconciler, LiveTransportContext,
};
use crate::session_runtime::runtime_state::ArchiveRuntimeCleanup;
use crate::{PersistentSessionService, StagedSessionRegistry};
pub struct MobOwnedOnlyIngress;
#[async_trait]
impl LiveSessionIngressReconciler for MobOwnedOnlyIngress {
async fn ensure_session_owned_live_ingress(
&self,
session_id: &SessionId,
) -> Result<(), LiveIngressError> {
Err(LiveIngressError::Internal(format!(
"member-host live open reached session-owned peer ingress for {session_id}; \
member sessions are mob-owned by construction"
)))
}
}
static MOB_OWNED_ONLY_INGRESS: MobOwnedOnlyIngress = MobOwnedOnlyIngress;
pub struct ServiceMemberLiveHostConfig {
pub service: Arc<PersistentSessionService<FactoryAgentBuilder>>,
pub runtime_adapter: Arc<MeerkatMachine>,
pub host: Arc<LiveAdapterHost>,
pub ws_state: Option<Arc<LiveWsState>>,
pub base_url: Option<String>,
pub session_factory: Arc<dyn RealtimeSessionFactory>,
pub realm_id: Option<RealmId>,
pub instance_id: Option<String>,
pub backend: Option<String>,
}
pub struct ServiceMemberLiveHost {
service: Arc<PersistentSessionService<FactoryAgentBuilder>>,
staged_sessions: Arc<StagedSessionRegistry>,
staged_capacity_admissions: StagedCapacityAdmissions,
runtime_adapter: Arc<MeerkatMachine>,
host: Arc<LiveAdapterHost>,
ws_state: Option<Arc<LiveWsState>>,
base_url: Option<String>,
#[cfg(feature = "live-webrtc")]
webrtc_state: Option<Arc<meerkat_live::LiveWebrtcState>>,
session_factory: Arc<dyn RealtimeSessionFactory>,
realm_id: Option<RealmId>,
instance_id: Option<String>,
backend: Option<String>,
}
impl ServiceMemberLiveHost {
#[must_use]
pub fn new(config: ServiceMemberLiveHostConfig) -> Self {
Self {
service: config.service,
staged_sessions: Arc::new(StagedSessionRegistry::new()),
staged_capacity_admissions: Arc::new(std::sync::Mutex::new(
std::collections::HashMap::new(),
)),
runtime_adapter: config.runtime_adapter,
host: config.host,
ws_state: config.ws_state,
base_url: config.base_url,
#[cfg(feature = "live-webrtc")]
webrtc_state: None,
session_factory: config.session_factory,
realm_id: config.realm_id,
instance_id: config.instance_id,
backend: config.backend,
}
}
fn orchestrator(&self) -> LiveOrchestrator<'_> {
LiveOrchestrator {
service: &self.service,
staged_sessions: &self.staged_sessions,
staged_capacity_admissions: &self.staged_capacity_admissions,
runtime_adapter: &self.runtime_adapter,
host: Some(Arc::clone(&self.host)),
config_runtime: None,
default_llm_client: None,
agent_llm_client_decorator: None,
external_tools: None,
archive_runtime_cleanup: ArchiveRuntimeCleanup {
runtime_adapter: Arc::clone(&self.runtime_adapter),
pending_session_event_streams: None,
mcp_state: None,
mob_state: None,
},
realm_id: self.realm_id.as_ref(),
instance_id: self.instance_id.as_deref(),
backend: self.backend.as_deref(),
ingress_reconciler: Some(&MOB_OWNED_ONLY_INGRESS),
}
}
#[cfg(feature = "live-webrtc")]
#[must_use]
pub fn with_webrtc_cleanup_state(mut self, state: Arc<meerkat_live::LiveWebrtcState>) -> Self {
self.webrtc_state = Some(state);
self
}
fn transport_context(&self) -> LiveTransportContext<'_> {
LiveTransportContext::new(self.ws_state.as_deref(), self.base_url.as_deref())
}
}
#[async_trait]
impl MemberLiveHost for ServiceMemberLiveHost {
async fn open(
&self,
session: &SessionId,
turning_mode: Option<RealtimeTurningMode>,
transport: Option<LiveOpenTransport>,
) -> Result<LiveOpenResult, MemberLiveError> {
let orchestrator = self.orchestrator();
let opened = tokio::time::timeout(
MEMBER_LIVE_OPEN_CEILING,
orchestrator.open_live_channel(
&self.host,
self.transport_context(),
Some(self.session_factory.as_ref()),
session,
turning_mode,
transport,
),
)
.await;
match opened {
Ok(result) => result.map_err(member_live_error_from_open),
Err(_elapsed) => {
if let Some(channel_id) = self
.runtime_adapter
.live_active_channel_for_session(session)
.await
{
self.orchestrator()
.close_live_channel_after_open_failure(&self.host, session, &channel_id)
.await;
}
Err(MemberLiveError::Unavailable {
reason: format!(
"live open exceeded the member ceiling of {}s and was aborted fail-closed",
MEMBER_LIVE_OPEN_CEILING.as_secs()
),
})
}
}
}
async fn close(
&self,
session: &SessionId,
channel_id: &str,
) -> Result<LiveCloseStatus, MemberLiveError> {
let channel = LiveChannelId::new(channel_id);
#[cfg(feature = "live-webrtc")]
if let Some(webrtc_state) = self.webrtc_state.as_ref() {
webrtc_state
.close_peer_checked(&channel)
.await
.map_err(|error| MemberLiveError::Unavailable {
reason: format!(
"WebRTC physical cleanup failed for live channel '{channel_id}': {error}"
),
})?;
}
self.orchestrator()
.close_live_channel(&self.host, &channel, Some(session))
.await
.map(|result| result.status)
.map_err(member_live_error_from_verb)
}
async fn status(
&self,
session: &SessionId,
channel_id: Option<String>,
) -> Result<MemberLiveStatus, MemberLiveError> {
let channel = match channel_id {
Some(channel_id) => LiveChannelId::new(channel_id),
None => self
.runtime_adapter
.live_active_channel_for_session(session)
.await
.ok_or(MemberLiveError::ChannelNotFound)?,
};
let status = self
.orchestrator()
.live_channel_status(&self.host, &channel, Some(session))
.await
.map_err(member_live_error_from_verb)?;
Ok(MemberLiveStatus {
channel_id: channel.to_string(),
status,
})
}
async fn control(
&self,
session: &SessionId,
channel_id: &str,
verb: BridgeLiveControlVerb,
) -> Result<BridgeLiveControlOutcome, MemberLiveError> {
let channel = LiveChannelId::new(channel_id);
self.orchestrator()
.control_live_channel(
&self.host,
self.transport_context(),
&channel,
Some(session),
verb,
)
.await
.map_err(member_live_error_from_verb)
}
}
#[must_use]
pub fn member_live_error_from_open(error: LiveOpenError) -> MemberLiveError {
use crate::session_runtime::errors::LiveOpenPrecheckError;
match error {
LiveOpenError::SessionNotFound { .. } => MemberLiveError::Unavailable {
reason: error.to_string(),
},
LiveOpenError::RealtimeFactoryMissing => MemberLiveError::TransportUnavailable,
LiveOpenError::NoTransportConfigured | LiveOpenError::WebsocketNotConfigured => {
MemberLiveError::TransportUnavailable
}
LiveOpenError::AdmissionRejectedAlreadyBound { .. } => MemberLiveError::ChannelAlreadyBound,
LiveOpenError::AdmissionRejectedLifecycleClosed => MemberLiveError::Unavailable {
reason: "session lifecycle is closed to live channel admission".to_string(),
},
LiveOpenError::Precheck(LiveOpenPrecheckError::ModelNotRealtime { model, provider }) => {
MemberLiveError::ModelNotRealtime {
model,
provider: provider.to_string(),
}
}
LiveOpenError::Precheck(LiveOpenPrecheckError::ProviderHasNoLiveAdapter { provider }) => {
MemberLiveError::AdapterUnavailable {
provider: provider.to_string(),
}
}
LiveOpenError::ProviderUnsupportedByFactory { provider } => {
MemberLiveError::AdapterUnavailable {
provider: provider.to_string(),
}
}
LiveOpenError::WebrtcNotConfigured | LiveOpenError::WebrtcNotCompiled => {
MemberLiveError::TransportUnsupported {
requested: "webrtc".to_string(),
}
}
LiveOpenError::UnsupportedTransport => MemberLiveError::TransportUnsupported {
requested: "unknown".to_string(),
},
error @ (LiveOpenError::SessionStateFault(_)
| LiveOpenError::OpenConfig(_)
| LiveOpenError::AdmissionAuthority(_)
| LiveOpenError::AdmissionRejectedChannelCollision { .. }
| LiveOpenError::AdmissionRejectedNoReason
| LiveOpenError::MissingHostHandoff
| LiveOpenError::HostOpenSessionAlreadyBound { .. }
| LiveOpenError::HostOpen(_)
| LiveOpenError::Precheck(LiveOpenPrecheckError::SessionLookup { .. })
| LiveOpenError::AdapterOpen(_)
| LiveOpenError::AdapterAttach(_)
| LiveOpenError::Ingress(_)
| LiveOpenError::TokenMint(_)
| LiveOpenError::AudioPolicyMissing
| LiveOpenError::AudioFormatUnmappable { .. }
| LiveOpenError::WebrtcClock(_)
| LiveOpenError::WebrtcTokenMint(_)) => MemberLiveError::Internal {
reason: error.to_string(),
},
}
}
#[must_use]
pub fn member_live_error_from_verb(error: LiveChannelVerbError) -> MemberLiveError {
use meerkat_runtime::meerkat_machine::dsl::{
LiveChannelRequestRejectionReason, LiveCommandRejectionReason,
};
match error {
LiveChannelVerbError::UnboundCommand { .. }
| LiveChannelVerbError::UnboundRequest { .. }
| LiveChannelVerbError::SessionPinMismatch { .. } => MemberLiveError::ChannelNotFound,
LiveChannelVerbError::CommandRejected {
authority, detail, ..
} => match authority.rejection {
LiveCommandRejectionReason::ChannelNotFound
| LiveCommandRejectionReason::NoAdapter
| LiveCommandRejectionReason::ChannelNotReady => {
MemberLiveError::Unavailable { reason: detail }
}
LiveCommandRejectionReason::UnsupportedCommand
| LiveCommandRejectionReason::AdapterError
| LiveCommandRejectionReason::InternalHostError => {
MemberLiveError::Internal { reason: detail }
}
},
LiveChannelVerbError::RequestRejected {
authority, detail, ..
} => {
let reason = detail.unwrap_or_else(|| "live channel request rejected".to_string());
match authority.rejection {
LiveChannelRequestRejectionReason::ChannelNotFound
| LiveChannelRequestRejectionReason::NoAdapter => {
MemberLiveError::Unavailable { reason }
}
LiveChannelRequestRejectionReason::InvalidToken
| LiveChannelRequestRejectionReason::InvalidPayload
| LiveChannelRequestRejectionReason::WebrtcAnswerError
| LiveChannelRequestRejectionReason::InternalHostError => {
MemberLiveError::Internal { reason }
}
}
}
error @ (LiveChannelVerbError::RejectionAuthorityFailed { .. }
| LiveChannelVerbError::ResultAuthority { .. }
| LiveChannelVerbError::CommitOmitted
| LiveChannelVerbError::HostCommit { .. }
| LiveChannelVerbError::ResultProjection { .. }
| LiveChannelVerbError::RefreshConfig(_)) => MemberLiveError::Internal {
reason: error.to_string(),
},
}
}
pub struct ServiceLiveToolDispatcher {
service: Arc<PersistentSessionService<FactoryAgentBuilder>>,
}
impl ServiceLiveToolDispatcher {
#[must_use]
pub fn new(service: Arc<PersistentSessionService<FactoryAgentBuilder>>) -> Self {
Self { service }
}
}
#[async_trait]
impl meerkat_live::LiveToolDispatcher for ServiceLiveToolDispatcher {
async fn dispatch_live_tool_call(
&self,
session_id: &SessionId,
call: meerkat_core::ToolCall,
) -> Result<meerkat_core::ops::ToolDispatchOutcome, meerkat_live::LiveToolDispatchError> {
self.service
.dispatch_external_tool_call_with_timeout_policy(
session_id,
call,
meerkat_core::ToolDispatchTimeoutPolicy::Disabled,
)
.await
.map_err(|err| meerkat_live::LiveToolDispatchError::from_session_error(session_id, err))
}
}