use super::state::{AgentState, reduce};
use anyhow::{Context, Result};
use ringo_agent::audio::{self, ToneAnalysis};
use ringo_agent::{AgentConfig, ProcessClient};
use ringo_core::account::{Account, BackendOptions};
use ringo_core::event::AppEvent;
use ringo_core::event::InviteHeaders;
use ringo_core::event::MediaStats;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use tokio::sync::watch;
const TRACE_POLL_INTERVAL: Duration = Duration::from_millis(150);
const RX_BUFFER_CAP_SAMPLES: usize = 48_000 * 30;
#[derive(Default)]
struct RxBuffer {
samples: Vec<i16>,
rate: u32,
}
pub struct AgentSession {
pub aor: String,
pub regint: u32,
client: ProcessClient,
state_rx: watch::Receiver<AgentState>,
rx_audio: Mutex<Option<Arc<Mutex<RxBuffer>>>>,
}
impl AgentSession {
pub async fn connect(name: &str, account: Account, options: &BackendOptions) -> Result<Self> {
let config = agent_config(name, &account, options);
let (client, events) =
ProcessClient::spawn(config).with_context(|| format!("spawn agent `{name}`"))?;
let (state_tx, state_rx) = watch::channel(AgentState::default());
let state_tx = Arc::new(state_tx);
let (async_event_tx, mut async_event_rx) = tokio::sync::mpsc::channel::<AppEvent>(64);
tokio::task::spawn_blocking(move || {
while let Ok(event) = events.recv() {
if async_event_tx.blocking_send(event).is_err() {
break;
}
}
});
let reader_tx = Arc::clone(&state_tx);
tokio::spawn(async move {
while let Some(event) = async_event_rx.recv().await {
reader_tx.send_modify(|s| reduce(s, &event));
}
});
let headers = client.headers_handle();
let trace_tx = Arc::clone(&state_tx);
tokio::spawn(async move {
let mut ticker = tokio::time::interval(TRACE_POLL_INTERVAL);
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
ticker.tick().await;
if trace_tx.is_closed() {
break;
}
let new = headers.lock().unwrap_or_else(|e| e.into_inner()).take();
if let Some(invites) = new {
merge_received_headers(&trace_tx, invites);
}
}
});
Ok(Self {
aor: format!("sip:{}@{}", account.username, account.domain),
regint: account.regint.unwrap_or(3600),
client,
state_rx,
rx_audio: Mutex::new(None),
})
}
pub fn domain(&self) -> &str {
self.aor.rsplit('@').next().unwrap_or("")
}
pub fn state(&self) -> watch::Receiver<AgentState> {
self.state_rx.clone()
}
pub fn set_audio_source(&self, spec: &str) {
self.client.set_audio_source(spec);
}
pub fn analyze_tone(&self, freq: u32, window: Duration) -> ToneAnalysis {
let buf = self.rx_audio_buffer();
let g = buf.lock().unwrap_or_else(|e| e.into_inner());
if g.rate == 0 {
return ToneAnalysis::default();
}
audio::analyze_tone_samples(&g.samples, g.rate, freq, window)
}
pub fn prime_received_audio(&self) {
let buf = self.rx_audio_buffer();
let mut g = buf.lock().unwrap_or_else(|e| e.into_inner());
g.samples.clear();
}
fn rx_audio_buffer(&self) -> Arc<Mutex<RxBuffer>> {
let mut slot = self.rx_audio.lock().unwrap_or_else(|e| e.into_inner());
if let Some(buf) = slot.as_ref() {
return Arc::clone(buf);
}
let buf = Arc::new(Mutex::new(RxBuffer::default()));
let rx = self.client.start_rx_audio();
let drain = Arc::clone(&buf);
std::thread::spawn(move || {
for frame in rx {
let mut g = drain.lock().unwrap_or_else(|e| e.into_inner());
g.rate = frame.rate;
g.samples.extend_from_slice(&frame.samples);
if g.samples.len() > RX_BUFFER_CAP_SAMPLES {
let excess = g.samples.len() - RX_BUFFER_CAP_SAMPLES;
g.samples.drain(0..excess);
}
}
});
*slot = Some(Arc::clone(&buf));
buf
}
pub fn save_audio(&self, prefix: &str) -> Vec<String> {
self.client.save_audio(prefix)
}
pub fn call_count(&self) -> u32 {
self.client.call_count()
}
pub fn request_shutdown(&self) {
self.client.request_shutdown();
}
pub fn media_stats(&self) -> Option<MediaStats> {
self.client.media_stats()
}
pub fn received_dtmf(&self) -> String {
self.client.received_dtmf()
}
pub fn register(&self) {
self.client.register(&self.aor, self.regint);
}
pub fn dial(&self, target: &str) {
self.client.dial(target);
}
pub fn accept(&self) {
self.client.accept();
}
pub fn hold(&self) {
self.client.hold();
}
pub fn resume(&self) {
self.client.resume();
}
pub fn mute(&self) {
self.client.mute();
}
pub fn send_dtmf(&self, digit: char) {
self.client.send_dtmf(digit);
}
pub fn add_header(&self, key: &str, value: &str) {
self.client.add_header(key, value);
}
pub fn hangup(&self) {
self.client.hangup();
}
pub fn hangup_all(&self) {
self.client.hangup_all();
}
pub fn transfer(&self, uri: &str) {
self.client.transfer(uri);
}
pub fn attended_transfer_start(&self, uri: &str) {
self.client.attended_transfer_start(uri);
}
pub fn attended_transfer_exec(&self) {
self.client.attended_transfer_exec();
}
pub fn attended_transfer_abort(&self) {
self.client.attended_transfer_abort();
}
pub fn deflect_incoming(&self, contact: &str, diversion: Option<&str>) {
self.client.deflect_incoming(contact, diversion);
}
pub fn arm_invite_response(&self, scode: u16, reason: &str, headers: Vec<String>) {
self.client.arm_invite_response(scode, reason, headers);
}
pub fn disarm_invite_response(&self) {
self.client.disarm_invite_response();
}
}
fn agent_config(name: &str, account: &Account, options: &BackendOptions) -> AgentConfig {
let mut account = account.clone();
account.catchall = true;
AgentConfig {
name: name.to_string(),
account,
options: options.clone(),
}
}
fn merge_received_headers(state_tx: &watch::Sender<AgentState>, invites: InviteHeaders) {
if invites.is_empty() {
return;
}
let has_new = {
let cur = state_tx.borrow();
invites
.keys()
.any(|k| !cur.received_headers.contains_key(k))
};
if has_new {
state_tx.send_modify(|s| {
for (call_id, headers) in invites {
s.received_headers.entry(call_id).or_insert(headers);
}
});
}
}