1use super::Command;
2use crate::{
3 CallOption, ReferOption,
4 call::state::{ActorMsg, CallProgress, CallRuntime, Extras, LegShared, build_callrecord},
5 event::{EventReceiver, EventSender, SessionEvent},
6 media::{
7 TrackId,
8 engine::StreamEngine,
9 recorder::RecorderOption,
10 stream::{MediaStream, MediaStreamBuilder, SERVER_SIDE_TRACK_ID},
11 track::{
12 Track, TrackConfig, forwarding::ForwardingTrack, media_pass::MediaPassTrack,
13 tts::SynthesisHandle, websocket::WebsocketBytesReceiver,
14 },
15 },
16 synthesis::SynthesisCommand,
17 transcription::TranscriptionOption,
18};
19use crate::{
20 app::AppState,
21 call::{CommandReceiver, CommandSender, sip::Invitation},
22 callrecord::{CallRecord, CallRecordEvent, CallRecordEventType, CallRecordHangupReason},
23};
24use anyhow::Result;
25use arc_swap::{ArcSwap, ArcSwapOption};
26use chrono::{DateTime, Utc};
27use rsipstack::dialog::invite_dialog::InviteDialog;
28use serde::{Deserialize, Serialize};
29use std::{
30 collections::HashMap,
31 path::Path,
32 sync::{
33 Arc,
34 atomic::{AtomicBool, Ordering},
35 },
36 time::Duration,
37};
38use tokio::{fs::File, select, sync::mpsc};
39use tokio_util::sync::CancellationToken;
40use tracing::{debug, info, warn};
41
42pub enum PendingCallerTrack {
44 StartedForEarlyMedia,
47 NotStarted(Box<dyn Track>),
50}
51
52#[cfg(test)]
53mod tests {
54 use super::*;
55 use crate::app::AppStateBuilder;
56 use crate::config::Config;
57 use crate::media::track::tts::SynthesisHandle;
58 use crate::synthesis::SynthesisCommand;
59 use tokio::sync::mpsc;
60
61 async fn make_active_call_with_option(option: CallOption) -> Arc<ActiveCall> {
62 let mut config = Config::default();
63 config.udp_port = 0; config.media_cache_path = "/tmp/mediacache".to_string();
65 let app_state = AppStateBuilder::new()
66 .with_config(config)
67 .with_stream_engine(Arc::new(StreamEngine::default()))
68 .build()
69 .await
70 .unwrap();
71 let active_call = Arc::new(ActiveCall::new(CallSpec {
72 call_type: ActiveCallType::Sip,
73 cancel_token: CancellationToken::new(),
74 session_id: "test-session".to_string(),
75 invitation: app_state.invitation.clone(),
76 app_state: app_state.clone(),
77 track_config: TrackConfig::default(),
78 audio_receiver: None,
79 dump_events: false,
80 server_side_track_id: None,
81 extras: None,
82 }));
83 active_call.set_option(option);
84 active_call
85 }
86
87 #[tokio::test]
88 async fn test_tts_ssrc_reuse_for_autohangup() -> Result<()> {
89 let mut option = crate::CallOption::default();
90 option.tts = Some(crate::synthesis::SynthesisOption::default());
91 let active_call = make_active_call_with_option(option).await;
92
93 let (tx, mut rx) = mpsc::unbounded_channel::<SynthesisCommand>();
94 let initial_ssrc = 12345;
95 let handle = SynthesisHandle::new(tx, Some("play_1".to_string()), initial_ssrc);
96
97 active_call.tts_handle.store(Some(Arc::new(handle)));
99 active_call.set_current_play(Some("play_1".to_string()));
100
101 active_call
103 .do_tts(Command::Tts {
104 text: "hangup now".to_string(),
105 speaker: None,
106 play_id: Some("play_1".to_string()),
107 auto_hangup: Some(true),
108 streaming: Some(false),
109 end_of_stream: Some(true),
110 option: None,
111 wait_input_timeout: None,
112 base64: Some(false),
113 cache_key: None,
114 })
115 .await?;
116
117 let cmd = rx.try_recv().expect("Should have received tts command");
119 assert_eq!(cmd.text, "hangup now");
120 assert_eq!(cmd.auto_hangup, Some(true));
121
122 Ok(())
123 }
124
125 #[tokio::test]
126 async fn test_tts_new_ssrc_for_different_play_id() -> Result<()> {
127 let mut tts_opt = crate::synthesis::SynthesisOption::default();
128 tts_opt.provider = Some(crate::synthesis::SynthesisType::Aliyun);
129 let mut option = crate::CallOption::default();
130 option.tts = Some(tts_opt);
131 let active_call = make_active_call_with_option(option).await;
132
133 let (tx, _rx) = mpsc::unbounded_channel();
134 let initial_ssrc = 111;
135 let handle = SynthesisHandle::new(tx, Some("play_1".to_string()), initial_ssrc);
136
137 active_call.tts_handle.store(Some(Arc::new(handle)));
138 active_call.set_current_play(Some("play_1".to_string()));
139
140 active_call
142 .do_tts(Command::Tts {
143 text: "new play".to_string(),
144 speaker: None,
145 play_id: Some("play_2".to_string()),
146 auto_hangup: Some(true),
147 streaming: Some(false),
148 end_of_stream: Some(true),
149 option: None,
150 wait_input_timeout: None,
151 base64: Some(false),
152 cache_key: None,
153 })
154 .await?;
155
156 {
159 let handle = active_call.tts_handle.load_full();
160 assert!(handle.is_some(), "new tts handle should be stored");
161 let handle = handle.unwrap();
162 assert_ne!(
163 handle.ssrc, initial_ssrc,
164 "Should use a new SSRC for different play_id"
165 );
166 }
167
168 Ok(())
169 }
170
171 #[tokio::test]
173 async fn test_hangup_refer_true_cancels_refer_only() -> Result<()> {
174 let active_call = make_active_call_with_option(crate::CallOption::default()).await;
175
176 let refer_token = active_call.cancel_token.child_token();
177 let refer_leg = LegShared::new(1, true, CallProgress::default());
178 active_call.set_refer_call_token(refer_token.clone());
179 active_call.set_refer_leg(Some(refer_leg.clone()));
180
181 active_call.do_hangup(None, None, None, Some(true)).await?;
182
183 assert!(
184 refer_token.is_cancelled(),
185 "refer token should be cancelled"
186 );
187 assert!(
188 !active_call.media_stream.cancel_token.is_cancelled(),
189 "media stream should NOT stop"
190 );
191 assert!(
192 refer_leg.progress.load_full().hangup_reason.is_some(),
193 "hangup_reason should be set on refer leg"
194 );
195 Ok(())
196 }
197
198 #[tokio::test]
200 async fn test_hangup_none_cancels_refer_too() -> Result<()> {
201 let active_call = make_active_call_with_option(crate::CallOption::default()).await;
202
203 let refer_token = active_call.cancel_token.child_token();
204 active_call.set_refer_call_token(refer_token.clone());
205
206 active_call.do_hangup(None, None, None, None).await?;
207
208 assert!(
209 refer_token.is_cancelled(),
210 "refer token should be cancelled"
211 );
212 assert!(
213 active_call.media_stream.cancel_token.is_cancelled(),
214 "media stream should stop"
215 );
216 Ok(())
217 }
218
219 struct MockCallerTrack {
234 id: TrackId,
235 config: crate::media::track::TrackConfig,
236 processor_chain: crate::media::processor::ProcessorChain,
237 }
238
239 impl MockCallerTrack {
240 fn new(id: TrackId) -> Self {
241 Self {
242 id,
243 config: crate::media::track::TrackConfig::default(),
244 processor_chain: crate::media::processor::ProcessorChain::new(16000),
245 }
246 }
247 }
248
249 #[async_trait::async_trait]
250 impl crate::media::track::Track for MockCallerTrack {
251 fn ssrc(&self) -> u32 {
252 0
253 }
254 fn id(&self) -> &TrackId {
255 &self.id
256 }
257 fn config(&self) -> &crate::media::track::TrackConfig {
258 &self.config
259 }
260 fn processor_chain(&mut self) -> &mut crate::media::processor::ProcessorChain {
261 &mut self.processor_chain
262 }
263 async fn handshake(
264 &mut self,
265 _o: String,
266 _t: Option<tokio::time::Duration>,
267 ) -> Result<String> {
268 Ok(String::new())
269 }
270 async fn update_remote_description(&mut self, _a: &String) -> Result<()> {
271 Ok(())
272 }
273 async fn start(
274 &mut self,
275 _e: crate::event::EventSender,
276 _p: crate::media::track::TrackPacketSender,
277 ) -> Result<()> {
278 Ok(())
279 }
280 async fn stop(&self) -> Result<()> {
281 Ok(())
282 }
283 async fn send_packet(&mut self, _f: &crate::media::AudioFrame) -> Result<()> {
284 Ok(())
285 }
286 }
287
288 struct MockAsrClient;
289
290 #[async_trait::async_trait]
291 impl crate::transcription::TranscriptionClient for MockAsrClient {
292 fn send_audio(
293 &self,
294 _s: &[crate::media::Sample],
295 _src: Option<&crate::media::SourcePacket>,
296 ) -> Result<()> {
297 Ok(())
298 }
299 }
300
301 async fn make_active_call_with_engine_and_option(
302 engine: Arc<StreamEngine>,
303 cache_dir: &str,
304 option: crate::CallOption,
305 ) -> Arc<ActiveCall> {
306 let mut config = Config::default();
307 config.udp_port = 0;
308 config.media_cache_path = cache_dir.to_string();
309 let app_state = AppStateBuilder::new()
310 .with_config(config)
311 .with_stream_engine(engine)
312 .build()
313 .await
314 .unwrap();
315 let session_id = format!("test-{}-{}", cache_dir, uuid::Uuid::new_v4());
316 let active_call = Arc::new(ActiveCall::new(CallSpec {
317 call_type: ActiveCallType::Sip,
318 cancel_token: CancellationToken::new(),
319 session_id: session_id.clone(),
320 invitation: app_state.invitation.clone(),
321 app_state: app_state.clone(),
322 track_config: TrackConfig::default(),
323 audio_receiver: None,
324 dump_events: false,
325 server_side_track_id: None,
326 extras: None,
327 }));
328 active_call.set_option(option);
329 active_call
330 }
331
332 #[tokio::test]
333 async fn test_setup_track_with_stream_builds_processors_from_accept_option() -> Result<()> {
334 let (asr_created_tx, mut asr_created_rx) = mpsc::channel::<()>(1);
335
336 let mock_provider =
337 crate::transcription::TranscriptionType::Other("mock-ringing-asr".to_string());
338
339 let mut engine = StreamEngine::new();
340 engine.register_asr(
341 mock_provider.clone(),
342 Box::new(move |_tid, _tok, _opt, _es| {
343 let tx = asr_created_tx.clone();
344 Box::pin(async move {
345 let _ = tx.send(()).await;
346 Ok(Box::new(MockAsrClient)
347 as Box<dyn crate::transcription::TranscriptionClient>)
348 })
349 }),
350 );
351 let engine = Arc::new(engine);
352
353 let accept_option = crate::CallOption {
354 asr: Some(crate::transcription::TranscriptionOption {
355 provider: Some(mock_provider),
356 ..Default::default()
357 }),
358 ..Default::default()
359 };
360 let active_call = make_active_call_with_engine_and_option(
361 engine,
362 "mediacache_ringing_accept_test",
363 accept_option,
364 )
365 .await;
366 let cancel_token = active_call.cancel_token.clone();
367
368 let mock_track = Box::new(MockCallerTrack::new(active_call.session_id.clone()));
371
372 let accept_option = active_call.progress.load_full().option.clone().unwrap();
375 active_call
376 .setup_track_with_stream(&accept_option, mock_track)
377 .await?;
378
379 let received =
382 tokio::time::timeout(std::time::Duration::from_secs(3), asr_created_rx.recv()).await;
383 assert!(
384 received.is_ok() && received.unwrap().is_some(),
385 "ASR processor was NOT created — setup_track_with_stream did not build \
386 processors from the accept option (regression: ringing-before-accept)"
387 );
388
389 cancel_token.cancel();
390 Ok(())
391 }
392
393 #[tokio::test]
407 async fn test_update_track_wrapper_does_not_build_asr_processor() -> Result<()> {
408 let (asr_created_tx, mut asr_created_rx) = mpsc::channel::<()>(1);
409
410 let mock_provider =
411 crate::transcription::TranscriptionType::Other("mock-ringing-asr".to_string());
412
413 let mut engine = StreamEngine::new();
414 engine.register_asr(
415 mock_provider.clone(),
416 Box::new(move |_tid, _tok, _opt, _es| {
417 let tx = asr_created_tx.clone();
418 Box::pin(async move {
419 let _ = tx.send(()).await;
420 Ok(Box::new(MockAsrClient)
421 as Box<dyn crate::transcription::TranscriptionClient>)
422 })
423 }),
424 );
425 let engine = Arc::new(engine);
426
427 let accept_option = crate::CallOption {
430 asr: Some(crate::transcription::TranscriptionOption {
431 provider: Some(mock_provider),
432 ..Default::default()
433 }),
434 ..Default::default()
435 };
436 let active_call = make_active_call_with_engine_and_option(
437 engine,
438 "mediacache_update_track_wrapper_test",
439 accept_option,
440 )
441 .await;
442 let cancel_token = active_call.cancel_token.clone();
443
444 let mock_track = Box::new(MockCallerTrack::new(active_call.session_id.clone()));
445 active_call.update_track_wrapper(mock_track, None).await;
446
447 let received =
450 tokio::time::timeout(std::time::Duration::from_millis(500), asr_created_rx.recv())
451 .await;
452 assert!(
453 received.is_err(),
454 "ASR builder fired during track preparation — double-ASR regression"
455 );
456
457 cancel_token.cancel();
458 Ok(())
459 }
460}
461
462#[derive(Deserialize)]
463#[serde(rename_all = "camelCase")]
464pub struct CallParams {
465 pub id: Option<String>,
466 #[serde(rename = "dump")]
467 pub dump_events: Option<bool>,
468 #[serde(rename = "ping")]
469 pub ping_interval: Option<u32>,
470 pub server_side_track: Option<String>,
471 #[serde(default)]
476 pub forward: Option<bool>,
477 #[serde(default)]
479 pub visited: Option<String>,
480}
481
482impl CallParams {
483 pub fn to_forward_query(&self) -> String {
486 let mut parts: Vec<String> = Vec::new();
487 if let Some(id) = &self.id {
488 parts.push(format!("id={}", urlencoding::encode(id)));
489 }
490 if let Some(dump) = self.dump_events {
491 parts.push(format!("dump={}", dump));
492 }
493 if let Some(ping) = self.ping_interval {
494 parts.push(format!("ping={}", ping));
495 }
496 if let Some(track) = &self.server_side_track {
497 parts.push(format!("server_side_track={}", urlencoding::encode(track)));
498 }
499 parts.push("forward=true".to_string());
500 parts.join("&")
501 }
502}
503
504#[derive(Debug, Serialize, Deserialize, Clone, Default, PartialEq, Eq)]
505#[serde(rename_all = "camelCase")]
506pub enum ActiveCallType {
507 Webrtc,
508 B2bua,
509 WebSocket,
510 #[default]
511 Sip,
512}
513
514pub type ActiveCallRef = Arc<ActiveCall>;
515
516pub struct ActiveCall {
520 pub cancel_token: CancellationToken,
521 pub call_type: ActiveCallType,
522 pub session_id: String,
523 pub start_time: DateTime<Utc>,
524 pub media_stream: Arc<MediaStream>,
525 pub track_config: TrackConfig,
526 pub event_sender: EventSender,
527 pub app_state: AppState,
528 pub invitation: Invitation,
529 pub cmd_sender: CommandSender,
530 pub dump_events: bool,
531 pub server_side_track_id: TrackId,
532
533 pub ssrc: u32,
535 pub bridge_paused: Arc<AtomicBool>,
537 pub progress: Arc<ArcSwap<CallProgress>>,
539 pub extras: Extras,
541 pub moh: ArcSwapOption<String>,
543 pub current_play_id: ArcSwapOption<String>,
545 pub tts_handle: ArcSwapOption<SynthesisHandle>,
547 pub refer_leg: ArcSwapOption<LegShared>,
549 pub ready_to_answer: ArcSwapOption<ReadyAnswer>,
552 pub pending_sip_answer: ArcSwapOption<PendingSipAnswer>,
557 pub refer_call_token: ArcSwapOption<CancellationToken>,
559 pub wait_input_timeout: ArcSwapOption<u32>,
561 pub pending_asr_resume: ArcSwapOption<(u32, TranscriptionOption)>,
563 pub audio_receiver: std::sync::Mutex<Option<WebsocketBytesReceiver>>,
565}
566
567impl ActiveCall {
568 pub fn leg(&self) -> LegShared {
570 LegShared {
571 ssrc: self.ssrc,
572 is_refer: false,
573 progress: self.progress.clone(),
574 extras: self.extras.clone(),
575 }
576 }
577
578 pub fn set_option(&self, option: CallOption) {
580 self.progress.rcu(|p| {
581 let mut p = CallProgress::clone(p);
582 p.option = Some(option.clone());
583 p
584 });
585 }
586
587 pub fn moh_path(&self) -> Option<String> {
588 self.moh.load_full().map(|s| s.to_string())
589 }
590
591 pub fn set_moh(&self, v: Option<String>) {
592 self.moh.store(v.map(Arc::new));
593 }
594
595 pub fn current_play(&self) -> Option<String> {
596 self.current_play_id.load_full().map(|s| s.to_string())
597 }
598
599 pub fn set_current_play(&self, v: Option<String>) {
600 self.current_play_id.store(v.map(Arc::new));
601 }
602
603 pub fn refer_leg_value(&self) -> Option<LegShared> {
604 self.refer_leg.load_full().map(|l| l.as_ref().clone())
605 }
606
607 pub fn set_refer_leg(&self, v: Option<LegShared>) {
608 self.refer_leg.store(v.map(Arc::new));
609 }
610
611 pub fn set_ready_to_answer(&self, ready: ReadyAnswer) {
613 self.ready_to_answer.store(Some(Arc::new(ready)));
614 }
615
616 pub fn take_ready_to_answer(&self) -> Option<Arc<ReadyAnswer>> {
617 self.ready_to_answer.swap(None)
618 }
619
620 pub fn has_ready_to_answer(&self) -> bool {
621 self.ready_to_answer.load().is_some()
622 }
623
624 pub fn set_pending_sip_answer(&self, pending: PendingSipAnswer) {
626 self.pending_sip_answer.store(Some(Arc::new(pending)));
627 }
628
629 pub fn take_pending_sip_answer(&self) -> Option<Arc<PendingSipAnswer>> {
630 self.pending_sip_answer.swap(None)
631 }
632
633 pub fn take_refer_call_token(&self) -> Option<CancellationToken> {
636 self.refer_call_token.swap(None).map(|t| (*t).clone())
637 }
638
639 pub fn set_refer_call_token(&self, token: CancellationToken) {
640 self.refer_call_token.store(Some(Arc::new(token)));
641 }
642
643 pub fn take_wait_input_timeout(&self) -> Option<u32> {
646 self.wait_input_timeout.swap(None).map(|t| *t)
647 }
648
649 pub fn set_wait_input_timeout(&self, v: Option<u32>) {
650 self.wait_input_timeout.store(v.map(Arc::new));
651 }
652
653 pub fn set_pending_asr_resume(&self, v: (u32, TranscriptionOption)) {
655 self.pending_asr_resume.store(Some(Arc::new(v)));
656 }
657
658 pub fn take_pending_asr_resume(&self) -> Option<(u32, TranscriptionOption)> {
659 self.pending_asr_resume.swap(None).map(|a| (*a).clone())
660 }
661
662 pub fn set_extra(&self, key: &str, value: serde_json::Value) {
664 self.leg().set_extra(key, value);
665 }
666
667 fn has_pending_invite(&self) -> bool {
669 self.invitation
670 .find_dialog_id_by_session_id(&self.session_id)
671 .is_some()
672 }
673
674 async fn hangup_now(&self, reason: Option<CallRecordHangupReason>) {
676 self.do_hangup(reason, None, None, None).await.ok();
677 }
678
679 async fn reject_now(&self, code: Option<rsipstack::rsip::StatusCode>, reason: Option<String>) {
681 self.do_reject(code, reason).await.ok();
682 }
683}
684
685pub struct ReadyAnswer {
688 pub answer: String,
689 pub track: PendingCallerTrack,
690 pub dialog: InviteDialog,
691}
692
693pub struct PendingSipAnswer {
697 pub answer: String,
698 pub dialog: InviteDialog,
699 pub track: Box<dyn Track>,
700}
701
702pub struct ActiveCallGuard {
703 pub call: ActiveCallRef,
704 pub active_calls: usize,
705}
706
707impl ActiveCallGuard {
708 pub fn new(call: ActiveCallRef) -> Self {
709 let active_calls = {
710 call.app_state
711 .total_calls
712 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
713 let mut calls = call.app_state.active_calls.lock().unwrap();
714 calls.insert(call.session_id.clone(), call.clone());
715 calls.len()
716 };
717 Self { call, active_calls }
718 }
719}
720
721impl Drop for ActiveCallGuard {
722 fn drop(&mut self) {
723 self.call
724 .app_state
725 .active_calls
726 .lock()
727 .unwrap()
728 .remove(&self.call.session_id);
729 self.call
732 .invitation
733 .unregister_session(&self.call.session_id);
734 }
735}
736
737pub struct ActiveCallReceiver {
738 pub cmd_receiver: CommandReceiver,
739 pub dump_cmd_receiver: CommandReceiver,
740 pub dump_event_receiver: EventReceiver,
741}
742
743pub struct CallSpec {
745 pub call_type: ActiveCallType,
746 pub cancel_token: CancellationToken,
747 pub session_id: String,
748 pub invitation: Invitation,
749 pub app_state: AppState,
750 pub track_config: TrackConfig,
751 pub audio_receiver: Option<WebsocketBytesReceiver>,
753 pub dump_events: bool,
754 pub server_side_track_id: Option<TrackId>,
756 pub extras: Option<HashMap<String, serde_json::Value>>,
759}
760
761impl ActiveCall {
762 pub fn new(spec: CallSpec) -> Self {
763 let CallSpec {
764 call_type,
765 cancel_token,
766 session_id,
767 invitation,
768 app_state,
769 track_config,
770 audio_receiver,
771 dump_events,
772 server_side_track_id,
773 extras,
774 } = spec;
775 let event_sender = crate::event::create_event_sender();
776 let cmd_sender = tokio::sync::broadcast::Sender::<Command>::new(32);
777 let server_side_track_id = server_side_track_id.unwrap_or(SERVER_SIDE_TRACK_ID.to_string());
778 let media_stream_builder = MediaStreamBuilder::new(event_sender.clone())
779 .with_id(session_id.clone())
780 .with_cancel_token(cancel_token.child_token());
781 let media_stream = Arc::new(media_stream_builder.build());
782 let start_time = Utc::now();
783 let call_type_str = match &call_type {
785 ActiveCallType::Sip => "sip",
786 ActiveCallType::WebSocket => "websocket",
787 ActiveCallType::Webrtc => "webrtc",
788 ActiveCallType::B2bua => "b2bua",
789 };
790 let mut extras = extras.unwrap_or_default();
791 extras
792 .entry(crate::playbook::BUILTIN_SESSION_ID.to_string())
793 .or_insert_with(|| serde_json::Value::String(session_id.clone()));
794 extras
795 .entry(crate::playbook::BUILTIN_CALL_TYPE.to_string())
796 .or_insert_with(|| serde_json::Value::String(call_type_str.to_string()));
797 extras
798 .entry(crate::playbook::BUILTIN_START_TIME.to_string())
799 .or_insert_with(|| serde_json::Value::String(start_time.to_rfc3339()));
800
801 let progress = CallProgress {
802 session_id: session_id.clone(),
803 start_time: Some(start_time),
804 ..Default::default()
805 };
806
807 Self {
808 cancel_token,
809 call_type,
810 session_id,
811 start_time,
812 media_stream,
813 track_config,
814 event_sender,
815 app_state,
816 invitation,
817 cmd_sender,
818 dump_events,
819 server_side_track_id,
820 ssrc: rand::random::<u32>(),
821 bridge_paused: Arc::new(AtomicBool::new(false)),
822 progress: Arc::new(ArcSwap::from_pointee(progress)),
823 extras: Arc::new(ArcSwap::from_pointee(extras)),
824 moh: ArcSwapOption::new(None),
825 current_play_id: ArcSwapOption::new(None),
826 tts_handle: ArcSwapOption::new(None),
827 refer_leg: ArcSwapOption::new(None),
828 ready_to_answer: ArcSwapOption::new(None),
829 pending_sip_answer: ArcSwapOption::new(None),
830 refer_call_token: ArcSwapOption::new(None),
831 wait_input_timeout: ArcSwapOption::new(None),
832 pending_asr_resume: ArcSwapOption::new(None),
833 audio_receiver: std::sync::Mutex::new(audio_receiver),
834 }
835 }
836
837 pub async fn enqueue_command(&self, command: Command) -> Result<()> {
838 self.cmd_sender
839 .send(command)
840 .map_err(|e| anyhow::anyhow!("Failed to send command: {}", e))?;
841 Ok(())
842 }
843
844 pub fn new_receiver(&self) -> ActiveCallReceiver {
848 ActiveCallReceiver {
849 cmd_receiver: self.cmd_sender.subscribe(),
850 dump_cmd_receiver: self.cmd_sender.subscribe(),
851 dump_event_receiver: self.event_sender.subscribe(),
852 }
853 }
854
855 pub async fn serve(self: Arc<Self>, receiver: ActiveCallReceiver) -> Result<()> {
863 let ActiveCallReceiver {
864 mut cmd_receiver,
865 dump_cmd_receiver,
866 dump_event_receiver,
867 } = receiver;
868
869 let mut event_receiver = self.event_sender.subscribe();
870 let (actor_tx, mut actor_rx) = mpsc::channel::<ActorMsg>(16);
871 let mut runtime = CallRuntime::new(actor_tx);
872 runtime.me = Some(self.clone());
873
874 self.app_state
875 .total_calls
876 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
877
878 let me = self.clone();
879 let actor = async move {
880 let _cancel_on_exit = CancelOnExit(&me.cancel_token);
884 let mut ticker = tokio::time::interval(Duration::from_millis(100));
885 let mut media_serve = Box::pin(me.media_stream.serve());
887 loop {
888 tokio::select! {
889 cmd = cmd_receiver.recv() => {
890 match cmd {
891 Ok(command) => {
892 if let Err(e) = Box::pin(me.dispatch(&mut runtime, command)).await {
896 warn!(session_id = me.session_id, "{}", e);
897 me.event_sender
898 .send(SessionEvent::Error {
899 track_id: me.session_id.clone(),
900 timestamp: crate::media::get_timestamp(),
901 sender: "command".to_string(),
902 error: e.to_string(),
903 code: None,
904 })
905 .ok();
906 }
907 }
908 Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue,
909 Err(_) => {
910 info!(session_id = me.session_id, "command loop done");
911 break;
912 }
913 }
914 }
915 ev = event_receiver.recv() => {
916 match ev {
917 Ok(event) => Box::pin(me.handle_event(&mut runtime, event)).await,
918 Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue,
919 Err(_) => {
920 info!(session_id = me.session_id, "event loop done");
921 break;
922 }
923 }
924 }
925 Some(msg) = actor_rx.recv() => {
926 if let Err(e) = Box::pin(me.handle_actor_msg(msg)).await {
927 warn!(session_id = me.session_id, "{}", e);
928 me.event_sender
929 .send(SessionEvent::Error {
930 track_id: me.session_id.clone(),
931 timestamp: crate::media::get_timestamp(),
932 sender: "command".to_string(),
933 error: e.to_string(),
934 code: None,
935 })
936 .ok();
937 }
938 }
939 _ = ticker.tick() => {
940 Box::pin(me.check_input_timeout(&mut runtime)).await;
941 }
942 _ = &mut media_serve => {
943 info!(session_id = me.session_id, "media stream loop done");
944 break;
945 }
946 _ = me.cancel_token.cancelled() => {
947 info!(session_id = me.session_id, "call cancelled - cleaning up resources");
948 break;
949 }
950 }
951 }
952 };
953
954 tokio::join!(
955 self.dump_loop(self.dump_events, dump_cmd_receiver, dump_event_receiver),
956 actor
957 );
958 Ok(())
959 }
960
961 async fn check_input_timeout(&self, runtime: &mut CallRuntime) {
963 let (start_time, expire) = runtime.input_timeout_expire;
964 if expire > 0 && crate::media::get_timestamp() >= start_time + expire as u64 {
965 info!(session_id = self.session_id, "wait input timeout reached");
966 runtime.input_timeout_expire = (0, 0);
967 self.event_sender
968 .send(SessionEvent::Silence {
969 track_id: self.server_side_track_id.clone(),
970 timestamp: crate::media::get_timestamp(),
971 start_time,
972 duration: expire as u64,
973 samples: None,
974 refer: Some(false),
975 })
976 .ok();
977 }
978 }
979
980 async fn handle_event(&self, runtime: &mut CallRuntime, event: SessionEvent) {
982 match event {
983 SessionEvent::Speaking { .. }
984 | SessionEvent::Dtmf { .. }
985 | SessionEvent::AsrDelta { .. }
986 | SessionEvent::AsrFinal { .. }
987 | SessionEvent::TrackStart { .. } => {
988 runtime.input_timeout_expire = (0, 0);
989 }
990 SessionEvent::TrackEnd {
991 track_id,
992 play_id,
993 ssrc,
994 auto_hangup,
995 ..
996 } => {
997 if track_id != self.server_side_track_id {
998 return;
999 }
1000
1001 if play_id != self.current_play() {
1002 debug!(
1003 session_id = self.session_id,
1004 ?play_id,
1005 current = ?self.current_play(),
1006 "ignoring interrupted track end"
1007 );
1008 return;
1009 }
1010 self.set_current_play(None);
1011 let moh_path = self.moh_path();
1012 let wait_timeout_val = self.take_wait_input_timeout();
1013
1014 if let Some(path) = moh_path {
1015 info!(session_id = self.session_id, "looping moh: {}", path);
1016 let ssrc = rand::random::<u32>();
1017 let file_track = self.make_file_track(path.clone(), ssrc);
1018 self.update_track_wrapper(Box::new(file_track), Some(path))
1019 .await;
1020 return;
1021 }
1022
1023 if let Some(hangup_reason) = auto_hangup {
1024 info!(
1025 session_id = self.session_id,
1026 ssrc, "auto hangup when track end track_id:{}", track_id
1027 );
1028 self.do_hangup(Some(hangup_reason), None, None, None)
1029 .await
1030 .ok();
1031 }
1032
1033 if let Some(timeout) = wait_timeout_val {
1034 runtime.input_timeout_expire = if timeout > 0 {
1035 (crate::media::get_timestamp(), timeout)
1036 } else {
1037 (0, 0)
1038 };
1039 }
1040 }
1041 SessionEvent::Interrupt { receiver } => {
1042 let track_id = receiver.unwrap_or_else(|| self.server_side_track_id.clone());
1043 if track_id == self.server_side_track_id {
1044 debug!(
1045 session_id = self.session_id,
1046 "received interrupt event, stopping playback"
1047 );
1048 self.do_interrupt(true).await.ok();
1049 }
1050 }
1051 SessionEvent::Inactivity { track_id, .. } => {
1052 info!(
1053 session_id = self.session_id,
1054 track_id, "inactivity timeout reached, hanging up"
1055 );
1056 self.do_hangup(
1057 Some(CallRecordHangupReason::InactivityTimeout),
1058 None,
1059 None,
1060 None,
1061 )
1062 .await
1063 .ok();
1064 }
1065 SessionEvent::Hangup { refer, .. } => {
1066 if refer == Some(true) {
1068 if let Some((refer_ssrc, asr_option)) = self.take_pending_asr_resume() {
1069 let is_refer_hangup = self
1071 .refer_leg
1072 .load_full()
1073 .map(|leg| leg.ssrc == refer_ssrc)
1074 .unwrap_or(false);
1075
1076 if is_refer_hangup {
1077 info!(
1078 session_id = self.session_id,
1079 "Refer call ended, resuming parent ASR"
1080 );
1081
1082 match self
1084 .app_state
1085 .stream_engine
1086 .create_asr_processor(
1087 self.server_side_track_id.clone(),
1088 self.cancel_token.child_token(),
1089 asr_option,
1090 self.event_sender.clone(),
1091 )
1092 .await
1093 {
1094 Ok(asr_processor) => {
1095 if let Err(e) = self
1096 .media_stream
1097 .append_processor(&self.server_side_track_id, asr_processor)
1098 .await
1099 {
1100 warn!(
1101 session_id = self.session_id,
1102 "Failed to resume ASR after refer: {}", e
1103 );
1104 }
1105 }
1106 Err(e) => {
1107 warn!(
1108 session_id = self.session_id,
1109 "Failed to create ASR processor for resume: {}", e
1110 );
1111 }
1112 }
1113 }
1114 }
1115 }
1116 }
1117 SessionEvent::Error { track_id, .. } => {
1118 if track_id != self.server_side_track_id {
1119 return;
1120 }
1121
1122 let moh_info = {
1123 let path = self.moh_path();
1124 path.map(|path| {
1125 let fallback = "./config/sounds/refer_moh.wav".to_string();
1126 if path != fallback && std::path::Path::new(&fallback).exists() {
1127 info!(
1128 session_id = self.session_id,
1129 "moh error, switching to fallback: {}", fallback
1130 );
1131 self.set_moh(Some(fallback.clone()));
1132 fallback
1133 } else {
1134 info!(
1135 session_id = self.session_id,
1136 "looping moh on error: {}", path
1137 );
1138 path
1139 }
1140 })
1141 };
1142
1143 if let Some(next_path) = moh_info {
1144 let ssrc = rand::random::<u32>();
1145 let file_track = self.make_file_track(next_path.clone(), ssrc);
1146 self.update_track_wrapper(Box::new(file_track), Some(next_path))
1147 .await;
1148 }
1149 }
1150 SessionEvent::Hold { on_hold, .. } => {
1151 self.bridge_paused.store(on_hold, Ordering::Relaxed);
1152 }
1153 _ => {}
1154 }
1155 }
1156
1157 async fn handle_actor_msg(&self, msg: ActorMsg) -> Result<()> {
1159 match msg {
1160 ActorMsg::ReferDone {
1161 track_id,
1162 forward_dtmf,
1163 result,
1164 } => match result {
1165 Ok(answer) => {
1166 self.media_stream
1167 .set_track_refer(&track_id, Some(true))
1168 .await;
1169 if !forward_dtmf {
1170 self.media_stream
1171 .set_track_dtmf_forward(&track_id, false)
1172 .await;
1173 }
1174 self.event_sender
1175 .send(SessionEvent::Answer {
1176 timestamp: crate::media::get_timestamp(),
1177 track_id,
1178 sdp: answer,
1179 refer: Some(true),
1180 })
1181 .ok();
1182 Ok(())
1183 }
1184 Err(e) => {
1185 warn!(
1186 session_id = self.session_id,
1187 "failed to create refer sip track: {}", e
1188 );
1189 self.emit_reject_from_rsip_error(track_id, true, &e);
1190 Err(e.into())
1191 }
1192 },
1193 }
1194 }
1195
1196 async fn dispatch(&self, runtime: &mut CallRuntime, command: Command) -> Result<()> {
1197 match command {
1198 Command::Invite { option } => self.do_invite(runtime, option).await,
1199 Command::Accept { option } => self.do_accept(option).await,
1200 Command::Reject { reason, code } => {
1201 self.do_reject(code.map(|c| (c as u16).into()), Some(reason))
1202 .await
1203 }
1204 Command::Ringing { .. } => self.do_ringing(command).await,
1205 Command::Tts { .. } => self.do_tts(command).await,
1206 Command::Play { .. } => self.do_play(command).await,
1207 Command::Hangup {
1208 reason,
1209 initiator,
1210 headers,
1211 refer,
1212 } => {
1213 let reason = reason.map(|r| {
1214 r.parse::<CallRecordHangupReason>()
1215 .unwrap_or(CallRecordHangupReason::BySystem)
1216 });
1217 self.do_hangup(reason, initiator, headers, refer).await
1218 }
1219 Command::Refer {
1220 caller,
1221 callee,
1222 options,
1223 } => self.do_refer(runtime, caller, callee, options).await,
1224 Command::Message {
1225 body,
1226 content_type,
1227 headers,
1228 refer,
1229 } => self.do_message(body, content_type, headers, refer).await,
1230 Command::Bridge { target_session_id } => self.do_bridge(target_session_id).await,
1231 Command::Unbridge { target_session_id } => self.do_unbridge(target_session_id).await,
1232 Command::Mute { track_id } => self.do_mute(track_id).await,
1233 Command::Unmute { track_id } => self.do_unmute(track_id).await,
1234 Command::Pause {} => self.do_pause().await,
1235 Command::Resume {} => self.do_resume().await,
1236 Command::Interrupt {
1237 graceful: passage,
1238 fade_out_ms: _,
1239 } => self.do_interrupt(passage.unwrap_or_default()).await,
1240 Command::History { speaker, text } => self.do_history(speaker, text).await,
1241 Command::Custom { sender, data } => self.do_custom(sender, data),
1242 Command::AddIceCandidate {
1243 candidate,
1244 sdp_mid,
1245 sdp_mline_index,
1246 } => {
1247 self.media_stream
1248 .add_ice_candidate(&candidate, sdp_mid.as_deref(), sdp_mline_index)
1249 .await
1250 }
1251 }
1252 }
1253
1254 fn build_record_option(&self, option: &CallOption) -> Option<RecorderOption> {
1255 if let Some(recorder_option) = &option.recorder {
1256 let mut recorder_file = recorder_option.recorder_file.clone();
1257 if recorder_file.contains("{id}") {
1258 recorder_file = recorder_file.replace("{id}", &self.session_id);
1259 }
1260
1261 let recorder_file = if recorder_file.is_empty() {
1262 self.app_state.get_recorder_file(&self.session_id)
1263 } else {
1264 let p = Path::new(&recorder_file);
1265 p.is_absolute()
1266 .then(|| recorder_file.clone())
1267 .unwrap_or_else(|| self.app_state.get_recorder_file(&recorder_file))
1268 };
1269 info!(
1270 session_id = self.session_id,
1271 recorder_file, "created recording file"
1272 );
1273
1274 let track_samplerate = self.track_config.samplerate;
1275 let recorder_samplerate = if track_samplerate > 0 {
1276 track_samplerate
1277 } else {
1278 recorder_option.samplerate
1279 };
1280 let recorder_ptime = if recorder_option.ptime == 0 {
1281 200
1282 } else {
1283 recorder_option.ptime
1284 };
1285 let requested_format = recorder_option
1286 .format
1287 .unwrap_or(self.app_state.config.recorder_format());
1288 let format = requested_format.effective();
1289 if requested_format != format {
1290 warn!(
1291 session_id = self.session_id,
1292 requested = requested_format.extension(),
1293 "Recorder format fallback to wav due to unsupported feature"
1294 );
1295 }
1296 let mut recorder_config = RecorderOption {
1297 recorder_file,
1298 samplerate: recorder_samplerate,
1299 ptime: recorder_ptime,
1300 format: Some(format),
1301 native_samplerate: Some(
1302 recorder_option.native_samplerate.unwrap_or(false)
1303 || self.app_state.config.recorder_native_samplerate(),
1304 ),
1305 };
1306 recorder_config.ensure_path_extension(format);
1307 Some(recorder_config)
1308 } else {
1309 None
1310 }
1311 }
1312
1313 async fn invite_or_accept(&self, mut option: CallOption, sender: String) -> Result<CallOption> {
1314 {
1316 let state = self.progress.load_full();
1317 option = state.merge_option(option);
1318 }
1319
1320 option.check_default();
1321 if let Some(opt) = self.build_record_option(&option) {
1322 self.media_stream.update_recorder_option(opt).await;
1323 }
1324 self.ensure_call_ambiance(&option).await;
1325
1326 if let Some(opt) = &option.media_pass {
1327 let track_id = self.server_side_track_id.clone();
1328 let cancel_token = self.cancel_token.child_token();
1329 let ssrc = rand::random::<u32>();
1330 let media_pass_track = MediaPassTrack::new(
1331 self.session_id.clone(),
1332 ssrc,
1333 track_id,
1334 cancel_token,
1335 opt.clone(),
1336 );
1337 self.update_track_wrapper(Box::new(media_pass_track), None)
1338 .await;
1339 }
1340
1341 info!(
1342 session_id = self.session_id,
1343 call_type = ?self.call_type,
1344 sender,
1345 ?option,
1346 "caller with option"
1347 );
1348
1349 match self.setup_caller_track(&option).await {
1350 Ok(_) => return Ok(option),
1351 Err(e) => {
1352 self.app_state
1353 .total_failed_calls
1354 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1355 let error_event = crate::event::SessionEvent::Error {
1356 track_id: self.session_id.clone(),
1357 timestamp: crate::media::get_timestamp(),
1358 sender,
1359 error: e.to_string(),
1360 code: None,
1361 };
1362 self.event_sender.send(error_event).ok();
1363 self.hangup_now(Some(CallRecordHangupReason::BySystem))
1364 .await;
1365 return Err(e);
1366 }
1367 }
1368 }
1369
1370 async fn do_invite(&self, runtime: &mut CallRuntime, option: CallOption) -> Result<()> {
1371 let me = runtime
1378 .me
1379 .clone()
1380 .ok_or_else(|| anyhow::anyhow!("invite is only supported inside serve()"))?;
1381 crate::spawn(async move {
1382 if let Err(e) = me.invite_or_accept(option, "invite".to_string()).await {
1383 warn!(session_id = me.session_id, "{}", e);
1384 me.event_sender
1385 .send(SessionEvent::Error {
1386 track_id: me.session_id.clone(),
1387 timestamp: crate::media::get_timestamp(),
1388 sender: "command".to_string(),
1389 error: e.to_string(),
1390 code: None,
1391 })
1392 .ok();
1393 }
1394 });
1395 Ok(())
1396 }
1397
1398 async fn do_accept(&self, mut option: CallOption) -> Result<()> {
1399 let has_pending = self.has_pending_invite();
1400 let ready_to_answer_val = !self.has_ready_to_answer();
1401
1402 if ready_to_answer_val {
1403 if !has_pending {
1404 warn!(session_id = self.session_id, "no pending call to accept");
1406 let rejet_event = crate::event::SessionEvent::Reject {
1407 track_id: self.session_id.clone(),
1408 timestamp: crate::media::get_timestamp(),
1409 reason: "no pending call".to_string(),
1410 refer: None,
1411 code: Some(486),
1412 };
1413 self.event_sender.send(rejet_event).ok();
1414 self.hangup_now(Some(CallRecordHangupReason::BySystem))
1415 .await;
1416 return Err(anyhow::anyhow!("no pending call to accept"));
1417 }
1418 option = self.invite_or_accept(option, "accept".to_string()).await?;
1419 } else {
1420 option.check_default();
1421 if let Some(opt) = self.build_record_option(&option) {
1422 self.media_stream.update_recorder_option(opt).await;
1423 }
1424 self.set_option(option.clone());
1425 self.ensure_call_ambiance(&option).await;
1426 }
1427 info!(session_id = self.session_id, ?option, "accepting call");
1428 let ready = self.take_ready_to_answer();
1429 if let Some(ready) = ready {
1430 let ReadyAnswer {
1433 answer,
1434 track: pending_track,
1435 dialog,
1436 } = match Arc::try_unwrap(ready) {
1437 Ok(ready) => ready,
1438 Err(_) => {
1439 warn!(
1440 session_id = self.session_id,
1441 "ready_to_answer held elsewhere; skipping accept"
1442 );
1443 return Ok(());
1444 }
1445 };
1446 info!(session_id = self.session_id, "ready to answer with track");
1447
1448 let headers = vec![rsipstack::rsip::Header::ContentType(
1449 "application/sdp".to_string().into(),
1450 )];
1451
1452 match dialog.accept(Some(headers), Some(answer.as_bytes().to_vec())) {
1453 Ok(_) => {
1454 self.leg().update_progress(|p| {
1455 p.answer = Some(answer.clone());
1456 p.answer_time.get_or_insert_with(Utc::now);
1457 });
1458 self.finish_caller_stack(&option, pending_track).await?;
1459 }
1460 Err(e) => {
1461 warn!(session_id = self.session_id, "failed to accept call: {}", e);
1462 return Err(anyhow::anyhow!("failed to accept call"));
1463 }
1464 }
1465 }
1466
1467 if let Some(pending) = self.take_pending_sip_answer() {
1472 let Ok(pending) = Arc::try_unwrap(pending) else {
1473 warn!(
1474 session_id = self.session_id,
1475 "pending sip answer held elsewhere; skipping sip accept"
1476 );
1477 return Ok(());
1478 };
1479 let headers = vec![rsipstack::rsip::Header::ContentType(
1480 "application/sdp".to_string().into(),
1481 )];
1482 match pending
1483 .dialog
1484 .accept(Some(headers), Some(pending.answer.as_bytes().to_vec()))
1485 {
1486 Ok(_) => {
1487 info!(
1488 session_id = self.session_id,
1489 "answered underlying sip dialog"
1490 );
1491 self.leg().update_progress(|p| {
1492 p.answer = Some(pending.answer.clone());
1493 p.answer_time.get_or_insert_with(Utc::now);
1494 });
1495 self.media_stream.update_track(pending.track, None).await;
1498 }
1499 Err(e) => {
1500 warn!(
1501 session_id = self.session_id,
1502 "failed to accept underlying sip dialog: {}", e
1503 );
1504 }
1505 }
1506 }
1507 Ok(())
1508 }
1509
1510 async fn do_reject(
1511 &self,
1512 code: Option<rsipstack::rsip::StatusCode>,
1513 reason: Option<String>,
1514 ) -> Result<()> {
1515 if let Some(pending) = self.take_pending_sip_answer() {
1516 info!(
1517 session_id = self.session_id,
1518 ?reason,
1519 ?code,
1520 "rejecting underlying sip dialog"
1521 );
1522 if let Ok(pending) = Arc::try_unwrap(pending) {
1523 pending.dialog.reject(code.clone(), reason.clone()).ok();
1524 self.invitation.dialog_layer.remove_dialog(&pending.dialog.id());
1525 }
1526 }
1527 match self
1528 .invitation
1529 .find_dialog_id_by_session_id(&self.session_id)
1530 {
1531 Some(id) => {
1532 info!(
1533 session_id = self.session_id,
1534 ?reason,
1535 ?code,
1536 "rejecting call"
1537 );
1538 let result = self.invitation.hangup(id, code, reason).await;
1539 if result.is_ok() {
1540 self.cancel_token.cancel();
1541 }
1542 result
1543 }
1544 None => {
1545 if let Some(ready) = self.take_ready_to_answer() {
1546 info!(
1547 session_id = self.session_id,
1548 ?reason,
1549 ?code,
1550 "rejecting call from ready_to_answer"
1551 );
1552 let dialog = &ready.dialog;
1553 let dialog_id = dialog.id();
1554 dialog.reject(code, reason).ok();
1555 self.invitation.dialog_layer.remove_dialog(&dialog_id);
1556 self.cancel_token.cancel();
1557 }
1558 Ok(())
1559 }
1560 }
1561 }
1562
1563 async fn do_ringing(&self, command: Command) -> Result<()> {
1564 let Command::Ringing {
1565 ringtone,
1566 recorder,
1567 early_media,
1568 } = command
1569 else {
1570 unreachable!("do_ringing called with non-Ringing command");
1571 };
1572
1573 if !self.has_ready_to_answer() {
1574 let option = CallOption {
1575 recorder,
1576 ..Default::default()
1577 };
1578 let _ = self.invite_or_accept(option, "ringing".to_string()).await?;
1579 }
1580
1581 if let Some(ready) = self.ready_to_answer.load_full() {
1582 let (headers, body) = if early_media.unwrap_or_default() || ringtone.is_some() {
1583 let headers = vec![rsipstack::rsip::Header::ContentType(
1584 "application/sdp".to_string().into(),
1585 )];
1586 (Some(headers), Some(ready.answer.as_bytes().to_vec()))
1587 } else {
1588 (None, None)
1589 };
1590
1591 ready.dialog.ringing(headers, body).ok();
1592 info!(
1593 session_id = self.session_id,
1594 ringtone, early_media, "playing ringtone"
1595 );
1596 if let Some(ringtone_url) = ringtone {
1597 self.do_play(Command::Play {
1598 url: ringtone_url,
1599 play_id: None,
1600 auto_hangup: None,
1601 wait_input_timeout: None,
1602 offset_ms: None,
1603 })
1604 .await
1605 .ok();
1606 } else {
1607 info!(session_id = self.session_id, "no ringtone to play");
1608 }
1609 }
1610 Ok(())
1611 }
1612
1613 async fn do_tts(&self, command: Command) -> Result<()> {
1614 let Command::Tts {
1615 text,
1616 speaker,
1617 play_id,
1618 auto_hangup,
1619 streaming,
1620 end_of_stream,
1621 option,
1622 wait_input_timeout,
1623 base64,
1624 cache_key,
1625 } = command
1626 else {
1627 unreachable!("do_tts called with non-Tts command");
1628 };
1629 let streaming = streaming.unwrap_or_default();
1630 let end_of_stream = end_of_stream.unwrap_or_default();
1631 let base64 = base64.unwrap_or_default();
1632
1633 let tts_option = {
1634 let call_state = self.progress.load_full();
1635 match call_state.option.clone().unwrap_or_default().tts {
1636 Some(opt) => opt.merge_with(option),
1637 None => {
1638 if let Some(opt) = option {
1639 opt
1640 } else {
1641 return Err(anyhow::anyhow!("no tts option available"));
1642 }
1643 }
1644 }
1645 };
1646 let speaker = match speaker {
1647 Some(s) => Some(s),
1648 None => tts_option.speaker.clone(),
1649 };
1650
1651 let mut play_command = SynthesisCommand {
1652 text,
1653 speaker,
1654 play_id: play_id.clone(),
1655 streaming,
1656 end_of_stream: if !streaming { true } else { end_of_stream },
1657 option: tts_option,
1658 base64,
1659 cache_key,
1660 auto_hangup,
1661 };
1662 info!(
1663 session_id = self.session_id,
1664 provider = ?play_command.option.provider,
1665 text = %play_command.text.chars().take(10).collect::<String>(),
1666 speaker = play_command.speaker.as_deref(),
1667 auto_hangup = auto_hangup.unwrap_or_default(),
1668 play_id = play_command.play_id.as_deref(),
1669 streaming = play_command.streaming,
1670 end_of_stream = play_command.end_of_stream,
1671 wait_input_timeout = wait_input_timeout.unwrap_or_default(),
1672 is_base64 = play_command.base64,
1673 cache_key = play_command.cache_key.as_deref(),
1674 "new synthesis"
1675 );
1676
1677 let ssrc = rand::random::<u32>();
1678 let (should_interrupt, picked_ssrc) = {
1679 let existing_handle = self.tts_handle.load_full();
1680 let current_play_id = self.current_play();
1681
1682 let (target_ssrc, changed) = if let Some(handle) = &existing_handle {
1683 if play_id.is_some() && current_play_id != play_id {
1684 (ssrc, true)
1685 } else {
1686 (handle.ssrc, false)
1687 }
1688 } else {
1689 (ssrc, false)
1690 };
1691
1692 self.set_wait_input_timeout(wait_input_timeout);
1695
1696 self.set_current_play(play_id.clone());
1697 (changed, target_ssrc)
1698 };
1699
1700 if should_interrupt {
1701 let _ = self.do_interrupt(false).await;
1702 }
1703
1704 let existing_handle = self.tts_handle.load_full();
1708 if let Some(tts_handle) = existing_handle {
1709 match tts_handle.try_send(play_command) {
1710 Ok(_) => return Ok(()),
1711 Err(e) => {
1712 play_command = e.0;
1713 }
1714 }
1715 }
1716
1717 let (new_handle, tts_track) = StreamEngine::create_tts_track(
1718 self.app_state.stream_engine.clone(),
1719 self.cancel_token.child_token(),
1720 self.session_id.clone(),
1721 self.server_side_track_id.clone(),
1722 picked_ssrc,
1723 play_id.clone(),
1724 streaming,
1725 &play_command.option,
1726 play_command.auto_hangup,
1727 )
1728 .await?;
1729
1730 new_handle.try_send(play_command)?;
1731 self.tts_handle.store(Some(Arc::new(new_handle)));
1732 self.update_track_wrapper(tts_track, play_id).await;
1733 Ok(())
1734 }
1735
1736 async fn do_play(&self, command: Command) -> Result<()> {
1737 let Command::Play {
1738 url,
1739 play_id,
1740 auto_hangup,
1741 wait_input_timeout,
1742 offset_ms,
1743 } = command
1744 else {
1745 unreachable!("do_play called with non-Play command");
1746 };
1747 let ssrc = rand::random::<u32>();
1748 info!(
1749 session_id = self.session_id,
1750 ssrc, url, play_id, auto_hangup, "play file track"
1751 );
1752
1753 let play_id = play_id.or(Some(url.clone()));
1754
1755 let mut file_track = self
1757 .make_file_track(url, ssrc)
1758 .with_play_id(play_id.clone())
1759 .with_auto_hangup(auto_hangup);
1760
1761 if let Some(offset) = offset_ms {
1762 file_track = file_track.with_offset_ms(offset);
1763 }
1764
1765 {
1766 self.tts_handle.store(None);
1767 self.set_wait_input_timeout(wait_input_timeout);
1768 }
1769
1770 self.update_track_wrapper(Box::new(file_track), play_id)
1771 .await;
1772 Ok(())
1773 }
1774
1775 async fn do_history(&self, speaker: String, text: String) -> Result<()> {
1776 self.event_sender
1777 .send(SessionEvent::AddHistory {
1778 sender: Some(self.session_id.clone()),
1779 timestamp: crate::media::get_timestamp(),
1780 speaker,
1781 text,
1782 })
1783 .map(|_| ())
1784 .map_err(Into::into)
1785 }
1786
1787 fn do_custom(&self, sender: Option<String>, data: serde_json::Value) -> Result<()> {
1788 self.event_sender
1789 .send(SessionEvent::Custom {
1790 timestamp: crate::media::get_timestamp(),
1791 sender,
1792 data,
1793 })
1794 .map(|_| ())
1795 .map_err(Into::into)
1796 }
1797
1798 async fn do_interrupt(&self, graceful: bool) -> Result<()> {
1799 {
1800 self.tts_handle.store(None);
1801 self.set_moh(None);
1802 }
1803 self.media_stream
1804 .remove_track(&self.server_side_track_id, graceful)
1805 .await;
1806 Ok(())
1807 }
1808 async fn do_pause(&self) -> Result<()> {
1809 self.media_stream
1810 .pause_playback(self.server_side_track_id.clone())
1811 .await?;
1812 Ok(())
1813 }
1814 async fn do_resume(&self) -> Result<()> {
1815 self.media_stream
1816 .resume_playback(self.server_side_track_id.clone())
1817 .await?;
1818 Ok(())
1819 }
1820 async fn do_hangup(
1821 &self,
1822 reason: Option<CallRecordHangupReason>,
1823 initiator: Option<String>,
1824 headers: Option<HashMap<String, String>>,
1825 refer: Option<bool>,
1826 ) -> Result<()> {
1827 info!(
1828 session_id = self.session_id,
1829 ?reason,
1830 ?initiator,
1831 ?headers,
1832 ?refer,
1833 "do_hangup"
1834 );
1835
1836 let hangup_reason = match initiator.as_deref() {
1837 Some("caller") => CallRecordHangupReason::ByCaller,
1838 Some("callee") => CallRecordHangupReason::ByCallee,
1839 Some("system") => CallRecordHangupReason::Autohangup,
1840 _ => reason.unwrap_or(CallRecordHangupReason::BySystem),
1841 };
1842
1843 match refer {
1844 Some(true) => {
1845 let refer_token = self.take_refer_call_token();
1847 let refer_leg = self.refer_leg_value();
1848 let has_refer_leg = refer_leg.is_some();
1849 if let Some(leg) = refer_leg {
1850 if let Some(headers) = headers {
1851 let h_val = serde_json::to_value(&headers).unwrap_or_default();
1852 leg.set_extra("_hangup_headers", h_val);
1853 }
1854 let reason = hangup_reason.clone();
1856 leg.update_progress(|p| p.set_hangup_reason(reason.clone()));
1857 }
1858 if let Some(token) = refer_token {
1859 token.cancel();
1860 }
1861 if has_refer_leg {
1862 self.media_stream
1863 .remove_track(&self.server_side_track_id, false)
1864 .await;
1865 }
1866 }
1867 _ => {
1868 if let Some(headers) = headers {
1869 let h_val = serde_json::to_value(&headers).unwrap_or_default();
1870 self.leg().set_extra("_hangup_headers", h_val);
1871 }
1872 self.leg()
1873 .update_progress(|p| p.set_hangup_reason(hangup_reason.clone()));
1874 let refer_token = self.take_refer_call_token();
1875 self.media_stream
1876 .stop(Some(hangup_reason.to_string()), initiator);
1877 if let Some(token) = refer_token {
1878 token.cancel();
1879 }
1880 }
1881 }
1882 tokio::task::yield_now().await;
1883 Ok(())
1884 }
1885
1886 async fn do_refer(
1892 &self,
1893 runtime: &mut CallRuntime,
1894 caller: String,
1895 callee: String,
1896 refer_option: Option<ReferOption>,
1897 ) -> Result<()> {
1898 self.do_interrupt(false).await.ok();
1899
1900 let pause_parent_asr = refer_option
1902 .as_ref()
1903 .and_then(|o| o.pause_parent_asr)
1904 .unwrap_or(false);
1905
1906 let original_asr_option = if pause_parent_asr {
1908 self.progress
1909 .load_full()
1910 .option
1911 .as_ref()
1912 .and_then(|o| o.asr.clone())
1913 } else {
1914 None
1915 };
1916
1917 if pause_parent_asr {
1919 info!(
1920 session_id = self.session_id,
1921 "Pausing parent call ASR during refer"
1922 );
1923 self.media_stream
1924 .remove_processor::<crate::media::asr_processor::AsrProcessor>(
1925 &self.server_side_track_id,
1926 )
1927 .await
1928 .ok();
1929 }
1930
1931 let mut moh = refer_option.as_ref().and_then(|o| o.moh.clone());
1932 if let Some(ref path) = moh {
1933 if !path.starts_with("http") && !std::path::Path::new(path).exists() {
1934 let fallback = "./config/sounds/refer_moh.wav";
1935 if std::path::Path::new(fallback).exists() {
1936 info!(
1937 session_id = self.session_id,
1938 "moh {} not found, using fallback {}", path, fallback
1939 );
1940 moh = Some(fallback.to_string());
1941 }
1942 }
1943 }
1944 let ref_call_id = refer_option
1945 .as_ref()
1946 .and_then(|o| o.call_id.clone())
1947 .unwrap_or_else(|| format!("ref-{}-{}", rand::random::<u32>(), self.session_id));
1948
1949 let session_id = self.session_id.clone();
1950 let track_id = self.server_side_track_id.clone();
1951
1952 let (recorder, parent_caller) = {
1953 let progress = self.progress.load_full();
1954 let option = progress.option.as_ref();
1955 (
1956 option.map(|o| o.recorder.clone()).unwrap_or_default(),
1957 option.and_then(|o| o.caller.clone()),
1958 )
1959 };
1960 let caller = if caller.trim().is_empty() {
1961 parent_caller.unwrap_or_default()
1962 } else {
1963 caller
1964 };
1965
1966 let mut call_option = CallOption {
1967 caller: Some(caller),
1968 callee: Some(callee.clone()),
1969 sip: refer_option.as_ref().and_then(|o| o.sip.clone()),
1970 vad: refer_option
1971 .as_ref()
1972 .and_then(|o| o.vad.clone())
1973 .map(|mut opts| {
1974 opts.refer = Some(true);
1975 opts
1976 }),
1977 asr: refer_option
1978 .as_ref()
1979 .and_then(|o| o.asr.clone())
1980 .map(|mut opts| {
1981 opts.refer = Some(true);
1982 opts
1983 }),
1984 denoise: refer_option.as_ref().and_then(|o| o.denoise.clone()),
1985 agc: refer_option.as_ref().and_then(|o| o.agc.clone()),
1986 recorder,
1987 ..Default::default()
1988 };
1989 call_option.check_default();
1990
1991 let mut invite_option = call_option.build_invite_option()?;
1992 invite_option.call_id = Some(ref_call_id.clone());
1993
1994 let headers = invite_option.headers.get_or_insert_with(|| Vec::new());
1995
1996 {
1997 let progress = self.progress.load_full();
1998 if let Some(opt) = progress.option.as_ref() {
1999 if let Some(callee) = opt.callee.as_ref() {
2000 headers.push(rsipstack::rsip::Header::Other(
2001 "X-Referred-To".to_string(),
2002 callee.clone(),
2003 ));
2004 }
2005 if let Some(caller) = opt.caller.as_ref() {
2006 headers.push(rsipstack::rsip::Header::Other(
2007 "X-Referred-From".to_string(),
2008 caller.clone(),
2009 ));
2010 }
2011 }
2012 }
2013
2014 headers.push(rsipstack::rsip::Header::Other(
2015 "X-Referred-Id".to_string(),
2016 self.session_id.clone(),
2017 ));
2018
2019 let ssrc = rand::random::<u32>();
2020 let refer_leg = LegShared::new(
2021 ssrc,
2022 true,
2023 CallProgress {
2024 session_id: ref_call_id.clone(),
2025 start_time: Some(Utc::now()),
2026 option: Some(call_option.clone()),
2027 ..Default::default()
2028 },
2029 );
2030 self.set_refer_leg(Some(refer_leg.clone()));
2031
2032 let auto_hangup_requested = refer_option
2033 .as_ref()
2034 .and_then(|o| o.auto_hangup)
2035 .unwrap_or(true);
2036
2037 if !auto_hangup_requested && pause_parent_asr && original_asr_option.is_some() {
2041 let asr_option = original_asr_option.unwrap();
2042 self.set_pending_asr_resume((ssrc, asr_option));
2043 }
2044
2045 let timeout_secs = refer_option.as_ref().and_then(|o| o.timeout).unwrap_or(30);
2046 let forward_dtmf = refer_option
2047 .as_ref()
2048 .and_then(|o| o.forward_dtmf)
2049 .unwrap_or(true);
2050
2051 info!(
2052 session_id = self.session_id,
2053 ssrc,
2054 auto_hangup = auto_hangup_requested,
2055 callee,
2056 timeout_secs,
2057 "do_refer"
2058 );
2059
2060 let refer_cancel_token = self.cancel_token.child_token();
2061 self.set_refer_call_token(refer_cancel_token.clone());
2062
2063 let me = runtime
2066 .me
2067 .clone()
2068 .ok_or_else(|| anyhow::anyhow!("refer is only supported inside serve()"))?;
2069 let actor_tx = runtime.actor_tx.clone();
2070 let event_sender = self.event_sender.clone();
2071 let log_session_id = session_id.clone();
2072 let reject_track_id = track_id.clone();
2073 crate::spawn(async move {
2074 let out = crate::call::tracks::OutgoingLeg {
2075 cancel_token: refer_cancel_token,
2076 leg: refer_leg,
2077 track_id: track_id.clone(),
2078 invite_option,
2079 call_option,
2080 moh,
2081 auto_hangup: auto_hangup_requested,
2082 };
2083 let result = match tokio::time::timeout(
2084 Duration::from_secs(timeout_secs as u64),
2085 me.create_outgoing_sip_track(out),
2086 )
2087 .await
2088 {
2089 Ok(res) => res,
2090 Err(_) => {
2091 warn!(
2092 session_id = log_session_id,
2093 "refer sip track creation timed out after {} seconds", timeout_secs
2094 );
2095 event_sender
2096 .send(SessionEvent::Reject {
2097 track_id: reject_track_id,
2098 timestamp: crate::media::get_timestamp(),
2099 reason: "Timeout when refer".into(),
2100 code: Some(408),
2101 refer: Some(true),
2102 })
2103 .ok();
2104 Err(rsipstack::Error::Error(
2105 "refer sip track creation timed out".to_string(),
2106 ))
2107 }
2108 };
2109 me.set_moh(None);
2110 actor_tx
2111 .send(ActorMsg::ReferDone {
2112 track_id,
2113 forward_dtmf,
2114 result,
2115 })
2116 .await
2117 .ok();
2118 });
2119
2120 Ok(())
2121 }
2122
2123 async fn do_message(
2124 &self,
2125 body: String,
2126 content_type: Option<String>,
2127 headers: Option<HashMap<String, String>>,
2128 refer: Option<bool>,
2129 ) -> Result<()> {
2130 if !matches!(self.call_type, ActiveCallType::Sip | ActiveCallType::B2bua) {
2131 return Err(anyhow::anyhow!(
2132 "message command is only supported for SIP calls"
2133 ));
2134 }
2135
2136 let dialog_key = if refer == Some(true) {
2137 self.refer_leg_value()
2138 .map(|leg| leg.progress.load_full().session_id.clone())
2139 } else {
2140 Some(self.progress.load_full().session_id.clone())
2141 };
2142
2143 let mut dialog = dialog_key
2144 .as_ref()
2145 .filter(|id| !id.is_empty())
2146 .and_then(|id| self.invitation.dialog_layer.get_dialog_with(id));
2147
2148 if dialog.is_none() {
2151 if let Some(target_id) = dialog_key.as_ref().filter(|id| !id.is_empty()) {
2152 dialog = self
2153 .invitation
2154 .find_dialog_id_by_session_id(target_id)
2155 .and_then(|dialog_id| self.invitation.dialog_layer.get_dialog(&dialog_id));
2156 }
2157 }
2158
2159 if dialog.is_none() {
2160 if let Some(target_id) = dialog_key.as_ref().filter(|id| !id.is_empty()) {
2161 dialog = self
2162 .invitation
2163 .dialog_layer
2164 .all_dialog_ids()
2165 .into_iter()
2166 .filter_map(|id| self.invitation.dialog_layer.get_dialog_with(&id))
2167 .find(|dialog| dialog.id().to_string() == *target_id);
2168 }
2169 }
2170
2171 if dialog.is_none() {
2174 let call_id = match (refer == Some(true), dialog_key.as_deref()) {
2175 (true, Some(id)) if !id.is_empty() => Some(id),
2176 (false, _) => Some(self.session_id.as_str()),
2177 _ => None,
2178 };
2179 if let Some(call_id) = call_id {
2180 dialog = self
2181 .invitation
2182 .dialog_layer
2183 .get_client_dialog_by_call_id(call_id)
2184 .into_iter()
2185 .find(|d| {
2186 matches!(
2187 d.state(),
2188 rsipstack::dialog::dialog::DialogState::Confirmed(_, _)
2189 )
2190 })
2191 .map(rsipstack::dialog::dialog::Dialog::Invite);
2192 }
2193 }
2194
2195 let dialog = dialog.ok_or_else(|| {
2196 anyhow::anyhow!(
2197 "no established SIP dialog found for message command, refer={}",
2198 refer.unwrap_or_default()
2199 )
2200 })?;
2201
2202 let mut sip_headers = vec![rsipstack::rsip::Header::ContentType(
2203 content_type
2204 .clone()
2205 .unwrap_or_else(|| "text/plain;charset=utf-8".to_string())
2206 .into(),
2207 )];
2208 if let Some(headers) = &headers {
2209 sip_headers.extend(crate::sip_util::sip_headers_from_map(headers));
2210 }
2211
2212 info!(
2213 session_id = self.session_id,
2214 dialog_id = %dialog.id(),
2215 content_type = content_type.as_deref().unwrap_or("text/plain;charset=utf-8"),
2216 refer = refer.unwrap_or_default(),
2217 body = %body.chars().take(64).collect::<String>(),
2218 "sending SIP MESSAGE"
2219 );
2220
2221 let response = dialog
2222 .message(Some(sip_headers), Some(body.into_bytes()))
2223 .await?;
2224 match response {
2225 Some(resp)
2226 if resp.status_code.kind() == rsipstack::rsip::StatusCodeKind::Successful =>
2227 {
2228 Ok(())
2229 }
2230 Some(resp) => Err(anyhow::anyhow!(
2231 "SIP MESSAGE rejected with status {}",
2232 resp.status_code
2233 )),
2234 None => Err(anyhow::anyhow!(
2235 "SIP MESSAGE was not sent because dialog is not confirmed"
2236 )),
2237 }
2238 }
2239
2240 fn bridge_track_id(source_session_id: &str, target_session_id: &str) -> TrackId {
2241 format!("bridge:{}:to:{}", source_session_id, target_session_id)
2242 }
2243
2244 async fn do_bridge(&self, target_session_id: String) -> Result<()> {
2245 let target = {
2246 let calls = self.app_state.active_calls.lock().unwrap();
2247 calls.get(&target_session_id).cloned()
2248 };
2249 let target = target.ok_or_else(|| {
2250 anyhow::anyhow!("bridge target session not found: {}", target_session_id)
2251 })?;
2252
2253 if target.session_id == self.session_id {
2254 return Err(anyhow::anyhow!("cannot bridge a call to itself").into());
2255 }
2256
2257 let self_bridge_track_id = Self::bridge_track_id(&self.session_id, &target.session_id);
2258 let target_bridge_track_id = Self::bridge_track_id(&target.session_id, &self.session_id);
2259
2260 self.media_stream
2261 .remove_track(&self_bridge_track_id, false)
2262 .await;
2263 target
2264 .media_stream
2265 .remove_track(&target_bridge_track_id, false)
2266 .await;
2267
2268 let (self_bridge_sender, self_bridge_receiver) = mpsc::channel(25);
2269 let (target_bridge_sender, target_bridge_receiver) = mpsc::channel(25);
2270
2271 let self_paused = self.bridge_paused.clone();
2272 let target_paused = target.bridge_paused.clone();
2273
2274 let self_forwarding_track = ForwardingTrack::new(
2275 self_bridge_track_id.clone(),
2276 self.session_id.clone(),
2277 target_bridge_sender,
2278 self_bridge_receiver,
2279 self.track_config.clone(),
2280 self.cancel_token.child_token(),
2281 rand::random::<u32>(),
2282 self_paused,
2283 );
2284
2285 let target_forwarding_track = ForwardingTrack::new(
2286 target_bridge_track_id.clone(),
2287 target.session_id.clone(),
2288 self_bridge_sender,
2289 target_bridge_receiver,
2290 target.track_config.clone(),
2291 target.cancel_token.child_token(),
2292 rand::random::<u32>(),
2293 target_paused,
2294 );
2295
2296 self.media_stream
2297 .update_track(Box::new(self_forwarding_track), None)
2298 .await;
2299 target
2300 .media_stream
2301 .update_track(Box::new(target_forwarding_track), None)
2302 .await;
2303
2304 info!(
2305 session_id = self.session_id,
2306 target = target_session_id,
2307 self_bridge_track_id,
2308 target_bridge_track_id,
2309 "audio bridge established"
2310 );
2311 Ok(())
2312 }
2313
2314 async fn do_unbridge(&self, target_session_id: String) -> Result<()> {
2315 let target = {
2316 let calls = self.app_state.active_calls.lock().unwrap();
2317 calls.get(&target_session_id).cloned()
2318 };
2319
2320 let self_bridge_track_id = Self::bridge_track_id(&self.session_id, &target_session_id);
2321 self.media_stream
2322 .remove_track(&self_bridge_track_id, false)
2323 .await;
2324
2325 if let Some(target) = target {
2326 let target_bridge_track_id =
2327 Self::bridge_track_id(&target.session_id, &self.session_id);
2328 target
2329 .media_stream
2330 .remove_track(&target_bridge_track_id, false)
2331 .await;
2332 info!(
2333 session_id = self.session_id,
2334 target = target.session_id,
2335 self_bridge_track_id,
2336 target_bridge_track_id,
2337 "audio bridge removed"
2338 );
2339 } else {
2340 info!(
2341 session_id = self.session_id,
2342 target = target_session_id,
2343 self_bridge_track_id,
2344 "audio bridge removed locally; target session not active"
2345 );
2346 }
2347
2348 Ok(())
2349 }
2350
2351 async fn do_mute(&self, track_id: Option<String>) -> Result<()> {
2352 self.media_stream.mute_track(track_id).await;
2353 Ok(())
2354 }
2355
2356 async fn do_unmute(&self, track_id: Option<String>) -> Result<()> {
2357 self.media_stream.unmute_track(track_id).await;
2358 Ok(())
2359 }
2360
2361 pub async fn cleanup(&self) -> Result<()> {
2362 if matches!(self.call_type, ActiveCallType::Sip | ActiveCallType::B2bua) {
2363 self.reject_now(
2364 Some(rsipstack::rsip::StatusCode::Decline),
2365 Some("handler disconnected".to_string()),
2366 )
2367 .await;
2368 }
2369 if let Some(pending) = self.take_pending_sip_answer() {
2372 if let Ok(pending) = Arc::try_unwrap(pending) {
2373 pending
2374 .dialog
2375 .reject(
2376 Some(rsipstack::rsip::StatusCode::Decline),
2377 Some("handler disconnected".to_string()),
2378 )
2379 .ok();
2380 self.invitation
2381 .dialog_layer
2382 .remove_dialog(&pending.dialog.id());
2383 }
2384 }
2385 self.tts_handle.store(None);
2386 self.media_stream.cleanup().await.ok();
2387 Ok(())
2388 }
2389
2390 pub fn get_callrecord(&self) -> Option<CallRecord> {
2393 let progress = self.progress.load_full();
2394 let extras = self.extras.load_full();
2395 let refer_leg = self.refer_leg_value();
2396 Some(build_callrecord(
2397 &progress,
2398 &extras,
2399 refer_leg.as_ref(),
2400 &self.app_state,
2401 self.session_id.clone(),
2402 self.call_type.clone(),
2403 ))
2404 }
2405
2406 async fn dump_to_file(
2407 &self,
2408 dump_file: &mut File,
2409 cmd_receiver: &mut CommandReceiver,
2410 event_receiver: &mut EventReceiver,
2411 ) {
2412 loop {
2413 select! {
2414 _ = self.cancel_token.cancelled() => {
2415 break;
2416 }
2417 Ok(cmd) = cmd_receiver.recv() => {
2418 CallRecordEvent::write(CallRecordEventType::Command, cmd, dump_file)
2419 .await;
2420 }
2421 Ok(event) = event_receiver.recv() => {
2422 if matches!(event, SessionEvent::Binary{..}) {
2423 continue;
2424 }
2425 CallRecordEvent::write(CallRecordEventType::Event, event, dump_file)
2426 .await;
2427 }
2428 };
2429 }
2430 }
2431
2432 async fn dump_loop(
2433 &self,
2434 dump_events: bool,
2435 mut dump_cmd_receiver: CommandReceiver,
2436 mut dump_event_receiver: EventReceiver,
2437 ) {
2438 if !dump_events {
2439 return;
2440 }
2441
2442 let file_name = self.app_state.get_dump_events_file(&self.session_id);
2443 let mut dump_file = match File::options()
2444 .create(true)
2445 .append(true)
2446 .open(&file_name)
2447 .await
2448 {
2449 Ok(file) => file,
2450 Err(e) => {
2451 warn!(
2452 session_id = self.session_id,
2453 file_name, "failed to open dump events file: {}", e
2454 );
2455 return;
2456 }
2457 };
2458 self.dump_to_file(
2459 &mut dump_file,
2460 &mut dump_cmd_receiver,
2461 &mut dump_event_receiver,
2462 )
2463 .await;
2464
2465 while let Ok(event) = dump_event_receiver.try_recv() {
2466 if matches!(event, SessionEvent::Binary { .. }) {
2467 continue;
2468 }
2469 CallRecordEvent::write(CallRecordEventType::Event, event, &mut dump_file).await;
2470 }
2471 }
2472}
2473
2474struct CancelOnExit<'a>(&'a CancellationToken);
2477
2478impl Drop for CancelOnExit<'_> {
2479 fn drop(&mut self) {
2480 self.0.cancel();
2481 }
2482}
2483
2484impl Drop for ActiveCall {
2485 fn drop(&mut self) {
2486 info!(session_id = self.session_id, "dropping active call");
2487 if let Some(sender) = self.app_state.callrecord_sender.as_ref() {
2488 if let Some(record) = self.get_callrecord() {
2489 if let Err(e) = sender.send(record) {
2490 warn!(
2491 session_id = self.session_id,
2492 "failed to send call record: {}", e
2493 );
2494 }
2495 }
2496 }
2497 }
2498}