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