1use super::Command;
2use crate::{
3 CallOption, ReferOption,
4 event::{EventReceiver, EventSender, SessionEvent},
5 media::{
6 TrackId,
7 ambiance::AmbianceProcessor,
8 engine::StreamEngine,
9 negotiate::strip_ipv6_candidates,
10 processor::SubscribeProcessor,
11 recorder::RecorderOption,
12 stream::{MediaStream, MediaStreamBuilder, SERVER_SIDE_TRACK_ID},
13 track::{
14 Track, TrackConfig,
15 file::FileTrack,
16 forwarding::ForwardingTrack,
17 media_pass::MediaPassTrack,
18 rtc::{RtcTrack, RtcTrackConfig},
19 tts::SynthesisHandle,
20 websocket::{WebsocketBytesReceiver, WebsocketTrack},
21 },
22 },
23 synthesis::{SynthesisCommand, SynthesisOption},
24 transcription::TranscriptionOption,
25};
26use crate::{
27 app::AppState,
28 call::{
29 CommandReceiver, CommandSender,
30 sip::{DialogStateReceiverGuard, Invitation, InviteDialogStates},
31 },
32 callrecord::{CallRecord, CallRecordEvent, CallRecordEventType, CallRecordHangupReason},
33 useragent::{
34 invitation::PendingDialog,
35 public_address::{
36 build_public_contact_uri, contact_needs_public_resolution, find_local_addr_for_uri,
37 },
38 },
39};
40use anyhow::Result;
41use audio_codec::CodecType;
42use chrono::{DateTime, Utc};
43use rsipstack::dialog::{invitation::InviteOption, server_dialog::ServerInviteDialog};
44use serde::{Deserialize, Serialize};
45use std::{
46 collections::HashMap,
47 path::Path,
48 sync::{
49 Arc,
50 atomic::{AtomicBool, Ordering},
51 },
52 time::Duration,
53};
54use tokio::{fs::File, select, sync::Mutex, sync::RwLock, sync::mpsc, time::sleep};
55use tokio_util::sync::CancellationToken;
56use tracing::{debug, info, warn};
57
58#[cfg(test)]
59mod tests {
60 use super::*;
61 use crate::app::AppStateBuilder;
62 use crate::callrecord::CallRecordHangupReason;
63 use crate::config::Config;
64 use crate::media::track::tts::SynthesisHandle;
65 use crate::synthesis::SynthesisCommand;
66 use tokio::sync::mpsc;
67
68 #[tokio::test]
69 async fn test_tts_ssrc_reuse_for_autohangup() -> Result<()> {
70 let mut config = Config::default();
71 config.udp_port = 0; config.media_cache_path = "/tmp/mediacache".to_string();
73 let stream_engine = Arc::new(StreamEngine::default());
74 let app_state = AppStateBuilder::new()
75 .with_config(config)
76 .with_stream_engine(stream_engine)
77 .build()
78 .await?;
79
80 let cancel_token = CancellationToken::new();
81 let session_id = "test-session".to_string();
82 let track_config = TrackConfig::default();
83
84 let mut option = crate::CallOption::default();
85 option.tts = Some(crate::synthesis::SynthesisOption::default());
86
87 let active_call = Arc::new(ActiveCall::new(
88 ActiveCallType::Sip,
89 cancel_token.clone(),
90 session_id.clone(),
91 app_state.invitation.clone(),
92 app_state.clone(),
93 track_config,
94 None,
95 false,
96 None,
97 None,
98 None,
99 ));
100
101 {
102 let mut state = active_call.call_state.write().await;
103 state.option = Some(option);
104 }
105
106 let (tx, mut rx) = mpsc::unbounded_channel::<SynthesisCommand>();
107 let initial_ssrc = 12345;
108 let handle = SynthesisHandle::new(tx, Some("play_1".to_string()), initial_ssrc);
109
110 {
112 let mut state = active_call.call_state.write().await;
113 state.tts_handle = Some(handle);
114 state.current_play_id = Some("play_1".to_string());
115 }
116
117 active_call
119 .do_tts(
120 "hangup now".to_string(),
121 None,
122 Some("play_1".to_string()),
123 Some(true),
124 false,
125 true,
126 None,
127 None,
128 false,
129 None,
130 )
131 .await?;
132
133 {
135 let state = active_call.call_state.read().await;
136 assert!(state.auto_hangup.is_some());
137 let (h_ssrc, reason) = state.auto_hangup.clone().unwrap();
138 assert_eq!(
139 h_ssrc, initial_ssrc,
140 "SSRC should be reused from existing handle"
141 );
142 assert_eq!(reason, CallRecordHangupReason::BySystem);
143 }
144
145 let cmd = rx.try_recv().expect("Should have received tts command");
147 assert_eq!(cmd.text, "hangup now");
148
149 Ok(())
150 }
151
152 #[tokio::test]
153 async fn test_tts_new_ssrc_for_different_play_id() -> Result<()> {
154 let mut config = Config::default();
155 config.udp_port = 0; config.media_cache_path = "/tmp/mediacache".to_string();
157 let stream_engine = Arc::new(StreamEngine::default());
158 let app_state = AppStateBuilder::new()
159 .with_config(config)
160 .with_stream_engine(stream_engine)
161 .build()
162 .await?;
163
164 let active_call = Arc::new(ActiveCall::new(
165 ActiveCallType::Sip,
166 CancellationToken::new(),
167 "test-session".to_string(),
168 app_state.invitation.clone(),
169 app_state.clone(),
170 TrackConfig::default(),
171 None,
172 false,
173 None,
174 None,
175 None,
176 ));
177
178 let mut tts_opt = crate::synthesis::SynthesisOption::default();
179 tts_opt.provider = Some(crate::synthesis::SynthesisType::Aliyun);
180 let mut option = crate::CallOption::default();
181 option.tts = Some(tts_opt);
182 {
183 let mut state = active_call.call_state.write().await;
184 state.option = Some(option);
185 }
186
187 let (tx, _rx) = mpsc::unbounded_channel();
188 let initial_ssrc = 111;
189 let handle = SynthesisHandle::new(tx, Some("play_1".to_string()), initial_ssrc);
190
191 {
192 let mut state = active_call.call_state.write().await;
193 state.tts_handle = Some(handle);
194 state.current_play_id = Some("play_1".to_string());
195 }
196
197 active_call
199 .do_tts(
200 "new play".to_string(),
201 None,
202 Some("play_2".to_string()),
203 Some(true),
204 false,
205 true,
206 None,
207 None,
208 false,
209 None,
210 )
211 .await?;
212
213 {
215 let state = active_call.call_state.read().await;
216 let (h_ssrc, _) = state.auto_hangup.clone().unwrap();
217 assert_ne!(
218 h_ssrc, initial_ssrc,
219 "Should use a new SSRC for different play_id"
220 );
221 }
222
223 Ok(())
224 }
225
226 async fn make_active_call() -> Arc<ActiveCall> {
227 let mut config = Config::default();
228 config.udp_port = 0;
229 config.media_cache_path = "/tmp/mediacache".to_string();
230 let app_state = AppStateBuilder::new()
231 .with_config(config)
232 .with_stream_engine(Arc::new(StreamEngine::default()))
233 .build()
234 .await
235 .unwrap();
236 Arc::new(ActiveCall::new(
237 ActiveCallType::Sip,
238 CancellationToken::new(),
239 "test-session".to_string(),
240 app_state.invitation.clone(),
241 app_state.clone(),
242 TrackConfig::default(),
243 None,
244 false,
245 None,
246 None,
247 None,
248 ))
249 }
250
251 #[tokio::test]
253 async fn test_hangup_refer_true_cancels_refer_only() -> Result<()> {
254 let active_call = make_active_call().await;
255
256 let refer_token = active_call.cancel_token.child_token();
257 let refer_state = Arc::new(RwLock::new(ActiveCallState {
258 ssrc: 1,
259 is_refer: true,
260 ..Default::default()
261 }));
262 {
263 let mut cs = active_call.call_state.write().await;
264 cs.refer_call_token = Some(refer_token.clone());
265 cs.refer_callstate = Some(refer_state.clone());
266 }
267
268 active_call.do_hangup(None, None, None, Some(true)).await?;
269
270 assert!(refer_token.is_cancelled(), "refer token should be cancelled");
271 assert!(
272 !active_call.media_stream.cancel_token.is_cancelled(),
273 "media stream should NOT stop"
274 );
275 assert!(
276 refer_state.read().await.hangup_reason.is_some(),
277 "hangup_reason should be set on refer state"
278 );
279 Ok(())
280 }
281
282 #[tokio::test]
284 async fn test_hangup_none_cancels_refer_too() -> Result<()> {
285 let active_call = make_active_call().await;
286
287 let refer_token = active_call.cancel_token.child_token();
288 {
289 let mut cs = active_call.call_state.write().await;
290 cs.refer_call_token = Some(refer_token.clone());
291 }
292
293 active_call.do_hangup(None, None, None, None).await?;
294
295 assert!(refer_token.is_cancelled(), "refer token should be cancelled");
296 assert!(
297 active_call.media_stream.cancel_token.is_cancelled(),
298 "media stream should stop"
299 );
300 Ok(())
301 }
302}
303
304#[derive(Deserialize)]
305#[serde(rename_all = "camelCase")]
306pub struct CallParams {
307 pub id: Option<String>,
308 #[serde(rename = "dump")]
309 pub dump_events: Option<bool>,
310 #[serde(rename = "ping")]
311 pub ping_interval: Option<u32>,
312 pub server_side_track: Option<String>,
313}
314
315#[derive(Debug, Serialize, Deserialize, Clone, Default, PartialEq, Eq)]
316#[serde(rename_all = "camelCase")]
317pub enum ActiveCallType {
318 Webrtc,
319 B2bua,
320 WebSocket,
321 #[default]
322 Sip,
323}
324
325#[derive(Default)]
326pub struct ActiveCallState {
327 pub session_id: String,
328 pub start_time: DateTime<Utc>,
329 pub ring_time: Option<DateTime<Utc>>,
330 pub answer_time: Option<DateTime<Utc>>,
331 pub hangup_reason: Option<CallRecordHangupReason>,
332 pub last_status_code: u16,
333 pub option: Option<CallOption>,
334 pub answer: Option<String>,
335 pub ssrc: u32,
336 pub refer_callstate: Option<ActiveCallStateRef>,
337 pub extras: Option<HashMap<String, serde_json::Value>>,
338 pub is_refer: bool,
339 pub sip_hangup_headers_template: Option<HashMap<String, String>>,
340
341 pub tts_handle: Option<SynthesisHandle>,
343 pub auto_hangup: Option<(u32, CallRecordHangupReason)>,
344 pub wait_input_timeout: Option<u32>,
345 pub moh: Option<String>,
346 pub current_play_id: Option<String>,
347 pub audio_receiver: Option<WebsocketBytesReceiver>,
348 pub ready_to_answer: Option<(String, Option<Box<dyn Track>>, ServerInviteDialog)>,
349 pub pending_asr_resume: Option<(u32, TranscriptionOption)>,
350 pub bridge_paused: Arc<AtomicBool>,
351 pub refer_call_token: Option<CancellationToken>,
353}
354
355pub type ActiveCallRef = Arc<ActiveCall>;
356pub type ActiveCallStateRef = Arc<RwLock<ActiveCallState>>;
357
358pub struct ActiveCall {
359 pub call_state: ActiveCallStateRef,
360 pub cancel_token: CancellationToken,
361 pub call_type: ActiveCallType,
362 pub session_id: String,
363 pub media_stream: Arc<MediaStream>,
364 pub track_config: TrackConfig,
365 pub event_sender: EventSender,
366 pub app_state: AppState,
367 pub invitation: Invitation,
368 pub cmd_sender: CommandSender,
369 pub dump_events: bool,
370 pub server_side_track_id: TrackId,
371}
372
373pub struct ActiveCallGuard {
374 pub call: ActiveCallRef,
375 pub active_calls: usize,
376}
377
378impl ActiveCallGuard {
379 pub fn new(call: ActiveCallRef) -> Self {
380 let active_calls = {
381 call.app_state
382 .total_calls
383 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
384 let mut calls = call.app_state.active_calls.lock().unwrap();
385 calls.insert(call.session_id.clone(), call.clone());
386 calls.len()
387 };
388 Self { call, active_calls }
389 }
390}
391
392impl Drop for ActiveCallGuard {
393 fn drop(&mut self) {
394 self.call
395 .app_state
396 .active_calls
397 .lock()
398 .unwrap()
399 .remove(&self.call.session_id);
400 }
401}
402
403pub struct ActiveCallReceiver {
404 pub cmd_receiver: CommandReceiver,
405 pub dump_cmd_receiver: CommandReceiver,
406 pub dump_event_receiver: EventReceiver,
407}
408
409impl ActiveCall {
410 pub fn new(
411 call_type: ActiveCallType,
412 cancel_token: CancellationToken,
413 session_id: String,
414 invitation: Invitation,
415 app_state: AppState,
416 track_config: TrackConfig,
417 audio_receiver: Option<WebsocketBytesReceiver>,
418 dump_events: bool,
419 server_side_track_id: Option<TrackId>,
420 extras: Option<HashMap<String, serde_json::Value>>,
421 sip_hangup_headers_template: Option<HashMap<String, String>>,
422 ) -> Self {
423 let event_sender = crate::event::create_event_sender();
424 let cmd_sender = tokio::sync::broadcast::Sender::<Command>::new(32);
425 let server_side_track_id =
426 server_side_track_id.unwrap_or(SERVER_SIDE_TRACK_ID.to_string());
427 let media_stream_builder = MediaStreamBuilder::new(event_sender.clone())
428 .with_id(session_id.clone())
429 .with_cancel_token(cancel_token.child_token());
430 let media_stream = Arc::new(media_stream_builder.build());
431 let start_time = Utc::now();
432 let call_type_str = match &call_type {
434 ActiveCallType::Sip => "sip",
435 ActiveCallType::WebSocket => "websocket",
436 ActiveCallType::Webrtc => "webrtc",
437 ActiveCallType::B2bua => "b2bua",
438 };
439 let extras = {
440 let mut e = extras.unwrap_or_default();
441 e.entry(crate::playbook::BUILTIN_SESSION_ID.to_string())
442 .or_insert_with(|| serde_json::Value::String(session_id.clone()));
443 e.entry(crate::playbook::BUILTIN_CALL_TYPE.to_string())
444 .or_insert_with(|| serde_json::Value::String(call_type_str.to_string()));
445 e.entry(crate::playbook::BUILTIN_START_TIME.to_string())
446 .or_insert_with(|| serde_json::Value::String(start_time.to_rfc3339()));
447 Some(e)
448 };
449 let call_state = Arc::new(RwLock::new(ActiveCallState {
450 session_id: session_id.clone(),
451 start_time,
452 ssrc: rand::random::<u32>(),
453 extras,
454 audio_receiver,
455 sip_hangup_headers_template,
456 ..Default::default()
457 }));
458 Self {
459 cancel_token,
460 call_type,
461 session_id,
462 call_state,
463 media_stream,
464 track_config,
465 event_sender,
466 app_state,
467 invitation,
468 cmd_sender,
469 dump_events,
470 server_side_track_id,
471 }
472 }
473
474 pub async fn enqueue_command(&self, command: Command) -> Result<()> {
475 self.cmd_sender
476 .send(command)
477 .map_err(|e| anyhow::anyhow!("Failed to send command: {}", e))?;
478 Ok(())
479 }
480
481 pub fn new_receiver(&self) -> ActiveCallReceiver {
485 ActiveCallReceiver {
486 cmd_receiver: self.cmd_sender.subscribe(),
487 dump_cmd_receiver: self.cmd_sender.subscribe(),
488 dump_event_receiver: self.event_sender.subscribe(),
489 }
490 }
491
492 pub async fn serve(&self, receiver: ActiveCallReceiver) -> Result<()> {
493 let ActiveCallReceiver {
494 mut cmd_receiver,
495 dump_cmd_receiver,
496 dump_event_receiver,
497 } = receiver;
498
499 let process_command_loop = async move {
500 while let Ok(command) = cmd_receiver.recv().await {
501 match Box::pin(self.dispatch(command)).await {
502 Ok(_) => (),
503 Err(e) => {
504 warn!(session_id = self.session_id, "{}", e);
505 self.event_sender
506 .send(SessionEvent::Error {
507 track_id: self.session_id.clone(),
508 timestamp: crate::media::get_timestamp(),
509 sender: "command".to_string(),
510 error: e.to_string(),
511 code: None,
512 })
513 .ok();
514 }
515 }
516 }
517 };
518 self.app_state
519 .total_calls
520 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
521
522 tokio::join!(
523 self.dump_loop(self.dump_events, dump_cmd_receiver, dump_event_receiver),
524 async {
525 select! {
526 _ = process_command_loop => {
527 info!(session_id = self.session_id, "command loop done");
528 }
529 _ = self.process() => {
530 info!(session_id = self.session_id, "call serve done");
531 }
532 _ = self.cancel_token.cancelled() => {
533 info!(session_id = self.session_id, "call cancelled - cleaning up resources");
534 }
535 }
536 self.cancel_token.cancel();
537 }
538 );
539 Ok(())
540 }
541
542 async fn process(&self) -> Result<()> {
543 let mut event_receiver = self.event_sender.subscribe();
544
545 let input_timeout_expire = Arc::new(Mutex::new((0u64, 0u32)));
546 let input_timeout_expire_ref = input_timeout_expire.clone();
547 let event_sender = self.event_sender.clone();
548 let wait_input_timeout_loop = async {
549 loop {
550 let (start_time, expire) = { *input_timeout_expire.lock().await };
551 if expire > 0 && crate::media::get_timestamp() >= start_time + expire as u64 {
552 info!(session_id = self.session_id, "wait input timeout reached");
553 *input_timeout_expire.lock().await = (0, 0);
554 let is_refer = self.call_state.read().await.is_refer;
555 event_sender
556 .send(SessionEvent::Silence {
557 track_id: self.server_side_track_id.clone(),
558 timestamp: crate::media::get_timestamp(),
559 start_time,
560 duration: expire as u64,
561 samples: None,
562 refer: Some(is_refer),
563 })
564 .ok();
565 }
566 sleep(Duration::from_millis(100)).await;
567 }
568 };
569 let server_side_track_id = self.server_side_track_id.clone();
570 let event_hook_loop = async move {
571 while let Ok(event) = event_receiver.recv().await {
572 match event {
573 SessionEvent::Speaking { .. }
574 | SessionEvent::Dtmf { .. }
575 | SessionEvent::AsrDelta { .. }
576 | SessionEvent::AsrFinal { .. }
577 | SessionEvent::TrackStart { .. } => {
578 *input_timeout_expire_ref.lock().await = (0, 0);
579 }
580 SessionEvent::TrackEnd {
581 track_id,
582 play_id,
583 ssrc,
584 ..
585 } => {
586 if track_id != server_side_track_id {
587 continue;
588 }
589
590 let (moh_path, auto_hangup, wait_timeout_val) = {
591 let mut state = self.call_state.write().await;
592 if play_id != state.current_play_id {
593 debug!(
594 session_id = self.session_id,
595 ?play_id,
596 current = ?state.current_play_id,
597 "ignoring interrupted track end"
598 );
599 continue;
600 }
601 state.current_play_id = None;
602 (
603 state.moh.clone(),
604 state.auto_hangup.clone(),
605 state.wait_input_timeout.take(),
606 )
607 };
608
609 if let Some(path) = moh_path {
610 info!(session_id = self.session_id, "looping moh: {}", path);
611 let ssrc = rand::random::<u32>();
612 let file_track = FileTrack::new(self.server_side_track_id.clone())
613 .with_play_id(Some(path.clone()))
614 .with_ssrc(ssrc)
615 .with_path(path.clone())
616 .with_cancel_token(self.cancel_token.child_token());
617 self.update_track_wrapper(Box::new(file_track), Some(path))
618 .await;
619 continue;
620 }
621
622 if let Some((hangup_ssrc, hangup_reason)) = auto_hangup {
623 if hangup_ssrc == ssrc {
624 info!(
625 session_id = self.session_id,
626 ssrc, "auto hangup when track end track_id:{}", track_id
627 );
628 self.do_hangup(Some(hangup_reason), None, None, None)
629 .await
630 .ok();
631 }
632 }
633
634 if let Some(timeout) = wait_timeout_val {
635 let expire = if timeout > 0 {
636 (crate::media::get_timestamp(), timeout)
637 } else {
638 (0, 0)
639 };
640 *input_timeout_expire_ref.lock().await = expire;
641 }
642 }
643 SessionEvent::Interrupt { receiver } => {
644 let track_id =
645 receiver.unwrap_or_else(|| self.server_side_track_id.clone());
646 if track_id == self.server_side_track_id {
647 debug!(
648 session_id = self.session_id,
649 "received interrupt event, stopping playback"
650 );
651 self.do_interrupt(true).await.ok();
652 }
653 }
654 SessionEvent::Inactivity { track_id, .. } => {
655 info!(
656 session_id = self.session_id,
657 track_id, "inactivity timeout reached, hanging up"
658 );
659 self.do_hangup(
660 Some(CallRecordHangupReason::InactivityTimeout),
661 None,
662 None,
663 None,
664 )
665 .await
666 .ok();
667 }
668 SessionEvent::Hangup { refer, .. } => {
669 if refer == Some(true) {
671 let mut cs = self.call_state.write().await;
672 if let Some((refer_ssrc, asr_option)) = cs.pending_asr_resume.take() {
673 let is_refer_hangup = cs
675 .refer_callstate
676 .as_ref()
677 .map(|rcs| {
678 rcs.try_read()
679 .map(|g| g.ssrc == refer_ssrc)
680 .unwrap_or(false)
681 })
682 .unwrap_or(false);
683
684 if is_refer_hangup {
685 drop(cs); info!(
687 session_id = self.session_id,
688 "Refer call ended, resuming parent ASR"
689 );
690
691 match self
693 .app_state
694 .stream_engine
695 .create_asr_processor(
696 self.server_side_track_id.clone(),
697 self.cancel_token.child_token(),
698 asr_option,
699 self.event_sender.clone(),
700 )
701 .await
702 {
703 Ok(asr_processor) => {
704 if let Err(e) = self
705 .media_stream
706 .append_processor(
707 &self.server_side_track_id,
708 asr_processor,
709 )
710 .await
711 {
712 warn!(
713 session_id = self.session_id,
714 "Failed to resume ASR after refer: {}", e
715 );
716 }
717 }
718 Err(e) => {
719 warn!(
720 session_id = self.session_id,
721 "Failed to create ASR processor for resume: {}", e
722 );
723 }
724 }
725 }
726 }
727 }
728 }
729 SessionEvent::Error { track_id, .. } => {
730 if track_id != server_side_track_id {
731 continue;
732 }
733
734 let moh_info = {
735 let mut state = self.call_state.write().await;
736 if let Some(path) = state.moh.clone() {
737 let fallback = "./config/sounds/refer_moh.wav".to_string();
738 let next_path = if path != fallback
739 && std::path::Path::new(&fallback).exists()
740 {
741 info!(
742 session_id = self.session_id,
743 "moh error, switching to fallback: {}", fallback
744 );
745 state.moh = Some(fallback.clone());
746 fallback
747 } else {
748 info!(
749 session_id = self.session_id,
750 "looping moh on error: {}", path
751 );
752 path
753 };
754 Some(next_path)
755 } else {
756 None
757 }
758 };
759
760 if let Some(next_path) = moh_info {
761 let ssrc = rand::random::<u32>();
762 let file_track = FileTrack::new(self.server_side_track_id.clone())
763 .with_play_id(Some(next_path.clone()))
764 .with_ssrc(ssrc)
765 .with_path(next_path.clone())
766 .with_cancel_token(self.cancel_token.child_token());
767 self.update_track_wrapper(Box::new(file_track), Some(next_path))
768 .await;
769 continue;
770 }
771 }
772 SessionEvent::Hold { on_hold, .. } => {
773 self.call_state
774 .read()
775 .await
776 .bridge_paused
777 .store(on_hold, Ordering::Relaxed);
778 }
779 _ => {}
780 }
781 }
782 };
783
784 select! {
785 _ = wait_input_timeout_loop=>{
786 info!(session_id = self.session_id, "wait input timeout loop done");
787 }
788 _ = self.media_stream.serve() => {
789 info!(session_id = self.session_id, "media stream loop done");
790 }
791 _ = event_hook_loop => {
792 info!(session_id = self.session_id, "event loop done");
793 }
794 }
795 Ok(())
796 }
797
798 async fn dispatch(&self, command: Command) -> Result<()> {
799 match command {
800 Command::Invite { option } => self.do_invite(option).await,
801 Command::Accept { option } => self.do_accept(option).await,
802 Command::Reject { reason, code } => {
803 self.do_reject(code.map(|c| (c as u16).into()), Some(reason))
804 .await
805 }
806 Command::Ringing {
807 ringtone,
808 recorder,
809 early_media,
810 } => self.do_ringing(ringtone, recorder, early_media).await,
811 Command::Tts {
812 text,
813 speaker,
814 play_id,
815 auto_hangup,
816 streaming,
817 end_of_stream,
818 option,
819 wait_input_timeout,
820 base64,
821 cache_key,
822 } => {
823 self.do_tts(
824 text,
825 speaker,
826 play_id,
827 auto_hangup,
828 streaming.unwrap_or_default(),
829 end_of_stream.unwrap_or_default(),
830 option,
831 wait_input_timeout,
832 base64.unwrap_or_default(),
833 cache_key,
834 )
835 .await
836 }
837 Command::Play {
838 url,
839 play_id,
840 auto_hangup,
841 wait_input_timeout,
842 offset_ms,
843 } => {
844 self.do_play(url, play_id, auto_hangup, wait_input_timeout, offset_ms)
845 .await
846 }
847 Command::Hangup {
848 reason,
849 initiator,
850 headers,
851 refer,
852 } => {
853 let reason = reason.map(|r| {
854 r.parse::<CallRecordHangupReason>()
855 .unwrap_or(CallRecordHangupReason::BySystem)
856 });
857 self.do_hangup(reason, initiator, headers, refer).await
858 }
859 Command::Refer {
860 caller,
861 callee,
862 options,
863 } => self.do_refer(caller, callee, options).await,
864 Command::Message {
865 body,
866 content_type,
867 headers,
868 refer,
869 } => self.do_message(body, content_type, headers, refer).await,
870 Command::Bridge { target_session_id } => self.do_bridge(target_session_id).await,
871 Command::Unbridge { target_session_id } => self.do_unbridge(target_session_id).await,
872 Command::Mute { track_id } => self.do_mute(track_id).await,
873 Command::Unmute { track_id } => self.do_unmute(track_id).await,
874 Command::Pause {} => self.do_pause().await,
875 Command::Resume {} => self.do_resume().await,
876 Command::Interrupt {
877 graceful: passage,
878 fade_out_ms: _,
879 } => self.do_interrupt(passage.unwrap_or_default()).await,
880 Command::History { speaker, text } => self.do_history(speaker, text).await,
881 Command::Custom { sender, data } => self.do_custom(sender, data),
882 }
883 }
884
885 fn build_record_option(&self, option: &CallOption) -> Option<RecorderOption> {
886 if let Some(recorder_option) = &option.recorder {
887 let mut recorder_file = recorder_option.recorder_file.clone();
888 if recorder_file.contains("{id}") {
889 recorder_file = recorder_file.replace("{id}", &self.session_id);
890 }
891
892 let recorder_file = if recorder_file.is_empty() {
893 self.app_state.get_recorder_file(&self.session_id)
894 } else {
895 let p = Path::new(&recorder_file);
896 p.is_absolute()
897 .then(|| recorder_file.clone())
898 .unwrap_or_else(|| self.app_state.get_recorder_file(&recorder_file))
899 };
900 info!(
901 session_id = self.session_id,
902 recorder_file, "created recording file"
903 );
904
905 let track_samplerate = self.track_config.samplerate;
906 let recorder_samplerate = if track_samplerate > 0 {
907 track_samplerate
908 } else {
909 recorder_option.samplerate
910 };
911 let recorder_ptime = if recorder_option.ptime == 0 {
912 200
913 } else {
914 recorder_option.ptime
915 };
916 let requested_format = recorder_option
917 .format
918 .unwrap_or(self.app_state.config.recorder_format());
919 let format = requested_format.effective();
920 if requested_format != format {
921 warn!(
922 session_id = self.session_id,
923 requested = requested_format.extension(),
924 "Recorder format fallback to wav due to unsupported feature"
925 );
926 }
927 let mut recorder_config = RecorderOption {
928 recorder_file,
929 samplerate: recorder_samplerate,
930 ptime: recorder_ptime,
931 format: Some(format),
932 };
933 recorder_config.ensure_path_extension(format);
934 Some(recorder_config)
935 } else {
936 None
937 }
938 }
939
940 async fn invite_or_accept(&self, mut option: CallOption, sender: String) -> Result<CallOption> {
941 {
943 let state = self.call_state.read().await;
944 option = state.merge_option(option);
945 }
946
947 option.check_default();
948 if let Some(opt) = self.build_record_option(&option) {
949 self.media_stream.update_recorder_option(opt).await;
950 }
951
952 if let Some(opt) = &option.media_pass {
953 let track_id = self.server_side_track_id.clone();
954 let cancel_token = self.cancel_token.child_token();
955 let ssrc = rand::random::<u32>();
956 let media_pass_track = MediaPassTrack::new(
957 self.session_id.clone(),
958 ssrc,
959 track_id,
960 cancel_token,
961 opt.clone(),
962 );
963 self.update_track_wrapper(Box::new(media_pass_track), None)
964 .await;
965 }
966
967 info!(
968 session_id = self.session_id,
969 call_type = ?self.call_type,
970 sender,
971 ?option,
972 "caller with option"
973 );
974
975 match self.setup_caller_track(&option).await {
976 Ok(_) => return Ok(option),
977 Err(e) => {
978 self.app_state
979 .total_failed_calls
980 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
981 let error_event = crate::event::SessionEvent::Error {
982 track_id: self.session_id.clone(),
983 timestamp: crate::media::get_timestamp(),
984 sender,
985 error: e.to_string(),
986 code: None,
987 };
988 self.event_sender.send(error_event).ok();
989 self.do_hangup(Some(CallRecordHangupReason::BySystem), None, None, None)
990 .await
991 .ok();
992 return Err(e);
993 }
994 }
995 }
996
997 async fn do_invite(&self, option: CallOption) -> Result<()> {
998 self.invite_or_accept(option, "invite".to_string())
999 .await
1000 .map(|_| ())
1001 }
1002
1003 async fn do_accept(&self, mut option: CallOption) -> Result<()> {
1004 let has_pending = self
1005 .invitation
1006 .find_dialog_id_by_session_id(&self.session_id)
1007 .is_some();
1008 let ready_to_answer_val = {
1009 let state = self.call_state.read().await;
1010 state.ready_to_answer.is_none()
1011 };
1012
1013 if ready_to_answer_val {
1014 if !has_pending {
1015 warn!(session_id = self.session_id, "no pending call to accept");
1017 let rejet_event = crate::event::SessionEvent::Reject {
1018 track_id: self.session_id.clone(),
1019 timestamp: crate::media::get_timestamp(),
1020 reason: "no pending call".to_string(),
1021 refer: None,
1022 code: Some(486),
1023 };
1024 self.event_sender.send(rejet_event).ok();
1025 self.do_hangup(Some(CallRecordHangupReason::BySystem), None, None, None)
1026 .await
1027 .ok();
1028 return Err(anyhow::anyhow!("no pending call to accept"));
1029 }
1030 option = self.invite_or_accept(option, "accept".to_string()).await?;
1031 } else {
1032 option.check_default();
1033 if let Some(opt) = self.build_record_option(&option) {
1034 self.media_stream.update_recorder_option(opt).await;
1035 }
1036 self.call_state.write().await.option = Some(option.clone());
1037 }
1038 info!(session_id = self.session_id, ?option, "accepting call");
1039 let ready = self.call_state.write().await.ready_to_answer.take();
1040 if let Some((answer, track, dialog)) = ready {
1041 info!(
1042 session_id = self.session_id,
1043 track_id = track.as_ref().map(|t| t.id()),
1044 "ready to answer with track"
1045 );
1046
1047 let headers = vec![rsipstack::rsip::Header::ContentType(
1048 "application/sdp".to_string().into(),
1049 )];
1050
1051 match dialog.accept(Some(headers), Some(answer.as_bytes().to_vec())) {
1052 Ok(_) => {
1053 {
1054 let mut state = self.call_state.write().await;
1055 state.answer = Some(answer);
1056 state.answer_time = Some(Utc::now());
1057 }
1058 self.finish_caller_stack(&option, track).await?;
1059 }
1060 Err(e) => {
1061 warn!(session_id = self.session_id, "failed to accept call: {}", e);
1062 return Err(anyhow::anyhow!("failed to accept call"));
1063 }
1064 }
1065 }
1066 return Ok(());
1067 }
1068
1069 async fn do_reject(
1070 &self,
1071 code: Option<rsipstack::rsip::StatusCode>,
1072 reason: Option<String>,
1073 ) -> Result<()> {
1074 match self
1075 .invitation
1076 .find_dialog_id_by_session_id(&self.session_id)
1077 {
1078 Some(id) => {
1079 info!(
1080 session_id = self.session_id,
1081 ?reason,
1082 ?code,
1083 "rejecting call"
1084 );
1085 let result = self.invitation.hangup(id, code, reason).await;
1086 if result.is_ok() {
1087 self.cancel_token.cancel();
1088 }
1089 result
1090 }
1091 None => {
1092 let ready = self.call_state.write().await.ready_to_answer.take();
1093 if let Some((_, _, dialog)) = ready {
1094 info!(
1095 session_id = self.session_id,
1096 ?reason,
1097 ?code,
1098 "rejecting call from ready_to_answer"
1099 );
1100 let dialog_id = dialog.id();
1101 dialog.reject(code, reason).ok();
1102 self.invitation.dialog_layer.remove_dialog(&dialog_id);
1103 self.cancel_token.cancel();
1104 }
1105 Ok(())
1106 }
1107 }
1108 }
1109
1110 async fn do_ringing(
1111 &self,
1112 ringtone: Option<String>,
1113 recorder: Option<RecorderOption>,
1114 early_media: Option<bool>,
1115 ) -> Result<()> {
1116 let ready_to_answer_val = self.call_state.read().await.ready_to_answer.is_none();
1117 if ready_to_answer_val {
1118 let option = CallOption {
1119 recorder,
1120 ..Default::default()
1121 };
1122 let _ = self.invite_or_accept(option, "ringing".to_string()).await?;
1123 }
1124
1125 let state = self.call_state.read().await;
1126 if let Some((answer, _, dialog)) = state.ready_to_answer.as_ref() {
1127 let (headers, body) = if early_media.unwrap_or_default() || ringtone.is_some() {
1128 let headers = vec![rsipstack::rsip::Header::ContentType(
1129 "application/sdp".to_string().into(),
1130 )];
1131 (Some(headers), Some(answer.as_bytes().to_vec()))
1132 } else {
1133 (None, None)
1134 };
1135
1136 dialog.ringing(headers, body).ok();
1137 info!(
1138 session_id = self.session_id,
1139 ringtone, early_media, "playing ringtone"
1140 );
1141 if let Some(ringtone_url) = ringtone {
1142 drop(state);
1143 self.do_play(ringtone_url, None, None, None, None)
1144 .await
1145 .ok();
1146 } else {
1147 info!(session_id = self.session_id, "no ringtone to play");
1148 }
1149 }
1150 Ok(())
1151 }
1152
1153 async fn do_tts(
1154 &self,
1155 text: String,
1156 speaker: Option<String>,
1157 play_id: Option<String>,
1158 auto_hangup: Option<bool>,
1159 streaming: bool,
1160 end_of_stream: bool,
1161 option: Option<SynthesisOption>,
1162 wait_input_timeout: Option<u32>,
1163 base64: bool,
1164 cache_key: Option<String>,
1165 ) -> Result<()> {
1166 let tts_option = {
1167 let call_state = self.call_state.read().await;
1168 match call_state.option.clone().unwrap_or_default().tts {
1169 Some(opt) => opt.merge_with(option),
1170 None => {
1171 if let Some(opt) = option {
1172 opt
1173 } else {
1174 return Err(anyhow::anyhow!("no tts option available"));
1175 }
1176 }
1177 }
1178 };
1179 let speaker = match speaker {
1180 Some(s) => Some(s),
1181 None => tts_option.speaker.clone(),
1182 };
1183
1184 let mut play_command = SynthesisCommand {
1185 text,
1186 speaker,
1187 play_id: play_id.clone(),
1188 streaming,
1189 end_of_stream: if !streaming { true } else { end_of_stream },
1190 option: tts_option,
1191 base64,
1192 cache_key,
1193 };
1194 info!(
1195 session_id = self.session_id,
1196 provider = ?play_command.option.provider,
1197 text = %play_command.text.chars().take(10).collect::<String>(),
1198 speaker = play_command.speaker.as_deref(),
1199 auto_hangup = auto_hangup.unwrap_or_default(),
1200 play_id = play_command.play_id.as_deref(),
1201 streaming = play_command.streaming,
1202 end_of_stream = play_command.end_of_stream,
1203 wait_input_timeout = wait_input_timeout.unwrap_or_default(),
1204 is_base64 = play_command.base64,
1205 cache_key = play_command.cache_key.as_deref(),
1206 "new synthesis"
1207 );
1208
1209 let ssrc = rand::random::<u32>();
1210 let (should_interrupt, picked_ssrc) = {
1211 let mut state = self.call_state.write().await;
1212
1213 let (target_ssrc, changed) = if let Some(handle) = &state.tts_handle {
1214 if play_id.is_some() && state.current_play_id != play_id {
1215 (ssrc, true)
1216 } else {
1217 (handle.ssrc, false)
1218 }
1219 } else {
1220 (ssrc, false)
1221 };
1222
1223 state.wait_input_timeout = wait_input_timeout;
1226
1227 state.current_play_id = play_id.clone();
1228 (changed, target_ssrc)
1229 };
1230
1231 if should_interrupt {
1232 let _ = self.do_interrupt(false).await;
1233 }
1234
1235 {
1239 let mut state = self.call_state.write().await;
1240 state.auto_hangup = match auto_hangup {
1241 Some(true) => Some((picked_ssrc, CallRecordHangupReason::BySystem)),
1242 _ => {
1243 if state.tts_handle.is_some() && !should_interrupt {
1247 state.auto_hangup.clone()
1248 } else {
1249 None
1250 }
1251 }
1252 };
1253 }
1254
1255 let existing_handle = self.call_state.read().await.tts_handle.clone();
1256 if let Some(tts_handle) = existing_handle {
1257 match tts_handle.try_send(play_command) {
1258 Ok(_) => return Ok(()),
1259 Err(e) => {
1260 play_command = e.0;
1261 }
1262 }
1263 }
1264
1265 let (new_handle, tts_track) = StreamEngine::create_tts_track(
1266 self.app_state.stream_engine.clone(),
1267 self.cancel_token.child_token(),
1268 self.session_id.clone(),
1269 self.server_side_track_id.clone(),
1270 picked_ssrc,
1271 play_id.clone(),
1272 streaming,
1273 &play_command.option,
1274 )
1275 .await?;
1276
1277 new_handle.try_send(play_command)?;
1278 self.call_state.write().await.tts_handle = Some(new_handle);
1279 self.update_track_wrapper(tts_track, play_id).await;
1280 Ok(())
1281 }
1282
1283 async fn do_play(
1284 &self,
1285 url: String,
1286 play_id: Option<String>,
1287 auto_hangup: Option<bool>,
1288 wait_input_timeout: Option<u32>,
1289 offset_ms: Option<u32>,
1290 ) -> Result<()> {
1291 let ssrc = rand::random::<u32>();
1292 info!(
1293 session_id = self.session_id,
1294 ssrc, url, play_id, auto_hangup, "play file track"
1295 );
1296
1297 let play_id = play_id.or(Some(url.clone()));
1298
1299 let mut file_track = FileTrack::new(self.server_side_track_id.clone())
1300 .with_play_id(play_id.clone())
1301 .with_ssrc(ssrc)
1302 .with_path(url)
1303 .with_cancel_token(self.cancel_token.child_token());
1304
1305 if let Some(offset) = offset_ms {
1306 file_track = file_track.with_offset_ms(offset);
1307 }
1308
1309 {
1310 let mut state = self.call_state.write().await;
1311 state.tts_handle = None;
1312 state.auto_hangup = match auto_hangup {
1313 Some(true) => Some((ssrc, CallRecordHangupReason::BySystem)),
1314 _ => None,
1315 };
1316 state.wait_input_timeout = wait_input_timeout;
1317 }
1318
1319 self.update_track_wrapper(Box::new(file_track), play_id)
1320 .await;
1321 Ok(())
1322 }
1323
1324 async fn do_history(&self, speaker: String, text: String) -> Result<()> {
1325 self.event_sender
1326 .send(SessionEvent::AddHistory {
1327 sender: Some(self.session_id.clone()),
1328 timestamp: crate::media::get_timestamp(),
1329 speaker,
1330 text,
1331 })
1332 .map(|_| ())
1333 .map_err(Into::into)
1334 }
1335
1336 fn do_custom(&self, sender: Option<String>, data: serde_json::Value) -> Result<()> {
1337 self.event_sender
1338 .send(SessionEvent::Custom {
1339 timestamp: crate::media::get_timestamp(),
1340 sender,
1341 data,
1342 })
1343 .map(|_| ())
1344 .map_err(Into::into)
1345 }
1346
1347 async fn do_interrupt(&self, graceful: bool) -> Result<()> {
1348 {
1349 let mut state = self.call_state.write().await;
1350 state.tts_handle = None;
1351 state.moh = None;
1352 state.auto_hangup = None;
1353 }
1354 self.media_stream
1355 .remove_track(&self.server_side_track_id, graceful)
1356 .await;
1357 Ok(())
1358 }
1359 async fn do_pause(&self) -> Result<()> {
1360 self.media_stream
1361 .pause_playback(self.server_side_track_id.clone())
1362 .await?;
1363 Ok(())
1364 }
1365 async fn do_resume(&self) -> Result<()> {
1366 self.media_stream
1367 .resume_playback(self.server_side_track_id.clone())
1368 .await?;
1369 Ok(())
1370 }
1371 async fn do_hangup(
1372 &self,
1373 reason: Option<CallRecordHangupReason>,
1374 initiator: Option<String>,
1375 headers: Option<HashMap<String, String>>,
1376 refer: Option<bool>,
1377 ) -> Result<()> {
1378 info!(
1379 session_id = self.session_id,
1380 ?reason,
1381 ?initiator,
1382 ?headers,
1383 ?refer,
1384 "do_hangup"
1385 );
1386
1387 let hangup_reason = match initiator.as_deref() {
1388 Some("caller") => CallRecordHangupReason::ByCaller,
1389 Some("callee") => CallRecordHangupReason::ByCallee,
1390 Some("system") => CallRecordHangupReason::Autohangup,
1391 _ => reason.unwrap_or(CallRecordHangupReason::BySystem),
1392 };
1393
1394 match refer {
1395 Some(true) => {
1396 let (refer_state, refer_token) = {
1398 let mut state = self.call_state.write().await;
1399 (state.refer_callstate.clone(), state.refer_call_token.take())
1400 };
1401 let mut has_refer_state = false;
1402 if let Some(refer_state) = refer_state {
1403 has_refer_state = true;
1404 let mut refer_state = refer_state.write().await;
1405 if let Some(headers) = headers {
1406 let h_val = serde_json::to_value(&headers).unwrap_or_default();
1407 let mut extras = refer_state.extras.take().unwrap_or_default();
1408 extras.insert("_hangup_headers".to_string(), h_val);
1409 refer_state.extras = Some(extras);
1410 }
1411 refer_state.set_hangup_reason(hangup_reason);
1413 }
1414 if let Some(token) = refer_token {
1415 token.cancel();
1416 }
1417 if has_refer_state {
1418 self.media_stream
1419 .remove_track(&self.server_side_track_id, false)
1420 .await;
1421 }
1422 }
1423 _ => {
1424 let refer_token = {
1425 let mut state = self.call_state.write().await;
1426 if let Some(headers) = headers {
1427 let h_val = serde_json::to_value(&headers).unwrap_or_default();
1428 let mut extras = state.extras.take().unwrap_or_default();
1429 extras.insert("_hangup_headers".to_string(), h_val);
1430 state.extras = Some(extras);
1431 }
1432 state.set_hangup_reason(hangup_reason.clone());
1433 state.refer_call_token.take()
1434 };
1435 self.media_stream
1436 .stop(Some(hangup_reason.to_string()), initiator);
1437 if let Some(token) = refer_token {
1438 token.cancel();
1439 }
1440 }
1441 }
1442 tokio::task::yield_now().await;
1443 Ok(())
1444 }
1445
1446 async fn do_refer(
1447 &self,
1448 caller: String,
1449 callee: String,
1450 refer_option: Option<ReferOption>,
1451 ) -> Result<()> {
1452 self.do_interrupt(false).await.ok();
1453
1454 let pause_parent_asr = refer_option
1456 .as_ref()
1457 .and_then(|o| o.pause_parent_asr)
1458 .unwrap_or(false);
1459
1460 let original_asr_option = if pause_parent_asr {
1462 let cs = self.call_state.read().await;
1463 cs.option.as_ref().and_then(|o| o.asr.clone())
1464 } else {
1465 None
1466 };
1467
1468 if pause_parent_asr {
1470 info!(
1471 session_id = self.session_id,
1472 "Pausing parent call ASR during refer"
1473 );
1474 self.media_stream
1475 .remove_processor::<crate::media::asr_processor::AsrProcessor>(
1476 &self.server_side_track_id,
1477 )
1478 .await
1479 .ok();
1480 }
1481
1482 let mut moh = refer_option.as_ref().and_then(|o| o.moh.clone());
1483 if let Some(ref path) = moh {
1484 if !path.starts_with("http") && !std::path::Path::new(path).exists() {
1485 let fallback = "./config/sounds/refer_moh.wav";
1486 if std::path::Path::new(fallback).exists() {
1487 info!(
1488 session_id = self.session_id,
1489 "moh {} not found, using fallback {}", path, fallback
1490 );
1491 moh = Some(fallback.to_string());
1492 }
1493 }
1494 }
1495 let ref_call_id = refer_option
1496 .as_ref()
1497 .and_then(|o| o.call_id.clone())
1498 .unwrap_or_else(|| format!("ref-{}-{}", rand::random::<u32>(), self.session_id));
1499
1500 let session_id = self.session_id.clone();
1501 let track_id = self.server_side_track_id.clone();
1502
1503 let (recorder, parent_caller) = {
1504 let cs = self.call_state.read().await;
1505 let option = cs.option.as_ref();
1506 (
1507 option.map(|o| o.recorder.clone()).unwrap_or_default(),
1508 option.and_then(|o| o.caller.clone()),
1509 )
1510 };
1511 let caller = if caller.trim().is_empty() {
1512 parent_caller.unwrap_or_default()
1513 } else {
1514 caller
1515 };
1516
1517 let mut call_option = CallOption {
1518 caller: Some(caller),
1519 callee: Some(callee.clone()),
1520 sip: refer_option.as_ref().and_then(|o| o.sip.clone()),
1521 vad: refer_option
1522 .as_ref()
1523 .and_then(|o| o.vad.clone())
1524 .map(|mut opts| {
1525 opts.refer = Some(true);
1526 opts
1527 }),
1528 asr: refer_option
1529 .as_ref()
1530 .and_then(|o| o.asr.clone())
1531 .map(|mut opts| {
1532 opts.refer = Some(true);
1533 opts
1534 }),
1535 denoise: refer_option.as_ref().and_then(|o| o.denoise.clone()),
1536 agc: refer_option.as_ref().and_then(|o| o.agc.clone()),
1537 recorder,
1538 ..Default::default()
1539 };
1540 call_option.check_default();
1541
1542 let mut invite_option = call_option.build_invite_option()?;
1543 invite_option.call_id = Some(ref_call_id.clone());
1544
1545 let headers = invite_option.headers.get_or_insert_with(|| Vec::new());
1546
1547 {
1548 let cs = self.call_state.read().await;
1549 if let Some(opt) = cs.option.as_ref() {
1550 if let Some(callee) = opt.callee.as_ref() {
1551 headers.push(rsipstack::rsip::Header::Other(
1552 "X-Referred-To".to_string(),
1553 callee.clone(),
1554 ));
1555 }
1556 if let Some(caller) = opt.caller.as_ref() {
1557 headers.push(rsipstack::rsip::Header::Other(
1558 "X-Referred-From".to_string(),
1559 caller.clone(),
1560 ));
1561 }
1562 }
1563 }
1564
1565 headers.push(rsipstack::rsip::Header::Other(
1566 "X-Referred-Id".to_string(),
1567 self.session_id.clone(),
1568 ));
1569
1570 let ssrc = rand::random::<u32>();
1571 let refer_call_state = Arc::new(RwLock::new(ActiveCallState {
1572 session_id: ref_call_id.clone(),
1573 start_time: Utc::now(),
1574 ssrc,
1575 option: Some(call_option.clone()),
1576 is_refer: true,
1577 ..Default::default()
1578 }));
1579
1580 {
1581 let mut cs = self.call_state.write().await;
1582 cs.refer_callstate.replace(refer_call_state.clone());
1583 }
1584
1585 let auto_hangup_requested = refer_option
1586 .as_ref()
1587 .and_then(|o| o.auto_hangup)
1588 .unwrap_or(true);
1589
1590 if auto_hangup_requested {
1591 self.call_state.write().await.auto_hangup =
1592 Some((ssrc, CallRecordHangupReason::ByRefer));
1593 } else {
1594 self.call_state.write().await.auto_hangup = None;
1595 }
1596
1597 if !auto_hangup_requested && pause_parent_asr && original_asr_option.is_some() {
1599 let asr_option = original_asr_option.unwrap();
1600 self.call_state.write().await.pending_asr_resume = Some((ssrc, asr_option));
1601 }
1602
1603 let timeout_secs = refer_option.as_ref().and_then(|o| o.timeout).unwrap_or(30);
1604
1605 info!(
1606 session_id = self.session_id,
1607 ssrc,
1608 auto_hangup = auto_hangup_requested,
1609 callee,
1610 timeout_secs,
1611 "do_refer"
1612 );
1613
1614 let refer_cancel_token = self.cancel_token.child_token();
1615 self.call_state.write().await.refer_call_token = Some(refer_cancel_token.clone());
1616
1617 let r = tokio::time::timeout(
1618 Duration::from_secs(timeout_secs as u64),
1619 self.create_outgoing_sip_track(
1620 refer_cancel_token,
1621 refer_call_state.clone(),
1622 &track_id,
1623 invite_option,
1624 &call_option,
1625 moh,
1626 auto_hangup_requested,
1627 ),
1628 )
1629 .await;
1630
1631 {
1632 self.call_state.write().await.moh = None;
1633 }
1634
1635 let result = match r {
1636 Ok(res) => res,
1637 Err(_) => {
1638 warn!(
1639 session_id = session_id,
1640 "refer sip track creation timed out after {} seconds", timeout_secs
1641 );
1642 self.event_sender
1643 .send(SessionEvent::Reject {
1644 track_id,
1645 timestamp: crate::media::get_timestamp(),
1646 reason: "Timeout when refer".into(),
1647 code: Some(408),
1648 refer: Some(true),
1649 })
1650 .ok();
1651 return Err(anyhow::anyhow!("refer sip track creation timed out").into());
1652 }
1653 };
1654
1655 match result {
1656 Ok(answer) => {
1657 self.media_stream.set_track_refer(&track_id, Some(true)).await;
1658 let forward_dtmf = refer_option.as_ref().and_then(|o| o.forward_dtmf).unwrap_or(true);
1659 if !forward_dtmf {
1660 self.media_stream.set_track_dtmf_forward(&track_id, false).await;
1661 }
1662 self.event_sender
1663 .send(SessionEvent::Answer {
1664 timestamp: crate::media::get_timestamp(),
1665 track_id,
1666 sdp: answer,
1667 refer: Some(true),
1668 })
1669 .ok();
1670 }
1671 Err(e) => {
1672 warn!(
1673 session_id = session_id,
1674 "failed to create refer sip track: {}", e
1675 );
1676 match &e {
1677 rsipstack::Error::DialogError(reason, _, code) => {
1678 self.event_sender
1679 .send(SessionEvent::Reject {
1680 track_id,
1681 timestamp: crate::media::get_timestamp(),
1682 reason: reason.clone(),
1683 code: Some(code.code() as u32),
1684 refer: Some(true),
1685 })
1686 .ok();
1687 }
1688 _ => {}
1689 }
1690 return Err(e.into());
1691 }
1692 }
1693 Ok(())
1694 }
1695
1696 async fn do_message(
1697 &self,
1698 body: String,
1699 content_type: Option<String>,
1700 headers: Option<HashMap<String, String>>,
1701 refer: Option<bool>,
1702 ) -> Result<()> {
1703 if !matches!(self.call_type, ActiveCallType::Sip | ActiveCallType::B2bua) {
1704 return Err(anyhow::anyhow!("message command is only supported for SIP calls"));
1705 }
1706
1707 let dialog_key = if refer == Some(true) {
1708 let refer_state = self.call_state.read().await.refer_callstate.clone();
1709 match refer_state {
1710 Some(state) => Some(state.read().await.session_id.clone()),
1711 None => None,
1712 }
1713 } else {
1714 Some(self.call_state.read().await.session_id.clone())
1715 };
1716
1717 let mut dialog = dialog_key
1718 .as_ref()
1719 .filter(|id| !id.is_empty())
1720 .and_then(|id| self.invitation.dialog_layer.get_dialog_with(id));
1721
1722 if dialog.is_none() {
1723 if let Some(target_id) = dialog_key.as_ref().filter(|id| !id.is_empty()) {
1724 dialog = self
1725 .invitation
1726 .dialog_layer
1727 .all_dialog_ids()
1728 .into_iter()
1729 .filter_map(|id| self.invitation.dialog_layer.get_dialog_with(&id))
1730 .find(|dialog| dialog.id().to_string() == *target_id);
1731 }
1732 }
1733
1734 if dialog.is_none() && refer != Some(true) {
1735 dialog = self
1736 .invitation
1737 .dialog_layer
1738 .get_client_dialog_by_call_id(&self.session_id)
1739 .into_iter()
1740 .find(|d| {
1741 matches!(
1742 d.state(),
1743 rsipstack::dialog::dialog::DialogState::Confirmed(_, _)
1744 )
1745 })
1746 .map(rsipstack::dialog::dialog::Dialog::ClientInvite);
1747 }
1748
1749 if dialog.is_none() && refer == Some(true) {
1750 if let Some(call_id) = dialog_key.as_ref().filter(|id| !id.is_empty()) {
1751 dialog = self
1752 .invitation
1753 .dialog_layer
1754 .get_client_dialog_by_call_id(call_id)
1755 .into_iter()
1756 .find(|d| {
1757 matches!(
1758 d.state(),
1759 rsipstack::dialog::dialog::DialogState::Confirmed(_, _)
1760 )
1761 })
1762 .map(rsipstack::dialog::dialog::Dialog::ClientInvite);
1763 }
1764 }
1765
1766 let dialog = dialog.ok_or_else(|| {
1767 anyhow::anyhow!(
1768 "no established SIP dialog found for message command, refer={}",
1769 refer.unwrap_or_default()
1770 )
1771 })?;
1772
1773 let mut sip_headers = vec![rsipstack::rsip::Header::ContentType(
1774 content_type
1775 .clone()
1776 .unwrap_or_else(|| "text/plain;charset=utf-8".to_string())
1777 .into(),
1778 )];
1779 if let Some(headers) = headers {
1780 sip_headers.extend(
1781 headers
1782 .into_iter()
1783 .map(|(k, v)| rsipstack::rsip::Header::Other(k.into(), v.into())),
1784 );
1785 }
1786
1787 info!(
1788 session_id = self.session_id,
1789 dialog_id = %dialog.id(),
1790 content_type = content_type.as_deref().unwrap_or("text/plain;charset=utf-8"),
1791 refer = refer.unwrap_or_default(),
1792 body = %body.chars().take(64).collect::<String>(),
1793 "sending SIP MESSAGE"
1794 );
1795
1796 let response = dialog
1797 .message(Some(sip_headers), Some(body.into_bytes()))
1798 .await?;
1799 match response {
1800 Some(resp)
1801 if resp.status_code.kind() == rsipstack::rsip::StatusCodeKind::Successful =>
1802 {
1803 Ok(())
1804 }
1805 Some(resp) => Err(anyhow::anyhow!(
1806 "SIP MESSAGE rejected with status {}",
1807 resp.status_code
1808 )),
1809 None => Err(anyhow::anyhow!(
1810 "SIP MESSAGE was not sent because dialog is not confirmed"
1811 )),
1812 }
1813 }
1814
1815 fn bridge_track_id(source_session_id: &str, target_session_id: &str) -> TrackId {
1816 format!("bridge:{}:to:{}", source_session_id, target_session_id)
1817 }
1818
1819 async fn do_bridge(&self, target_session_id: String) -> Result<()> {
1820 let target = {
1821 let calls = self.app_state.active_calls.lock().unwrap();
1822 calls.get(&target_session_id).cloned()
1823 };
1824 let target = target.ok_or_else(|| {
1825 anyhow::anyhow!("bridge target session not found: {}", target_session_id)
1826 })?;
1827
1828 if target.session_id == self.session_id {
1829 return Err(anyhow::anyhow!("cannot bridge a call to itself").into());
1830 }
1831
1832 let self_bridge_track_id = Self::bridge_track_id(&self.session_id, &target.session_id);
1833 let target_bridge_track_id = Self::bridge_track_id(&target.session_id, &self.session_id);
1834
1835 self.media_stream
1836 .remove_track(&self_bridge_track_id, false)
1837 .await;
1838 target
1839 .media_stream
1840 .remove_track(&target_bridge_track_id, false)
1841 .await;
1842
1843 let (self_bridge_sender, self_bridge_receiver) = mpsc::channel(25);
1844 let (target_bridge_sender, target_bridge_receiver) = mpsc::channel(25);
1845
1846 let self_paused = self.call_state.read().await.bridge_paused.clone();
1847 let target_paused = target.call_state.read().await.bridge_paused.clone();
1848
1849 let self_forwarding_track = ForwardingTrack::new(
1850 self_bridge_track_id.clone(),
1851 self.session_id.clone(),
1852 target_bridge_sender,
1853 self_bridge_receiver,
1854 self.track_config.clone(),
1855 self.cancel_token.child_token(),
1856 rand::random::<u32>(),
1857 self_paused,
1858 );
1859
1860 let target_forwarding_track = ForwardingTrack::new(
1861 target_bridge_track_id.clone(),
1862 target.session_id.clone(),
1863 self_bridge_sender,
1864 target_bridge_receiver,
1865 target.track_config.clone(),
1866 target.cancel_token.child_token(),
1867 rand::random::<u32>(),
1868 target_paused,
1869 );
1870
1871 self.media_stream
1872 .update_track(Box::new(self_forwarding_track), None)
1873 .await;
1874 target
1875 .media_stream
1876 .update_track(Box::new(target_forwarding_track), None)
1877 .await;
1878
1879 info!(
1880 session_id = self.session_id,
1881 target = target_session_id,
1882 self_bridge_track_id,
1883 target_bridge_track_id,
1884 "audio bridge established"
1885 );
1886 Ok(())
1887 }
1888
1889 async fn do_unbridge(&self, target_session_id: String) -> Result<()> {
1890 let target = {
1891 let calls = self.app_state.active_calls.lock().unwrap();
1892 calls.get(&target_session_id).cloned()
1893 };
1894
1895 let self_bridge_track_id = Self::bridge_track_id(&self.session_id, &target_session_id);
1896 self.media_stream
1897 .remove_track(&self_bridge_track_id, false)
1898 .await;
1899
1900 if let Some(target) = target {
1901 let target_bridge_track_id =
1902 Self::bridge_track_id(&target.session_id, &self.session_id);
1903 target
1904 .media_stream
1905 .remove_track(&target_bridge_track_id, false)
1906 .await;
1907 info!(
1908 session_id = self.session_id,
1909 target = target.session_id,
1910 self_bridge_track_id,
1911 target_bridge_track_id,
1912 "audio bridge removed"
1913 );
1914 } else {
1915 info!(
1916 session_id = self.session_id,
1917 target = target_session_id,
1918 self_bridge_track_id,
1919 "audio bridge removed locally; target session not active"
1920 );
1921 }
1922
1923 Ok(())
1924 }
1925
1926 async fn do_mute(&self, track_id: Option<String>) -> Result<()> {
1927 self.media_stream.mute_track(track_id).await;
1928 Ok(())
1929 }
1930
1931 async fn do_unmute(&self, track_id: Option<String>) -> Result<()> {
1932 self.media_stream.unmute_track(track_id).await;
1933 Ok(())
1934 }
1935
1936 pub async fn cleanup(&self) -> Result<()> {
1937 if matches!(self.call_type, ActiveCallType::Sip | ActiveCallType::B2bua) {
1938 self.do_reject(
1939 Some(rsipstack::rsip::StatusCode::Decline),
1940 Some("handler disconnected".to_string()),
1941 )
1942 .await
1943 .ok();
1944 }
1945 self.call_state.write().await.tts_handle = None;
1946 self.media_stream.cleanup().await.ok();
1947 Ok(())
1948 }
1949
1950 pub fn get_callrecord(&self) -> Option<CallRecord> {
1951 self.call_state.try_read().ok().map(|call_state| {
1952 call_state.build_callrecord(
1953 self.app_state.clone(),
1954 self.session_id.clone(),
1955 self.call_type.clone(),
1956 )
1957 })
1958 }
1959
1960 async fn dump_to_file(
1961 &self,
1962 dump_file: &mut File,
1963 cmd_receiver: &mut CommandReceiver,
1964 event_receiver: &mut EventReceiver,
1965 ) {
1966 loop {
1967 select! {
1968 _ = self.cancel_token.cancelled() => {
1969 break;
1970 }
1971 Ok(cmd) = cmd_receiver.recv() => {
1972 CallRecordEvent::write(CallRecordEventType::Command, cmd, dump_file)
1973 .await;
1974 }
1975 Ok(event) = event_receiver.recv() => {
1976 if matches!(event, SessionEvent::Binary{..}) {
1977 continue;
1978 }
1979 CallRecordEvent::write(CallRecordEventType::Event, event, dump_file)
1980 .await;
1981 }
1982 };
1983 }
1984 }
1985
1986 async fn dump_loop(
1987 &self,
1988 dump_events: bool,
1989 mut dump_cmd_receiver: CommandReceiver,
1990 mut dump_event_receiver: EventReceiver,
1991 ) {
1992 if !dump_events {
1993 return;
1994 }
1995
1996 let file_name = self.app_state.get_dump_events_file(&self.session_id);
1997 let mut dump_file = match File::options()
1998 .create(true)
1999 .append(true)
2000 .open(&file_name)
2001 .await
2002 {
2003 Ok(file) => file,
2004 Err(e) => {
2005 warn!(
2006 session_id = self.session_id,
2007 file_name, "failed to open dump events file: {}", e
2008 );
2009 return;
2010 }
2011 };
2012 self.dump_to_file(
2013 &mut dump_file,
2014 &mut dump_cmd_receiver,
2015 &mut dump_event_receiver,
2016 )
2017 .await;
2018
2019 while let Ok(event) = dump_event_receiver.try_recv() {
2020 if matches!(event, SessionEvent::Binary { .. }) {
2021 continue;
2022 }
2023 CallRecordEvent::write(CallRecordEventType::Event, event, &mut dump_file).await;
2024 }
2025 }
2026
2027 pub async fn create_rtp_track(
2028 &self,
2029 track_id: TrackId,
2030 ssrc: u32,
2031 enable_srtp: Option<bool>,
2032 ) -> Result<RtcTrack> {
2033 let mut rtc_config = RtcTrackConfig::default();
2034 let use_srtp = enable_srtp
2036 .or(self.app_state.config.enable_srtp)
2037 .unwrap_or(false);
2038 rtc_config.mode = if use_srtp {
2039 rustrtc::TransportMode::Srtp
2040 } else {
2041 rustrtc::TransportMode::Rtp
2042 };
2043
2044 if let Some(codecs) = &self.app_state.config.codecs {
2045 let mut codec_types = Vec::new();
2046 for c in codecs {
2047 match c.to_lowercase().as_str() {
2048 "pcmu" => codec_types.push(CodecType::PCMU),
2049 "pcma" => codec_types.push(CodecType::PCMA),
2050 "g722" => codec_types.push(CodecType::G722),
2051 "g729" => codec_types.push(CodecType::G729),
2052 "opus" => codec_types.push(CodecType::Opus),
2053 "dtmf" | "2833" | "telephone_event" => {
2054 codec_types.push(CodecType::TelephoneEvent)
2055 }
2056 _ => {}
2057 }
2058 }
2059 if !codec_types.is_empty() {
2060 rtc_config.preferred_codec = Some(codec_types[0].clone());
2061 rtc_config.codecs = codec_types;
2062 }
2063 }
2064
2065 if rtc_config.preferred_codec.is_none() {
2066 rtc_config.preferred_codec = Some(self.track_config.codec.clone());
2067 }
2068
2069 rtc_config.rtp_port_range = self
2070 .app_state
2071 .config
2072 .rtp_start_port
2073 .zip(self.app_state.config.rtp_end_port);
2074
2075 if let Some(ref external_ip) = self.app_state.config.external_ip {
2076 rtc_config.external_ip = Some(external_ip.clone());
2077 }
2078 if let Some(ref bind_ip) = self.app_state.config.rtp_bind_ip {
2079 rtc_config.bind_ip = Some(bind_ip.clone());
2080 }
2081
2082 rtc_config.enable_latching = self.app_state.config.enable_rtp_latching;
2083 rtc_config.enable_ice_lite = self
2084 .call_state
2085 .read()
2086 .await
2087 .option
2088 .as_ref()
2089 .and_then(|o| o.enable_ice_lite)
2090 .or(self.app_state.config.enable_ice_lite);
2091
2092 let mut track = RtcTrack::new(
2093 self.cancel_token.child_token(),
2094 track_id,
2095 self.track_config.clone(),
2096 rtc_config,
2097 )
2098 .with_ssrc(ssrc);
2099
2100 track.create().await?;
2101
2102 Ok(track)
2103 }
2104
2105 async fn setup_caller_track(&self, option: &CallOption) -> Result<()> {
2106 let hangup_headers = option
2107 .sip
2108 .as_ref()
2109 .and_then(|s| s.hangup_headers.as_ref())
2110 .map(|headers_map| {
2111 headers_map
2112 .iter()
2113 .map(|(k, v)| rsipstack::rsip::Header::Other(k.clone(), v.clone()))
2114 .collect::<Vec<rsipstack::rsip::Header>>()
2115 });
2116 self.call_state.write().await.option = Some(option.clone());
2117 info!(
2118 session_id = self.session_id,
2119 call_type = ?self.call_type,
2120 "setup caller track"
2121 );
2122
2123 let track = match self.call_type {
2124 ActiveCallType::Webrtc => Some(self.create_webrtc_track().await?),
2125 ActiveCallType::WebSocket => {
2126 let audio_receiver = self.call_state.write().await.audio_receiver.take();
2127 if let Some(receiver) = audio_receiver {
2128 Some(self.create_websocket_track(receiver).await?)
2129 } else {
2130 None
2131 }
2132 }
2133 ActiveCallType::Sip => {
2134 if let Some(dialog_id) = self
2135 .invitation
2136 .find_dialog_id_by_session_id(&self.session_id)
2137 {
2138 if let Some(pending_dialog) = self.invitation.get_pending_call(&dialog_id) {
2139 return self
2140 .prepare_incoming_sip_track(
2141 self.cancel_token.clone(),
2142 self.call_state.clone(),
2143 &self.session_id,
2144 pending_dialog,
2145 hangup_headers,
2146 )
2147 .await;
2148 }
2149 }
2150
2151 let mut option = option.clone();
2153 if option.sip.is_none()
2154 || option
2155 .sip
2156 .as_ref()
2157 .and_then(|s| s.username.as_ref())
2158 .is_none()
2159 {
2160 if let Some(callee) = &option.callee {
2161 if let Some(cred) = self.app_state.find_credentials_for_callee(callee) {
2162 if option.sip.is_none() {
2163 option.sip = Some(crate::SipOption {
2164 username: Some(cred.username.clone()),
2165 password: Some(cred.password.clone()),
2166 realm: cred.realm.clone(),
2167 ..Default::default()
2168 });
2169 }
2170 }
2171 }
2172 }
2173
2174 let mut invite_option = option.build_invite_option()?;
2175 invite_option.call_id = Some(self.session_id.clone());
2176
2177 match self
2178 .create_outgoing_sip_track(
2179 self.cancel_token.clone(),
2180 self.call_state.clone(),
2181 &self.session_id,
2182 invite_option,
2183 &option,
2184 None,
2185 false,
2186 )
2187 .await
2188 {
2189 Ok(answer) => {
2190 self.event_sender
2191 .send(SessionEvent::Answer {
2192 timestamp: crate::media::get_timestamp(),
2193 track_id: self.session_id.clone(),
2194 sdp: answer,
2195 refer: Some(false),
2196 })
2197 .ok();
2198 return Ok(());
2199 }
2200 Err(e) => {
2201 warn!(
2202 session_id = self.session_id,
2203 "failed to create sip track: {}", e
2204 );
2205 match &e {
2206 rsipstack::Error::DialogError(reason, _, code) => {
2207 self.event_sender
2208 .send(SessionEvent::Reject {
2209 track_id: self.session_id.clone(),
2210 timestamp: crate::media::get_timestamp(),
2211 reason: reason.clone(),
2212 code: Some(code.code() as u32),
2213 refer: Some(false),
2214 })
2215 .ok();
2216 }
2217 _ => {}
2218 }
2219 return Err(e.into());
2220 }
2221 }
2222 }
2223 ActiveCallType::B2bua => {
2224 if let Some(dialog_id) = self
2225 .invitation
2226 .find_dialog_id_by_session_id(&self.session_id)
2227 {
2228 if let Some(pending_dialog) = self.invitation.get_pending_call(&dialog_id) {
2229 return self
2230 .prepare_incoming_sip_track(
2231 self.cancel_token.clone(),
2232 self.call_state.clone(),
2233 &self.session_id,
2234 pending_dialog,
2235 hangup_headers,
2236 )
2237 .await;
2238 }
2239 }
2240
2241 warn!(
2242 session_id = self.session_id,
2243 "no pending dialog found for B2BUA call"
2244 );
2245 return Err(anyhow::anyhow!(
2246 "no pending dialog found for session_id: {}",
2247 self.session_id
2248 ));
2249 }
2250 };
2251 match track {
2252 Some(track) => {
2253 self.finish_caller_stack(&option, Some(track)).await?;
2254 }
2255 None => {
2256 warn!(session_id = self.session_id, "no track created for caller");
2257 return Err(anyhow::anyhow!("no track created for caller"));
2258 }
2259 }
2260 Ok(())
2261 }
2262
2263 async fn finish_caller_stack(
2264 &self,
2265 option: &CallOption,
2266 track: Option<Box<dyn Track>>,
2267 ) -> Result<()> {
2268 if let Some(track) = track {
2269 self.setup_track_with_stream(&option, track).await?;
2270 }
2271
2272 {
2273 let call_state = self.call_state.read().await;
2274 if let Some(ref answer) = call_state.answer {
2275 info!(
2276 session_id = self.session_id,
2277 "sending answer event: {}", answer,
2278 );
2279 self.event_sender
2280 .send(SessionEvent::Answer {
2281 timestamp: crate::media::get_timestamp(),
2282 track_id: self.session_id.clone(),
2283 sdp: answer.clone(),
2284 refer: Some(false),
2285 })
2286 .ok();
2287 } else {
2288 warn!(
2289 session_id = self.session_id,
2290 "no answer in state to send event"
2291 );
2292 }
2293 }
2294 Ok(())
2295 }
2296
2297 pub async fn setup_track_with_stream(
2298 &self,
2299 option: &CallOption,
2300 mut track: Box<dyn Track>,
2301 ) -> Result<()> {
2302 let processors = match StreamEngine::create_processors(
2303 self.app_state.stream_engine.clone(),
2304 track.as_ref(),
2305 self.cancel_token.child_token(),
2306 self.event_sender.clone(),
2307 self.media_stream.packet_sender.clone(),
2308 option,
2309 )
2310 .await
2311 {
2312 Ok(processors) => processors,
2313 Err(e) => {
2314 warn!(
2315 session_id = self.session_id,
2316 "failed to prepare stream processors: {}", e
2317 );
2318 vec![]
2319 }
2320 };
2321
2322 for processor in processors {
2324 track.append_processor(processor);
2325 }
2326
2327 self.update_track_wrapper(track, None).await;
2328 Ok(())
2329 }
2330
2331 pub async fn update_track_wrapper(&self, mut track: Box<dyn Track>, play_id: Option<String>) {
2332 let (ambiance_opt, subscribe) = {
2333 let state = self.call_state.read().await;
2334 let mut opt = state
2335 .option
2336 .as_ref()
2337 .and_then(|o| o.ambiance.clone())
2338 .unwrap_or_default();
2339
2340 if let Some(global) = &self.app_state.config.ambiance {
2341 opt.merge(global);
2342 }
2343
2344 let subscribe = state
2345 .option
2346 .as_ref()
2347 .and_then(|o| o.subscribe)
2348 .unwrap_or_default();
2349
2350 (opt, subscribe)
2351 };
2352 if track.id() == &self.server_side_track_id && ambiance_opt.path.is_some() {
2353 match AmbianceProcessor::new(ambiance_opt).await {
2354 Ok(ambiance) => {
2355 info!(session_id = self.session_id, "loaded ambiance processor");
2356 track.append_processor(Box::new(ambiance));
2357 }
2358 Err(e) => {
2359 tracing::error!("failed to load ambiance wav {}", e);
2360 }
2361 }
2362 }
2363
2364 if subscribe && self.call_type != ActiveCallType::WebSocket {
2365 let (track_index, sub_track_id) = if track.id() == &self.server_side_track_id {
2366 (0, self.server_side_track_id.clone())
2367 } else {
2368 (1, self.session_id.clone())
2369 };
2370 let sub_processor =
2371 SubscribeProcessor::new(self.event_sender.clone(), sub_track_id, track_index);
2372 track.append_processor(Box::new(sub_processor));
2373 }
2374
2375 self.call_state.write().await.current_play_id = play_id.clone();
2376 self.media_stream.update_track(track, play_id).await;
2377 }
2378
2379 pub async fn create_websocket_track(
2380 &self,
2381 audio_receiver: WebsocketBytesReceiver,
2382 ) -> Result<Box<dyn Track>> {
2383 let (ssrc, codec) = {
2384 let call_state = self.call_state.read().await;
2385 (
2386 call_state.ssrc,
2387 call_state
2388 .option
2389 .as_ref()
2390 .map(|o| o.codec.clone())
2391 .unwrap_or_default(),
2392 )
2393 };
2394
2395 let ws_track = WebsocketTrack::new(
2396 self.cancel_token.child_token(),
2397 self.session_id.clone(),
2398 self.track_config.clone(),
2399 self.event_sender.clone(),
2400 audio_receiver,
2401 codec,
2402 ssrc,
2403 );
2404
2405 {
2406 let mut call_state = self.call_state.write().await;
2407 call_state.answer_time = Some(Utc::now());
2408 call_state.answer = Some("".to_string());
2409 call_state.last_status_code = 200;
2410 }
2411
2412 Ok(Box::new(ws_track))
2413 }
2414
2415 pub(super) async fn create_webrtc_track(&self) -> Result<Box<dyn Track>> {
2416 let (ssrc, option) = {
2417 let call_state = self.call_state.read().await;
2418 (
2419 call_state.ssrc,
2420 call_state.option.clone().unwrap_or_default(),
2421 )
2422 };
2423
2424 let mut rtc_config = RtcTrackConfig::default();
2425 rtc_config.mode = rustrtc::TransportMode::WebRtc; rtc_config.ice_servers = self.app_state.config.ice_servers.clone();
2427
2428 if let Some(codecs) = &self.app_state.config.codecs {
2429 let mut codec_types = Vec::new();
2430 for c in codecs {
2431 match c.to_lowercase().as_str() {
2432 "pcmu" => codec_types.push(CodecType::PCMU),
2433 "pcma" => codec_types.push(CodecType::PCMA),
2434 "g722" => codec_types.push(CodecType::G722),
2435 "g729" => codec_types.push(CodecType::G729),
2436 "opus" => codec_types.push(CodecType::Opus),
2437 "dtmf" | "2833" | "telephone_event" => {
2438 codec_types.push(CodecType::TelephoneEvent)
2439 }
2440 _ => {}
2441 }
2442 }
2443 if !codec_types.is_empty() {
2444 rtc_config.preferred_codec = Some(codec_types[0].clone());
2445 rtc_config.codecs = codec_types;
2446 }
2447 }
2448
2449 if let Some(ref external_ip) = self.app_state.config.external_ip {
2450 rtc_config.external_ip = Some(external_ip.clone());
2451 }
2452 if let Some(ref bind_ip) = self.app_state.config.rtp_bind_ip {
2453 rtc_config.bind_ip = Some(bind_ip.clone());
2454 }
2455
2456 let mut webrtc_track = RtcTrack::new(
2457 self.cancel_token.child_token(),
2458 self.session_id.clone(),
2459 self.track_config.clone(),
2460 rtc_config,
2461 )
2462 .with_ssrc(ssrc);
2463
2464 let timeout = option.handshake_timeout.map(|t| Duration::from_secs(t));
2465 let offer = match option.enable_ipv6 {
2466 Some(false) | None => {
2467 strip_ipv6_candidates(option.offer.as_ref().unwrap_or(&"".to_string()))
2468 }
2469 _ => option.offer.clone().unwrap_or("".to_string()),
2470 };
2471 let answer: Option<String>;
2472 match webrtc_track.handshake(offer, timeout).await {
2473 Ok(answer_sdp) => {
2474 answer = match option.enable_ipv6 {
2475 Some(false) | None => Some(strip_ipv6_candidates(&answer_sdp)),
2476 Some(true) => Some(answer_sdp.to_string()),
2477 };
2478 }
2479 Err(e) => {
2480 warn!(session_id = self.session_id, "failed to setup track: {}", e);
2481 return Err(anyhow::anyhow!("Failed to setup track: {}", e));
2482 }
2483 }
2484
2485 {
2486 let mut call_state = self.call_state.write().await;
2487 call_state.answer_time = Some(Utc::now());
2488 call_state.answer = answer;
2489 call_state.last_status_code = 200;
2490 }
2491 Ok(Box::new(webrtc_track))
2492 }
2493
2494 async fn create_outgoing_sip_track(
2495 &self,
2496 cancel_token: CancellationToken,
2497 call_state_ref: ActiveCallStateRef,
2498 track_id: &String,
2499 mut invite_option: InviteOption,
2500 call_option: &CallOption,
2501 moh: Option<String>,
2502 auto_hangup: bool,
2503 ) -> Result<String, rsipstack::Error> {
2504 let ssrc = call_state_ref.read().await.ssrc;
2505 let per_call_srtp = call_option.sip.as_ref().and_then(|s| s.enable_srtp);
2506 let rtp_track = self
2507 .create_rtp_track(track_id.clone(), ssrc, per_call_srtp)
2508 .await
2509 .map_err(|e| rsipstack::Error::Error(e.to_string()))?;
2510
2511 let offer = Some(
2512 rtp_track
2513 .local_description()
2514 .await
2515 .map_err(|e| rsipstack::Error::Error(e.to_string()))?,
2516 );
2517
2518 {
2519 let mut cs = call_state_ref.write().await;
2520 if let Some(o) = cs.option.as_mut() {
2521 o.offer = offer.clone();
2522 }
2523 cs.start_time = Utc::now();
2524 };
2525
2526 invite_option.offer = offer.clone().map(|s| s.into());
2527
2528 let needs_contact = contact_needs_public_resolution(&invite_option.contact);
2531
2532 if needs_contact {
2533 let addrs = self.invitation.dialog_layer.endpoint.get_addrs();
2534 if let Some(addr) = find_local_addr_for_uri(&addrs, &invite_option.callee) {
2535 let contact_username = invite_option
2536 .contact
2537 .auth
2538 .as_ref()
2539 .map(|auth| auth.user.as_str())
2540 .or_else(|| {
2541 invite_option
2542 .caller
2543 .auth
2544 .as_ref()
2545 .map(|auth| auth.user.as_str())
2546 });
2547 invite_option.contact = build_public_contact_uri(
2548 &self.app_state.learned_public_address,
2549 self.app_state.auto_learn_public_address_enabled(),
2550 &addr,
2551 contact_username,
2552 Some(&invite_option.contact),
2553 );
2554 } else {
2555 return Err(rsipstack::Error::Error(format!(
2556 "missing local SIP address for callee transport: {}",
2557 invite_option.callee
2558 )));
2559 }
2560 }
2561
2562 let mut rtp_track_to_setup = Some(Box::new(rtp_track) as Box<dyn Track>);
2563
2564 if let Some(moh) = moh {
2565 let ssrc_and_moh = {
2566 let mut state = call_state_ref.write().await;
2567 state.moh = Some(moh.clone());
2568 if state.current_play_id.is_none() {
2569 let ssrc = rand::random::<u32>();
2570 Some((ssrc, moh.clone()))
2571 } else {
2572 info!(
2573 session_id = self.session_id,
2574 "Something is playing, MOH will start after it ends"
2575 );
2576 None
2577 }
2578 };
2579
2580 if let Some((ssrc, moh_path)) = ssrc_and_moh {
2581 let file_track = FileTrack::new(self.server_side_track_id.clone())
2582 .with_play_id(Some(moh_path.clone()))
2583 .with_ssrc(ssrc)
2584 .with_path(moh_path.clone())
2585 .with_cancel_token(self.cancel_token.child_token());
2586 self.update_track_wrapper(Box::new(file_track), Some(moh_path))
2587 .await;
2588 }
2589 } else {
2590 let track = rtp_track_to_setup.take().unwrap();
2591 self.setup_track_with_stream(&call_option, track)
2592 .await
2593 .map_err(|e| rsipstack::Error::Error(e.to_string()))?;
2594 }
2595
2596 info!(
2597 session_id = self.session_id,
2598 track_id,
2599 contact = %invite_option.contact,
2600 "invite {} -> {} offer: \n{}",
2601 invite_option.caller,
2602 invite_option.callee,
2603 offer.as_ref().map(|s| s.as_str()).unwrap_or("<NO OFFER>")
2604 );
2605
2606 let (dlg_state_sender, dlg_state_receiver) =
2607 self.invitation.dialog_layer.new_dialog_state_channel();
2608
2609 let states = InviteDialogStates {
2610 is_client: true,
2611 session_id: self.session_id.clone(),
2612 track_id: track_id.clone(),
2613 event_sender: self.event_sender.clone(),
2614 media_stream: self.media_stream.clone(),
2615 call_state: call_state_ref.clone(),
2616 cancel_token,
2617 terminated_reason: None,
2618 has_early_media: false,
2619 };
2620
2621 let hangup_headers = call_option
2622 .sip
2623 .as_ref()
2624 .and_then(|s| s.hangup_headers.as_ref())
2625 .map(|headers_map| {
2626 headers_map
2627 .iter()
2628 .map(|(k, v)| rsipstack::rsip::Header::Other(k.clone(), v.clone()))
2629 .collect::<Vec<rsipstack::rsip::Header>>()
2630 });
2631
2632 let mut client_dialog_handler = DialogStateReceiverGuard::new(
2633 self.invitation.dialog_layer.clone(),
2634 dlg_state_receiver,
2635 hangup_headers,
2636 );
2637
2638 crate::spawn(async move {
2639 client_dialog_handler.process_dialog(states).await;
2640 });
2641
2642 let (dialog_id, answer) = self
2643 .invitation
2644 .invite(invite_option, dlg_state_sender)
2645 .await?;
2646
2647 self.call_state.write().await.moh = None;
2648
2649 if let Some(track) = rtp_track_to_setup {
2650 info!(
2651 session_id = self.session_id,
2652 track_id, "Stopping MOH and setting up RTP track"
2653 );
2654 self.media_stream
2655 .remove_track(&self.server_side_track_id, false)
2656 .await;
2657
2658 self.setup_track_with_stream(&call_option, track)
2659 .await
2660 .map_err(|e| rsipstack::Error::Error(e.to_string()))?;
2661 }
2662
2663 let answer = match answer {
2664 Some(answer) => {
2665 let s = String::from_utf8_lossy(&answer).to_string();
2666 if s.trim().is_empty() {
2667 let cs = call_state_ref.read().await;
2671 match cs.answer.clone() {
2672 Some(early_sdp) if !early_sdp.is_empty() => {
2673 info!(
2674 session_id = self.session_id,
2675 "200 OK has empty body; using early-media SDP from 183"
2676 );
2677 (early_sdp, true )
2678 }
2679 _ => {
2680 warn!(
2681 session_id = self.session_id,
2682 "200 OK has empty body and no early-media SDP available"
2683 );
2684 (s, false)
2685 }
2686 }
2687 } else {
2688 (s, false)
2689 }
2690 }
2691 None => {
2692 let cs = call_state_ref.read().await;
2694 match cs.answer.clone() {
2695 Some(early_sdp) if !early_sdp.is_empty() => {
2696 info!(
2697 session_id = self.session_id,
2698 "200 OK had no answer; using early-media SDP from 183"
2699 );
2700 (early_sdp, true )
2701 }
2702 _ => {
2703 warn!(session_id = self.session_id, "no answer received");
2704 return Err(rsipstack::Error::DialogError(
2705 "No answer received".to_string(),
2706 dialog_id,
2707 rsipstack::rsip::StatusCode::NotAcceptableHere,
2708 ));
2709 }
2710 }
2711 }
2712 };
2713 let (answer, remote_description_already_applied) = answer;
2714
2715 {
2716 let mut cs = call_state_ref.write().await;
2717 if cs.answer.is_none() {
2718 cs.answer = Some(answer.clone());
2719 }
2720 if auto_hangup {
2721 cs.auto_hangup = Some((ssrc, CallRecordHangupReason::ByRefer));
2722 }
2723 }
2724 if !remote_description_already_applied {
2725 self.media_stream
2726 .update_remote_description(&track_id, &answer)
2727 .await
2728 .ok();
2729 }
2730
2731 Ok(answer)
2732 }
2733
2734 pub fn is_webrtc_sdp(sdp: &str) -> bool {
2736 (sdp.contains("a=ice-ufrag:") || sdp.contains("a=ice-pwd:"))
2737 && sdp.contains("a=fingerprint:")
2738 }
2739
2740 pub async fn setup_answer_track(
2741 &self,
2742 ssrc: u32,
2743 option: &CallOption,
2744 offer: String,
2745 ) -> Result<(String, Box<dyn Track>)> {
2746 let offer = match option.enable_ipv6 {
2747 Some(false) | None => strip_ipv6_candidates(&offer),
2748 _ => offer.clone(),
2749 };
2750
2751 let timeout = option.handshake_timeout.map(|t| Duration::from_secs(t));
2752
2753 let mut media_track = if Self::is_webrtc_sdp(&offer) {
2754 let mut rtc_config = RtcTrackConfig::default();
2755 rtc_config.mode = rustrtc::TransportMode::WebRtc;
2756 rtc_config.ice_servers = self.app_state.config.ice_servers.clone();
2757 if let Some(ref external_ip) = self.app_state.config.external_ip {
2758 rtc_config.external_ip = Some(external_ip.clone());
2759 }
2760 if let Some(ref bind_ip) = self.app_state.config.rtp_bind_ip {
2761 rtc_config.bind_ip = Some(bind_ip.clone());
2762 }
2763 rtc_config.enable_latching = self.app_state.config.enable_rtp_latching;
2764 rtc_config.enable_ice_lite = self
2765 .call_state
2766 .read()
2767 .await
2768 .option
2769 .as_ref()
2770 .and_then(|o| o.enable_ice_lite)
2771 .or(self.app_state.config.enable_ice_lite);
2772
2773 let webrtc_track = RtcTrack::new(
2774 self.cancel_token.child_token(),
2775 self.session_id.clone(),
2776 self.track_config.clone(),
2777 rtc_config,
2778 )
2779 .with_ssrc(ssrc);
2780
2781 Box::new(webrtc_track) as Box<dyn Track>
2782 } else {
2783 let per_call_srtp = option.sip.as_ref().and_then(|s| s.enable_srtp);
2784 let rtp_track = self
2785 .create_rtp_track(self.session_id.clone(), ssrc, per_call_srtp)
2786 .await?;
2787 Box::new(rtp_track) as Box<dyn Track>
2788 };
2789
2790 let answer = match media_track.handshake(offer.clone(), timeout).await {
2791 Ok(answer) => answer,
2792 Err(e) => {
2793 return Err(anyhow::anyhow!("handshake failed: {e}"));
2794 }
2795 };
2796
2797 return Ok((answer, media_track));
2798 }
2799
2800 pub async fn prepare_incoming_sip_track(
2801 &self,
2802 cancel_token: CancellationToken,
2803 call_state_ref: ActiveCallStateRef,
2804 track_id: &String,
2805 pending_dialog: PendingDialog,
2806 hangup_headers: Option<Vec<rsipstack::rsip::Header>>,
2807 ) -> Result<()> {
2808 let state_receiver = pending_dialog.state_receiver;
2809 let states = InviteDialogStates {
2812 is_client: false,
2813 session_id: self.session_id.clone(),
2814 track_id: track_id.clone(),
2815 event_sender: self.event_sender.clone(),
2816 media_stream: self.media_stream.clone(),
2817 call_state: self.call_state.clone(),
2818 cancel_token,
2819 terminated_reason: None,
2820 has_early_media: false,
2821 };
2822
2823 let initial_request = pending_dialog.dialog.initial_request();
2824 let offer = String::from_utf8_lossy(&initial_request.body).to_string();
2825
2826 let (ssrc, option) = {
2827 let call_state = call_state_ref.read().await;
2828 (
2829 call_state.ssrc,
2830 call_state.option.clone().unwrap_or_default(),
2831 )
2832 };
2833
2834 match self.setup_answer_track(ssrc, &option, offer).await {
2835 Ok((offer, track)) => {
2836 self.setup_track_with_stream(&option, track).await?;
2837 {
2838 let mut state = self.call_state.write().await;
2839 state.ready_to_answer = Some((offer, None, pending_dialog.dialog));
2840 }
2841 }
2842 Err(e) => {
2843 return Err(anyhow::anyhow!("error creating track: {}", e));
2844 }
2845 }
2846
2847 let mut client_dialog_handler = DialogStateReceiverGuard::new(
2848 self.invitation.dialog_layer.clone(),
2849 state_receiver,
2850 hangup_headers,
2851 );
2852
2853 crate::spawn(async move {
2854 client_dialog_handler.process_dialog(states).await;
2855 });
2856 Ok(())
2857 }
2858}
2859
2860impl Drop for ActiveCall {
2861 fn drop(&mut self) {
2862 info!(session_id = self.session_id, "dropping active call");
2863 if let Some(sender) = self.app_state.callrecord_sender.as_ref() {
2864 if let Some(record) = self.get_callrecord() {
2865 if let Err(e) = sender.send(record) {
2866 warn!(
2867 session_id = self.session_id,
2868 "failed to send call record: {}", e
2869 );
2870 }
2871 }
2872 }
2873 }
2874}
2875
2876impl ActiveCallState {
2877 pub fn merge_option(&self, mut option: CallOption) -> CallOption {
2878 if let Some(existing) = &self.option {
2879 if option.asr.is_none() {
2880 option.asr = existing.asr.clone();
2881 }
2882 if option.tts.is_none() {
2883 option.tts = existing.tts.clone();
2884 }
2885 if option.vad.is_none() {
2886 option.vad = existing.vad.clone();
2887 }
2888 if option.denoise.is_none() {
2889 option.denoise = existing.denoise;
2890 }
2891 if option.agc.is_none() {
2892 option.agc = existing.agc.clone();
2893 }
2894 if option.recorder.is_none() {
2895 option.recorder = existing.recorder.clone();
2896 }
2897 if option.eou.is_none() {
2898 option.eou = existing.eou.clone();
2899 }
2900 if option.extra.is_none() {
2901 option.extra = existing.extra.clone();
2902 }
2903 if option.ambiance.is_none() {
2904 option.ambiance = existing.ambiance.clone();
2905 }
2906 if option.ringback_detection.is_none() {
2907 option.ringback_detection = existing.ringback_detection.clone();
2908 }
2909 }
2910 option
2911 }
2912
2913 pub fn set_hangup_reason(&mut self, reason: CallRecordHangupReason) {
2914 if self.hangup_reason.is_none() {
2915 self.hangup_reason = Some(reason);
2916 }
2917 }
2918
2919 pub fn build_hangup_event(
2920 &self,
2921 track_id: TrackId,
2922 initiator: Option<String>,
2923 ) -> crate::event::SessionEvent {
2924 let from = self.option.as_ref().and_then(|o| o.caller.as_ref());
2925 let to = self.option.as_ref().and_then(|o| o.callee.as_ref());
2926 let extra = self.extras.clone();
2927
2928 crate::event::SessionEvent::Hangup {
2929 track_id,
2930 timestamp: crate::media::get_timestamp(),
2931 reason: Some(format!("{:?}", self.hangup_reason)),
2932 initiator,
2933 start_time: self.start_time.to_rfc3339(),
2934 answer_time: self.answer_time.map(|t| t.to_rfc3339()),
2935 ringing_time: self.ring_time.map(|t| t.to_rfc3339()),
2936 hangup_time: Utc::now().to_rfc3339(),
2937 extra,
2938 from: from.map(|f| f.into()),
2939 to: to.map(|f| f.into()),
2940 refer: Some(self.is_refer),
2941 }
2942 }
2943
2944 pub fn build_callrecord(
2945 &self,
2946 app_state: AppState,
2947 session_id: String,
2948 call_type: ActiveCallType,
2949 ) -> CallRecord {
2950 let option = self.option.clone().unwrap_or_default();
2951 let recorder = if option.recorder.is_some() {
2952 let recorder_file = app_state.get_recorder_file(&session_id);
2953 if std::path::Path::new(&recorder_file).exists() {
2954 let file_size = std::fs::metadata(&recorder_file)
2955 .map(|m| m.len())
2956 .unwrap_or(0);
2957 vec![crate::callrecord::CallRecordMedia {
2958 track_id: session_id.clone(),
2959 path: recorder_file,
2960 size: file_size,
2961 extra: None,
2962 }]
2963 } else {
2964 vec![]
2965 }
2966 } else {
2967 vec![]
2968 };
2969
2970 let dump_event_file = app_state.get_dump_events_file(&session_id);
2971 let dump_event_file = if std::path::Path::new(&dump_event_file).exists() {
2972 Some(dump_event_file)
2973 } else {
2974 None
2975 };
2976
2977 let refer_callrecord = self.refer_callstate.as_ref().and_then(|rc| {
2978 if let Ok(rc) = rc.try_read() {
2979 Some(Box::new(rc.build_callrecord(
2980 app_state.clone(),
2981 rc.session_id.clone(),
2982 ActiveCallType::B2bua,
2983 )))
2984 } else {
2985 None
2986 }
2987 });
2988
2989 let caller = option.caller.clone().unwrap_or_default();
2990 let callee = option.callee.clone().unwrap_or_default();
2991
2992 CallRecord {
2993 option: Some(option),
2994 call_id: session_id,
2995 call_type,
2996 start_time: self.start_time,
2997 ring_time: self.ring_time.clone(),
2998 answer_time: self.answer_time.clone(),
2999 end_time: Utc::now(),
3000 caller,
3001 callee,
3002 hangup_reason: self.hangup_reason.clone(),
3003 hangup_messages: Vec::new(),
3004 status_code: self.last_status_code,
3005 extras: self.extras.clone(),
3006 dump_event_file,
3007 recorder,
3008 refer_callrecord,
3009 }
3010 }
3011}