use super::{executor::ResponseStateInput, StateMachine};
use crate::{
errors::{Result, SessionError},
session_registry::SessionRegistryHandle,
state_table::types::{EventType, Role},
types::{CallState, SessionId, SessionInfo},
};
use std::collections::HashMap;
use std::sync::Arc;
use tokio::sync::RwLock;
type SessionSubscriber = Box<dyn Fn(SessionEvent) + Send + Sync>;
type SessionSubscribers = Arc<RwLock<HashMap<SessionId, Vec<SessionSubscriber>>>>;
pub struct StateMachineHelpers {
pub state_machine: Arc<StateMachine>,
subscribers: SessionSubscribers,
}
#[derive(Clone)]
pub enum SessionEvent {
StateChanged { from: CallState, to: CallState },
CallEstablished,
CallTerminated { reason: String },
MediaReady,
IncomingCall { from: String },
}
impl std::fmt::Debug for SessionEvent {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::StateChanged { from, to } => formatter
.debug_struct("StateChanged")
.field("from", from)
.field("to", to)
.finish(),
Self::CallEstablished => formatter.write_str("CallEstablished"),
Self::CallTerminated { reason } => formatter
.debug_struct("CallTerminated")
.field("reason_bytes", &reason.len())
.finish(),
Self::MediaReady => formatter.write_str("MediaReady"),
Self::IncomingCall { from } => formatter
.debug_struct("IncomingCall")
.field("from_bytes", &from.len())
.finish(),
}
}
}
impl StateMachineHelpers {
pub fn new(state_machine: Arc<StateMachine>) -> Self {
Self {
state_machine,
subscribers: Arc::new(RwLock::new(HashMap::new())),
}
}
pub async fn create_session(
&self,
session_id: SessionId,
from: String,
to: String,
role: Role,
) -> Result<()> {
self.state_machine
.store
.create_session_initialized(
session_id.clone(),
role,
true, |session| {
session.local_uri = Some(from.clone());
session.remote_uri = Some(to.clone());
},
)
.await?;
Ok(())
}
pub async fn make_transfer_leg(
&self,
from: &str,
to: &str,
transferor_session_id: &SessionId,
) -> Result<SessionId> {
self.make_call_inner(
from,
to,
None,
Some(transferor_session_id.clone()),
None,
Vec::new(),
)
.await
}
pub async fn set_transferor_session(
&self,
leg_session_id: &SessionId,
transferor_session_id: &SessionId,
) -> Result<()> {
self.state_machine
.set_transferor_session(leg_session_id, transferor_session_id.clone())
.await?;
Ok(())
}
async fn make_call_inner(
&self,
from: &str,
to: &str,
credentials: Option<crate::types::Credentials>,
transferor_session_id: Option<SessionId>,
pai_uri: Option<String>,
extra_headers: Vec<rvoip_sip_core::types::TypedHeader>,
) -> Result<SessionId> {
let transferor_lifecycle_handle = match transferor_session_id.as_ref() {
Some(transferor) => Some(
self.state_machine
.store
.lifecycle_handle(transferor)
.ok_or_else(|| SessionError::SessionNotFound(transferor.to_string()))?,
),
None => None,
};
let session_id = SessionId::new();
self.create_session(
session_id.clone(),
from.to_string(),
to.to_string(),
Role::UAC,
)
.await?;
self.state_machine
.process_event_with_outbound_session_input(
&session_id,
EventType::MakeCall {
target: to.to_string(),
},
super::executor::OutboundSessionStateInput::new(
credentials,
None,
pai_uri,
transferor_session_id,
transferor_lifecycle_handle,
extra_headers,
),
)
.await?;
Ok(session_id)
}
pub async fn accept_call(&self, session_id: &SessionId) -> Result<()> {
self.state_machine
.process_event(session_id, EventType::AcceptCall)
.await?;
Ok(())
}
pub(crate) async fn accept_call_exact(&self, handle: &SessionRegistryHandle) -> Result<()> {
self.state_machine
.process_event_exact(handle, EventType::AcceptCall)
.await?;
Ok(())
}
pub async fn accept_call_with_sdp(&self, session_id: &SessionId, sdp: String) -> Result<()> {
self.state_machine
.process_event_with_local_sdp(session_id, EventType::AcceptCall, sdp)
.await?;
Ok(())
}
pub(crate) async fn accept_call_with_sdp_exact(
&self,
handle: &SessionRegistryHandle,
sdp: String,
) -> Result<()> {
self.state_machine
.process_event_with_response_input_exact(
handle,
EventType::AcceptCall,
ResponseStateInput::accept(Some(sdp), Vec::new()),
)
.await?;
Ok(())
}
pub(crate) async fn accept_call_with_response(
&self,
session_id: &SessionId,
handle: &SessionRegistryHandle,
sdp: Option<String>,
extras: Vec<rvoip_sip_core::types::TypedHeader>,
) -> Result<()> {
if handle.session_id() != session_id {
return Err(SessionError::InvalidInput(
"response lifecycle handle does not match its session".to_string(),
));
}
self.state_machine
.process_event_with_response_input_exact(
handle,
EventType::AcceptCall,
ResponseStateInput::accept(sdp, extras),
)
.await?;
Ok(())
}
pub async fn send_early_media(
&self,
session_id: &SessionId,
sdp: Option<String>,
) -> Result<()> {
self.state_machine
.process_event(session_id, EventType::SendEarlyMedia { sdp })
.await?;
Ok(())
}
pub(crate) async fn send_early_media_exact(
&self,
handle: &SessionRegistryHandle,
sdp: Option<String>,
) -> Result<()> {
self.state_machine
.process_event_exact(handle, EventType::SendEarlyMedia { sdp })
.await?;
Ok(())
}
pub(crate) async fn send_provisional_with_response(
&self,
session_id: &SessionId,
handle: &SessionRegistryHandle,
status: u16,
sdp: Option<String>,
extras: Vec<rvoip_sip_core::types::TypedHeader>,
) -> Result<()> {
if handle.session_id() != session_id {
return Err(SessionError::InvalidInput(
"response lifecycle handle does not match its session".to_string(),
));
}
self.state_machine
.process_event_with_response_input_exact(
handle,
EventType::SendEarlyMedia { sdp },
ResponseStateInput::provisional(status, extras),
)
.await?;
Ok(())
}
pub async fn reject_call(
&self,
session_id: &SessionId,
status: u16,
reason: &str,
) -> Result<()> {
self.state_machine
.process_event(
session_id,
EventType::RejectCall {
status,
reason: reason.to_string(),
},
)
.await?;
Ok(())
}
pub(crate) async fn reject_call_exact(
&self,
handle: &SessionRegistryHandle,
status: u16,
reason: &str,
) -> Result<()> {
self.state_machine
.process_event_exact(
handle,
EventType::RejectCall {
status,
reason: reason.to_string(),
},
)
.await?;
Ok(())
}
pub(crate) async fn reject_call_with_extras_exact(
&self,
handle: &SessionRegistryHandle,
status: u16,
reason: &str,
extras: Vec<rvoip_sip_core::types::TypedHeader>,
) -> Result<()> {
self.state_machine
.process_event_with_response_input_exact(
handle,
EventType::RejectCall {
status,
reason: reason.to_string(),
},
ResponseStateInput::headers(extras),
)
.await?;
Ok(())
}
pub async fn redirect_call(
&self,
session_id: &SessionId,
status: u16,
contacts: Vec<String>,
) -> Result<()> {
self.state_machine
.process_event(session_id, EventType::RedirectCall { status, contacts })
.await?;
Ok(())
}
pub(crate) async fn redirect_call_exact(
&self,
handle: &SessionRegistryHandle,
status: u16,
contacts: Vec<String>,
) -> Result<()> {
self.state_machine
.process_event_exact(handle, EventType::RedirectCall { status, contacts })
.await?;
Ok(())
}
pub(crate) async fn redirect_call_with_extras_exact(
&self,
handle: &SessionRegistryHandle,
status: u16,
contacts: Vec<String>,
extras: Vec<rvoip_sip_core::types::TypedHeader>,
) -> Result<()> {
self.state_machine
.process_event_with_response_input_exact(
handle,
EventType::RedirectCall { status, contacts },
ResponseStateInput::headers(extras),
)
.await?;
Ok(())
}
pub async fn hangup(&self, session_id: &SessionId) -> Result<()> {
let handle = self
.state_machine
.store
.lifecycle_handle(session_id)
.ok_or_else(|| SessionError::SessionNotFound(session_id.to_string()))?;
self.hangup_exact(&handle).await
}
pub(crate) async fn hangup_exact(&self, handle: &SessionRegistryHandle) -> Result<()> {
self.state_machine
.process_event_exact(handle, EventType::HangupCall)
.await?;
Ok(())
}
pub async fn create_conference(&self, session_id: &SessionId, name: &str) -> Result<()> {
self.state_machine
.process_event(
session_id,
EventType::CreateConference {
name: name.to_string(),
},
)
.await?;
Ok(())
}
pub async fn add_to_conference(
&self,
host_session_id: &SessionId,
participant_session_id: &SessionId,
) -> Result<()> {
self.state_machine
.process_event(
host_session_id,
EventType::AddParticipant {
session_id: participant_session_id.to_string(),
},
)
.await?;
Ok(())
}
pub async fn get_session_info(&self, session_id: &SessionId) -> Result<SessionInfo> {
let session = self
.state_machine
.store
.with_session(session_id, Clone::clone)?;
Ok(session_info_from_state(session))
}
pub(crate) async fn get_session_info_exact(
&self,
handle: &SessionRegistryHandle,
) -> Result<SessionInfo> {
let session = self.state_machine.store.get_session_exact(handle).await?;
Ok(session_info_from_state(session))
}
pub async fn list_sessions(&self) -> Vec<SessionInfo> {
let sessions = self.state_machine.store.get_all_sessions().await;
sessions.into_iter().map(session_info_from_state).collect()
}
#[cfg(feature = "perf-tests")]
pub async fn perf_diagnostic_counts(&self) -> serde_json::Value {
let active_sessions = self.state_machine.store.get_all_sessions().await.len();
let subscribers = self.subscribers.read().await.len();
serde_json::json!({
"active_sessions": active_sessions,
"subscriber_sessions": subscribers,
})
}
pub async fn get_state(&self, session_id: &SessionId) -> Result<CallState> {
Ok(self
.state_machine
.store
.with_session(session_id, |session| session.call_state)?)
}
pub(crate) async fn get_state_exact(
&self,
handle: &SessionRegistryHandle,
) -> Result<CallState> {
Ok(self
.state_machine
.store
.get_session_snapshot_exact(handle)?
.call_state)
}
pub(crate) async fn negotiated_media_config(
&self,
session_id: &SessionId,
) -> Result<Option<(crate::session_store::state::NegotiatedConfig, u8)>> {
let (negotiated_config, payload_type, sdp_negotiated) = self
.state_machine
.store
.with_session(session_id, |session| {
(
session.negotiated_config.clone(),
session.negotiated_payload_type(),
session.sdp_negotiated,
)
})?;
match negotiated_config.zip(payload_type) {
Some(config) => Ok(Some(config)),
None if sdp_negotiated => Err(crate::errors::SessionError::MediaError(
"SDP was supplied without an anchored negotiated media configuration".to_string(),
)),
None => Ok(None),
}
}
pub async fn is_in_conference(&self, session_id: &SessionId) -> Result<bool> {
let _ = session_id;
Ok(false)
}
pub async fn subscribe<F>(&self, session_id: SessionId, callback: F)
where
F: Fn(SessionEvent) + Send + Sync + 'static,
{
self.subscribers
.write()
.await
.entry(session_id)
.or_insert_with(Vec::new)
.push(Box::new(callback));
}
pub async fn unsubscribe(&self, session_id: &SessionId) {
self.subscribers.write().await.remove(session_id);
}
pub(crate) async fn notify_subscribers(&self, session_id: &SessionId, event: SessionEvent) {
if let Some(callbacks) = self.subscribers.read().await.get(session_id) {
for callback in callbacks {
callback(event.clone());
}
}
}
pub(crate) async fn cleanup_session(&self, session_id: &SessionId) {
self.subscribers.write().await.remove(session_id);
}
}
fn session_info_from_state(session: crate::session_store::SessionState) -> SessionInfo {
let start_time = std::time::SystemTime::now()
.checked_sub(session.session_duration())
.unwrap_or(std::time::SystemTime::UNIX_EPOCH);
SessionInfo {
session_id: session.session_id,
from: session.local_uri.unwrap_or_default(),
to: session.remote_uri.unwrap_or_default(),
state: session.call_state,
start_time,
media_active: session.media_session_id.is_some(),
}
}
#[cfg(test)]
mod single_session_view_tests {
use super::*;
#[test]
fn session_info_is_projected_from_the_canonical_session_state() {
let session_id = SessionId::new();
let mut session = crate::session_store::SessionState::new(session_id.clone(), Role::UAC);
session.local_uri = Some("sip:alice@example.test".to_string());
session.remote_uri = Some("sip:bob@example.test".to_string());
session.call_state = CallState::Active;
session.media_session_id = Some(crate::types::MediaSessionId::new("media-exact"));
let info = session_info_from_state(session);
assert_eq!(info.session_id, session_id);
assert_eq!(info.from, "sip:alice@example.test");
assert_eq!(info.to, "sip:bob@example.test");
assert_eq!(info.state, CallState::Active);
assert!(info.media_active);
}
#[test]
fn helpers_have_no_second_mutable_active_session_view() {
let production = include_str!("helpers.rs")
.split("#[cfg(test)]")
.next()
.expect("production helper source");
assert!(!production.contains("active_sessions:"));
assert!(!production.contains("self.active_sessions"));
assert!(production.contains("with_session(session_id, Clone::clone)"));
}
}