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 refer_call_token: ArcSwapOption<CancellationToken>,
554 pub wait_input_timeout: ArcSwapOption<u32>,
556 pub pending_asr_resume: ArcSwapOption<(u32, TranscriptionOption)>,
558 pub audio_receiver: std::sync::Mutex<Option<WebsocketBytesReceiver>>,
560}
561
562impl ActiveCall {
563 pub fn leg(&self) -> LegShared {
565 LegShared {
566 ssrc: self.ssrc,
567 is_refer: false,
568 progress: self.progress.clone(),
569 extras: self.extras.clone(),
570 }
571 }
572
573 pub fn set_option(&self, option: CallOption) {
575 self.progress.rcu(|p| {
576 let mut p = CallProgress::clone(p);
577 p.option = Some(option.clone());
578 p
579 });
580 }
581
582 pub fn moh_path(&self) -> Option<String> {
583 self.moh.load_full().map(|s| s.to_string())
584 }
585
586 pub fn set_moh(&self, v: Option<String>) {
587 self.moh.store(v.map(Arc::new));
588 }
589
590 pub fn current_play(&self) -> Option<String> {
591 self.current_play_id.load_full().map(|s| s.to_string())
592 }
593
594 pub fn set_current_play(&self, v: Option<String>) {
595 self.current_play_id.store(v.map(Arc::new));
596 }
597
598 pub fn refer_leg_value(&self) -> Option<LegShared> {
599 self.refer_leg.load_full().map(|l| l.as_ref().clone())
600 }
601
602 pub fn set_refer_leg(&self, v: Option<LegShared>) {
603 self.refer_leg.store(v.map(Arc::new));
604 }
605
606 pub fn set_ready_to_answer(&self, ready: ReadyAnswer) {
608 self.ready_to_answer.store(Some(Arc::new(ready)));
609 }
610
611 pub fn take_ready_to_answer(&self) -> Option<Arc<ReadyAnswer>> {
612 self.ready_to_answer.swap(None)
613 }
614
615 pub fn has_ready_to_answer(&self) -> bool {
616 self.ready_to_answer.load().is_some()
617 }
618
619 pub fn take_refer_call_token(&self) -> Option<CancellationToken> {
622 self.refer_call_token.swap(None).map(|t| (*t).clone())
623 }
624
625 pub fn set_refer_call_token(&self, token: CancellationToken) {
626 self.refer_call_token.store(Some(Arc::new(token)));
627 }
628
629 pub fn take_wait_input_timeout(&self) -> Option<u32> {
632 self.wait_input_timeout.swap(None).map(|t| *t)
633 }
634
635 pub fn set_wait_input_timeout(&self, v: Option<u32>) {
636 self.wait_input_timeout.store(v.map(Arc::new));
637 }
638
639 pub fn set_pending_asr_resume(&self, v: (u32, TranscriptionOption)) {
641 self.pending_asr_resume.store(Some(Arc::new(v)));
642 }
643
644 pub fn take_pending_asr_resume(&self) -> Option<(u32, TranscriptionOption)> {
645 self.pending_asr_resume.swap(None).map(|a| (*a).clone())
646 }
647
648 pub fn set_extra(&self, key: &str, value: serde_json::Value) {
650 self.leg().set_extra(key, value);
651 }
652
653 fn has_pending_invite(&self) -> bool {
655 self.invitation
656 .find_dialog_id_by_session_id(&self.session_id)
657 .is_some()
658 }
659
660 async fn hangup_now(&self, reason: Option<CallRecordHangupReason>) {
662 self.do_hangup(reason, None, None, None).await.ok();
663 }
664
665 async fn reject_now(&self, code: Option<rsipstack::rsip::StatusCode>, reason: Option<String>) {
667 self.do_reject(code, reason).await.ok();
668 }
669}
670
671pub struct ReadyAnswer {
674 pub answer: String,
675 pub track: PendingCallerTrack,
676 pub dialog: InviteDialog,
677}
678
679pub struct ActiveCallGuard {
680 pub call: ActiveCallRef,
681 pub active_calls: usize,
682}
683
684impl ActiveCallGuard {
685 pub fn new(call: ActiveCallRef) -> Self {
686 let active_calls = {
687 call.app_state
688 .total_calls
689 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
690 let mut calls = call.app_state.active_calls.lock().unwrap();
691 calls.insert(call.session_id.clone(), call.clone());
692 calls.len()
693 };
694 Self { call, active_calls }
695 }
696}
697
698impl Drop for ActiveCallGuard {
699 fn drop(&mut self) {
700 self.call
701 .app_state
702 .active_calls
703 .lock()
704 .unwrap()
705 .remove(&self.call.session_id);
706 }
707}
708
709pub struct ActiveCallReceiver {
710 pub cmd_receiver: CommandReceiver,
711 pub dump_cmd_receiver: CommandReceiver,
712 pub dump_event_receiver: EventReceiver,
713}
714
715pub struct CallSpec {
717 pub call_type: ActiveCallType,
718 pub cancel_token: CancellationToken,
719 pub session_id: String,
720 pub invitation: Invitation,
721 pub app_state: AppState,
722 pub track_config: TrackConfig,
723 pub audio_receiver: Option<WebsocketBytesReceiver>,
725 pub dump_events: bool,
726 pub server_side_track_id: Option<TrackId>,
728 pub extras: Option<HashMap<String, serde_json::Value>>,
731}
732
733impl ActiveCall {
734 pub fn new(spec: CallSpec) -> Self {
735 let CallSpec {
736 call_type,
737 cancel_token,
738 session_id,
739 invitation,
740 app_state,
741 track_config,
742 audio_receiver,
743 dump_events,
744 server_side_track_id,
745 extras,
746 } = spec;
747 let event_sender = crate::event::create_event_sender();
748 let cmd_sender = tokio::sync::broadcast::Sender::<Command>::new(32);
749 let server_side_track_id = server_side_track_id.unwrap_or(SERVER_SIDE_TRACK_ID.to_string());
750 let media_stream_builder = MediaStreamBuilder::new(event_sender.clone())
751 .with_id(session_id.clone())
752 .with_cancel_token(cancel_token.child_token());
753 let media_stream = Arc::new(media_stream_builder.build());
754 let start_time = Utc::now();
755 let call_type_str = match &call_type {
757 ActiveCallType::Sip => "sip",
758 ActiveCallType::WebSocket => "websocket",
759 ActiveCallType::Webrtc => "webrtc",
760 ActiveCallType::B2bua => "b2bua",
761 };
762 let mut extras = extras.unwrap_or_default();
763 extras
764 .entry(crate::playbook::BUILTIN_SESSION_ID.to_string())
765 .or_insert_with(|| serde_json::Value::String(session_id.clone()));
766 extras
767 .entry(crate::playbook::BUILTIN_CALL_TYPE.to_string())
768 .or_insert_with(|| serde_json::Value::String(call_type_str.to_string()));
769 extras
770 .entry(crate::playbook::BUILTIN_START_TIME.to_string())
771 .or_insert_with(|| serde_json::Value::String(start_time.to_rfc3339()));
772
773 let progress = CallProgress {
774 session_id: session_id.clone(),
775 start_time: Some(start_time),
776 ..Default::default()
777 };
778
779 Self {
780 cancel_token,
781 call_type,
782 session_id,
783 start_time,
784 media_stream,
785 track_config,
786 event_sender,
787 app_state,
788 invitation,
789 cmd_sender,
790 dump_events,
791 server_side_track_id,
792 ssrc: rand::random::<u32>(),
793 bridge_paused: Arc::new(AtomicBool::new(false)),
794 progress: Arc::new(ArcSwap::from_pointee(progress)),
795 extras: Arc::new(ArcSwap::from_pointee(extras)),
796 moh: ArcSwapOption::new(None),
797 current_play_id: ArcSwapOption::new(None),
798 tts_handle: ArcSwapOption::new(None),
799 refer_leg: ArcSwapOption::new(None),
800 ready_to_answer: ArcSwapOption::new(None),
801 refer_call_token: ArcSwapOption::new(None),
802 wait_input_timeout: ArcSwapOption::new(None),
803 pending_asr_resume: ArcSwapOption::new(None),
804 audio_receiver: std::sync::Mutex::new(audio_receiver),
805 }
806 }
807
808 pub async fn enqueue_command(&self, command: Command) -> Result<()> {
809 self.cmd_sender
810 .send(command)
811 .map_err(|e| anyhow::anyhow!("Failed to send command: {}", e))?;
812 Ok(())
813 }
814
815 pub fn new_receiver(&self) -> ActiveCallReceiver {
819 ActiveCallReceiver {
820 cmd_receiver: self.cmd_sender.subscribe(),
821 dump_cmd_receiver: self.cmd_sender.subscribe(),
822 dump_event_receiver: self.event_sender.subscribe(),
823 }
824 }
825
826 pub async fn serve(self: Arc<Self>, receiver: ActiveCallReceiver) -> Result<()> {
834 let ActiveCallReceiver {
835 mut cmd_receiver,
836 dump_cmd_receiver,
837 dump_event_receiver,
838 } = receiver;
839
840 let mut event_receiver = self.event_sender.subscribe();
841 let (actor_tx, mut actor_rx) = mpsc::channel::<ActorMsg>(16);
842 let mut runtime = CallRuntime::new(actor_tx);
843 runtime.me = Some(self.clone());
844
845 self.app_state
846 .total_calls
847 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
848
849 let me = self.clone();
850 let actor = async move {
851 let _cancel_on_exit = CancelOnExit(&me.cancel_token);
855 let mut ticker = tokio::time::interval(Duration::from_millis(100));
856 let mut media_serve = Box::pin(me.media_stream.serve());
858 loop {
859 tokio::select! {
860 cmd = cmd_receiver.recv() => {
861 match cmd {
862 Ok(command) => {
863 if let Err(e) = Box::pin(me.dispatch(&mut runtime, command)).await {
867 warn!(session_id = me.session_id, "{}", e);
868 me.event_sender
869 .send(SessionEvent::Error {
870 track_id: me.session_id.clone(),
871 timestamp: crate::media::get_timestamp(),
872 sender: "command".to_string(),
873 error: e.to_string(),
874 code: None,
875 })
876 .ok();
877 }
878 }
879 Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue,
880 Err(_) => {
881 info!(session_id = me.session_id, "command loop done");
882 break;
883 }
884 }
885 }
886 ev = event_receiver.recv() => {
887 match ev {
888 Ok(event) => Box::pin(me.handle_event(&mut runtime, event)).await,
889 Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue,
890 Err(_) => {
891 info!(session_id = me.session_id, "event loop done");
892 break;
893 }
894 }
895 }
896 Some(msg) = actor_rx.recv() => {
897 if let Err(e) = Box::pin(me.handle_actor_msg(msg)).await {
898 warn!(session_id = me.session_id, "{}", e);
899 me.event_sender
900 .send(SessionEvent::Error {
901 track_id: me.session_id.clone(),
902 timestamp: crate::media::get_timestamp(),
903 sender: "command".to_string(),
904 error: e.to_string(),
905 code: None,
906 })
907 .ok();
908 }
909 }
910 _ = ticker.tick() => {
911 Box::pin(me.check_input_timeout(&mut runtime)).await;
912 }
913 _ = &mut media_serve => {
914 info!(session_id = me.session_id, "media stream loop done");
915 break;
916 }
917 _ = me.cancel_token.cancelled() => {
918 info!(session_id = me.session_id, "call cancelled - cleaning up resources");
919 break;
920 }
921 }
922 }
923 };
924
925 tokio::join!(
926 self.dump_loop(self.dump_events, dump_cmd_receiver, dump_event_receiver),
927 actor
928 );
929 Ok(())
930 }
931
932 async fn check_input_timeout(&self, runtime: &mut CallRuntime) {
934 let (start_time, expire) = runtime.input_timeout_expire;
935 if expire > 0 && crate::media::get_timestamp() >= start_time + expire as u64 {
936 info!(session_id = self.session_id, "wait input timeout reached");
937 runtime.input_timeout_expire = (0, 0);
938 self.event_sender
939 .send(SessionEvent::Silence {
940 track_id: self.server_side_track_id.clone(),
941 timestamp: crate::media::get_timestamp(),
942 start_time,
943 duration: expire as u64,
944 samples: None,
945 refer: Some(false),
946 })
947 .ok();
948 }
949 }
950
951 async fn handle_event(&self, runtime: &mut CallRuntime, event: SessionEvent) {
953 match event {
954 SessionEvent::Speaking { .. }
955 | SessionEvent::Dtmf { .. }
956 | SessionEvent::AsrDelta { .. }
957 | SessionEvent::AsrFinal { .. }
958 | SessionEvent::TrackStart { .. } => {
959 runtime.input_timeout_expire = (0, 0);
960 }
961 SessionEvent::TrackEnd {
962 track_id,
963 play_id,
964 ssrc,
965 auto_hangup,
966 ..
967 } => {
968 if track_id != self.server_side_track_id {
969 return;
970 }
971
972 if play_id != self.current_play() {
973 debug!(
974 session_id = self.session_id,
975 ?play_id,
976 current = ?self.current_play(),
977 "ignoring interrupted track end"
978 );
979 return;
980 }
981 self.set_current_play(None);
982 let moh_path = self.moh_path();
983 let wait_timeout_val = self.take_wait_input_timeout();
984
985 if let Some(path) = moh_path {
986 info!(session_id = self.session_id, "looping moh: {}", path);
987 let ssrc = rand::random::<u32>();
988 let file_track = self.make_file_track(path.clone(), ssrc);
989 self.update_track_wrapper(Box::new(file_track), Some(path))
990 .await;
991 return;
992 }
993
994 if let Some(hangup_reason) = auto_hangup {
995 info!(
996 session_id = self.session_id,
997 ssrc, "auto hangup when track end track_id:{}", track_id
998 );
999 self.do_hangup(Some(hangup_reason), None, None, None)
1000 .await
1001 .ok();
1002 }
1003
1004 if let Some(timeout) = wait_timeout_val {
1005 runtime.input_timeout_expire = if timeout > 0 {
1006 (crate::media::get_timestamp(), timeout)
1007 } else {
1008 (0, 0)
1009 };
1010 }
1011 }
1012 SessionEvent::Interrupt { receiver } => {
1013 let track_id = receiver.unwrap_or_else(|| self.server_side_track_id.clone());
1014 if track_id == self.server_side_track_id {
1015 debug!(
1016 session_id = self.session_id,
1017 "received interrupt event, stopping playback"
1018 );
1019 self.do_interrupt(true).await.ok();
1020 }
1021 }
1022 SessionEvent::Inactivity { track_id, .. } => {
1023 info!(
1024 session_id = self.session_id,
1025 track_id, "inactivity timeout reached, hanging up"
1026 );
1027 self.do_hangup(
1028 Some(CallRecordHangupReason::InactivityTimeout),
1029 None,
1030 None,
1031 None,
1032 )
1033 .await
1034 .ok();
1035 }
1036 SessionEvent::Hangup { refer, .. } => {
1037 if refer == Some(true) {
1039 if let Some((refer_ssrc, asr_option)) = self.take_pending_asr_resume() {
1040 let is_refer_hangup = self
1042 .refer_leg
1043 .load_full()
1044 .map(|leg| leg.ssrc == refer_ssrc)
1045 .unwrap_or(false);
1046
1047 if is_refer_hangup {
1048 info!(
1049 session_id = self.session_id,
1050 "Refer call ended, resuming parent ASR"
1051 );
1052
1053 match self
1055 .app_state
1056 .stream_engine
1057 .create_asr_processor(
1058 self.server_side_track_id.clone(),
1059 self.cancel_token.child_token(),
1060 asr_option,
1061 self.event_sender.clone(),
1062 )
1063 .await
1064 {
1065 Ok(asr_processor) => {
1066 if let Err(e) = self
1067 .media_stream
1068 .append_processor(&self.server_side_track_id, asr_processor)
1069 .await
1070 {
1071 warn!(
1072 session_id = self.session_id,
1073 "Failed to resume ASR after refer: {}", e
1074 );
1075 }
1076 }
1077 Err(e) => {
1078 warn!(
1079 session_id = self.session_id,
1080 "Failed to create ASR processor for resume: {}", e
1081 );
1082 }
1083 }
1084 }
1085 }
1086 }
1087 }
1088 SessionEvent::Error { track_id, .. } => {
1089 if track_id != self.server_side_track_id {
1090 return;
1091 }
1092
1093 let moh_info = {
1094 let path = self.moh_path();
1095 path.map(|path| {
1096 let fallback = "./config/sounds/refer_moh.wav".to_string();
1097 if path != fallback && std::path::Path::new(&fallback).exists() {
1098 info!(
1099 session_id = self.session_id,
1100 "moh error, switching to fallback: {}", fallback
1101 );
1102 self.set_moh(Some(fallback.clone()));
1103 fallback
1104 } else {
1105 info!(
1106 session_id = self.session_id,
1107 "looping moh on error: {}", path
1108 );
1109 path
1110 }
1111 })
1112 };
1113
1114 if let Some(next_path) = moh_info {
1115 let ssrc = rand::random::<u32>();
1116 let file_track = self.make_file_track(next_path.clone(), ssrc);
1117 self.update_track_wrapper(Box::new(file_track), Some(next_path))
1118 .await;
1119 }
1120 }
1121 SessionEvent::Hold { on_hold, .. } => {
1122 self.bridge_paused.store(on_hold, Ordering::Relaxed);
1123 }
1124 _ => {}
1125 }
1126 }
1127
1128 async fn handle_actor_msg(&self, msg: ActorMsg) -> Result<()> {
1130 match msg {
1131 ActorMsg::ReferDone {
1132 track_id,
1133 forward_dtmf,
1134 result,
1135 } => match result {
1136 Ok(answer) => {
1137 self.media_stream
1138 .set_track_refer(&track_id, Some(true))
1139 .await;
1140 if !forward_dtmf {
1141 self.media_stream
1142 .set_track_dtmf_forward(&track_id, false)
1143 .await;
1144 }
1145 self.event_sender
1146 .send(SessionEvent::Answer {
1147 timestamp: crate::media::get_timestamp(),
1148 track_id,
1149 sdp: answer,
1150 refer: Some(true),
1151 })
1152 .ok();
1153 Ok(())
1154 }
1155 Err(e) => {
1156 warn!(
1157 session_id = self.session_id,
1158 "failed to create refer sip track: {}", e
1159 );
1160 self.emit_reject_from_rsip_error(track_id, true, &e);
1161 Err(e.into())
1162 }
1163 },
1164 }
1165 }
1166
1167 async fn dispatch(&self, runtime: &mut CallRuntime, command: Command) -> Result<()> {
1168 match command {
1169 Command::Invite { option } => self.do_invite(runtime, option).await,
1170 Command::Accept { option } => self.do_accept(option).await,
1171 Command::Reject { reason, code } => {
1172 self.do_reject(code.map(|c| (c as u16).into()), Some(reason))
1173 .await
1174 }
1175 Command::Ringing { .. } => self.do_ringing(command).await,
1176 Command::Tts { .. } => self.do_tts(command).await,
1177 Command::Play { .. } => self.do_play(command).await,
1178 Command::Hangup {
1179 reason,
1180 initiator,
1181 headers,
1182 refer,
1183 } => {
1184 let reason = reason.map(|r| {
1185 r.parse::<CallRecordHangupReason>()
1186 .unwrap_or(CallRecordHangupReason::BySystem)
1187 });
1188 self.do_hangup(reason, initiator, headers, refer).await
1189 }
1190 Command::Refer {
1191 caller,
1192 callee,
1193 options,
1194 } => self.do_refer(runtime, caller, callee, options).await,
1195 Command::Message {
1196 body,
1197 content_type,
1198 headers,
1199 refer,
1200 } => self.do_message(body, content_type, headers, refer).await,
1201 Command::Bridge { target_session_id } => self.do_bridge(target_session_id).await,
1202 Command::Unbridge { target_session_id } => self.do_unbridge(target_session_id).await,
1203 Command::Mute { track_id } => self.do_mute(track_id).await,
1204 Command::Unmute { track_id } => self.do_unmute(track_id).await,
1205 Command::Pause {} => self.do_pause().await,
1206 Command::Resume {} => self.do_resume().await,
1207 Command::Interrupt {
1208 graceful: passage,
1209 fade_out_ms: _,
1210 } => self.do_interrupt(passage.unwrap_or_default()).await,
1211 Command::History { speaker, text } => self.do_history(speaker, text).await,
1212 Command::Custom { sender, data } => self.do_custom(sender, data),
1213 Command::AddIceCandidate {
1214 candidate,
1215 sdp_mid,
1216 sdp_mline_index,
1217 } => {
1218 self.media_stream
1219 .add_ice_candidate(&candidate, sdp_mid.as_deref(), sdp_mline_index)
1220 .await
1221 }
1222 }
1223 }
1224
1225 fn build_record_option(&self, option: &CallOption) -> Option<RecorderOption> {
1226 if let Some(recorder_option) = &option.recorder {
1227 let mut recorder_file = recorder_option.recorder_file.clone();
1228 if recorder_file.contains("{id}") {
1229 recorder_file = recorder_file.replace("{id}", &self.session_id);
1230 }
1231
1232 let recorder_file = if recorder_file.is_empty() {
1233 self.app_state.get_recorder_file(&self.session_id)
1234 } else {
1235 let p = Path::new(&recorder_file);
1236 p.is_absolute()
1237 .then(|| recorder_file.clone())
1238 .unwrap_or_else(|| self.app_state.get_recorder_file(&recorder_file))
1239 };
1240 info!(
1241 session_id = self.session_id,
1242 recorder_file, "created recording file"
1243 );
1244
1245 let track_samplerate = self.track_config.samplerate;
1246 let recorder_samplerate = if track_samplerate > 0 {
1247 track_samplerate
1248 } else {
1249 recorder_option.samplerate
1250 };
1251 let recorder_ptime = if recorder_option.ptime == 0 {
1252 200
1253 } else {
1254 recorder_option.ptime
1255 };
1256 let requested_format = recorder_option
1257 .format
1258 .unwrap_or(self.app_state.config.recorder_format());
1259 let format = requested_format.effective();
1260 if requested_format != format {
1261 warn!(
1262 session_id = self.session_id,
1263 requested = requested_format.extension(),
1264 "Recorder format fallback to wav due to unsupported feature"
1265 );
1266 }
1267 let mut recorder_config = RecorderOption {
1268 recorder_file,
1269 samplerate: recorder_samplerate,
1270 ptime: recorder_ptime,
1271 format: Some(format),
1272 };
1273 recorder_config.ensure_path_extension(format);
1274 Some(recorder_config)
1275 } else {
1276 None
1277 }
1278 }
1279
1280 async fn invite_or_accept(&self, mut option: CallOption, sender: String) -> Result<CallOption> {
1281 {
1283 let state = self.progress.load_full();
1284 option = state.merge_option(option);
1285 }
1286
1287 option.check_default();
1288 if let Some(opt) = self.build_record_option(&option) {
1289 self.media_stream.update_recorder_option(opt).await;
1290 }
1291 self.ensure_call_ambiance(&option).await;
1292
1293 if let Some(opt) = &option.media_pass {
1294 let track_id = self.server_side_track_id.clone();
1295 let cancel_token = self.cancel_token.child_token();
1296 let ssrc = rand::random::<u32>();
1297 let media_pass_track = MediaPassTrack::new(
1298 self.session_id.clone(),
1299 ssrc,
1300 track_id,
1301 cancel_token,
1302 opt.clone(),
1303 );
1304 self.update_track_wrapper(Box::new(media_pass_track), None)
1305 .await;
1306 }
1307
1308 info!(
1309 session_id = self.session_id,
1310 call_type = ?self.call_type,
1311 sender,
1312 ?option,
1313 "caller with option"
1314 );
1315
1316 match self.setup_caller_track(&option).await {
1317 Ok(_) => return Ok(option),
1318 Err(e) => {
1319 self.app_state
1320 .total_failed_calls
1321 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1322 let error_event = crate::event::SessionEvent::Error {
1323 track_id: self.session_id.clone(),
1324 timestamp: crate::media::get_timestamp(),
1325 sender,
1326 error: e.to_string(),
1327 code: None,
1328 };
1329 self.event_sender.send(error_event).ok();
1330 self.hangup_now(Some(CallRecordHangupReason::BySystem))
1331 .await;
1332 return Err(e);
1333 }
1334 }
1335 }
1336
1337 async fn do_invite(&self, runtime: &mut CallRuntime, option: CallOption) -> Result<()> {
1338 let me = runtime
1345 .me
1346 .clone()
1347 .ok_or_else(|| anyhow::anyhow!("invite is only supported inside serve()"))?;
1348 crate::spawn(async move {
1349 if let Err(e) = me.invite_or_accept(option, "invite".to_string()).await {
1350 warn!(session_id = me.session_id, "{}", e);
1351 me.event_sender
1352 .send(SessionEvent::Error {
1353 track_id: me.session_id.clone(),
1354 timestamp: crate::media::get_timestamp(),
1355 sender: "command".to_string(),
1356 error: e.to_string(),
1357 code: None,
1358 })
1359 .ok();
1360 }
1361 });
1362 Ok(())
1363 }
1364
1365 async fn do_accept(&self, mut option: CallOption) -> Result<()> {
1366 let has_pending = self.has_pending_invite();
1367 let ready_to_answer_val = !self.has_ready_to_answer();
1368
1369 if ready_to_answer_val {
1370 if !has_pending {
1371 warn!(session_id = self.session_id, "no pending call to accept");
1373 let rejet_event = crate::event::SessionEvent::Reject {
1374 track_id: self.session_id.clone(),
1375 timestamp: crate::media::get_timestamp(),
1376 reason: "no pending call".to_string(),
1377 refer: None,
1378 code: Some(486),
1379 };
1380 self.event_sender.send(rejet_event).ok();
1381 self.hangup_now(Some(CallRecordHangupReason::BySystem))
1382 .await;
1383 return Err(anyhow::anyhow!("no pending call to accept"));
1384 }
1385 option = self.invite_or_accept(option, "accept".to_string()).await?;
1386 } else {
1387 option.check_default();
1388 if let Some(opt) = self.build_record_option(&option) {
1389 self.media_stream.update_recorder_option(opt).await;
1390 }
1391 self.set_option(option.clone());
1392 self.ensure_call_ambiance(&option).await;
1393 }
1394 info!(session_id = self.session_id, ?option, "accepting call");
1395 let ready = self.take_ready_to_answer();
1396 if let Some(ready) = ready {
1397 let ReadyAnswer {
1400 answer,
1401 track: pending_track,
1402 dialog,
1403 } = match Arc::try_unwrap(ready) {
1404 Ok(ready) => ready,
1405 Err(_) => {
1406 warn!(
1407 session_id = self.session_id,
1408 "ready_to_answer held elsewhere; skipping accept"
1409 );
1410 return Ok(());
1411 }
1412 };
1413 info!(session_id = self.session_id, "ready to answer with track");
1414
1415 let headers = vec![rsipstack::rsip::Header::ContentType(
1416 "application/sdp".to_string().into(),
1417 )];
1418
1419 match dialog.accept(Some(headers), Some(answer.as_bytes().to_vec())) {
1420 Ok(_) => {
1421 self.leg().update_progress(|p| {
1422 p.answer = Some(answer.clone());
1423 p.answer_time.get_or_insert_with(Utc::now);
1424 });
1425 self.finish_caller_stack(&option, pending_track).await?;
1426 }
1427 Err(e) => {
1428 warn!(session_id = self.session_id, "failed to accept call: {}", e);
1429 return Err(anyhow::anyhow!("failed to accept call"));
1430 }
1431 }
1432 }
1433 return Ok(());
1434 }
1435
1436 async fn do_reject(
1437 &self,
1438 code: Option<rsipstack::rsip::StatusCode>,
1439 reason: Option<String>,
1440 ) -> Result<()> {
1441 match self
1442 .invitation
1443 .find_dialog_id_by_session_id(&self.session_id)
1444 {
1445 Some(id) => {
1446 info!(
1447 session_id = self.session_id,
1448 ?reason,
1449 ?code,
1450 "rejecting call"
1451 );
1452 let result = self.invitation.hangup(id, code, reason).await;
1453 if result.is_ok() {
1454 self.cancel_token.cancel();
1455 }
1456 result
1457 }
1458 None => {
1459 if let Some(ready) = self.take_ready_to_answer() {
1460 info!(
1461 session_id = self.session_id,
1462 ?reason,
1463 ?code,
1464 "rejecting call from ready_to_answer"
1465 );
1466 let dialog = &ready.dialog;
1467 let dialog_id = dialog.id();
1468 dialog.reject(code, reason).ok();
1469 self.invitation.dialog_layer.remove_dialog(&dialog_id);
1470 self.cancel_token.cancel();
1471 }
1472 Ok(())
1473 }
1474 }
1475 }
1476
1477 async fn do_ringing(&self, command: Command) -> Result<()> {
1478 let Command::Ringing {
1479 ringtone,
1480 recorder,
1481 early_media,
1482 } = command
1483 else {
1484 unreachable!("do_ringing called with non-Ringing command");
1485 };
1486
1487 if !self.has_ready_to_answer() {
1488 let option = CallOption {
1489 recorder,
1490 ..Default::default()
1491 };
1492 let _ = self.invite_or_accept(option, "ringing".to_string()).await?;
1493 }
1494
1495 if let Some(ready) = self.ready_to_answer.load_full() {
1496 let (headers, body) = if early_media.unwrap_or_default() || ringtone.is_some() {
1497 let headers = vec![rsipstack::rsip::Header::ContentType(
1498 "application/sdp".to_string().into(),
1499 )];
1500 (Some(headers), Some(ready.answer.as_bytes().to_vec()))
1501 } else {
1502 (None, None)
1503 };
1504
1505 ready.dialog.ringing(headers, body).ok();
1506 info!(
1507 session_id = self.session_id,
1508 ringtone, early_media, "playing ringtone"
1509 );
1510 if let Some(ringtone_url) = ringtone {
1511 self.do_play(Command::Play {
1512 url: ringtone_url,
1513 play_id: None,
1514 auto_hangup: None,
1515 wait_input_timeout: None,
1516 offset_ms: None,
1517 })
1518 .await
1519 .ok();
1520 } else {
1521 info!(session_id = self.session_id, "no ringtone to play");
1522 }
1523 }
1524 Ok(())
1525 }
1526
1527 async fn do_tts(&self, command: Command) -> Result<()> {
1528 let Command::Tts {
1529 text,
1530 speaker,
1531 play_id,
1532 auto_hangup,
1533 streaming,
1534 end_of_stream,
1535 option,
1536 wait_input_timeout,
1537 base64,
1538 cache_key,
1539 } = command
1540 else {
1541 unreachable!("do_tts called with non-Tts command");
1542 };
1543 let streaming = streaming.unwrap_or_default();
1544 let end_of_stream = end_of_stream.unwrap_or_default();
1545 let base64 = base64.unwrap_or_default();
1546
1547 let tts_option = {
1548 let call_state = self.progress.load_full();
1549 match call_state.option.clone().unwrap_or_default().tts {
1550 Some(opt) => opt.merge_with(option),
1551 None => {
1552 if let Some(opt) = option {
1553 opt
1554 } else {
1555 return Err(anyhow::anyhow!("no tts option available"));
1556 }
1557 }
1558 }
1559 };
1560 let speaker = match speaker {
1561 Some(s) => Some(s),
1562 None => tts_option.speaker.clone(),
1563 };
1564
1565 let mut play_command = SynthesisCommand {
1566 text,
1567 speaker,
1568 play_id: play_id.clone(),
1569 streaming,
1570 end_of_stream: if !streaming { true } else { end_of_stream },
1571 option: tts_option,
1572 base64,
1573 cache_key,
1574 auto_hangup,
1575 };
1576 info!(
1577 session_id = self.session_id,
1578 provider = ?play_command.option.provider,
1579 text = %play_command.text.chars().take(10).collect::<String>(),
1580 speaker = play_command.speaker.as_deref(),
1581 auto_hangup = auto_hangup.unwrap_or_default(),
1582 play_id = play_command.play_id.as_deref(),
1583 streaming = play_command.streaming,
1584 end_of_stream = play_command.end_of_stream,
1585 wait_input_timeout = wait_input_timeout.unwrap_or_default(),
1586 is_base64 = play_command.base64,
1587 cache_key = play_command.cache_key.as_deref(),
1588 "new synthesis"
1589 );
1590
1591 let ssrc = rand::random::<u32>();
1592 let (should_interrupt, picked_ssrc) = {
1593 let existing_handle = self.tts_handle.load_full();
1594 let current_play_id = self.current_play();
1595
1596 let (target_ssrc, changed) = if let Some(handle) = &existing_handle {
1597 if play_id.is_some() && current_play_id != play_id {
1598 (ssrc, true)
1599 } else {
1600 (handle.ssrc, false)
1601 }
1602 } else {
1603 (ssrc, false)
1604 };
1605
1606 self.set_wait_input_timeout(wait_input_timeout);
1609
1610 self.set_current_play(play_id.clone());
1611 (changed, target_ssrc)
1612 };
1613
1614 if should_interrupt {
1615 let _ = self.do_interrupt(false).await;
1616 }
1617
1618 let existing_handle = self.tts_handle.load_full();
1622 if let Some(tts_handle) = existing_handle {
1623 match tts_handle.try_send(play_command) {
1624 Ok(_) => return Ok(()),
1625 Err(e) => {
1626 play_command = e.0;
1627 }
1628 }
1629 }
1630
1631 let (new_handle, tts_track) = StreamEngine::create_tts_track(
1632 self.app_state.stream_engine.clone(),
1633 self.cancel_token.child_token(),
1634 self.session_id.clone(),
1635 self.server_side_track_id.clone(),
1636 picked_ssrc,
1637 play_id.clone(),
1638 streaming,
1639 &play_command.option,
1640 play_command.auto_hangup,
1641 )
1642 .await?;
1643
1644 new_handle.try_send(play_command)?;
1645 self.tts_handle.store(Some(Arc::new(new_handle)));
1646 self.update_track_wrapper(tts_track, play_id).await;
1647 Ok(())
1648 }
1649
1650 async fn do_play(&self, command: Command) -> Result<()> {
1651 let Command::Play {
1652 url,
1653 play_id,
1654 auto_hangup,
1655 wait_input_timeout,
1656 offset_ms,
1657 } = command
1658 else {
1659 unreachable!("do_play called with non-Play command");
1660 };
1661 let ssrc = rand::random::<u32>();
1662 info!(
1663 session_id = self.session_id,
1664 ssrc, url, play_id, auto_hangup, "play file track"
1665 );
1666
1667 let play_id = play_id.or(Some(url.clone()));
1668
1669 let mut file_track = self
1671 .make_file_track(url, ssrc)
1672 .with_play_id(play_id.clone())
1673 .with_auto_hangup(auto_hangup);
1674
1675 if let Some(offset) = offset_ms {
1676 file_track = file_track.with_offset_ms(offset);
1677 }
1678
1679 {
1680 self.tts_handle.store(None);
1681 self.set_wait_input_timeout(wait_input_timeout);
1682 }
1683
1684 self.update_track_wrapper(Box::new(file_track), play_id)
1685 .await;
1686 Ok(())
1687 }
1688
1689 async fn do_history(&self, speaker: String, text: String) -> Result<()> {
1690 self.event_sender
1691 .send(SessionEvent::AddHistory {
1692 sender: Some(self.session_id.clone()),
1693 timestamp: crate::media::get_timestamp(),
1694 speaker,
1695 text,
1696 })
1697 .map(|_| ())
1698 .map_err(Into::into)
1699 }
1700
1701 fn do_custom(&self, sender: Option<String>, data: serde_json::Value) -> Result<()> {
1702 self.event_sender
1703 .send(SessionEvent::Custom {
1704 timestamp: crate::media::get_timestamp(),
1705 sender,
1706 data,
1707 })
1708 .map(|_| ())
1709 .map_err(Into::into)
1710 }
1711
1712 async fn do_interrupt(&self, graceful: bool) -> Result<()> {
1713 {
1714 self.tts_handle.store(None);
1715 self.set_moh(None);
1716 }
1717 self.media_stream
1718 .remove_track(&self.server_side_track_id, graceful)
1719 .await;
1720 Ok(())
1721 }
1722 async fn do_pause(&self) -> Result<()> {
1723 self.media_stream
1724 .pause_playback(self.server_side_track_id.clone())
1725 .await?;
1726 Ok(())
1727 }
1728 async fn do_resume(&self) -> Result<()> {
1729 self.media_stream
1730 .resume_playback(self.server_side_track_id.clone())
1731 .await?;
1732 Ok(())
1733 }
1734 async fn do_hangup(
1735 &self,
1736 reason: Option<CallRecordHangupReason>,
1737 initiator: Option<String>,
1738 headers: Option<HashMap<String, String>>,
1739 refer: Option<bool>,
1740 ) -> Result<()> {
1741 info!(
1742 session_id = self.session_id,
1743 ?reason,
1744 ?initiator,
1745 ?headers,
1746 ?refer,
1747 "do_hangup"
1748 );
1749
1750 let hangup_reason = match initiator.as_deref() {
1751 Some("caller") => CallRecordHangupReason::ByCaller,
1752 Some("callee") => CallRecordHangupReason::ByCallee,
1753 Some("system") => CallRecordHangupReason::Autohangup,
1754 _ => reason.unwrap_or(CallRecordHangupReason::BySystem),
1755 };
1756
1757 match refer {
1758 Some(true) => {
1759 let refer_token = self.take_refer_call_token();
1761 let refer_leg = self.refer_leg_value();
1762 let has_refer_leg = refer_leg.is_some();
1763 if let Some(leg) = refer_leg {
1764 if let Some(headers) = headers {
1765 let h_val = serde_json::to_value(&headers).unwrap_or_default();
1766 leg.set_extra("_hangup_headers", h_val);
1767 }
1768 let reason = hangup_reason.clone();
1770 leg.update_progress(|p| p.set_hangup_reason(reason.clone()));
1771 }
1772 if let Some(token) = refer_token {
1773 token.cancel();
1774 }
1775 if has_refer_leg {
1776 self.media_stream
1777 .remove_track(&self.server_side_track_id, false)
1778 .await;
1779 }
1780 }
1781 _ => {
1782 if let Some(headers) = headers {
1783 let h_val = serde_json::to_value(&headers).unwrap_or_default();
1784 self.leg().set_extra("_hangup_headers", h_val);
1785 }
1786 self.leg()
1787 .update_progress(|p| p.set_hangup_reason(hangup_reason.clone()));
1788 let refer_token = self.take_refer_call_token();
1789 self.media_stream
1790 .stop(Some(hangup_reason.to_string()), initiator);
1791 if let Some(token) = refer_token {
1792 token.cancel();
1793 }
1794 }
1795 }
1796 tokio::task::yield_now().await;
1797 Ok(())
1798 }
1799
1800 async fn do_refer(
1806 &self,
1807 runtime: &mut CallRuntime,
1808 caller: String,
1809 callee: String,
1810 refer_option: Option<ReferOption>,
1811 ) -> Result<()> {
1812 self.do_interrupt(false).await.ok();
1813
1814 let pause_parent_asr = refer_option
1816 .as_ref()
1817 .and_then(|o| o.pause_parent_asr)
1818 .unwrap_or(false);
1819
1820 let original_asr_option = if pause_parent_asr {
1822 self.progress
1823 .load_full()
1824 .option
1825 .as_ref()
1826 .and_then(|o| o.asr.clone())
1827 } else {
1828 None
1829 };
1830
1831 if pause_parent_asr {
1833 info!(
1834 session_id = self.session_id,
1835 "Pausing parent call ASR during refer"
1836 );
1837 self.media_stream
1838 .remove_processor::<crate::media::asr_processor::AsrProcessor>(
1839 &self.server_side_track_id,
1840 )
1841 .await
1842 .ok();
1843 }
1844
1845 let mut moh = refer_option.as_ref().and_then(|o| o.moh.clone());
1846 if let Some(ref path) = moh {
1847 if !path.starts_with("http") && !std::path::Path::new(path).exists() {
1848 let fallback = "./config/sounds/refer_moh.wav";
1849 if std::path::Path::new(fallback).exists() {
1850 info!(
1851 session_id = self.session_id,
1852 "moh {} not found, using fallback {}", path, fallback
1853 );
1854 moh = Some(fallback.to_string());
1855 }
1856 }
1857 }
1858 let ref_call_id = refer_option
1859 .as_ref()
1860 .and_then(|o| o.call_id.clone())
1861 .unwrap_or_else(|| format!("ref-{}-{}", rand::random::<u32>(), self.session_id));
1862
1863 let session_id = self.session_id.clone();
1864 let track_id = self.server_side_track_id.clone();
1865
1866 let (recorder, parent_caller) = {
1867 let progress = self.progress.load_full();
1868 let option = progress.option.as_ref();
1869 (
1870 option.map(|o| o.recorder.clone()).unwrap_or_default(),
1871 option.and_then(|o| o.caller.clone()),
1872 )
1873 };
1874 let caller = if caller.trim().is_empty() {
1875 parent_caller.unwrap_or_default()
1876 } else {
1877 caller
1878 };
1879
1880 let mut call_option = CallOption {
1881 caller: Some(caller),
1882 callee: Some(callee.clone()),
1883 sip: refer_option.as_ref().and_then(|o| o.sip.clone()),
1884 vad: refer_option
1885 .as_ref()
1886 .and_then(|o| o.vad.clone())
1887 .map(|mut opts| {
1888 opts.refer = Some(true);
1889 opts
1890 }),
1891 asr: refer_option
1892 .as_ref()
1893 .and_then(|o| o.asr.clone())
1894 .map(|mut opts| {
1895 opts.refer = Some(true);
1896 opts
1897 }),
1898 denoise: refer_option.as_ref().and_then(|o| o.denoise.clone()),
1899 agc: refer_option.as_ref().and_then(|o| o.agc.clone()),
1900 recorder,
1901 ..Default::default()
1902 };
1903 call_option.check_default();
1904
1905 let mut invite_option = call_option.build_invite_option()?;
1906 invite_option.call_id = Some(ref_call_id.clone());
1907
1908 let headers = invite_option.headers.get_or_insert_with(|| Vec::new());
1909
1910 {
1911 let progress = self.progress.load_full();
1912 if let Some(opt) = progress.option.as_ref() {
1913 if let Some(callee) = opt.callee.as_ref() {
1914 headers.push(rsipstack::rsip::Header::Other(
1915 "X-Referred-To".to_string(),
1916 callee.clone(),
1917 ));
1918 }
1919 if let Some(caller) = opt.caller.as_ref() {
1920 headers.push(rsipstack::rsip::Header::Other(
1921 "X-Referred-From".to_string(),
1922 caller.clone(),
1923 ));
1924 }
1925 }
1926 }
1927
1928 headers.push(rsipstack::rsip::Header::Other(
1929 "X-Referred-Id".to_string(),
1930 self.session_id.clone(),
1931 ));
1932
1933 let ssrc = rand::random::<u32>();
1934 let refer_leg = LegShared::new(
1935 ssrc,
1936 true,
1937 CallProgress {
1938 session_id: ref_call_id.clone(),
1939 start_time: Some(Utc::now()),
1940 option: Some(call_option.clone()),
1941 ..Default::default()
1942 },
1943 );
1944 self.set_refer_leg(Some(refer_leg.clone()));
1945
1946 let auto_hangup_requested = refer_option
1947 .as_ref()
1948 .and_then(|o| o.auto_hangup)
1949 .unwrap_or(true);
1950
1951 if !auto_hangup_requested && pause_parent_asr && original_asr_option.is_some() {
1955 let asr_option = original_asr_option.unwrap();
1956 self.set_pending_asr_resume((ssrc, asr_option));
1957 }
1958
1959 let timeout_secs = refer_option.as_ref().and_then(|o| o.timeout).unwrap_or(30);
1960 let forward_dtmf = refer_option
1961 .as_ref()
1962 .and_then(|o| o.forward_dtmf)
1963 .unwrap_or(true);
1964
1965 info!(
1966 session_id = self.session_id,
1967 ssrc,
1968 auto_hangup = auto_hangup_requested,
1969 callee,
1970 timeout_secs,
1971 "do_refer"
1972 );
1973
1974 let refer_cancel_token = self.cancel_token.child_token();
1975 self.set_refer_call_token(refer_cancel_token.clone());
1976
1977 let me = runtime
1980 .me
1981 .clone()
1982 .ok_or_else(|| anyhow::anyhow!("refer is only supported inside serve()"))?;
1983 let actor_tx = runtime.actor_tx.clone();
1984 let event_sender = self.event_sender.clone();
1985 let log_session_id = session_id.clone();
1986 let reject_track_id = track_id.clone();
1987 crate::spawn(async move {
1988 let out = crate::call::tracks::OutgoingLeg {
1989 cancel_token: refer_cancel_token,
1990 leg: refer_leg,
1991 track_id: track_id.clone(),
1992 invite_option,
1993 call_option,
1994 moh,
1995 auto_hangup: auto_hangup_requested,
1996 };
1997 let result = match tokio::time::timeout(
1998 Duration::from_secs(timeout_secs as u64),
1999 me.create_outgoing_sip_track(out),
2000 )
2001 .await
2002 {
2003 Ok(res) => res,
2004 Err(_) => {
2005 warn!(
2006 session_id = log_session_id,
2007 "refer sip track creation timed out after {} seconds", timeout_secs
2008 );
2009 event_sender
2010 .send(SessionEvent::Reject {
2011 track_id: reject_track_id,
2012 timestamp: crate::media::get_timestamp(),
2013 reason: "Timeout when refer".into(),
2014 code: Some(408),
2015 refer: Some(true),
2016 })
2017 .ok();
2018 Err(rsipstack::Error::Error(
2019 "refer sip track creation timed out".to_string(),
2020 ))
2021 }
2022 };
2023 me.set_moh(None);
2024 actor_tx
2025 .send(ActorMsg::ReferDone {
2026 track_id,
2027 forward_dtmf,
2028 result,
2029 })
2030 .await
2031 .ok();
2032 });
2033
2034 Ok(())
2035 }
2036
2037 async fn do_message(
2038 &self,
2039 body: String,
2040 content_type: Option<String>,
2041 headers: Option<HashMap<String, String>>,
2042 refer: Option<bool>,
2043 ) -> Result<()> {
2044 if !matches!(self.call_type, ActiveCallType::Sip | ActiveCallType::B2bua) {
2045 return Err(anyhow::anyhow!(
2046 "message command is only supported for SIP calls"
2047 ));
2048 }
2049
2050 let dialog_key = if refer == Some(true) {
2051 self.refer_leg_value()
2052 .map(|leg| leg.progress.load_full().session_id.clone())
2053 } else {
2054 Some(self.progress.load_full().session_id.clone())
2055 };
2056
2057 let mut dialog = dialog_key
2058 .as_ref()
2059 .filter(|id| !id.is_empty())
2060 .and_then(|id| self.invitation.dialog_layer.get_dialog_with(id));
2061
2062 if dialog.is_none() {
2063 if let Some(target_id) = dialog_key.as_ref().filter(|id| !id.is_empty()) {
2064 dialog = self
2065 .invitation
2066 .dialog_layer
2067 .all_dialog_ids()
2068 .into_iter()
2069 .filter_map(|id| self.invitation.dialog_layer.get_dialog_with(&id))
2070 .find(|dialog| dialog.id().to_string() == *target_id);
2071 }
2072 }
2073
2074 if dialog.is_none() {
2077 let call_id = match (refer == Some(true), dialog_key.as_deref()) {
2078 (true, Some(id)) if !id.is_empty() => Some(id),
2079 (false, _) => Some(self.session_id.as_str()),
2080 _ => None,
2081 };
2082 if let Some(call_id) = call_id {
2083 dialog = self
2084 .invitation
2085 .dialog_layer
2086 .get_client_dialog_by_call_id(call_id)
2087 .into_iter()
2088 .find(|d| {
2089 matches!(
2090 d.state(),
2091 rsipstack::dialog::dialog::DialogState::Confirmed(_, _)
2092 )
2093 })
2094 .map(rsipstack::dialog::dialog::Dialog::Invite);
2095 }
2096 }
2097
2098 let dialog = dialog.ok_or_else(|| {
2099 anyhow::anyhow!(
2100 "no established SIP dialog found for message command, refer={}",
2101 refer.unwrap_or_default()
2102 )
2103 })?;
2104
2105 let mut sip_headers = vec![rsipstack::rsip::Header::ContentType(
2106 content_type
2107 .clone()
2108 .unwrap_or_else(|| "text/plain;charset=utf-8".to_string())
2109 .into(),
2110 )];
2111 if let Some(headers) = &headers {
2112 sip_headers.extend(crate::sip_util::sip_headers_from_map(headers));
2113 }
2114
2115 info!(
2116 session_id = self.session_id,
2117 dialog_id = %dialog.id(),
2118 content_type = content_type.as_deref().unwrap_or("text/plain;charset=utf-8"),
2119 refer = refer.unwrap_or_default(),
2120 body = %body.chars().take(64).collect::<String>(),
2121 "sending SIP MESSAGE"
2122 );
2123
2124 let response = dialog
2125 .message(Some(sip_headers), Some(body.into_bytes()))
2126 .await?;
2127 match response {
2128 Some(resp)
2129 if resp.status_code.kind() == rsipstack::rsip::StatusCodeKind::Successful =>
2130 {
2131 Ok(())
2132 }
2133 Some(resp) => Err(anyhow::anyhow!(
2134 "SIP MESSAGE rejected with status {}",
2135 resp.status_code
2136 )),
2137 None => Err(anyhow::anyhow!(
2138 "SIP MESSAGE was not sent because dialog is not confirmed"
2139 )),
2140 }
2141 }
2142
2143 fn bridge_track_id(source_session_id: &str, target_session_id: &str) -> TrackId {
2144 format!("bridge:{}:to:{}", source_session_id, target_session_id)
2145 }
2146
2147 async fn do_bridge(&self, target_session_id: String) -> Result<()> {
2148 let target = {
2149 let calls = self.app_state.active_calls.lock().unwrap();
2150 calls.get(&target_session_id).cloned()
2151 };
2152 let target = target.ok_or_else(|| {
2153 anyhow::anyhow!("bridge target session not found: {}", target_session_id)
2154 })?;
2155
2156 if target.session_id == self.session_id {
2157 return Err(anyhow::anyhow!("cannot bridge a call to itself").into());
2158 }
2159
2160 let self_bridge_track_id = Self::bridge_track_id(&self.session_id, &target.session_id);
2161 let target_bridge_track_id = Self::bridge_track_id(&target.session_id, &self.session_id);
2162
2163 self.media_stream
2164 .remove_track(&self_bridge_track_id, false)
2165 .await;
2166 target
2167 .media_stream
2168 .remove_track(&target_bridge_track_id, false)
2169 .await;
2170
2171 let (self_bridge_sender, self_bridge_receiver) = mpsc::channel(25);
2172 let (target_bridge_sender, target_bridge_receiver) = mpsc::channel(25);
2173
2174 let self_paused = self.bridge_paused.clone();
2175 let target_paused = target.bridge_paused.clone();
2176
2177 let self_forwarding_track = ForwardingTrack::new(
2178 self_bridge_track_id.clone(),
2179 self.session_id.clone(),
2180 target_bridge_sender,
2181 self_bridge_receiver,
2182 self.track_config.clone(),
2183 self.cancel_token.child_token(),
2184 rand::random::<u32>(),
2185 self_paused,
2186 );
2187
2188 let target_forwarding_track = ForwardingTrack::new(
2189 target_bridge_track_id.clone(),
2190 target.session_id.clone(),
2191 self_bridge_sender,
2192 target_bridge_receiver,
2193 target.track_config.clone(),
2194 target.cancel_token.child_token(),
2195 rand::random::<u32>(),
2196 target_paused,
2197 );
2198
2199 self.media_stream
2200 .update_track(Box::new(self_forwarding_track), None)
2201 .await;
2202 target
2203 .media_stream
2204 .update_track(Box::new(target_forwarding_track), None)
2205 .await;
2206
2207 info!(
2208 session_id = self.session_id,
2209 target = target_session_id,
2210 self_bridge_track_id,
2211 target_bridge_track_id,
2212 "audio bridge established"
2213 );
2214 Ok(())
2215 }
2216
2217 async fn do_unbridge(&self, target_session_id: String) -> Result<()> {
2218 let target = {
2219 let calls = self.app_state.active_calls.lock().unwrap();
2220 calls.get(&target_session_id).cloned()
2221 };
2222
2223 let self_bridge_track_id = Self::bridge_track_id(&self.session_id, &target_session_id);
2224 self.media_stream
2225 .remove_track(&self_bridge_track_id, false)
2226 .await;
2227
2228 if let Some(target) = target {
2229 let target_bridge_track_id =
2230 Self::bridge_track_id(&target.session_id, &self.session_id);
2231 target
2232 .media_stream
2233 .remove_track(&target_bridge_track_id, false)
2234 .await;
2235 info!(
2236 session_id = self.session_id,
2237 target = target.session_id,
2238 self_bridge_track_id,
2239 target_bridge_track_id,
2240 "audio bridge removed"
2241 );
2242 } else {
2243 info!(
2244 session_id = self.session_id,
2245 target = target_session_id,
2246 self_bridge_track_id,
2247 "audio bridge removed locally; target session not active"
2248 );
2249 }
2250
2251 Ok(())
2252 }
2253
2254 async fn do_mute(&self, track_id: Option<String>) -> Result<()> {
2255 self.media_stream.mute_track(track_id).await;
2256 Ok(())
2257 }
2258
2259 async fn do_unmute(&self, track_id: Option<String>) -> Result<()> {
2260 self.media_stream.unmute_track(track_id).await;
2261 Ok(())
2262 }
2263
2264 pub async fn cleanup(&self) -> Result<()> {
2265 if matches!(self.call_type, ActiveCallType::Sip | ActiveCallType::B2bua) {
2266 self.reject_now(
2267 Some(rsipstack::rsip::StatusCode::Decline),
2268 Some("handler disconnected".to_string()),
2269 )
2270 .await;
2271 }
2272 self.tts_handle.store(None);
2273 self.media_stream.cleanup().await.ok();
2274 Ok(())
2275 }
2276
2277 pub fn get_callrecord(&self) -> Option<CallRecord> {
2280 let progress = self.progress.load_full();
2281 let extras = self.extras.load_full();
2282 let refer_leg = self.refer_leg_value();
2283 Some(build_callrecord(
2284 &progress,
2285 &extras,
2286 refer_leg.as_ref(),
2287 &self.app_state,
2288 self.session_id.clone(),
2289 self.call_type.clone(),
2290 ))
2291 }
2292
2293 async fn dump_to_file(
2294 &self,
2295 dump_file: &mut File,
2296 cmd_receiver: &mut CommandReceiver,
2297 event_receiver: &mut EventReceiver,
2298 ) {
2299 loop {
2300 select! {
2301 _ = self.cancel_token.cancelled() => {
2302 break;
2303 }
2304 Ok(cmd) = cmd_receiver.recv() => {
2305 CallRecordEvent::write(CallRecordEventType::Command, cmd, dump_file)
2306 .await;
2307 }
2308 Ok(event) = event_receiver.recv() => {
2309 if matches!(event, SessionEvent::Binary{..}) {
2310 continue;
2311 }
2312 CallRecordEvent::write(CallRecordEventType::Event, event, dump_file)
2313 .await;
2314 }
2315 };
2316 }
2317 }
2318
2319 async fn dump_loop(
2320 &self,
2321 dump_events: bool,
2322 mut dump_cmd_receiver: CommandReceiver,
2323 mut dump_event_receiver: EventReceiver,
2324 ) {
2325 if !dump_events {
2326 return;
2327 }
2328
2329 let file_name = self.app_state.get_dump_events_file(&self.session_id);
2330 let mut dump_file = match File::options()
2331 .create(true)
2332 .append(true)
2333 .open(&file_name)
2334 .await
2335 {
2336 Ok(file) => file,
2337 Err(e) => {
2338 warn!(
2339 session_id = self.session_id,
2340 file_name, "failed to open dump events file: {}", e
2341 );
2342 return;
2343 }
2344 };
2345 self.dump_to_file(
2346 &mut dump_file,
2347 &mut dump_cmd_receiver,
2348 &mut dump_event_receiver,
2349 )
2350 .await;
2351
2352 while let Ok(event) = dump_event_receiver.try_recv() {
2353 if matches!(event, SessionEvent::Binary { .. }) {
2354 continue;
2355 }
2356 CallRecordEvent::write(CallRecordEventType::Event, event, &mut dump_file).await;
2357 }
2358 }
2359}
2360
2361struct CancelOnExit<'a>(&'a CancellationToken);
2364
2365impl Drop for CancelOnExit<'_> {
2366 fn drop(&mut self) {
2367 self.0.cancel();
2368 }
2369}
2370
2371impl Drop for ActiveCall {
2372 fn drop(&mut self) {
2373 info!(session_id = self.session_id, "dropping active call");
2374 if let Some(sender) = self.app_state.callrecord_sender.as_ref() {
2375 if let Some(record) = self.get_callrecord() {
2376 if let Err(e) = sender.send(record) {
2377 warn!(
2378 session_id = self.session_id,
2379 "failed to send call record: {}", e
2380 );
2381 }
2382 }
2383 }
2384 }
2385}