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