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