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