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::Bridge { target_session_id } => self.do_bridge(target_session_id).await,
865 Command::Unbridge { target_session_id } => self.do_unbridge(target_session_id).await,
866 Command::Mute { track_id } => self.do_mute(track_id).await,
867 Command::Unmute { track_id } => self.do_unmute(track_id).await,
868 Command::Pause {} => self.do_pause().await,
869 Command::Resume {} => self.do_resume().await,
870 Command::Interrupt {
871 graceful: passage,
872 fade_out_ms: _,
873 } => self.do_interrupt(passage.unwrap_or_default()).await,
874 Command::History { speaker, text } => self.do_history(speaker, text).await,
875 Command::Custom { sender, data } => self.do_custom(sender, data),
876 }
877 }
878
879 fn build_record_option(&self, option: &CallOption) -> Option<RecorderOption> {
880 if let Some(recorder_option) = &option.recorder {
881 let mut recorder_file = recorder_option.recorder_file.clone();
882 if recorder_file.contains("{id}") {
883 recorder_file = recorder_file.replace("{id}", &self.session_id);
884 }
885
886 let recorder_file = if recorder_file.is_empty() {
887 self.app_state.get_recorder_file(&self.session_id)
888 } else {
889 let p = Path::new(&recorder_file);
890 p.is_absolute()
891 .then(|| recorder_file.clone())
892 .unwrap_or_else(|| self.app_state.get_recorder_file(&recorder_file))
893 };
894 info!(
895 session_id = self.session_id,
896 recorder_file, "created recording file"
897 );
898
899 let track_samplerate = self.track_config.samplerate;
900 let recorder_samplerate = if track_samplerate > 0 {
901 track_samplerate
902 } else {
903 recorder_option.samplerate
904 };
905 let recorder_ptime = if recorder_option.ptime == 0 {
906 200
907 } else {
908 recorder_option.ptime
909 };
910 let requested_format = recorder_option
911 .format
912 .unwrap_or(self.app_state.config.recorder_format());
913 let format = requested_format.effective();
914 if requested_format != format {
915 warn!(
916 session_id = self.session_id,
917 requested = requested_format.extension(),
918 "Recorder format fallback to wav due to unsupported feature"
919 );
920 }
921 let mut recorder_config = RecorderOption {
922 recorder_file,
923 samplerate: recorder_samplerate,
924 ptime: recorder_ptime,
925 format: Some(format),
926 };
927 recorder_config.ensure_path_extension(format);
928 Some(recorder_config)
929 } else {
930 None
931 }
932 }
933
934 async fn invite_or_accept(&self, mut option: CallOption, sender: String) -> Result<CallOption> {
935 {
937 let state = self.call_state.read().await;
938 option = state.merge_option(option);
939 }
940
941 option.check_default();
942 if let Some(opt) = self.build_record_option(&option) {
943 self.media_stream.update_recorder_option(opt).await;
944 }
945
946 if let Some(opt) = &option.media_pass {
947 let track_id = self.server_side_track_id.clone();
948 let cancel_token = self.cancel_token.child_token();
949 let ssrc = rand::random::<u32>();
950 let media_pass_track = MediaPassTrack::new(
951 self.session_id.clone(),
952 ssrc,
953 track_id,
954 cancel_token,
955 opt.clone(),
956 );
957 self.update_track_wrapper(Box::new(media_pass_track), None)
958 .await;
959 }
960
961 info!(
962 session_id = self.session_id,
963 call_type = ?self.call_type,
964 sender,
965 ?option,
966 "caller with option"
967 );
968
969 match self.setup_caller_track(&option).await {
970 Ok(_) => return Ok(option),
971 Err(e) => {
972 self.app_state
973 .total_failed_calls
974 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
975 let error_event = crate::event::SessionEvent::Error {
976 track_id: self.session_id.clone(),
977 timestamp: crate::media::get_timestamp(),
978 sender,
979 error: e.to_string(),
980 code: None,
981 };
982 self.event_sender.send(error_event).ok();
983 self.do_hangup(Some(CallRecordHangupReason::BySystem), None, None, None)
984 .await
985 .ok();
986 return Err(e);
987 }
988 }
989 }
990
991 async fn do_invite(&self, option: CallOption) -> Result<()> {
992 self.invite_or_accept(option, "invite".to_string())
993 .await
994 .map(|_| ())
995 }
996
997 async fn do_accept(&self, mut option: CallOption) -> Result<()> {
998 let has_pending = self
999 .invitation
1000 .find_dialog_id_by_session_id(&self.session_id)
1001 .is_some();
1002 let ready_to_answer_val = {
1003 let state = self.call_state.read().await;
1004 state.ready_to_answer.is_none()
1005 };
1006
1007 if ready_to_answer_val {
1008 if !has_pending {
1009 warn!(session_id = self.session_id, "no pending call to accept");
1011 let rejet_event = crate::event::SessionEvent::Reject {
1012 track_id: self.session_id.clone(),
1013 timestamp: crate::media::get_timestamp(),
1014 reason: "no pending call".to_string(),
1015 refer: None,
1016 code: Some(486),
1017 };
1018 self.event_sender.send(rejet_event).ok();
1019 self.do_hangup(Some(CallRecordHangupReason::BySystem), None, None, None)
1020 .await
1021 .ok();
1022 return Err(anyhow::anyhow!("no pending call to accept"));
1023 }
1024 option = self.invite_or_accept(option, "accept".to_string()).await?;
1025 } else {
1026 option.check_default();
1027 if let Some(opt) = self.build_record_option(&option) {
1028 self.media_stream.update_recorder_option(opt).await;
1029 }
1030 self.call_state.write().await.option = Some(option.clone());
1031 }
1032 info!(session_id = self.session_id, ?option, "accepting call");
1033 let ready = self.call_state.write().await.ready_to_answer.take();
1034 if let Some((answer, track, dialog)) = ready {
1035 info!(
1036 session_id = self.session_id,
1037 track_id = track.as_ref().map(|t| t.id()),
1038 "ready to answer with track"
1039 );
1040
1041 let headers = vec![rsipstack::rsip::Header::ContentType(
1042 "application/sdp".to_string().into(),
1043 )];
1044
1045 match dialog.accept(Some(headers), Some(answer.as_bytes().to_vec())) {
1046 Ok(_) => {
1047 {
1048 let mut state = self.call_state.write().await;
1049 state.answer = Some(answer);
1050 state.answer_time = Some(Utc::now());
1051 }
1052 self.finish_caller_stack(&option, track).await?;
1053 }
1054 Err(e) => {
1055 warn!(session_id = self.session_id, "failed to accept call: {}", e);
1056 return Err(anyhow::anyhow!("failed to accept call"));
1057 }
1058 }
1059 }
1060 return Ok(());
1061 }
1062
1063 async fn do_reject(
1064 &self,
1065 code: Option<rsipstack::rsip::StatusCode>,
1066 reason: Option<String>,
1067 ) -> Result<()> {
1068 match self
1069 .invitation
1070 .find_dialog_id_by_session_id(&self.session_id)
1071 {
1072 Some(id) => {
1073 info!(
1074 session_id = self.session_id,
1075 ?reason,
1076 ?code,
1077 "rejecting call"
1078 );
1079 let result = self.invitation.hangup(id, code, reason).await;
1080 if result.is_ok() {
1081 self.cancel_token.cancel();
1082 }
1083 result
1084 }
1085 None => {
1086 let ready = self.call_state.write().await.ready_to_answer.take();
1087 if let Some((_, _, dialog)) = ready {
1088 info!(
1089 session_id = self.session_id,
1090 ?reason,
1091 ?code,
1092 "rejecting call from ready_to_answer"
1093 );
1094 let dialog_id = dialog.id();
1095 dialog.reject(code, reason).ok();
1096 self.invitation.dialog_layer.remove_dialog(&dialog_id);
1097 self.cancel_token.cancel();
1098 }
1099 Ok(())
1100 }
1101 }
1102 }
1103
1104 async fn do_ringing(
1105 &self,
1106 ringtone: Option<String>,
1107 recorder: Option<RecorderOption>,
1108 early_media: Option<bool>,
1109 ) -> Result<()> {
1110 let ready_to_answer_val = self.call_state.read().await.ready_to_answer.is_none();
1111 if ready_to_answer_val {
1112 let option = CallOption {
1113 recorder,
1114 ..Default::default()
1115 };
1116 let _ = self.invite_or_accept(option, "ringing".to_string()).await?;
1117 }
1118
1119 let state = self.call_state.read().await;
1120 if let Some((answer, _, dialog)) = state.ready_to_answer.as_ref() {
1121 let (headers, body) = if early_media.unwrap_or_default() || ringtone.is_some() {
1122 let headers = vec![rsipstack::rsip::Header::ContentType(
1123 "application/sdp".to_string().into(),
1124 )];
1125 (Some(headers), Some(answer.as_bytes().to_vec()))
1126 } else {
1127 (None, None)
1128 };
1129
1130 dialog.ringing(headers, body).ok();
1131 info!(
1132 session_id = self.session_id,
1133 ringtone, early_media, "playing ringtone"
1134 );
1135 if let Some(ringtone_url) = ringtone {
1136 drop(state);
1137 self.do_play(ringtone_url, None, None, None, None)
1138 .await
1139 .ok();
1140 } else {
1141 info!(session_id = self.session_id, "no ringtone to play");
1142 }
1143 }
1144 Ok(())
1145 }
1146
1147 async fn do_tts(
1148 &self,
1149 text: String,
1150 speaker: Option<String>,
1151 play_id: Option<String>,
1152 auto_hangup: Option<bool>,
1153 streaming: bool,
1154 end_of_stream: bool,
1155 option: Option<SynthesisOption>,
1156 wait_input_timeout: Option<u32>,
1157 base64: bool,
1158 cache_key: Option<String>,
1159 ) -> Result<()> {
1160 let tts_option = {
1161 let call_state = self.call_state.read().await;
1162 match call_state.option.clone().unwrap_or_default().tts {
1163 Some(opt) => opt.merge_with(option),
1164 None => {
1165 if let Some(opt) = option {
1166 opt
1167 } else {
1168 return Err(anyhow::anyhow!("no tts option available"));
1169 }
1170 }
1171 }
1172 };
1173 let speaker = match speaker {
1174 Some(s) => Some(s),
1175 None => tts_option.speaker.clone(),
1176 };
1177
1178 let mut play_command = SynthesisCommand {
1179 text,
1180 speaker,
1181 play_id: play_id.clone(),
1182 streaming,
1183 end_of_stream: if !streaming { true } else { end_of_stream },
1184 option: tts_option,
1185 base64,
1186 cache_key,
1187 };
1188 info!(
1189 session_id = self.session_id,
1190 provider = ?play_command.option.provider,
1191 text = %play_command.text.chars().take(10).collect::<String>(),
1192 speaker = play_command.speaker.as_deref(),
1193 auto_hangup = auto_hangup.unwrap_or_default(),
1194 play_id = play_command.play_id.as_deref(),
1195 streaming = play_command.streaming,
1196 end_of_stream = play_command.end_of_stream,
1197 wait_input_timeout = wait_input_timeout.unwrap_or_default(),
1198 is_base64 = play_command.base64,
1199 cache_key = play_command.cache_key.as_deref(),
1200 "new synthesis"
1201 );
1202
1203 let ssrc = rand::random::<u32>();
1204 let (should_interrupt, picked_ssrc) = {
1205 let mut state = self.call_state.write().await;
1206
1207 let (target_ssrc, changed) = if let Some(handle) = &state.tts_handle {
1208 if play_id.is_some() && state.current_play_id != play_id {
1209 (ssrc, true)
1210 } else {
1211 (handle.ssrc, false)
1212 }
1213 } else {
1214 (ssrc, false)
1215 };
1216
1217 state.wait_input_timeout = wait_input_timeout;
1220
1221 state.current_play_id = play_id.clone();
1222 (changed, target_ssrc)
1223 };
1224
1225 if should_interrupt {
1226 let _ = self.do_interrupt(false).await;
1227 }
1228
1229 {
1233 let mut state = self.call_state.write().await;
1234 state.auto_hangup = match auto_hangup {
1235 Some(true) => Some((picked_ssrc, CallRecordHangupReason::BySystem)),
1236 _ => {
1237 if state.tts_handle.is_some() && !should_interrupt {
1241 state.auto_hangup.clone()
1242 } else {
1243 None
1244 }
1245 }
1246 };
1247 }
1248
1249 let existing_handle = self.call_state.read().await.tts_handle.clone();
1250 if let Some(tts_handle) = existing_handle {
1251 match tts_handle.try_send(play_command) {
1252 Ok(_) => return Ok(()),
1253 Err(e) => {
1254 play_command = e.0;
1255 }
1256 }
1257 }
1258
1259 let (new_handle, tts_track) = StreamEngine::create_tts_track(
1260 self.app_state.stream_engine.clone(),
1261 self.cancel_token.child_token(),
1262 self.session_id.clone(),
1263 self.server_side_track_id.clone(),
1264 picked_ssrc,
1265 play_id.clone(),
1266 streaming,
1267 &play_command.option,
1268 )
1269 .await?;
1270
1271 new_handle.try_send(play_command)?;
1272 self.call_state.write().await.tts_handle = Some(new_handle);
1273 self.update_track_wrapper(tts_track, play_id).await;
1274 Ok(())
1275 }
1276
1277 async fn do_play(
1278 &self,
1279 url: String,
1280 play_id: Option<String>,
1281 auto_hangup: Option<bool>,
1282 wait_input_timeout: Option<u32>,
1283 offset_ms: Option<u32>,
1284 ) -> Result<()> {
1285 let ssrc = rand::random::<u32>();
1286 info!(
1287 session_id = self.session_id,
1288 ssrc, url, play_id, auto_hangup, "play file track"
1289 );
1290
1291 let play_id = play_id.or(Some(url.clone()));
1292
1293 let mut file_track = FileTrack::new(self.server_side_track_id.clone())
1294 .with_play_id(play_id.clone())
1295 .with_ssrc(ssrc)
1296 .with_path(url)
1297 .with_cancel_token(self.cancel_token.child_token());
1298
1299 if let Some(offset) = offset_ms {
1300 file_track = file_track.with_offset_ms(offset);
1301 }
1302
1303 {
1304 let mut state = self.call_state.write().await;
1305 state.tts_handle = None;
1306 state.auto_hangup = match auto_hangup {
1307 Some(true) => Some((ssrc, CallRecordHangupReason::BySystem)),
1308 _ => None,
1309 };
1310 state.wait_input_timeout = wait_input_timeout;
1311 }
1312
1313 self.update_track_wrapper(Box::new(file_track), play_id)
1314 .await;
1315 Ok(())
1316 }
1317
1318 async fn do_history(&self, speaker: String, text: String) -> Result<()> {
1319 self.event_sender
1320 .send(SessionEvent::AddHistory {
1321 sender: Some(self.session_id.clone()),
1322 timestamp: crate::media::get_timestamp(),
1323 speaker,
1324 text,
1325 })
1326 .map(|_| ())
1327 .map_err(Into::into)
1328 }
1329
1330 fn do_custom(&self, sender: Option<String>, data: serde_json::Value) -> Result<()> {
1331 self.event_sender
1332 .send(SessionEvent::Custom {
1333 timestamp: crate::media::get_timestamp(),
1334 sender,
1335 data,
1336 })
1337 .map(|_| ())
1338 .map_err(Into::into)
1339 }
1340
1341 async fn do_interrupt(&self, graceful: bool) -> Result<()> {
1342 {
1343 let mut state = self.call_state.write().await;
1344 state.tts_handle = None;
1345 state.moh = None;
1346 state.auto_hangup = None;
1347 }
1348 self.media_stream
1349 .remove_track(&self.server_side_track_id, graceful)
1350 .await;
1351 Ok(())
1352 }
1353 async fn do_pause(&self) -> Result<()> {
1354 self.media_stream
1355 .pause_playback(self.server_side_track_id.clone())
1356 .await?;
1357 Ok(())
1358 }
1359 async fn do_resume(&self) -> Result<()> {
1360 self.media_stream
1361 .resume_playback(self.server_side_track_id.clone())
1362 .await?;
1363 Ok(())
1364 }
1365 async fn do_hangup(
1366 &self,
1367 reason: Option<CallRecordHangupReason>,
1368 initiator: Option<String>,
1369 headers: Option<HashMap<String, String>>,
1370 refer: Option<bool>,
1371 ) -> Result<()> {
1372 info!(
1373 session_id = self.session_id,
1374 ?reason,
1375 ?initiator,
1376 ?headers,
1377 ?refer,
1378 "do_hangup"
1379 );
1380
1381 let hangup_reason = match initiator.as_deref() {
1382 Some("caller") => CallRecordHangupReason::ByCaller,
1383 Some("callee") => CallRecordHangupReason::ByCallee,
1384 Some("system") => CallRecordHangupReason::Autohangup,
1385 _ => reason.unwrap_or(CallRecordHangupReason::BySystem),
1386 };
1387
1388 match refer {
1389 Some(true) => {
1390 let (refer_state, refer_token) = {
1392 let mut state = self.call_state.write().await;
1393 (state.refer_callstate.clone(), state.refer_call_token.take())
1394 };
1395 let mut has_refer_state = false;
1396 if let Some(refer_state) = refer_state {
1397 has_refer_state = true;
1398 let mut refer_state = refer_state.write().await;
1399 if let Some(headers) = headers {
1400 let h_val = serde_json::to_value(&headers).unwrap_or_default();
1401 let mut extras = refer_state.extras.take().unwrap_or_default();
1402 extras.insert("_hangup_headers".to_string(), h_val);
1403 refer_state.extras = Some(extras);
1404 }
1405 refer_state.set_hangup_reason(hangup_reason);
1407 }
1408 if let Some(token) = refer_token {
1409 token.cancel();
1410 }
1411 if has_refer_state {
1412 self.media_stream
1413 .remove_track(&self.server_side_track_id, false)
1414 .await;
1415 }
1416 }
1417 _ => {
1418 let refer_token = {
1419 let mut state = self.call_state.write().await;
1420 if let Some(headers) = headers {
1421 let h_val = serde_json::to_value(&headers).unwrap_or_default();
1422 let mut extras = state.extras.take().unwrap_or_default();
1423 extras.insert("_hangup_headers".to_string(), h_val);
1424 state.extras = Some(extras);
1425 }
1426 state.set_hangup_reason(hangup_reason.clone());
1427 state.refer_call_token.take()
1428 };
1429 self.media_stream
1430 .stop(Some(hangup_reason.to_string()), initiator);
1431 if let Some(token) = refer_token {
1432 token.cancel();
1433 }
1434 }
1435 }
1436 tokio::task::yield_now().await;
1437 Ok(())
1438 }
1439
1440 async fn do_refer(
1441 &self,
1442 caller: String,
1443 callee: String,
1444 refer_option: Option<ReferOption>,
1445 ) -> Result<()> {
1446 self.do_interrupt(false).await.ok();
1447
1448 let pause_parent_asr = refer_option
1450 .as_ref()
1451 .and_then(|o| o.pause_parent_asr)
1452 .unwrap_or(false);
1453
1454 let original_asr_option = if pause_parent_asr {
1456 let cs = self.call_state.read().await;
1457 cs.option.as_ref().and_then(|o| o.asr.clone())
1458 } else {
1459 None
1460 };
1461
1462 if pause_parent_asr {
1464 info!(
1465 session_id = self.session_id,
1466 "Pausing parent call ASR during refer"
1467 );
1468 self.media_stream
1469 .remove_processor::<crate::media::asr_processor::AsrProcessor>(
1470 &self.server_side_track_id,
1471 )
1472 .await
1473 .ok();
1474 }
1475
1476 let mut moh = refer_option.as_ref().and_then(|o| o.moh.clone());
1477 if let Some(ref path) = moh {
1478 if !path.starts_with("http") && !std::path::Path::new(path).exists() {
1479 let fallback = "./config/sounds/refer_moh.wav";
1480 if std::path::Path::new(fallback).exists() {
1481 info!(
1482 session_id = self.session_id,
1483 "moh {} not found, using fallback {}", path, fallback
1484 );
1485 moh = Some(fallback.to_string());
1486 }
1487 }
1488 }
1489 let ref_call_id = refer_option
1490 .as_ref()
1491 .and_then(|o| o.call_id.clone())
1492 .unwrap_or_else(|| format!("ref-{}-{}", rand::random::<u32>(), self.session_id));
1493
1494 let session_id = self.session_id.clone();
1495 let track_id = self.server_side_track_id.clone();
1496
1497 let (recorder, parent_caller) = {
1498 let cs = self.call_state.read().await;
1499 let option = cs.option.as_ref();
1500 (
1501 option.map(|o| o.recorder.clone()).unwrap_or_default(),
1502 option.and_then(|o| o.caller.clone()),
1503 )
1504 };
1505 let caller = if caller.trim().is_empty() {
1506 parent_caller.unwrap_or_default()
1507 } else {
1508 caller
1509 };
1510
1511 let mut call_option = CallOption {
1512 caller: Some(caller),
1513 callee: Some(callee.clone()),
1514 sip: refer_option.as_ref().and_then(|o| o.sip.clone()),
1515 vad: refer_option
1516 .as_ref()
1517 .and_then(|o| o.vad.clone())
1518 .map(|mut opts| {
1519 opts.refer = Some(true);
1520 opts
1521 }),
1522 asr: refer_option
1523 .as_ref()
1524 .and_then(|o| o.asr.clone())
1525 .map(|mut opts| {
1526 opts.refer = Some(true);
1527 opts
1528 }),
1529 denoise: refer_option.as_ref().and_then(|o| o.denoise.clone()),
1530 recorder,
1531 ..Default::default()
1532 };
1533 call_option.check_default();
1534
1535 let mut invite_option = call_option.build_invite_option()?;
1536 invite_option.call_id = Some(ref_call_id);
1537
1538 let headers = invite_option.headers.get_or_insert_with(|| Vec::new());
1539
1540 {
1541 let cs = self.call_state.read().await;
1542 if let Some(opt) = cs.option.as_ref() {
1543 if let Some(callee) = opt.callee.as_ref() {
1544 headers.push(rsipstack::rsip::Header::Other(
1545 "X-Referred-To".to_string(),
1546 callee.clone(),
1547 ));
1548 }
1549 if let Some(caller) = opt.caller.as_ref() {
1550 headers.push(rsipstack::rsip::Header::Other(
1551 "X-Referred-From".to_string(),
1552 caller.clone(),
1553 ));
1554 }
1555 }
1556 }
1557
1558 headers.push(rsipstack::rsip::Header::Other(
1559 "X-Referred-Id".to_string(),
1560 self.session_id.clone(),
1561 ));
1562
1563 let ssrc = rand::random::<u32>();
1564 let refer_call_state = Arc::new(RwLock::new(ActiveCallState {
1565 start_time: Utc::now(),
1566 ssrc,
1567 option: Some(call_option.clone()),
1568 is_refer: true,
1569 ..Default::default()
1570 }));
1571
1572 {
1573 let mut cs = self.call_state.write().await;
1574 cs.refer_callstate.replace(refer_call_state.clone());
1575 }
1576
1577 let auto_hangup_requested = refer_option
1578 .as_ref()
1579 .and_then(|o| o.auto_hangup)
1580 .unwrap_or(true);
1581
1582 if auto_hangup_requested {
1583 self.call_state.write().await.auto_hangup =
1584 Some((ssrc, CallRecordHangupReason::ByRefer));
1585 } else {
1586 self.call_state.write().await.auto_hangup = None;
1587 }
1588
1589 if !auto_hangup_requested && pause_parent_asr && original_asr_option.is_some() {
1591 let asr_option = original_asr_option.unwrap();
1592 self.call_state.write().await.pending_asr_resume = Some((ssrc, asr_option));
1593 }
1594
1595 let timeout_secs = refer_option.as_ref().and_then(|o| o.timeout).unwrap_or(30);
1596
1597 info!(
1598 session_id = self.session_id,
1599 ssrc,
1600 auto_hangup = auto_hangup_requested,
1601 callee,
1602 timeout_secs,
1603 "do_refer"
1604 );
1605
1606 let refer_cancel_token = self.cancel_token.child_token();
1607 self.call_state.write().await.refer_call_token = Some(refer_cancel_token.clone());
1608
1609 let r = tokio::time::timeout(
1610 Duration::from_secs(timeout_secs as u64),
1611 self.create_outgoing_sip_track(
1612 refer_cancel_token,
1613 refer_call_state.clone(),
1614 &track_id,
1615 invite_option,
1616 &call_option,
1617 moh,
1618 auto_hangup_requested,
1619 ),
1620 )
1621 .await;
1622
1623 {
1624 self.call_state.write().await.moh = None;
1625 }
1626
1627 let result = match r {
1628 Ok(res) => res,
1629 Err(_) => {
1630 warn!(
1631 session_id = session_id,
1632 "refer sip track creation timed out after {} seconds", timeout_secs
1633 );
1634 self.event_sender
1635 .send(SessionEvent::Reject {
1636 track_id,
1637 timestamp: crate::media::get_timestamp(),
1638 reason: "Timeout when refer".into(),
1639 code: Some(408),
1640 refer: Some(true),
1641 })
1642 .ok();
1643 return Err(anyhow::anyhow!("refer sip track creation timed out").into());
1644 }
1645 };
1646
1647 match result {
1648 Ok(answer) => {
1649 self.media_stream.set_track_refer(&track_id, Some(true)).await;
1650 let forward_dtmf = refer_option.as_ref().and_then(|o| o.forward_dtmf).unwrap_or(true);
1651 if !forward_dtmf {
1652 self.media_stream.set_track_dtmf_forward(&track_id, false).await;
1653 }
1654 self.event_sender
1655 .send(SessionEvent::Answer {
1656 timestamp: crate::media::get_timestamp(),
1657 track_id,
1658 sdp: answer,
1659 refer: Some(true),
1660 })
1661 .ok();
1662 }
1663 Err(e) => {
1664 warn!(
1665 session_id = session_id,
1666 "failed to create refer sip track: {}", e
1667 );
1668 match &e {
1669 rsipstack::Error::DialogError(reason, _, code) => {
1670 self.event_sender
1671 .send(SessionEvent::Reject {
1672 track_id,
1673 timestamp: crate::media::get_timestamp(),
1674 reason: reason.clone(),
1675 code: Some(code.code() as u32),
1676 refer: Some(true),
1677 })
1678 .ok();
1679 }
1680 _ => {}
1681 }
1682 return Err(e.into());
1683 }
1684 }
1685 Ok(())
1686 }
1687
1688 fn bridge_track_id(source_session_id: &str, target_session_id: &str) -> TrackId {
1689 format!("bridge:{}:to:{}", source_session_id, target_session_id)
1690 }
1691
1692 async fn do_bridge(&self, target_session_id: String) -> Result<()> {
1693 let target = {
1694 let calls = self.app_state.active_calls.lock().unwrap();
1695 calls.get(&target_session_id).cloned()
1696 };
1697 let target = target.ok_or_else(|| {
1698 anyhow::anyhow!("bridge target session not found: {}", target_session_id)
1699 })?;
1700
1701 if target.session_id == self.session_id {
1702 return Err(anyhow::anyhow!("cannot bridge a call to itself").into());
1703 }
1704
1705 let self_bridge_track_id = Self::bridge_track_id(&self.session_id, &target.session_id);
1706 let target_bridge_track_id = Self::bridge_track_id(&target.session_id, &self.session_id);
1707
1708 self.media_stream
1709 .remove_track(&self_bridge_track_id, false)
1710 .await;
1711 target
1712 .media_stream
1713 .remove_track(&target_bridge_track_id, false)
1714 .await;
1715
1716 let (self_bridge_sender, self_bridge_receiver) = mpsc::channel(25);
1717 let (target_bridge_sender, target_bridge_receiver) = mpsc::channel(25);
1718
1719 let self_paused = self.call_state.read().await.bridge_paused.clone();
1720 let target_paused = target.call_state.read().await.bridge_paused.clone();
1721
1722 let self_forwarding_track = ForwardingTrack::new(
1723 self_bridge_track_id.clone(),
1724 self.session_id.clone(),
1725 target_bridge_sender,
1726 self_bridge_receiver,
1727 self.track_config.clone(),
1728 self.cancel_token.child_token(),
1729 rand::random::<u32>(),
1730 self_paused,
1731 );
1732
1733 let target_forwarding_track = ForwardingTrack::new(
1734 target_bridge_track_id.clone(),
1735 target.session_id.clone(),
1736 self_bridge_sender,
1737 target_bridge_receiver,
1738 target.track_config.clone(),
1739 target.cancel_token.child_token(),
1740 rand::random::<u32>(),
1741 target_paused,
1742 );
1743
1744 self.media_stream
1745 .update_track(Box::new(self_forwarding_track), None)
1746 .await;
1747 target
1748 .media_stream
1749 .update_track(Box::new(target_forwarding_track), None)
1750 .await;
1751
1752 info!(
1753 session_id = self.session_id,
1754 target = target_session_id,
1755 self_bridge_track_id,
1756 target_bridge_track_id,
1757 "audio bridge established"
1758 );
1759 Ok(())
1760 }
1761
1762 async fn do_unbridge(&self, target_session_id: String) -> Result<()> {
1763 let target = {
1764 let calls = self.app_state.active_calls.lock().unwrap();
1765 calls.get(&target_session_id).cloned()
1766 };
1767
1768 let self_bridge_track_id = Self::bridge_track_id(&self.session_id, &target_session_id);
1769 self.media_stream
1770 .remove_track(&self_bridge_track_id, false)
1771 .await;
1772
1773 if let Some(target) = target {
1774 let target_bridge_track_id =
1775 Self::bridge_track_id(&target.session_id, &self.session_id);
1776 target
1777 .media_stream
1778 .remove_track(&target_bridge_track_id, false)
1779 .await;
1780 info!(
1781 session_id = self.session_id,
1782 target = target.session_id,
1783 self_bridge_track_id,
1784 target_bridge_track_id,
1785 "audio bridge removed"
1786 );
1787 } else {
1788 info!(
1789 session_id = self.session_id,
1790 target = target_session_id,
1791 self_bridge_track_id,
1792 "audio bridge removed locally; target session not active"
1793 );
1794 }
1795
1796 Ok(())
1797 }
1798
1799 async fn do_mute(&self, track_id: Option<String>) -> Result<()> {
1800 self.media_stream.mute_track(track_id).await;
1801 Ok(())
1802 }
1803
1804 async fn do_unmute(&self, track_id: Option<String>) -> Result<()> {
1805 self.media_stream.unmute_track(track_id).await;
1806 Ok(())
1807 }
1808
1809 pub async fn cleanup(&self) -> Result<()> {
1810 if matches!(self.call_type, ActiveCallType::Sip | ActiveCallType::B2bua) {
1811 self.do_reject(
1812 Some(rsipstack::rsip::StatusCode::Decline),
1813 Some("handler disconnected".to_string()),
1814 )
1815 .await
1816 .ok();
1817 }
1818 self.call_state.write().await.tts_handle = None;
1819 self.media_stream.cleanup().await.ok();
1820 Ok(())
1821 }
1822
1823 pub fn get_callrecord(&self) -> Option<CallRecord> {
1824 self.call_state.try_read().ok().map(|call_state| {
1825 call_state.build_callrecord(
1826 self.app_state.clone(),
1827 self.session_id.clone(),
1828 self.call_type.clone(),
1829 )
1830 })
1831 }
1832
1833 async fn dump_to_file(
1834 &self,
1835 dump_file: &mut File,
1836 cmd_receiver: &mut CommandReceiver,
1837 event_receiver: &mut EventReceiver,
1838 ) {
1839 loop {
1840 select! {
1841 _ = self.cancel_token.cancelled() => {
1842 break;
1843 }
1844 Ok(cmd) = cmd_receiver.recv() => {
1845 CallRecordEvent::write(CallRecordEventType::Command, cmd, dump_file)
1846 .await;
1847 }
1848 Ok(event) = event_receiver.recv() => {
1849 if matches!(event, SessionEvent::Binary{..}) {
1850 continue;
1851 }
1852 CallRecordEvent::write(CallRecordEventType::Event, event, dump_file)
1853 .await;
1854 }
1855 };
1856 }
1857 }
1858
1859 async fn dump_loop(
1860 &self,
1861 dump_events: bool,
1862 mut dump_cmd_receiver: CommandReceiver,
1863 mut dump_event_receiver: EventReceiver,
1864 ) {
1865 if !dump_events {
1866 return;
1867 }
1868
1869 let file_name = self.app_state.get_dump_events_file(&self.session_id);
1870 let mut dump_file = match File::options()
1871 .create(true)
1872 .append(true)
1873 .open(&file_name)
1874 .await
1875 {
1876 Ok(file) => file,
1877 Err(e) => {
1878 warn!(
1879 session_id = self.session_id,
1880 file_name, "failed to open dump events file: {}", e
1881 );
1882 return;
1883 }
1884 };
1885 self.dump_to_file(
1886 &mut dump_file,
1887 &mut dump_cmd_receiver,
1888 &mut dump_event_receiver,
1889 )
1890 .await;
1891
1892 while let Ok(event) = dump_event_receiver.try_recv() {
1893 if matches!(event, SessionEvent::Binary { .. }) {
1894 continue;
1895 }
1896 CallRecordEvent::write(CallRecordEventType::Event, event, &mut dump_file).await;
1897 }
1898 }
1899
1900 pub async fn create_rtp_track(
1901 &self,
1902 track_id: TrackId,
1903 ssrc: u32,
1904 enable_srtp: Option<bool>,
1905 ) -> Result<RtcTrack> {
1906 let mut rtc_config = RtcTrackConfig::default();
1907 let use_srtp = enable_srtp
1909 .or(self.app_state.config.enable_srtp)
1910 .unwrap_or(false);
1911 rtc_config.mode = if use_srtp {
1912 rustrtc::TransportMode::Srtp
1913 } else {
1914 rustrtc::TransportMode::Rtp
1915 };
1916
1917 if let Some(codecs) = &self.app_state.config.codecs {
1918 let mut codec_types = Vec::new();
1919 for c in codecs {
1920 match c.to_lowercase().as_str() {
1921 "pcmu" => codec_types.push(CodecType::PCMU),
1922 "pcma" => codec_types.push(CodecType::PCMA),
1923 "g722" => codec_types.push(CodecType::G722),
1924 "g729" => codec_types.push(CodecType::G729),
1925 #[cfg(feature = "opus")]
1926 "opus" => codec_types.push(CodecType::Opus),
1927 "dtmf" | "2833" | "telephone_event" => {
1928 codec_types.push(CodecType::TelephoneEvent)
1929 }
1930 _ => {}
1931 }
1932 }
1933 if !codec_types.is_empty() {
1934 rtc_config.preferred_codec = Some(codec_types[0].clone());
1935 rtc_config.codecs = codec_types;
1936 }
1937 }
1938
1939 if rtc_config.preferred_codec.is_none() {
1940 rtc_config.preferred_codec = Some(self.track_config.codec.clone());
1941 }
1942
1943 rtc_config.rtp_port_range = self
1944 .app_state
1945 .config
1946 .rtp_start_port
1947 .zip(self.app_state.config.rtp_end_port);
1948
1949 if let Some(ref external_ip) = self.app_state.config.external_ip {
1950 rtc_config.external_ip = Some(external_ip.clone());
1951 }
1952 if let Some(ref bind_ip) = self.app_state.config.rtp_bind_ip {
1953 rtc_config.bind_ip = Some(bind_ip.clone());
1954 }
1955
1956 rtc_config.enable_latching = self.app_state.config.enable_rtp_latching;
1957 rtc_config.enable_ice_lite = self
1958 .call_state
1959 .read()
1960 .await
1961 .option
1962 .as_ref()
1963 .and_then(|o| o.enable_ice_lite)
1964 .or(self.app_state.config.enable_ice_lite);
1965
1966 let mut track = RtcTrack::new(
1967 self.cancel_token.child_token(),
1968 track_id,
1969 self.track_config.clone(),
1970 rtc_config,
1971 )
1972 .with_ssrc(ssrc);
1973
1974 track.create().await?;
1975
1976 Ok(track)
1977 }
1978
1979 async fn setup_caller_track(&self, option: &CallOption) -> Result<()> {
1980 let hangup_headers = option
1981 .sip
1982 .as_ref()
1983 .and_then(|s| s.hangup_headers.as_ref())
1984 .map(|headers_map| {
1985 headers_map
1986 .iter()
1987 .map(|(k, v)| rsipstack::rsip::Header::Other(k.clone(), v.clone()))
1988 .collect::<Vec<rsipstack::rsip::Header>>()
1989 });
1990 self.call_state.write().await.option = Some(option.clone());
1991 info!(
1992 session_id = self.session_id,
1993 call_type = ?self.call_type,
1994 "setup caller track"
1995 );
1996
1997 let track = match self.call_type {
1998 ActiveCallType::Webrtc => Some(self.create_webrtc_track().await?),
1999 ActiveCallType::WebSocket => {
2000 let audio_receiver = self.call_state.write().await.audio_receiver.take();
2001 if let Some(receiver) = audio_receiver {
2002 Some(self.create_websocket_track(receiver).await?)
2003 } else {
2004 None
2005 }
2006 }
2007 ActiveCallType::Sip => {
2008 if let Some(dialog_id) = self
2009 .invitation
2010 .find_dialog_id_by_session_id(&self.session_id)
2011 {
2012 if let Some(pending_dialog) = self.invitation.get_pending_call(&dialog_id) {
2013 return self
2014 .prepare_incoming_sip_track(
2015 self.cancel_token.clone(),
2016 self.call_state.clone(),
2017 &self.session_id,
2018 pending_dialog,
2019 hangup_headers,
2020 )
2021 .await;
2022 }
2023 }
2024
2025 let mut option = option.clone();
2027 if option.sip.is_none()
2028 || option
2029 .sip
2030 .as_ref()
2031 .and_then(|s| s.username.as_ref())
2032 .is_none()
2033 {
2034 if let Some(callee) = &option.callee {
2035 if let Some(cred) = self.app_state.find_credentials_for_callee(callee) {
2036 if option.sip.is_none() {
2037 option.sip = Some(crate::SipOption {
2038 username: Some(cred.username.clone()),
2039 password: Some(cred.password.clone()),
2040 realm: cred.realm.clone(),
2041 ..Default::default()
2042 });
2043 }
2044 }
2045 }
2046 }
2047
2048 let mut invite_option = option.build_invite_option()?;
2049 invite_option.call_id = Some(self.session_id.clone());
2050
2051 match self
2052 .create_outgoing_sip_track(
2053 self.cancel_token.clone(),
2054 self.call_state.clone(),
2055 &self.session_id,
2056 invite_option,
2057 &option,
2058 None,
2059 false,
2060 )
2061 .await
2062 {
2063 Ok(answer) => {
2064 self.event_sender
2065 .send(SessionEvent::Answer {
2066 timestamp: crate::media::get_timestamp(),
2067 track_id: self.session_id.clone(),
2068 sdp: answer,
2069 refer: Some(false),
2070 })
2071 .ok();
2072 return Ok(());
2073 }
2074 Err(e) => {
2075 warn!(
2076 session_id = self.session_id,
2077 "failed to create sip track: {}", e
2078 );
2079 match &e {
2080 rsipstack::Error::DialogError(reason, _, code) => {
2081 self.event_sender
2082 .send(SessionEvent::Reject {
2083 track_id: self.session_id.clone(),
2084 timestamp: crate::media::get_timestamp(),
2085 reason: reason.clone(),
2086 code: Some(code.code() as u32),
2087 refer: Some(false),
2088 })
2089 .ok();
2090 }
2091 _ => {}
2092 }
2093 return Err(e.into());
2094 }
2095 }
2096 }
2097 ActiveCallType::B2bua => {
2098 if let Some(dialog_id) = self
2099 .invitation
2100 .find_dialog_id_by_session_id(&self.session_id)
2101 {
2102 if let Some(pending_dialog) = self.invitation.get_pending_call(&dialog_id) {
2103 return self
2104 .prepare_incoming_sip_track(
2105 self.cancel_token.clone(),
2106 self.call_state.clone(),
2107 &self.session_id,
2108 pending_dialog,
2109 hangup_headers,
2110 )
2111 .await;
2112 }
2113 }
2114
2115 warn!(
2116 session_id = self.session_id,
2117 "no pending dialog found for B2BUA call"
2118 );
2119 return Err(anyhow::anyhow!(
2120 "no pending dialog found for session_id: {}",
2121 self.session_id
2122 ));
2123 }
2124 };
2125 match track {
2126 Some(track) => {
2127 self.finish_caller_stack(&option, Some(track)).await?;
2128 }
2129 None => {
2130 warn!(session_id = self.session_id, "no track created for caller");
2131 return Err(anyhow::anyhow!("no track created for caller"));
2132 }
2133 }
2134 Ok(())
2135 }
2136
2137 async fn finish_caller_stack(
2138 &self,
2139 option: &CallOption,
2140 track: Option<Box<dyn Track>>,
2141 ) -> Result<()> {
2142 if let Some(track) = track {
2143 self.setup_track_with_stream(&option, track).await?;
2144 }
2145
2146 {
2147 let call_state = self.call_state.read().await;
2148 if let Some(ref answer) = call_state.answer {
2149 info!(
2150 session_id = self.session_id,
2151 "sending answer event: {}", answer,
2152 );
2153 self.event_sender
2154 .send(SessionEvent::Answer {
2155 timestamp: crate::media::get_timestamp(),
2156 track_id: self.session_id.clone(),
2157 sdp: answer.clone(),
2158 refer: Some(false),
2159 })
2160 .ok();
2161 } else {
2162 warn!(
2163 session_id = self.session_id,
2164 "no answer in state to send event"
2165 );
2166 }
2167 }
2168 Ok(())
2169 }
2170
2171 pub async fn setup_track_with_stream(
2172 &self,
2173 option: &CallOption,
2174 mut track: Box<dyn Track>,
2175 ) -> Result<()> {
2176 let processors = match StreamEngine::create_processors(
2177 self.app_state.stream_engine.clone(),
2178 track.as_ref(),
2179 self.cancel_token.child_token(),
2180 self.event_sender.clone(),
2181 self.media_stream.packet_sender.clone(),
2182 option,
2183 )
2184 .await
2185 {
2186 Ok(processors) => processors,
2187 Err(e) => {
2188 warn!(
2189 session_id = self.session_id,
2190 "failed to prepare stream processors: {}", e
2191 );
2192 vec![]
2193 }
2194 };
2195
2196 for processor in processors {
2198 track.append_processor(processor);
2199 }
2200
2201 self.update_track_wrapper(track, None).await;
2202 Ok(())
2203 }
2204
2205 pub async fn update_track_wrapper(&self, mut track: Box<dyn Track>, play_id: Option<String>) {
2206 let (ambiance_opt, subscribe) = {
2207 let state = self.call_state.read().await;
2208 let mut opt = state
2209 .option
2210 .as_ref()
2211 .and_then(|o| o.ambiance.clone())
2212 .unwrap_or_default();
2213
2214 if let Some(global) = &self.app_state.config.ambiance {
2215 opt.merge(global);
2216 }
2217
2218 let subscribe = state
2219 .option
2220 .as_ref()
2221 .and_then(|o| o.subscribe)
2222 .unwrap_or_default();
2223
2224 (opt, subscribe)
2225 };
2226 if track.id() == &self.server_side_track_id && ambiance_opt.path.is_some() {
2227 match AmbianceProcessor::new(ambiance_opt).await {
2228 Ok(ambiance) => {
2229 info!(session_id = self.session_id, "loaded ambiance processor");
2230 track.append_processor(Box::new(ambiance));
2231 }
2232 Err(e) => {
2233 tracing::error!("failed to load ambiance wav {}", e);
2234 }
2235 }
2236 }
2237
2238 if subscribe && self.call_type != ActiveCallType::WebSocket {
2239 let (track_index, sub_track_id) = if track.id() == &self.server_side_track_id {
2240 (0, self.server_side_track_id.clone())
2241 } else {
2242 (1, self.session_id.clone())
2243 };
2244 let sub_processor =
2245 SubscribeProcessor::new(self.event_sender.clone(), sub_track_id, track_index);
2246 track.append_processor(Box::new(sub_processor));
2247 }
2248
2249 self.call_state.write().await.current_play_id = play_id.clone();
2250 self.media_stream.update_track(track, play_id).await;
2251 }
2252
2253 pub async fn create_websocket_track(
2254 &self,
2255 audio_receiver: WebsocketBytesReceiver,
2256 ) -> Result<Box<dyn Track>> {
2257 let (ssrc, codec) = {
2258 let call_state = self.call_state.read().await;
2259 (
2260 call_state.ssrc,
2261 call_state
2262 .option
2263 .as_ref()
2264 .map(|o| o.codec.clone())
2265 .unwrap_or_default(),
2266 )
2267 };
2268
2269 let ws_track = WebsocketTrack::new(
2270 self.cancel_token.child_token(),
2271 self.session_id.clone(),
2272 self.track_config.clone(),
2273 self.event_sender.clone(),
2274 audio_receiver,
2275 codec,
2276 ssrc,
2277 );
2278
2279 {
2280 let mut call_state = self.call_state.write().await;
2281 call_state.answer_time = Some(Utc::now());
2282 call_state.answer = Some("".to_string());
2283 call_state.last_status_code = 200;
2284 }
2285
2286 Ok(Box::new(ws_track))
2287 }
2288
2289 pub(super) async fn create_webrtc_track(&self) -> Result<Box<dyn Track>> {
2290 let (ssrc, option) = {
2291 let call_state = self.call_state.read().await;
2292 (
2293 call_state.ssrc,
2294 call_state.option.clone().unwrap_or_default(),
2295 )
2296 };
2297
2298 let mut rtc_config = RtcTrackConfig::default();
2299 rtc_config.mode = rustrtc::TransportMode::WebRtc; rtc_config.ice_servers = self.app_state.config.ice_servers.clone();
2301
2302 if let Some(codecs) = &self.app_state.config.codecs {
2303 let mut codec_types = Vec::new();
2304 for c in codecs {
2305 match c.to_lowercase().as_str() {
2306 "pcmu" => codec_types.push(CodecType::PCMU),
2307 "pcma" => codec_types.push(CodecType::PCMA),
2308 "g722" => codec_types.push(CodecType::G722),
2309 "g729" => codec_types.push(CodecType::G729),
2310 #[cfg(feature = "opus")]
2311 "opus" => codec_types.push(CodecType::Opus),
2312 "dtmf" | "2833" | "telephone_event" => {
2313 codec_types.push(CodecType::TelephoneEvent)
2314 }
2315 _ => {}
2316 }
2317 }
2318 if !codec_types.is_empty() {
2319 rtc_config.preferred_codec = Some(codec_types[0].clone());
2320 rtc_config.codecs = codec_types;
2321 }
2322 }
2323
2324 if let Some(ref external_ip) = self.app_state.config.external_ip {
2325 rtc_config.external_ip = Some(external_ip.clone());
2326 }
2327 if let Some(ref bind_ip) = self.app_state.config.rtp_bind_ip {
2328 rtc_config.bind_ip = Some(bind_ip.clone());
2329 }
2330
2331 let mut webrtc_track = RtcTrack::new(
2332 self.cancel_token.child_token(),
2333 self.session_id.clone(),
2334 self.track_config.clone(),
2335 rtc_config,
2336 )
2337 .with_ssrc(ssrc);
2338
2339 let timeout = option.handshake_timeout.map(|t| Duration::from_secs(t));
2340 let offer = match option.enable_ipv6 {
2341 Some(false) | None => {
2342 strip_ipv6_candidates(option.offer.as_ref().unwrap_or(&"".to_string()))
2343 }
2344 _ => option.offer.clone().unwrap_or("".to_string()),
2345 };
2346 let answer: Option<String>;
2347 match webrtc_track.handshake(offer, timeout).await {
2348 Ok(answer_sdp) => {
2349 answer = match option.enable_ipv6 {
2350 Some(false) | None => Some(strip_ipv6_candidates(&answer_sdp)),
2351 Some(true) => Some(answer_sdp.to_string()),
2352 };
2353 }
2354 Err(e) => {
2355 warn!(session_id = self.session_id, "failed to setup track: {}", e);
2356 return Err(anyhow::anyhow!("Failed to setup track: {}", e));
2357 }
2358 }
2359
2360 {
2361 let mut call_state = self.call_state.write().await;
2362 call_state.answer_time = Some(Utc::now());
2363 call_state.answer = answer;
2364 call_state.last_status_code = 200;
2365 }
2366 Ok(Box::new(webrtc_track))
2367 }
2368
2369 async fn create_outgoing_sip_track(
2370 &self,
2371 cancel_token: CancellationToken,
2372 call_state_ref: ActiveCallStateRef,
2373 track_id: &String,
2374 mut invite_option: InviteOption,
2375 call_option: &CallOption,
2376 moh: Option<String>,
2377 auto_hangup: bool,
2378 ) -> Result<String, rsipstack::Error> {
2379 let ssrc = call_state_ref.read().await.ssrc;
2380 let per_call_srtp = call_option.sip.as_ref().and_then(|s| s.enable_srtp);
2381 let rtp_track = self
2382 .create_rtp_track(track_id.clone(), ssrc, per_call_srtp)
2383 .await
2384 .map_err(|e| rsipstack::Error::Error(e.to_string()))?;
2385
2386 let offer = Some(
2387 rtp_track
2388 .local_description()
2389 .await
2390 .map_err(|e| rsipstack::Error::Error(e.to_string()))?,
2391 );
2392
2393 {
2394 let mut cs = call_state_ref.write().await;
2395 if let Some(o) = cs.option.as_mut() {
2396 o.offer = offer.clone();
2397 }
2398 cs.start_time = Utc::now();
2399 };
2400
2401 invite_option.offer = offer.clone().map(|s| s.into());
2402
2403 let needs_contact = contact_needs_public_resolution(&invite_option.contact);
2406
2407 if needs_contact {
2408 let addrs = self.invitation.dialog_layer.endpoint.get_addrs();
2409 if let Some(addr) = find_local_addr_for_uri(&addrs, &invite_option.callee) {
2410 let contact_username = invite_option
2411 .contact
2412 .auth
2413 .as_ref()
2414 .map(|auth| auth.user.as_str())
2415 .or_else(|| {
2416 invite_option
2417 .caller
2418 .auth
2419 .as_ref()
2420 .map(|auth| auth.user.as_str())
2421 });
2422 invite_option.contact = build_public_contact_uri(
2423 &self.app_state.learned_public_address,
2424 self.app_state.auto_learn_public_address_enabled(),
2425 &addr,
2426 contact_username,
2427 Some(&invite_option.contact),
2428 );
2429 } else {
2430 return Err(rsipstack::Error::Error(format!(
2431 "missing local SIP address for callee transport: {}",
2432 invite_option.callee
2433 )));
2434 }
2435 }
2436
2437 let mut rtp_track_to_setup = Some(Box::new(rtp_track) as Box<dyn Track>);
2438
2439 if let Some(moh) = moh {
2440 let ssrc_and_moh = {
2441 let mut state = call_state_ref.write().await;
2442 state.moh = Some(moh.clone());
2443 if state.current_play_id.is_none() {
2444 let ssrc = rand::random::<u32>();
2445 Some((ssrc, moh.clone()))
2446 } else {
2447 info!(
2448 session_id = self.session_id,
2449 "Something is playing, MOH will start after it ends"
2450 );
2451 None
2452 }
2453 };
2454
2455 if let Some((ssrc, moh_path)) = ssrc_and_moh {
2456 let file_track = FileTrack::new(self.server_side_track_id.clone())
2457 .with_play_id(Some(moh_path.clone()))
2458 .with_ssrc(ssrc)
2459 .with_path(moh_path.clone())
2460 .with_cancel_token(self.cancel_token.child_token());
2461 self.update_track_wrapper(Box::new(file_track), Some(moh_path))
2462 .await;
2463 }
2464 } else {
2465 let track = rtp_track_to_setup.take().unwrap();
2466 self.setup_track_with_stream(&call_option, track)
2467 .await
2468 .map_err(|e| rsipstack::Error::Error(e.to_string()))?;
2469 }
2470
2471 info!(
2472 session_id = self.session_id,
2473 track_id,
2474 contact = %invite_option.contact,
2475 "invite {} -> {} offer: \n{}",
2476 invite_option.caller,
2477 invite_option.callee,
2478 offer.as_ref().map(|s| s.as_str()).unwrap_or("<NO OFFER>")
2479 );
2480
2481 let (dlg_state_sender, dlg_state_receiver) =
2482 self.invitation.dialog_layer.new_dialog_state_channel();
2483
2484 let states = InviteDialogStates {
2485 is_client: true,
2486 session_id: self.session_id.clone(),
2487 track_id: track_id.clone(),
2488 event_sender: self.event_sender.clone(),
2489 media_stream: self.media_stream.clone(),
2490 call_state: call_state_ref.clone(),
2491 cancel_token,
2492 terminated_reason: None,
2493 has_early_media: false,
2494 };
2495
2496 let hangup_headers = call_option
2497 .sip
2498 .as_ref()
2499 .and_then(|s| s.hangup_headers.as_ref())
2500 .map(|headers_map| {
2501 headers_map
2502 .iter()
2503 .map(|(k, v)| rsipstack::rsip::Header::Other(k.clone(), v.clone()))
2504 .collect::<Vec<rsipstack::rsip::Header>>()
2505 });
2506
2507 let mut client_dialog_handler = DialogStateReceiverGuard::new(
2508 self.invitation.dialog_layer.clone(),
2509 dlg_state_receiver,
2510 hangup_headers,
2511 );
2512
2513 crate::spawn(async move {
2514 client_dialog_handler.process_dialog(states).await;
2515 });
2516
2517 let (dialog_id, answer) = self
2518 .invitation
2519 .invite(invite_option, dlg_state_sender)
2520 .await?;
2521
2522 self.call_state.write().await.moh = None;
2523
2524 if let Some(track) = rtp_track_to_setup {
2525 info!(
2526 session_id = self.session_id,
2527 track_id, "Stopping MOH and setting up RTP track"
2528 );
2529 self.media_stream
2530 .remove_track(&self.server_side_track_id, false)
2531 .await;
2532
2533 self.setup_track_with_stream(&call_option, track)
2534 .await
2535 .map_err(|e| rsipstack::Error::Error(e.to_string()))?;
2536 }
2537
2538 let answer = match answer {
2539 Some(answer) => {
2540 let s = String::from_utf8_lossy(&answer).to_string();
2541 if s.trim().is_empty() {
2542 let cs = call_state_ref.read().await;
2546 match cs.answer.clone() {
2547 Some(early_sdp) if !early_sdp.is_empty() => {
2548 info!(
2549 session_id = self.session_id,
2550 "200 OK has empty body; using early-media SDP from 183"
2551 );
2552 (early_sdp, true )
2553 }
2554 _ => {
2555 warn!(
2556 session_id = self.session_id,
2557 "200 OK has empty body and no early-media SDP available"
2558 );
2559 (s, false)
2560 }
2561 }
2562 } else {
2563 (s, false)
2564 }
2565 }
2566 None => {
2567 let cs = call_state_ref.read().await;
2569 match cs.answer.clone() {
2570 Some(early_sdp) if !early_sdp.is_empty() => {
2571 info!(
2572 session_id = self.session_id,
2573 "200 OK had no answer; using early-media SDP from 183"
2574 );
2575 (early_sdp, true )
2576 }
2577 _ => {
2578 warn!(session_id = self.session_id, "no answer received");
2579 return Err(rsipstack::Error::DialogError(
2580 "No answer received".to_string(),
2581 dialog_id,
2582 rsipstack::rsip::StatusCode::NotAcceptableHere,
2583 ));
2584 }
2585 }
2586 }
2587 };
2588 let (answer, remote_description_already_applied) = answer;
2589
2590 {
2591 let mut cs = call_state_ref.write().await;
2592 if cs.answer.is_none() {
2593 cs.answer = Some(answer.clone());
2594 }
2595 if auto_hangup {
2596 cs.auto_hangup = Some((ssrc, CallRecordHangupReason::ByRefer));
2597 }
2598 }
2599 if !remote_description_already_applied {
2600 self.media_stream
2601 .update_remote_description(&track_id, &answer)
2602 .await
2603 .ok();
2604 }
2605
2606 Ok(answer)
2607 }
2608
2609 pub fn is_webrtc_sdp(sdp: &str) -> bool {
2611 (sdp.contains("a=ice-ufrag:") || sdp.contains("a=ice-pwd:"))
2612 && sdp.contains("a=fingerprint:")
2613 }
2614
2615 pub async fn setup_answer_track(
2616 &self,
2617 ssrc: u32,
2618 option: &CallOption,
2619 offer: String,
2620 ) -> Result<(String, Box<dyn Track>)> {
2621 let offer = match option.enable_ipv6 {
2622 Some(false) | None => strip_ipv6_candidates(&offer),
2623 _ => offer.clone(),
2624 };
2625
2626 let timeout = option.handshake_timeout.map(|t| Duration::from_secs(t));
2627
2628 let mut media_track = if Self::is_webrtc_sdp(&offer) {
2629 let mut rtc_config = RtcTrackConfig::default();
2630 rtc_config.mode = rustrtc::TransportMode::WebRtc;
2631 rtc_config.ice_servers = self.app_state.config.ice_servers.clone();
2632 if let Some(ref external_ip) = self.app_state.config.external_ip {
2633 rtc_config.external_ip = Some(external_ip.clone());
2634 }
2635 if let Some(ref bind_ip) = self.app_state.config.rtp_bind_ip {
2636 rtc_config.bind_ip = Some(bind_ip.clone());
2637 }
2638 rtc_config.enable_latching = self.app_state.config.enable_rtp_latching;
2639 rtc_config.enable_ice_lite = self
2640 .call_state
2641 .read()
2642 .await
2643 .option
2644 .as_ref()
2645 .and_then(|o| o.enable_ice_lite)
2646 .or(self.app_state.config.enable_ice_lite);
2647
2648 let webrtc_track = RtcTrack::new(
2649 self.cancel_token.child_token(),
2650 self.session_id.clone(),
2651 self.track_config.clone(),
2652 rtc_config,
2653 )
2654 .with_ssrc(ssrc);
2655
2656 Box::new(webrtc_track) as Box<dyn Track>
2657 } else {
2658 let per_call_srtp = option.sip.as_ref().and_then(|s| s.enable_srtp);
2659 let rtp_track = self
2660 .create_rtp_track(self.session_id.clone(), ssrc, per_call_srtp)
2661 .await?;
2662 Box::new(rtp_track) as Box<dyn Track>
2663 };
2664
2665 let answer = match media_track.handshake(offer.clone(), timeout).await {
2666 Ok(answer) => answer,
2667 Err(e) => {
2668 return Err(anyhow::anyhow!("handshake failed: {e}"));
2669 }
2670 };
2671
2672 return Ok((answer, media_track));
2673 }
2674
2675 pub async fn prepare_incoming_sip_track(
2676 &self,
2677 cancel_token: CancellationToken,
2678 call_state_ref: ActiveCallStateRef,
2679 track_id: &String,
2680 pending_dialog: PendingDialog,
2681 hangup_headers: Option<Vec<rsipstack::rsip::Header>>,
2682 ) -> Result<()> {
2683 let state_receiver = pending_dialog.state_receiver;
2684 let states = InviteDialogStates {
2687 is_client: false,
2688 session_id: self.session_id.clone(),
2689 track_id: track_id.clone(),
2690 event_sender: self.event_sender.clone(),
2691 media_stream: self.media_stream.clone(),
2692 call_state: self.call_state.clone(),
2693 cancel_token,
2694 terminated_reason: None,
2695 has_early_media: false,
2696 };
2697
2698 let initial_request = pending_dialog.dialog.initial_request();
2699 let offer = String::from_utf8_lossy(&initial_request.body).to_string();
2700
2701 let (ssrc, option) = {
2702 let call_state = call_state_ref.read().await;
2703 (
2704 call_state.ssrc,
2705 call_state.option.clone().unwrap_or_default(),
2706 )
2707 };
2708
2709 match self.setup_answer_track(ssrc, &option, offer).await {
2710 Ok((offer, track)) => {
2711 self.setup_track_with_stream(&option, track).await?;
2712 {
2713 let mut state = self.call_state.write().await;
2714 state.ready_to_answer = Some((offer, None, pending_dialog.dialog));
2715 }
2716 }
2717 Err(e) => {
2718 return Err(anyhow::anyhow!("error creating track: {}", e));
2719 }
2720 }
2721
2722 let mut client_dialog_handler = DialogStateReceiverGuard::new(
2723 self.invitation.dialog_layer.clone(),
2724 state_receiver,
2725 hangup_headers,
2726 );
2727
2728 crate::spawn(async move {
2729 client_dialog_handler.process_dialog(states).await;
2730 });
2731 Ok(())
2732 }
2733}
2734
2735impl Drop for ActiveCall {
2736 fn drop(&mut self) {
2737 info!(session_id = self.session_id, "dropping active call");
2738 if let Some(sender) = self.app_state.callrecord_sender.as_ref() {
2739 if let Some(record) = self.get_callrecord() {
2740 if let Err(e) = sender.send(record) {
2741 warn!(
2742 session_id = self.session_id,
2743 "failed to send call record: {}", e
2744 );
2745 }
2746 }
2747 }
2748 }
2749}
2750
2751impl ActiveCallState {
2752 pub fn merge_option(&self, mut option: CallOption) -> CallOption {
2753 if let Some(existing) = &self.option {
2754 if option.asr.is_none() {
2755 option.asr = existing.asr.clone();
2756 }
2757 if option.tts.is_none() {
2758 option.tts = existing.tts.clone();
2759 }
2760 if option.vad.is_none() {
2761 option.vad = existing.vad.clone();
2762 }
2763 if option.denoise.is_none() {
2764 option.denoise = existing.denoise;
2765 }
2766 if option.recorder.is_none() {
2767 option.recorder = existing.recorder.clone();
2768 }
2769 if option.eou.is_none() {
2770 option.eou = existing.eou.clone();
2771 }
2772 if option.extra.is_none() {
2773 option.extra = existing.extra.clone();
2774 }
2775 if option.ambiance.is_none() {
2776 option.ambiance = existing.ambiance.clone();
2777 }
2778 if option.ringback_detection.is_none() {
2779 option.ringback_detection = existing.ringback_detection.clone();
2780 }
2781 }
2782 option
2783 }
2784
2785 pub fn set_hangup_reason(&mut self, reason: CallRecordHangupReason) {
2786 if self.hangup_reason.is_none() {
2787 self.hangup_reason = Some(reason);
2788 }
2789 }
2790
2791 pub fn build_hangup_event(
2792 &self,
2793 track_id: TrackId,
2794 initiator: Option<String>,
2795 ) -> crate::event::SessionEvent {
2796 let from = self.option.as_ref().and_then(|o| o.caller.as_ref());
2797 let to = self.option.as_ref().and_then(|o| o.callee.as_ref());
2798 let extra = self.extras.clone();
2799
2800 crate::event::SessionEvent::Hangup {
2801 track_id,
2802 timestamp: crate::media::get_timestamp(),
2803 reason: Some(format!("{:?}", self.hangup_reason)),
2804 initiator,
2805 start_time: self.start_time.to_rfc3339(),
2806 answer_time: self.answer_time.map(|t| t.to_rfc3339()),
2807 ringing_time: self.ring_time.map(|t| t.to_rfc3339()),
2808 hangup_time: Utc::now().to_rfc3339(),
2809 extra,
2810 from: from.map(|f| f.into()),
2811 to: to.map(|f| f.into()),
2812 refer: Some(self.is_refer),
2813 }
2814 }
2815
2816 pub fn build_callrecord(
2817 &self,
2818 app_state: AppState,
2819 session_id: String,
2820 call_type: ActiveCallType,
2821 ) -> CallRecord {
2822 let option = self.option.clone().unwrap_or_default();
2823 let recorder = if option.recorder.is_some() {
2824 let recorder_file = app_state.get_recorder_file(&session_id);
2825 if std::path::Path::new(&recorder_file).exists() {
2826 let file_size = std::fs::metadata(&recorder_file)
2827 .map(|m| m.len())
2828 .unwrap_or(0);
2829 vec![crate::callrecord::CallRecordMedia {
2830 track_id: session_id.clone(),
2831 path: recorder_file,
2832 size: file_size,
2833 extra: None,
2834 }]
2835 } else {
2836 vec![]
2837 }
2838 } else {
2839 vec![]
2840 };
2841
2842 let dump_event_file = app_state.get_dump_events_file(&session_id);
2843 let dump_event_file = if std::path::Path::new(&dump_event_file).exists() {
2844 Some(dump_event_file)
2845 } else {
2846 None
2847 };
2848
2849 let refer_callrecord = self.refer_callstate.as_ref().and_then(|rc| {
2850 if let Ok(rc) = rc.try_read() {
2851 Some(Box::new(rc.build_callrecord(
2852 app_state.clone(),
2853 rc.session_id.clone(),
2854 ActiveCallType::B2bua,
2855 )))
2856 } else {
2857 None
2858 }
2859 });
2860
2861 let caller = option.caller.clone().unwrap_or_default();
2862 let callee = option.callee.clone().unwrap_or_default();
2863
2864 CallRecord {
2865 option: Some(option),
2866 call_id: session_id,
2867 call_type,
2868 start_time: self.start_time,
2869 ring_time: self.ring_time.clone(),
2870 answer_time: self.answer_time.clone(),
2871 end_time: Utc::now(),
2872 caller,
2873 callee,
2874 hangup_reason: self.hangup_reason.clone(),
2875 hangup_messages: Vec::new(),
2876 status_code: self.last_status_code,
2877 extras: self.extras.clone(),
2878 dump_event_file,
2879 recorder,
2880 refer_callrecord,
2881 }
2882 }
2883}