1use super::Command;
2use crate::{
3 CallOption, ReferOption,
4 event::{EventReceiver, EventSender, SessionEvent},
5 media::{
6 TrackId,
7 ambiance::AmbianceProcessor,
8 engine::StreamEngine,
9 negotiate::strip_ipv6_candidates,
10 processor::SubscribeProcessor,
11 recorder::RecorderOption,
12 stream::{MediaStream, MediaStreamBuilder},
13 track::{
14 Track, TrackConfig,
15 file::FileTrack,
16 media_pass::MediaPassTrack,
17 rtc::{RtcTrack, RtcTrackConfig},
18 tts::SynthesisHandle,
19 websocket::{WebsocketBytesReceiver, WebsocketTrack},
20 },
21 },
22 synthesis::{SynthesisCommand, SynthesisOption},
23 transcription::TranscriptionOption,
24};
25use crate::{
26 app::AppState,
27 call::{
28 CommandReceiver, CommandSender,
29 sip::{DialogStateReceiverGuard, Invitation, InviteDialogStates},
30 },
31 callrecord::{CallRecord, CallRecordEvent, CallRecordEventType, CallRecordHangupReason},
32 useragent::{
33 invitation::PendingDialog,
34 public_address::{
35 build_public_contact_uri, contact_needs_public_resolution, find_local_addr_for_uri,
36 },
37 },
38};
39use anyhow::Result;
40use audio_codec::CodecType;
41use chrono::{DateTime, Utc};
42use rsipstack::dialog::{invitation::InviteOption, server_dialog::ServerInviteDialog};
43use serde::{Deserialize, Serialize};
44use std::{collections::HashMap, path::Path, sync::Arc, time::Duration};
45use tokio::{fs::File, select, sync::Mutex, sync::RwLock, time::sleep};
46use tokio_util::sync::CancellationToken;
47use tracing::{debug, info, warn};
48
49#[cfg(test)]
50mod tests {
51 use super::*;
52 use crate::app::AppStateBuilder;
53 use crate::callrecord::CallRecordHangupReason;
54 use crate::config::Config;
55 use crate::media::track::tts::SynthesisHandle;
56 use crate::synthesis::SynthesisCommand;
57 use tokio::sync::mpsc;
58
59 #[tokio::test]
60 async fn test_tts_ssrc_reuse_for_autohangup() -> Result<()> {
61 let mut config = Config::default();
62 config.udp_port = 0; config.media_cache_path = "/tmp/mediacache".to_string();
64 let stream_engine = Arc::new(StreamEngine::default());
65 let app_state = AppStateBuilder::new()
66 .with_config(config)
67 .with_stream_engine(stream_engine)
68 .build()
69 .await?;
70
71 let cancel_token = CancellationToken::new();
72 let session_id = "test-session".to_string();
73 let track_config = TrackConfig::default();
74
75 let mut option = crate::CallOption::default();
76 option.tts = Some(crate::synthesis::SynthesisOption::default());
77
78 let active_call = Arc::new(ActiveCall::new(
79 ActiveCallType::Sip,
80 cancel_token.clone(),
81 session_id.clone(),
82 app_state.invitation.clone(),
83 app_state.clone(),
84 track_config,
85 None,
86 false,
87 None,
88 None,
89 None,
90 ));
91
92 {
93 let mut state = active_call.call_state.write().await;
94 state.option = Some(option);
95 }
96
97 let (tx, mut rx) = mpsc::unbounded_channel::<SynthesisCommand>();
98 let initial_ssrc = 12345;
99 let handle = SynthesisHandle::new(tx, Some("play_1".to_string()), initial_ssrc);
100
101 {
103 let mut state = active_call.call_state.write().await;
104 state.tts_handle = Some(handle);
105 state.current_play_id = Some("play_1".to_string());
106 }
107
108 active_call
110 .do_tts(
111 "hangup now".to_string(),
112 None,
113 Some("play_1".to_string()),
114 Some(true),
115 false,
116 true,
117 None,
118 None,
119 false,
120 None,
121 )
122 .await?;
123
124 {
126 let state = active_call.call_state.read().await;
127 assert!(state.auto_hangup.is_some());
128 let (h_ssrc, reason) = state.auto_hangup.clone().unwrap();
129 assert_eq!(
130 h_ssrc, initial_ssrc,
131 "SSRC should be reused from existing handle"
132 );
133 assert_eq!(reason, CallRecordHangupReason::BySystem);
134 }
135
136 let cmd = rx.try_recv().expect("Should have received tts command");
138 assert_eq!(cmd.text, "hangup now");
139
140 Ok(())
141 }
142
143 #[tokio::test]
144 async fn test_tts_new_ssrc_for_different_play_id() -> Result<()> {
145 let mut config = Config::default();
146 config.udp_port = 0; config.media_cache_path = "/tmp/mediacache".to_string();
148 let stream_engine = Arc::new(StreamEngine::default());
149 let app_state = AppStateBuilder::new()
150 .with_config(config)
151 .with_stream_engine(stream_engine)
152 .build()
153 .await?;
154
155 let active_call = Arc::new(ActiveCall::new(
156 ActiveCallType::Sip,
157 CancellationToken::new(),
158 "test-session".to_string(),
159 app_state.invitation.clone(),
160 app_state.clone(),
161 TrackConfig::default(),
162 None,
163 false,
164 None,
165 None,
166 None,
167 ));
168
169 let mut tts_opt = crate::synthesis::SynthesisOption::default();
170 tts_opt.provider = Some(crate::synthesis::SynthesisType::Aliyun);
171 let mut option = crate::CallOption::default();
172 option.tts = Some(tts_opt);
173 {
174 let mut state = active_call.call_state.write().await;
175 state.option = Some(option);
176 }
177
178 let (tx, _rx) = mpsc::unbounded_channel();
179 let initial_ssrc = 111;
180 let handle = SynthesisHandle::new(tx, Some("play_1".to_string()), initial_ssrc);
181
182 {
183 let mut state = active_call.call_state.write().await;
184 state.tts_handle = Some(handle);
185 state.current_play_id = Some("play_1".to_string());
186 }
187
188 active_call
190 .do_tts(
191 "new play".to_string(),
192 None,
193 Some("play_2".to_string()),
194 Some(true),
195 false,
196 true,
197 None,
198 None,
199 false,
200 None,
201 )
202 .await?;
203
204 {
206 let state = active_call.call_state.read().await;
207 let (h_ssrc, _) = state.auto_hangup.clone().unwrap();
208 assert_ne!(
209 h_ssrc, initial_ssrc,
210 "Should use a new SSRC for different play_id"
211 );
212 }
213
214 Ok(())
215 }
216}
217
218#[derive(Deserialize)]
219#[serde(rename_all = "camelCase")]
220pub struct CallParams {
221 pub id: Option<String>,
222 #[serde(rename = "dump")]
223 pub dump_events: Option<bool>,
224 #[serde(rename = "ping")]
225 pub ping_interval: Option<u32>,
226 pub server_side_track: Option<String>,
227}
228
229#[derive(Debug, Serialize, Deserialize, Clone, Default, PartialEq, Eq)]
230#[serde(rename_all = "camelCase")]
231pub enum ActiveCallType {
232 Webrtc,
233 B2bua,
234 WebSocket,
235 #[default]
236 Sip,
237}
238
239#[derive(Default)]
240pub struct ActiveCallState {
241 pub session_id: String,
242 pub start_time: DateTime<Utc>,
243 pub ring_time: Option<DateTime<Utc>>,
244 pub answer_time: Option<DateTime<Utc>>,
245 pub hangup_reason: Option<CallRecordHangupReason>,
246 pub last_status_code: u16,
247 pub option: Option<CallOption>,
248 pub answer: Option<String>,
249 pub ssrc: u32,
250 pub refer_callstate: Option<ActiveCallStateRef>,
251 pub extras: Option<HashMap<String, serde_json::Value>>,
252 pub is_refer: bool,
253 pub sip_hangup_headers_template: Option<HashMap<String, String>>,
254
255 pub tts_handle: Option<SynthesisHandle>,
257 pub auto_hangup: Option<(u32, CallRecordHangupReason)>,
258 pub wait_input_timeout: Option<u32>,
259 pub moh: Option<String>,
260 pub current_play_id: Option<String>,
261 pub audio_receiver: Option<WebsocketBytesReceiver>,
262 pub ready_to_answer: Option<(String, Option<Box<dyn Track>>, ServerInviteDialog)>,
263 pub pending_asr_resume: Option<(u32, TranscriptionOption)>,
264}
265
266pub type ActiveCallRef = Arc<ActiveCall>;
267pub type ActiveCallStateRef = Arc<RwLock<ActiveCallState>>;
268
269pub struct ActiveCall {
270 pub call_state: ActiveCallStateRef,
271 pub cancel_token: CancellationToken,
272 pub call_type: ActiveCallType,
273 pub session_id: String,
274 pub media_stream: Arc<MediaStream>,
275 pub track_config: TrackConfig,
276 pub event_sender: EventSender,
277 pub app_state: AppState,
278 pub invitation: Invitation,
279 pub cmd_sender: CommandSender,
280 pub dump_events: bool,
281 pub server_side_track_id: TrackId,
282}
283
284pub struct ActiveCallGuard {
285 pub call: ActiveCallRef,
286 pub active_calls: usize,
287}
288
289impl ActiveCallGuard {
290 pub fn new(call: ActiveCallRef) -> Self {
291 let active_calls = {
292 call.app_state
293 .total_calls
294 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
295 let mut calls = call.app_state.active_calls.lock().unwrap();
296 calls.insert(call.session_id.clone(), call.clone());
297 calls.len()
298 };
299 Self { call, active_calls }
300 }
301}
302
303impl Drop for ActiveCallGuard {
304 fn drop(&mut self) {
305 self.call
306 .app_state
307 .active_calls
308 .lock()
309 .unwrap()
310 .remove(&self.call.session_id);
311 }
312}
313
314pub struct ActiveCallReceiver {
315 pub cmd_receiver: CommandReceiver,
316 pub dump_cmd_receiver: CommandReceiver,
317 pub dump_event_receiver: EventReceiver,
318}
319
320impl ActiveCall {
321 pub fn new(
322 call_type: ActiveCallType,
323 cancel_token: CancellationToken,
324 session_id: String,
325 invitation: Invitation,
326 app_state: AppState,
327 track_config: TrackConfig,
328 audio_receiver: Option<WebsocketBytesReceiver>,
329 dump_events: bool,
330 server_side_track_id: Option<TrackId>,
331 extras: Option<HashMap<String, serde_json::Value>>,
332 sip_hangup_headers_template: Option<HashMap<String, String>>,
333 ) -> Self {
334 let event_sender = crate::event::create_event_sender();
335 let cmd_sender = tokio::sync::broadcast::Sender::<Command>::new(32);
336 let media_stream_builder = MediaStreamBuilder::new(event_sender.clone())
337 .with_id(session_id.clone())
338 .with_cancel_token(cancel_token.child_token());
339 let media_stream = Arc::new(media_stream_builder.build());
340 let start_time = Utc::now();
341 let call_type_str = match &call_type {
343 ActiveCallType::Sip => "sip",
344 ActiveCallType::WebSocket => "websocket",
345 ActiveCallType::Webrtc => "webrtc",
346 ActiveCallType::B2bua => "b2bua",
347 };
348 let extras = {
349 let mut e = extras.unwrap_or_default();
350 e.entry(crate::playbook::BUILTIN_SESSION_ID.to_string())
351 .or_insert_with(|| serde_json::Value::String(session_id.clone()));
352 e.entry(crate::playbook::BUILTIN_CALL_TYPE.to_string())
353 .or_insert_with(|| serde_json::Value::String(call_type_str.to_string()));
354 e.entry(crate::playbook::BUILTIN_START_TIME.to_string())
355 .or_insert_with(|| serde_json::Value::String(start_time.to_rfc3339()));
356 Some(e)
357 };
358 let call_state = Arc::new(RwLock::new(ActiveCallState {
359 session_id: session_id.clone(),
360 start_time,
361 ssrc: rand::random::<u32>(),
362 extras,
363 audio_receiver,
364 sip_hangup_headers_template,
365 ..Default::default()
366 }));
367 Self {
368 cancel_token,
369 call_type,
370 session_id,
371 call_state,
372 media_stream,
373 track_config,
374 event_sender,
375 app_state,
376 invitation,
377 cmd_sender,
378 dump_events,
379 server_side_track_id: server_side_track_id.unwrap_or("server-side-track".to_string()),
380 }
381 }
382
383 pub async fn enqueue_command(&self, command: Command) -> Result<()> {
384 self.cmd_sender
385 .send(command)
386 .map_err(|e| anyhow::anyhow!("Failed to send command: {}", e))?;
387 Ok(())
388 }
389
390 pub fn new_receiver(&self) -> ActiveCallReceiver {
394 ActiveCallReceiver {
395 cmd_receiver: self.cmd_sender.subscribe(),
396 dump_cmd_receiver: self.cmd_sender.subscribe(),
397 dump_event_receiver: self.event_sender.subscribe(),
398 }
399 }
400
401 pub async fn serve(&self, receiver: ActiveCallReceiver) -> Result<()> {
402 let ActiveCallReceiver {
403 mut cmd_receiver,
404 dump_cmd_receiver,
405 dump_event_receiver,
406 } = receiver;
407
408 let process_command_loop = async move {
409 while let Ok(command) = cmd_receiver.recv().await {
410 match Box::pin(self.dispatch(command)).await {
411 Ok(_) => (),
412 Err(e) => {
413 warn!(session_id = self.session_id, "{}", e);
414 self.event_sender
415 .send(SessionEvent::Error {
416 track_id: self.session_id.clone(),
417 timestamp: crate::media::get_timestamp(),
418 sender: "command".to_string(),
419 error: e.to_string(),
420 code: None,
421 })
422 .ok();
423 }
424 }
425 }
426 };
427 self.app_state
428 .total_calls
429 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
430
431 tokio::join!(
432 self.dump_loop(self.dump_events, dump_cmd_receiver, dump_event_receiver),
433 async {
434 select! {
435 _ = process_command_loop => {
436 info!(session_id = self.session_id, "command loop done");
437 }
438 _ = self.process() => {
439 info!(session_id = self.session_id, "call serve done");
440 }
441 _ = self.cancel_token.cancelled() => {
442 info!(session_id = self.session_id, "call cancelled - cleaning up resources");
443 }
444 }
445 self.cancel_token.cancel();
446 }
447 );
448 Ok(())
449 }
450
451 async fn process(&self) -> Result<()> {
452 let mut event_receiver = self.event_sender.subscribe();
453
454 let input_timeout_expire = Arc::new(Mutex::new((0u64, 0u32)));
455 let input_timeout_expire_ref = input_timeout_expire.clone();
456 let event_sender = self.event_sender.clone();
457 let wait_input_timeout_loop = async {
458 loop {
459 let (start_time, expire) = { *input_timeout_expire.lock().await };
460 if expire > 0 && crate::media::get_timestamp() >= start_time + expire as u64 {
461 info!(session_id = self.session_id, "wait input timeout reached");
462 *input_timeout_expire.lock().await = (0, 0);
463 event_sender
464 .send(SessionEvent::Silence {
465 track_id: self.server_side_track_id.clone(),
466 timestamp: crate::media::get_timestamp(),
467 start_time,
468 duration: expire as u64,
469 samples: None,
470 })
471 .ok();
472 }
473 sleep(Duration::from_millis(100)).await;
474 }
475 };
476 let server_side_track_id = self.server_side_track_id.clone();
477 let event_hook_loop = async move {
478 while let Ok(event) = event_receiver.recv().await {
479 match event {
480 SessionEvent::Speaking { .. }
481 | SessionEvent::Dtmf { .. }
482 | SessionEvent::AsrDelta { .. }
483 | SessionEvent::AsrFinal { .. }
484 | SessionEvent::TrackStart { .. } => {
485 *input_timeout_expire_ref.lock().await = (0, 0);
486 }
487 SessionEvent::TrackEnd {
488 track_id,
489 play_id,
490 ssrc,
491 ..
492 } => {
493 if track_id != server_side_track_id {
494 continue;
495 }
496
497 let (moh_path, auto_hangup, wait_timeout_val) = {
498 let mut state = self.call_state.write().await;
499 if play_id != state.current_play_id {
500 debug!(
501 session_id = self.session_id,
502 ?play_id,
503 current = ?state.current_play_id,
504 "ignoring interrupted track end"
505 );
506 continue;
507 }
508 state.current_play_id = None;
509 (
510 state.moh.clone(),
511 state.auto_hangup.clone(),
512 state.wait_input_timeout.take(),
513 )
514 };
515
516 if let Some(path) = moh_path {
517 info!(session_id = self.session_id, "looping moh: {}", path);
518 let ssrc = rand::random::<u32>();
519 let file_track = FileTrack::new(self.server_side_track_id.clone())
520 .with_play_id(Some(path.clone()))
521 .with_ssrc(ssrc)
522 .with_path(path.clone())
523 .with_cancel_token(self.cancel_token.child_token());
524 self.update_track_wrapper(Box::new(file_track), Some(path))
525 .await;
526 continue;
527 }
528
529 if let Some((hangup_ssrc, hangup_reason)) = auto_hangup {
530 if hangup_ssrc == ssrc {
531 info!(
532 session_id = self.session_id,
533 ssrc, "auto hangup when track end track_id:{}", track_id
534 );
535 self.do_hangup(Some(hangup_reason), None, None).await.ok();
536 }
537 }
538
539 if let Some(timeout) = wait_timeout_val {
540 let expire = if timeout > 0 {
541 (crate::media::get_timestamp(), timeout)
542 } else {
543 (0, 0)
544 };
545 *input_timeout_expire_ref.lock().await = expire;
546 }
547 }
548 SessionEvent::Interrupt { receiver } => {
549 let track_id =
550 receiver.unwrap_or_else(|| self.server_side_track_id.clone());
551 if track_id == self.server_side_track_id {
552 debug!(
553 session_id = self.session_id,
554 "received interrupt event, stopping playback"
555 );
556 self.do_interrupt(true).await.ok();
557 }
558 }
559 SessionEvent::Inactivity { track_id, .. } => {
560 info!(
561 session_id = self.session_id,
562 track_id, "inactivity timeout reached, hanging up"
563 );
564 self.do_hangup(Some(CallRecordHangupReason::InactivityTimeout), None, None)
565 .await
566 .ok();
567 }
568 SessionEvent::Hangup { refer, .. } => {
569 if refer == Some(true) {
571 let mut cs = self.call_state.write().await;
572 if let Some((refer_ssrc, asr_option)) = cs.pending_asr_resume.take() {
573 let is_refer_hangup = cs
575 .refer_callstate
576 .as_ref()
577 .map(|rcs| {
578 rcs.try_read()
579 .map(|g| g.ssrc == refer_ssrc)
580 .unwrap_or(false)
581 })
582 .unwrap_or(false);
583
584 if is_refer_hangup {
585 drop(cs); info!(
587 session_id = self.session_id,
588 "Refer call ended, resuming parent ASR"
589 );
590
591 match self
593 .app_state
594 .stream_engine
595 .create_asr_processor(
596 self.server_side_track_id.clone(),
597 self.cancel_token.child_token(),
598 asr_option,
599 self.event_sender.clone(),
600 )
601 .await
602 {
603 Ok(asr_processor) => {
604 if let Err(e) = self
605 .media_stream
606 .append_processor(
607 &self.server_side_track_id,
608 asr_processor,
609 )
610 .await
611 {
612 warn!(
613 session_id = self.session_id,
614 "Failed to resume ASR after refer: {}", e
615 );
616 }
617 }
618 Err(e) => {
619 warn!(
620 session_id = self.session_id,
621 "Failed to create ASR processor for resume: {}", e
622 );
623 }
624 }
625 }
626 }
627 }
628 }
629 SessionEvent::Error { track_id, .. } => {
630 if track_id != server_side_track_id {
631 continue;
632 }
633
634 let moh_info = {
635 let mut state = self.call_state.write().await;
636 if let Some(path) = state.moh.clone() {
637 let fallback = "./config/sounds/refer_moh.wav".to_string();
638 let next_path = if path != fallback
639 && std::path::Path::new(&fallback).exists()
640 {
641 info!(
642 session_id = self.session_id,
643 "moh error, switching to fallback: {}", fallback
644 );
645 state.moh = Some(fallback.clone());
646 fallback
647 } else {
648 info!(
649 session_id = self.session_id,
650 "looping moh on error: {}", path
651 );
652 path
653 };
654 Some(next_path)
655 } else {
656 None
657 }
658 };
659
660 if let Some(next_path) = moh_info {
661 let ssrc = rand::random::<u32>();
662 let file_track = FileTrack::new(self.server_side_track_id.clone())
663 .with_play_id(Some(next_path.clone()))
664 .with_ssrc(ssrc)
665 .with_path(next_path.clone())
666 .with_cancel_token(self.cancel_token.child_token());
667 self.update_track_wrapper(Box::new(file_track), Some(next_path))
668 .await;
669 continue;
670 }
671 }
672 _ => {}
673 }
674 }
675 };
676
677 select! {
678 _ = wait_input_timeout_loop=>{
679 info!(session_id = self.session_id, "wait input timeout loop done");
680 }
681 _ = self.media_stream.serve() => {
682 info!(session_id = self.session_id, "media stream loop done");
683 }
684 _ = event_hook_loop => {
685 info!(session_id = self.session_id, "event loop done");
686 }
687 }
688 Ok(())
689 }
690
691 async fn dispatch(&self, command: Command) -> Result<()> {
692 match command {
693 Command::Invite { option } => self.do_invite(option).await,
694 Command::Accept { option } => self.do_accept(option).await,
695 Command::Reject { reason, code } => {
696 self.do_reject(code.map(|c| (c as u16).into()), Some(reason))
697 .await
698 }
699 Command::Ringing {
700 ringtone,
701 recorder,
702 early_media,
703 } => self.do_ringing(ringtone, recorder, early_media).await,
704 Command::Tts {
705 text,
706 speaker,
707 play_id,
708 auto_hangup,
709 streaming,
710 end_of_stream,
711 option,
712 wait_input_timeout,
713 base64,
714 cache_key,
715 } => {
716 self.do_tts(
717 text,
718 speaker,
719 play_id,
720 auto_hangup,
721 streaming.unwrap_or_default(),
722 end_of_stream.unwrap_or_default(),
723 option,
724 wait_input_timeout,
725 base64.unwrap_or_default(),
726 cache_key,
727 )
728 .await
729 }
730 Command::Play {
731 url,
732 play_id,
733 auto_hangup,
734 wait_input_timeout,
735 offset_ms,
736 } => {
737 self.do_play(url, play_id, auto_hangup, wait_input_timeout, offset_ms)
738 .await
739 }
740 Command::Hangup {
741 reason,
742 initiator,
743 headers,
744 } => {
745 let reason = reason.map(|r| {
746 r.parse::<CallRecordHangupReason>()
747 .unwrap_or(CallRecordHangupReason::BySystem)
748 });
749 self.do_hangup(reason, initiator, headers).await
750 }
751 Command::Refer {
752 caller,
753 callee,
754 options,
755 } => self.do_refer(caller, callee, options).await,
756 Command::Mute { track_id } => self.do_mute(track_id).await,
757 Command::Unmute { track_id } => self.do_unmute(track_id).await,
758 Command::Pause {} => self.do_pause().await,
759 Command::Resume {} => self.do_resume().await,
760 Command::Interrupt {
761 graceful: passage,
762 fade_out_ms: _,
763 } => self.do_interrupt(passage.unwrap_or_default()).await,
764 Command::History { speaker, text } => self.do_history(speaker, text).await,
765 Command::Custom { sender, data } => self.do_custom(sender, data),
766 }
767 }
768
769 fn build_record_option(&self, option: &CallOption) -> Option<RecorderOption> {
770 if let Some(recorder_option) = &option.recorder {
771 let mut recorder_file = recorder_option.recorder_file.clone();
772 if recorder_file.contains("{id}") {
773 recorder_file = recorder_file.replace("{id}", &self.session_id);
774 }
775
776 let recorder_file = if recorder_file.is_empty() {
777 self.app_state.get_recorder_file(&self.session_id)
778 } else {
779 let p = Path::new(&recorder_file);
780 p.is_absolute()
781 .then(|| recorder_file.clone())
782 .unwrap_or_else(|| self.app_state.get_recorder_file(&recorder_file))
783 };
784 info!(
785 session_id = self.session_id,
786 recorder_file, "created recording file"
787 );
788
789 let track_samplerate = self.track_config.samplerate;
790 let recorder_samplerate = if track_samplerate > 0 {
791 track_samplerate
792 } else {
793 recorder_option.samplerate
794 };
795 let recorder_ptime = if recorder_option.ptime == 0 {
796 200
797 } else {
798 recorder_option.ptime
799 };
800 let requested_format = recorder_option
801 .format
802 .unwrap_or(self.app_state.config.recorder_format());
803 let format = requested_format.effective();
804 if requested_format != format {
805 warn!(
806 session_id = self.session_id,
807 requested = requested_format.extension(),
808 "Recorder format fallback to wav due to unsupported feature"
809 );
810 }
811 let mut recorder_config = RecorderOption {
812 recorder_file,
813 samplerate: recorder_samplerate,
814 ptime: recorder_ptime,
815 format: Some(format),
816 };
817 recorder_config.ensure_path_extension(format);
818 Some(recorder_config)
819 } else {
820 None
821 }
822 }
823
824 async fn invite_or_accept(&self, mut option: CallOption, sender: String) -> Result<CallOption> {
825 {
827 let state = self.call_state.read().await;
828 option = state.merge_option(option);
829 }
830
831 option.check_default();
832 if let Some(opt) = self.build_record_option(&option) {
833 self.media_stream.update_recorder_option(opt).await;
834 }
835
836 if let Some(opt) = &option.media_pass {
837 let track_id = self.server_side_track_id.clone();
838 let cancel_token = self.cancel_token.child_token();
839 let ssrc = rand::random::<u32>();
840 let media_pass_track = MediaPassTrack::new(
841 self.session_id.clone(),
842 ssrc,
843 track_id,
844 cancel_token,
845 opt.clone(),
846 );
847 self.update_track_wrapper(Box::new(media_pass_track), None)
848 .await;
849 }
850
851 info!(
852 session_id = self.session_id,
853 call_type = ?self.call_type,
854 sender,
855 ?option,
856 "caller with option"
857 );
858
859 match self.setup_caller_track(&option).await {
860 Ok(_) => return Ok(option),
861 Err(e) => {
862 self.app_state
863 .total_failed_calls
864 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
865 let error_event = crate::event::SessionEvent::Error {
866 track_id: self.session_id.clone(),
867 timestamp: crate::media::get_timestamp(),
868 sender,
869 error: e.to_string(),
870 code: None,
871 };
872 self.event_sender.send(error_event).ok();
873 self.do_hangup(Some(CallRecordHangupReason::BySystem), None, None)
874 .await
875 .ok();
876 return Err(e);
877 }
878 }
879 }
880
881 async fn do_invite(&self, option: CallOption) -> Result<()> {
882 self.invite_or_accept(option, "invite".to_string())
883 .await
884 .map(|_| ())
885 }
886
887 async fn do_accept(&self, mut option: CallOption) -> Result<()> {
888 let has_pending = self
889 .invitation
890 .find_dialog_id_by_session_id(&self.session_id)
891 .is_some();
892 let ready_to_answer_val = {
893 let state = self.call_state.read().await;
894 state.ready_to_answer.is_none()
895 };
896
897 if ready_to_answer_val {
898 if !has_pending {
899 warn!(session_id = self.session_id, "no pending call to accept");
901 let rejet_event = crate::event::SessionEvent::Reject {
902 track_id: self.session_id.clone(),
903 timestamp: crate::media::get_timestamp(),
904 reason: "no pending call".to_string(),
905 refer: None,
906 code: Some(486),
907 };
908 self.event_sender.send(rejet_event).ok();
909 self.do_hangup(Some(CallRecordHangupReason::BySystem), None, None)
910 .await
911 .ok();
912 return Err(anyhow::anyhow!("no pending call to accept"));
913 }
914 option = self.invite_or_accept(option, "accept".to_string()).await?;
915 } else {
916 option.check_default();
917 self.call_state.write().await.option = Some(option.clone());
918 }
919 info!(session_id = self.session_id, ?option, "accepting call");
920 let ready = self.call_state.write().await.ready_to_answer.take();
921 if let Some((answer, track, dialog)) = ready {
922 info!(
923 session_id = self.session_id,
924 track_id = track.as_ref().map(|t| t.id()),
925 "ready to answer with track"
926 );
927
928 let headers = vec![rsipstack::rsip::Header::ContentType(
929 "application/sdp".to_string().into(),
930 )];
931
932 match dialog.accept(Some(headers), Some(answer.as_bytes().to_vec())) {
933 Ok(_) => {
934 {
935 let mut state = self.call_state.write().await;
936 state.answer = Some(answer);
937 state.answer_time = Some(Utc::now());
938 }
939 self.finish_caller_stack(&option, track).await?;
940 }
941 Err(e) => {
942 warn!(session_id = self.session_id, "failed to accept call: {}", e);
943 return Err(anyhow::anyhow!("failed to accept call"));
944 }
945 }
946 }
947 return Ok(());
948 }
949
950 async fn do_reject(
951 &self,
952 code: Option<rsipstack::rsip::StatusCode>,
953 reason: Option<String>,
954 ) -> Result<()> {
955 match self
956 .invitation
957 .find_dialog_id_by_session_id(&self.session_id)
958 {
959 Some(id) => {
960 info!(
961 session_id = self.session_id,
962 ?reason,
963 ?code,
964 "rejecting call"
965 );
966 let result = self.invitation.hangup(id, code, reason).await;
967 if result.is_ok() {
968 self.cancel_token.cancel();
969 }
970 result
971 }
972 None => {
973 let ready = self.call_state.write().await.ready_to_answer.take();
974 if let Some((_, _, dialog)) = ready {
975 info!(
976 session_id = self.session_id,
977 ?reason,
978 ?code,
979 "rejecting call from ready_to_answer"
980 );
981 let dialog_id = dialog.id();
982 dialog.reject(code, reason).ok();
983 self.invitation.dialog_layer.remove_dialog(&dialog_id);
984 self.cancel_token.cancel();
985 }
986 Ok(())
987 }
988 }
989 }
990
991 async fn do_ringing(
992 &self,
993 ringtone: Option<String>,
994 recorder: Option<RecorderOption>,
995 early_media: Option<bool>,
996 ) -> Result<()> {
997 let ready_to_answer_val = self.call_state.read().await.ready_to_answer.is_none();
998 if ready_to_answer_val {
999 let option = CallOption {
1000 recorder,
1001 ..Default::default()
1002 };
1003 let _ = self.invite_or_accept(option, "ringing".to_string()).await?;
1004 }
1005
1006 let state = self.call_state.read().await;
1007 if let Some((answer, _, dialog)) = state.ready_to_answer.as_ref() {
1008 let (headers, body) = if early_media.unwrap_or_default() || ringtone.is_some() {
1009 let headers = vec![rsipstack::rsip::Header::ContentType(
1010 "application/sdp".to_string().into(),
1011 )];
1012 (Some(headers), Some(answer.as_bytes().to_vec()))
1013 } else {
1014 (None, None)
1015 };
1016
1017 dialog.ringing(headers, body).ok();
1018 info!(
1019 session_id = self.session_id,
1020 ringtone, early_media, "playing ringtone"
1021 );
1022 if let Some(ringtone_url) = ringtone {
1023 drop(state);
1024 self.do_play(ringtone_url, None, None, None, None).await.ok();
1025 } else {
1026 info!(session_id = self.session_id, "no ringtone to play");
1027 }
1028 }
1029 Ok(())
1030 }
1031
1032 async fn do_tts(
1033 &self,
1034 text: String,
1035 speaker: Option<String>,
1036 play_id: Option<String>,
1037 auto_hangup: Option<bool>,
1038 streaming: bool,
1039 end_of_stream: bool,
1040 option: Option<SynthesisOption>,
1041 wait_input_timeout: Option<u32>,
1042 base64: bool,
1043 cache_key: Option<String>,
1044 ) -> Result<()> {
1045 let tts_option = {
1046 let call_state = self.call_state.read().await;
1047 match call_state.option.clone().unwrap_or_default().tts {
1048 Some(opt) => opt.merge_with(option),
1049 None => {
1050 if let Some(opt) = option {
1051 opt
1052 } else {
1053 return Err(anyhow::anyhow!("no tts option available"));
1054 }
1055 }
1056 }
1057 };
1058 let speaker = match speaker {
1059 Some(s) => Some(s),
1060 None => tts_option.speaker.clone(),
1061 };
1062
1063 let mut play_command = SynthesisCommand {
1064 text,
1065 speaker,
1066 play_id: play_id.clone(),
1067 streaming,
1068 end_of_stream: if !streaming { true } else { end_of_stream },
1069 option: tts_option,
1070 base64,
1071 cache_key,
1072 };
1073 info!(
1074 session_id = self.session_id,
1075 provider = ?play_command.option.provider,
1076 text = %play_command.text.chars().take(10).collect::<String>(),
1077 speaker = play_command.speaker.as_deref(),
1078 auto_hangup = auto_hangup.unwrap_or_default(),
1079 play_id = play_command.play_id.as_deref(),
1080 streaming = play_command.streaming,
1081 end_of_stream = play_command.end_of_stream,
1082 wait_input_timeout = wait_input_timeout.unwrap_or_default(),
1083 is_base64 = play_command.base64,
1084 cache_key = play_command.cache_key.as_deref(),
1085 "new synthesis"
1086 );
1087
1088 let ssrc = rand::random::<u32>();
1089 let (should_interrupt, picked_ssrc) = {
1090 let mut state = self.call_state.write().await;
1091
1092 let (target_ssrc, changed) = if let Some(handle) = &state.tts_handle {
1093 if play_id.is_some() && state.current_play_id != play_id {
1094 (ssrc, true)
1095 } else {
1096 (handle.ssrc, false)
1097 }
1098 } else {
1099 (ssrc, false)
1100 };
1101
1102 state.wait_input_timeout = wait_input_timeout;
1105
1106 state.current_play_id = play_id.clone();
1107 (changed, target_ssrc)
1108 };
1109
1110 if should_interrupt {
1111 let _ = self.do_interrupt(false).await;
1112 }
1113
1114 {
1118 let mut state = self.call_state.write().await;
1119 state.auto_hangup = match auto_hangup {
1120 Some(true) => Some((picked_ssrc, CallRecordHangupReason::BySystem)),
1121 _ => {
1122 if state.tts_handle.is_some() && !should_interrupt {
1126 state.auto_hangup.clone()
1127 } else {
1128 None
1129 }
1130 }
1131 };
1132 }
1133
1134 let existing_handle = self.call_state.read().await.tts_handle.clone();
1135 if let Some(tts_handle) = existing_handle {
1136 match tts_handle.try_send(play_command) {
1137 Ok(_) => return Ok(()),
1138 Err(e) => {
1139 play_command = e.0;
1140 }
1141 }
1142 }
1143
1144 let (new_handle, tts_track) = StreamEngine::create_tts_track(
1145 self.app_state.stream_engine.clone(),
1146 self.cancel_token.child_token(),
1147 self.session_id.clone(),
1148 self.server_side_track_id.clone(),
1149 picked_ssrc,
1150 play_id.clone(),
1151 streaming,
1152 &play_command.option,
1153 )
1154 .await?;
1155
1156 new_handle.try_send(play_command)?;
1157 self.call_state.write().await.tts_handle = Some(new_handle);
1158 self.update_track_wrapper(tts_track, play_id).await;
1159 Ok(())
1160 }
1161
1162 async fn do_play(
1163 &self,
1164 url: String,
1165 play_id: Option<String>,
1166 auto_hangup: Option<bool>,
1167 wait_input_timeout: Option<u32>,
1168 offset_ms: Option<u32>,
1169 ) -> Result<()> {
1170 let ssrc = rand::random::<u32>();
1171 info!(
1172 session_id = self.session_id,
1173 ssrc, url, play_id, auto_hangup, "play file track"
1174 );
1175
1176 let play_id = play_id.or(Some(url.clone()));
1177
1178 let mut file_track = FileTrack::new(self.server_side_track_id.clone())
1179 .with_play_id(play_id.clone())
1180 .with_ssrc(ssrc)
1181 .with_path(url)
1182 .with_cancel_token(self.cancel_token.child_token());
1183
1184 if let Some(offset) = offset_ms {
1185 file_track = file_track.with_offset_ms(offset);
1186 }
1187
1188 {
1189 let mut state = self.call_state.write().await;
1190 state.tts_handle = None;
1191 state.auto_hangup = match auto_hangup {
1192 Some(true) => Some((ssrc, CallRecordHangupReason::BySystem)),
1193 _ => None,
1194 };
1195 state.wait_input_timeout = wait_input_timeout;
1196 }
1197
1198 self.update_track_wrapper(Box::new(file_track), play_id)
1199 .await;
1200 Ok(())
1201 }
1202
1203 async fn do_history(&self, speaker: String, text: String) -> Result<()> {
1204 self.event_sender
1205 .send(SessionEvent::AddHistory {
1206 sender: Some(self.session_id.clone()),
1207 timestamp: crate::media::get_timestamp(),
1208 speaker,
1209 text,
1210 })
1211 .map(|_| ())
1212 .map_err(Into::into)
1213 }
1214
1215 fn do_custom(&self, sender: Option<String>, data: serde_json::Value) -> Result<()> {
1216 self.event_sender
1217 .send(SessionEvent::Custom {
1218 timestamp: crate::media::get_timestamp(),
1219 sender,
1220 data,
1221 })
1222 .map(|_| ())
1223 .map_err(Into::into)
1224 }
1225
1226 async fn do_interrupt(&self, graceful: bool) -> Result<()> {
1227 {
1228 let mut state = self.call_state.write().await;
1229 state.tts_handle = None;
1230 state.moh = None;
1231 state.auto_hangup = None;
1232 }
1233 self.media_stream
1234 .remove_track(&self.server_side_track_id, graceful)
1235 .await;
1236 Ok(())
1237 }
1238 async fn do_pause(&self) -> Result<()> {
1239 Ok(())
1240 }
1241 async fn do_resume(&self) -> Result<()> {
1242 Ok(())
1243 }
1244 async fn do_hangup(
1245 &self,
1246 reason: Option<CallRecordHangupReason>,
1247 initiator: Option<String>,
1248 headers: Option<HashMap<String, String>>,
1249 ) -> Result<()> {
1250 info!(
1251 session_id = self.session_id,
1252 ?reason,
1253 ?initiator,
1254 ?headers,
1255 "do_hangup"
1256 );
1257
1258 if let Some(headers) = headers {
1260 let h_val = serde_json::to_value(&headers).unwrap_or_default();
1261
1262 {
1263 let mut state = self.call_state.write().await;
1264 let mut extras = state.extras.take().unwrap_or_default();
1265 extras.insert("_hangup_headers".to_string(), h_val.clone());
1266 state.extras = Some(extras);
1267 }
1268 } else {
1269 }
1270
1271 let hangup_reason = match initiator.as_deref() {
1272 Some("caller") => CallRecordHangupReason::ByCaller,
1273 Some("callee") => CallRecordHangupReason::ByCallee,
1274 Some("system") => CallRecordHangupReason::Autohangup,
1275 _ => reason.unwrap_or(CallRecordHangupReason::BySystem),
1276 };
1277
1278 self.media_stream
1279 .stop(Some(hangup_reason.to_string()), initiator);
1280
1281 self.call_state
1282 .write()
1283 .await
1284 .set_hangup_reason(hangup_reason);
1285 Ok(())
1286 }
1287
1288 async fn do_refer(
1289 &self,
1290 caller: String,
1291 callee: String,
1292 refer_option: Option<ReferOption>,
1293 ) -> Result<()> {
1294 self.do_interrupt(false).await.ok();
1295
1296 let pause_parent_asr = refer_option
1298 .as_ref()
1299 .and_then(|o| o.pause_parent_asr)
1300 .unwrap_or(false);
1301
1302 let original_asr_option = if pause_parent_asr {
1304 let cs = self.call_state.read().await;
1305 cs.option.as_ref().and_then(|o| o.asr.clone())
1306 } else {
1307 None
1308 };
1309
1310 if pause_parent_asr {
1312 info!(
1313 session_id = self.session_id,
1314 "Pausing parent call ASR during refer"
1315 );
1316 self.media_stream
1317 .remove_processor::<crate::media::asr_processor::AsrProcessor>(
1318 &self.server_side_track_id,
1319 )
1320 .await
1321 .ok();
1322 }
1323
1324 let mut moh = refer_option.as_ref().and_then(|o| o.moh.clone());
1325 if let Some(ref path) = moh {
1326 if !path.starts_with("http") && !std::path::Path::new(path).exists() {
1327 let fallback = "./config/sounds/refer_moh.wav";
1328 if std::path::Path::new(fallback).exists() {
1329 info!(
1330 session_id = self.session_id,
1331 "moh {} not found, using fallback {}", path, fallback
1332 );
1333 moh = Some(fallback.to_string());
1334 }
1335 }
1336 }
1337 let ref_call_id = refer_option
1338 .as_ref()
1339 .and_then(|o| o.call_id.clone())
1340 .unwrap_or_else(|| format!("ref-{}-{}", rand::random::<u32>(), self.session_id));
1341
1342 let session_id = self.session_id.clone();
1343 let track_id = self.server_side_track_id.clone();
1344
1345 let (recorder, parent_caller) = {
1346 let cs = self.call_state.read().await;
1347 let option = cs.option.as_ref();
1348 (
1349 option.map(|o| o.recorder.clone()).unwrap_or_default(),
1350 option.and_then(|o| o.caller.clone()),
1351 )
1352 };
1353 let caller = if caller.trim().is_empty() {
1354 parent_caller.unwrap_or_default()
1355 } else {
1356 caller
1357 };
1358
1359 let call_option = CallOption {
1360 caller: Some(caller),
1361 callee: Some(callee.clone()),
1362 sip: refer_option.as_ref().and_then(|o| o.sip.clone()),
1363 asr: refer_option.as_ref().and_then(|o| o.asr.clone()),
1364 denoise: refer_option.as_ref().and_then(|o| o.denoise.clone()),
1365 recorder,
1366 ..Default::default()
1367 };
1368
1369 let mut invite_option = call_option.build_invite_option()?;
1370 invite_option.call_id = Some(ref_call_id);
1371
1372 let headers = invite_option.headers.get_or_insert_with(|| Vec::new());
1373
1374 {
1375 let cs = self.call_state.read().await;
1376 if let Some(opt) = cs.option.as_ref() {
1377 if let Some(callee) = opt.callee.as_ref() {
1378 headers.push(rsipstack::rsip::Header::Other(
1379 "X-Referred-To".to_string(),
1380 callee.clone(),
1381 ));
1382 }
1383 if let Some(caller) = opt.caller.as_ref() {
1384 headers.push(rsipstack::rsip::Header::Other(
1385 "X-Referred-From".to_string(),
1386 caller.clone(),
1387 ));
1388 }
1389 }
1390 }
1391
1392 headers.push(rsipstack::rsip::Header::Other(
1393 "X-Referred-Id".to_string(),
1394 self.session_id.clone(),
1395 ));
1396
1397 let ssrc = rand::random::<u32>();
1398 let refer_call_state = Arc::new(RwLock::new(ActiveCallState {
1399 start_time: Utc::now(),
1400 ssrc,
1401 option: Some(call_option.clone()),
1402 is_refer: true,
1403 ..Default::default()
1404 }));
1405
1406 {
1407 let mut cs = self.call_state.write().await;
1408 cs.refer_callstate.replace(refer_call_state.clone());
1409 }
1410
1411 let auto_hangup_requested = refer_option
1412 .as_ref()
1413 .and_then(|o| o.auto_hangup)
1414 .unwrap_or(true);
1415
1416 if auto_hangup_requested {
1417 self.call_state.write().await.auto_hangup =
1418 Some((ssrc, CallRecordHangupReason::ByRefer));
1419 } else {
1420 self.call_state.write().await.auto_hangup = None;
1421 }
1422
1423 if !auto_hangup_requested && pause_parent_asr && original_asr_option.is_some() {
1425 let asr_option = original_asr_option.unwrap();
1426 self.call_state.write().await.pending_asr_resume = Some((ssrc, asr_option));
1427 }
1428
1429 let timeout_secs = refer_option.as_ref().and_then(|o| o.timeout).unwrap_or(30);
1430
1431 info!(
1432 session_id = self.session_id,
1433 ssrc,
1434 auto_hangup = auto_hangup_requested,
1435 callee,
1436 timeout_secs,
1437 "do_refer"
1438 );
1439
1440 let r = tokio::time::timeout(
1441 Duration::from_secs(timeout_secs as u64),
1442 self.create_outgoing_sip_track(
1443 self.cancel_token.child_token(),
1444 refer_call_state.clone(),
1445 &track_id,
1446 invite_option,
1447 &call_option,
1448 moh,
1449 auto_hangup_requested,
1450 ),
1451 )
1452 .await;
1453
1454 {
1455 self.call_state.write().await.moh = None;
1456 }
1457
1458 let result = match r {
1459 Ok(res) => res,
1460 Err(_) => {
1461 warn!(
1462 session_id = session_id,
1463 "refer sip track creation timed out after {} seconds", timeout_secs
1464 );
1465 self.event_sender
1466 .send(SessionEvent::Reject {
1467 track_id,
1468 timestamp: crate::media::get_timestamp(),
1469 reason: "Timeout when refer".into(),
1470 code: Some(408),
1471 refer: Some(true),
1472 })
1473 .ok();
1474 return Err(anyhow::anyhow!("refer sip track creation timed out").into());
1475 }
1476 };
1477
1478 match result {
1479 Ok(answer) => {
1480 self.event_sender
1481 .send(SessionEvent::Answer {
1482 timestamp: crate::media::get_timestamp(),
1483 track_id,
1484 sdp: answer,
1485 refer: Some(true),
1486 })
1487 .ok();
1488 }
1489 Err(e) => {
1490 warn!(
1491 session_id = session_id,
1492 "failed to create refer sip track: {}", e
1493 );
1494 match &e {
1495 rsipstack::Error::DialogError(reason, _, code) => {
1496 self.event_sender
1497 .send(SessionEvent::Reject {
1498 track_id,
1499 timestamp: crate::media::get_timestamp(),
1500 reason: reason.clone(),
1501 code: Some(code.code() as u32),
1502 refer: Some(true),
1503 })
1504 .ok();
1505 }
1506 _ => {}
1507 }
1508 return Err(e.into());
1509 }
1510 }
1511 Ok(())
1512 }
1513
1514 async fn do_mute(&self, track_id: Option<String>) -> Result<()> {
1515 self.media_stream.mute_track(track_id).await;
1516 Ok(())
1517 }
1518
1519 async fn do_unmute(&self, track_id: Option<String>) -> Result<()> {
1520 self.media_stream.unmute_track(track_id).await;
1521 Ok(())
1522 }
1523
1524 pub async fn cleanup(&self) -> Result<()> {
1525 if matches!(self.call_type, ActiveCallType::Sip | ActiveCallType::B2bua) {
1526 self.do_reject(
1527 Some(rsipstack::rsip::StatusCode::Decline),
1528 Some("handler disconnected".to_string()),
1529 )
1530 .await
1531 .ok();
1532 }
1533 self.call_state.write().await.tts_handle = None;
1534 self.media_stream.cleanup().await.ok();
1535 Ok(())
1536 }
1537
1538 pub fn get_callrecord(&self) -> Option<CallRecord> {
1539 self.call_state.try_read().ok().map(|call_state| {
1540 call_state.build_callrecord(
1541 self.app_state.clone(),
1542 self.session_id.clone(),
1543 self.call_type.clone(),
1544 )
1545 })
1546 }
1547
1548 async fn dump_to_file(
1549 &self,
1550 dump_file: &mut File,
1551 cmd_receiver: &mut CommandReceiver,
1552 event_receiver: &mut EventReceiver,
1553 ) {
1554 loop {
1555 select! {
1556 _ = self.cancel_token.cancelled() => {
1557 break;
1558 }
1559 Ok(cmd) = cmd_receiver.recv() => {
1560 CallRecordEvent::write(CallRecordEventType::Command, cmd, dump_file)
1561 .await;
1562 }
1563 Ok(event) = event_receiver.recv() => {
1564 if matches!(event, SessionEvent::Binary{..}) {
1565 continue;
1566 }
1567 CallRecordEvent::write(CallRecordEventType::Event, event, dump_file)
1568 .await;
1569 }
1570 };
1571 }
1572 }
1573
1574 async fn dump_loop(
1575 &self,
1576 dump_events: bool,
1577 mut dump_cmd_receiver: CommandReceiver,
1578 mut dump_event_receiver: EventReceiver,
1579 ) {
1580 if !dump_events {
1581 return;
1582 }
1583
1584 let file_name = self.app_state.get_dump_events_file(&self.session_id);
1585 let mut dump_file = match File::options()
1586 .create(true)
1587 .append(true)
1588 .open(&file_name)
1589 .await
1590 {
1591 Ok(file) => file,
1592 Err(e) => {
1593 warn!(
1594 session_id = self.session_id,
1595 file_name, "failed to open dump events file: {}", e
1596 );
1597 return;
1598 }
1599 };
1600 self.dump_to_file(
1601 &mut dump_file,
1602 &mut dump_cmd_receiver,
1603 &mut dump_event_receiver,
1604 )
1605 .await;
1606
1607 while let Ok(event) = dump_event_receiver.try_recv() {
1608 if matches!(event, SessionEvent::Binary { .. }) {
1609 continue;
1610 }
1611 CallRecordEvent::write(CallRecordEventType::Event, event, &mut dump_file).await;
1612 }
1613 }
1614
1615 pub async fn create_rtp_track(
1616 &self,
1617 track_id: TrackId,
1618 ssrc: u32,
1619 enable_srtp: Option<bool>,
1620 ) -> Result<RtcTrack> {
1621 let mut rtc_config = RtcTrackConfig::default();
1622 let use_srtp = enable_srtp
1624 .or(self.app_state.config.enable_srtp)
1625 .unwrap_or(false);
1626 rtc_config.mode = if use_srtp {
1627 rustrtc::TransportMode::Srtp
1628 } else {
1629 rustrtc::TransportMode::Rtp
1630 };
1631
1632 if let Some(codecs) = &self.app_state.config.codecs {
1633 let mut codec_types = Vec::new();
1634 for c in codecs {
1635 match c.to_lowercase().as_str() {
1636 "pcmu" => codec_types.push(CodecType::PCMU),
1637 "pcma" => codec_types.push(CodecType::PCMA),
1638 "g722" => codec_types.push(CodecType::G722),
1639 "g729" => codec_types.push(CodecType::G729),
1640 #[cfg(feature = "opus")]
1641 "opus" => codec_types.push(CodecType::Opus),
1642 "dtmf" | "2833" | "telephone_event" => {
1643 codec_types.push(CodecType::TelephoneEvent)
1644 }
1645 _ => {}
1646 }
1647 }
1648 if !codec_types.is_empty() {
1649 rtc_config.preferred_codec = Some(codec_types[0].clone());
1650 rtc_config.codecs = codec_types;
1651 }
1652 }
1653
1654 if rtc_config.preferred_codec.is_none() {
1655 rtc_config.preferred_codec = Some(self.track_config.codec.clone());
1656 }
1657
1658 rtc_config.rtp_port_range = self
1659 .app_state
1660 .config
1661 .rtp_start_port
1662 .zip(self.app_state.config.rtp_end_port);
1663
1664 if let Some(ref external_ip) = self.app_state.config.external_ip {
1665 rtc_config.external_ip = Some(external_ip.clone());
1666 }
1667 if let Some(ref bind_ip) = self.app_state.config.rtp_bind_ip {
1668 rtc_config.bind_ip = Some(bind_ip.clone());
1669 }
1670
1671 rtc_config.enable_latching = self.app_state.config.enable_rtp_latching;
1672 rtc_config.enable_ice_lite = self.app_state.config.enable_ice_lite;
1673
1674 let mut track = RtcTrack::new(
1675 self.cancel_token.child_token(),
1676 track_id,
1677 self.track_config.clone(),
1678 rtc_config,
1679 )
1680 .with_ssrc(ssrc);
1681
1682 track.create().await?;
1683
1684 Ok(track)
1685 }
1686
1687 async fn setup_caller_track(&self, option: &CallOption) -> Result<()> {
1688 let hangup_headers = option
1689 .sip
1690 .as_ref()
1691 .and_then(|s| s.hangup_headers.as_ref())
1692 .map(|headers_map| {
1693 headers_map
1694 .iter()
1695 .map(|(k, v)| rsipstack::rsip::Header::Other(k.clone(), v.clone()))
1696 .collect::<Vec<rsipstack::rsip::Header>>()
1697 });
1698 self.call_state.write().await.option = Some(option.clone());
1699 info!(
1700 session_id = self.session_id,
1701 call_type = ?self.call_type,
1702 "setup caller track"
1703 );
1704
1705 let track = match self.call_type {
1706 ActiveCallType::Webrtc => Some(self.create_webrtc_track().await?),
1707 ActiveCallType::WebSocket => {
1708 let audio_receiver = self.call_state.write().await.audio_receiver.take();
1709 if let Some(receiver) = audio_receiver {
1710 Some(self.create_websocket_track(receiver).await?)
1711 } else {
1712 None
1713 }
1714 }
1715 ActiveCallType::Sip => {
1716 if let Some(dialog_id) = self
1717 .invitation
1718 .find_dialog_id_by_session_id(&self.session_id)
1719 {
1720 if let Some(pending_dialog) = self.invitation.get_pending_call(&dialog_id) {
1721 return self
1722 .prepare_incoming_sip_track(
1723 self.cancel_token.clone(),
1724 self.call_state.clone(),
1725 &self.session_id,
1726 pending_dialog,
1727 hangup_headers,
1728 )
1729 .await;
1730 }
1731 }
1732
1733 let mut option = option.clone();
1735 if option.sip.is_none()
1736 || option
1737 .sip
1738 .as_ref()
1739 .and_then(|s| s.username.as_ref())
1740 .is_none()
1741 {
1742 if let Some(callee) = &option.callee {
1743 if let Some(cred) = self.app_state.find_credentials_for_callee(callee) {
1744 if option.sip.is_none() {
1745 option.sip = Some(crate::SipOption {
1746 username: Some(cred.username.clone()),
1747 password: Some(cred.password.clone()),
1748 realm: cred.realm.clone(),
1749 ..Default::default()
1750 });
1751 }
1752 }
1753 }
1754 }
1755
1756 let mut invite_option = option.build_invite_option()?;
1757 invite_option.call_id = Some(self.session_id.clone());
1758
1759 match self
1760 .create_outgoing_sip_track(
1761 self.cancel_token.clone(),
1762 self.call_state.clone(),
1763 &self.session_id,
1764 invite_option,
1765 &option,
1766 None,
1767 false,
1768 )
1769 .await
1770 {
1771 Ok(answer) => {
1772 self.event_sender
1773 .send(SessionEvent::Answer {
1774 timestamp: crate::media::get_timestamp(),
1775 track_id: self.session_id.clone(),
1776 sdp: answer,
1777 refer: Some(false),
1778 })
1779 .ok();
1780 return Ok(());
1781 }
1782 Err(e) => {
1783 warn!(
1784 session_id = self.session_id,
1785 "failed to create sip track: {}", e
1786 );
1787 match &e {
1788 rsipstack::Error::DialogError(reason, _, code) => {
1789 self.event_sender
1790 .send(SessionEvent::Reject {
1791 track_id: self.session_id.clone(),
1792 timestamp: crate::media::get_timestamp(),
1793 reason: reason.clone(),
1794 code: Some(code.code() as u32),
1795 refer: Some(false),
1796 })
1797 .ok();
1798 }
1799 _ => {}
1800 }
1801 return Err(e.into());
1802 }
1803 }
1804 }
1805 ActiveCallType::B2bua => {
1806 if let Some(dialog_id) = self
1807 .invitation
1808 .find_dialog_id_by_session_id(&self.session_id)
1809 {
1810 if let Some(pending_dialog) = self.invitation.get_pending_call(&dialog_id) {
1811 return self
1812 .prepare_incoming_sip_track(
1813 self.cancel_token.clone(),
1814 self.call_state.clone(),
1815 &self.session_id,
1816 pending_dialog,
1817 hangup_headers,
1818 )
1819 .await;
1820 }
1821 }
1822
1823 warn!(
1824 session_id = self.session_id,
1825 "no pending dialog found for B2BUA call"
1826 );
1827 return Err(anyhow::anyhow!(
1828 "no pending dialog found for session_id: {}",
1829 self.session_id
1830 ));
1831 }
1832 };
1833 match track {
1834 Some(track) => {
1835 self.finish_caller_stack(&option, Some(track)).await?;
1836 }
1837 None => {
1838 warn!(session_id = self.session_id, "no track created for caller");
1839 return Err(anyhow::anyhow!("no track created for caller"));
1840 }
1841 }
1842 Ok(())
1843 }
1844
1845 async fn finish_caller_stack(
1846 &self,
1847 option: &CallOption,
1848 track: Option<Box<dyn Track>>,
1849 ) -> Result<()> {
1850 if let Some(track) = track {
1851 self.setup_track_with_stream(&option, track).await?;
1852 }
1853
1854 {
1855 let call_state = self.call_state.read().await;
1856 if let Some(ref answer) = call_state.answer {
1857 info!(
1858 session_id = self.session_id,
1859 "sending answer event: {}", answer,
1860 );
1861 self.event_sender
1862 .send(SessionEvent::Answer {
1863 timestamp: crate::media::get_timestamp(),
1864 track_id: self.session_id.clone(),
1865 sdp: answer.clone(),
1866 refer: Some(false),
1867 })
1868 .ok();
1869 } else {
1870 warn!(
1871 session_id = self.session_id,
1872 "no answer in state to send event"
1873 );
1874 }
1875 }
1876 Ok(())
1877 }
1878
1879 pub async fn setup_track_with_stream(
1880 &self,
1881 option: &CallOption,
1882 mut track: Box<dyn Track>,
1883 ) -> Result<()> {
1884 let processors = match StreamEngine::create_processors(
1885 self.app_state.stream_engine.clone(),
1886 track.as_ref(),
1887 self.cancel_token.child_token(),
1888 self.event_sender.clone(),
1889 self.media_stream.packet_sender.clone(),
1890 option,
1891 )
1892 .await
1893 {
1894 Ok(processors) => processors,
1895 Err(e) => {
1896 warn!(
1897 session_id = self.session_id,
1898 "failed to prepare stream processors: {}", e
1899 );
1900 vec![]
1901 }
1902 };
1903
1904 for processor in processors {
1906 track.append_processor(processor);
1907 }
1908
1909 self.update_track_wrapper(track, None).await;
1910 Ok(())
1911 }
1912
1913 pub async fn update_track_wrapper(&self, mut track: Box<dyn Track>, play_id: Option<String>) {
1914 let (ambiance_opt, subscribe) = {
1915 let state = self.call_state.read().await;
1916 let mut opt = state
1917 .option
1918 .as_ref()
1919 .and_then(|o| o.ambiance.clone())
1920 .unwrap_or_default();
1921
1922 if let Some(global) = &self.app_state.config.ambiance {
1923 opt.merge(global);
1924 }
1925
1926 let subscribe = state
1927 .option
1928 .as_ref()
1929 .and_then(|o| o.subscribe)
1930 .unwrap_or_default();
1931
1932 (opt, subscribe)
1933 };
1934 if track.id() == &self.server_side_track_id && ambiance_opt.path.is_some() {
1935 match AmbianceProcessor::new(ambiance_opt).await {
1936 Ok(ambiance) => {
1937 info!(session_id = self.session_id, "loaded ambiance processor");
1938 track.append_processor(Box::new(ambiance));
1939 }
1940 Err(e) => {
1941 tracing::error!("failed to load ambiance wav {}", e);
1942 }
1943 }
1944 }
1945
1946 if subscribe && self.call_type != ActiveCallType::WebSocket {
1947 let (track_index, sub_track_id) = if track.id() == &self.server_side_track_id {
1948 (0, self.server_side_track_id.clone())
1949 } else {
1950 (1, self.session_id.clone())
1951 };
1952 let sub_processor =
1953 SubscribeProcessor::new(self.event_sender.clone(), sub_track_id, track_index);
1954 track.append_processor(Box::new(sub_processor));
1955 }
1956
1957 self.call_state.write().await.current_play_id = play_id.clone();
1958 self.media_stream.update_track(track, play_id).await;
1959 }
1960
1961 pub async fn create_websocket_track(
1962 &self,
1963 audio_receiver: WebsocketBytesReceiver,
1964 ) -> Result<Box<dyn Track>> {
1965 let (ssrc, codec) = {
1966 let call_state = self.call_state.read().await;
1967 (
1968 call_state.ssrc,
1969 call_state
1970 .option
1971 .as_ref()
1972 .map(|o| o.codec.clone())
1973 .unwrap_or_default(),
1974 )
1975 };
1976
1977 let ws_track = WebsocketTrack::new(
1978 self.cancel_token.child_token(),
1979 self.session_id.clone(),
1980 self.track_config.clone(),
1981 self.event_sender.clone(),
1982 audio_receiver,
1983 codec,
1984 ssrc,
1985 );
1986
1987 {
1988 let mut call_state = self.call_state.write().await;
1989 call_state.answer_time = Some(Utc::now());
1990 call_state.answer = Some("".to_string());
1991 call_state.last_status_code = 200;
1992 }
1993
1994 Ok(Box::new(ws_track))
1995 }
1996
1997 pub(super) async fn create_webrtc_track(&self) -> Result<Box<dyn Track>> {
1998 let (ssrc, option) = {
1999 let call_state = self.call_state.read().await;
2000 (
2001 call_state.ssrc,
2002 call_state.option.clone().unwrap_or_default(),
2003 )
2004 };
2005
2006 let mut rtc_config = RtcTrackConfig::default();
2007 rtc_config.mode = rustrtc::TransportMode::WebRtc; rtc_config.ice_servers = self.app_state.config.ice_servers.clone();
2009
2010 if let Some(codecs) = &self.app_state.config.codecs {
2011 let mut codec_types = Vec::new();
2012 for c in codecs {
2013 match c.to_lowercase().as_str() {
2014 "pcmu" => codec_types.push(CodecType::PCMU),
2015 "pcma" => codec_types.push(CodecType::PCMA),
2016 "g722" => codec_types.push(CodecType::G722),
2017 "g729" => codec_types.push(CodecType::G729),
2018 #[cfg(feature = "opus")]
2019 "opus" => codec_types.push(CodecType::Opus),
2020 "dtmf" | "2833" | "telephone_event" => {
2021 codec_types.push(CodecType::TelephoneEvent)
2022 }
2023 _ => {}
2024 }
2025 }
2026 if !codec_types.is_empty() {
2027 rtc_config.preferred_codec = Some(codec_types[0].clone());
2028 rtc_config.codecs = codec_types;
2029 }
2030 }
2031
2032 if let Some(ref external_ip) = self.app_state.config.external_ip {
2033 rtc_config.external_ip = Some(external_ip.clone());
2034 }
2035 if let Some(ref bind_ip) = self.app_state.config.rtp_bind_ip {
2036 rtc_config.bind_ip = Some(bind_ip.clone());
2037 }
2038
2039 let mut webrtc_track = RtcTrack::new(
2040 self.cancel_token.child_token(),
2041 self.session_id.clone(),
2042 self.track_config.clone(),
2043 rtc_config,
2044 )
2045 .with_ssrc(ssrc);
2046
2047 let timeout = option.handshake_timeout.map(|t| Duration::from_secs(t));
2048 let offer = match option.enable_ipv6 {
2049 Some(false) | None => {
2050 strip_ipv6_candidates(option.offer.as_ref().unwrap_or(&"".to_string()))
2051 }
2052 _ => option.offer.clone().unwrap_or("".to_string()),
2053 };
2054 let answer: Option<String>;
2055 match webrtc_track.handshake(offer, timeout).await {
2056 Ok(answer_sdp) => {
2057 answer = match option.enable_ipv6 {
2058 Some(false) | None => Some(strip_ipv6_candidates(&answer_sdp)),
2059 Some(true) => Some(answer_sdp.to_string()),
2060 };
2061 }
2062 Err(e) => {
2063 warn!(session_id = self.session_id, "failed to setup track: {}", e);
2064 return Err(anyhow::anyhow!("Failed to setup track: {}", e));
2065 }
2066 }
2067
2068 {
2069 let mut call_state = self.call_state.write().await;
2070 call_state.answer_time = Some(Utc::now());
2071 call_state.answer = answer;
2072 call_state.last_status_code = 200;
2073 }
2074 Ok(Box::new(webrtc_track))
2075 }
2076
2077 async fn create_outgoing_sip_track(
2078 &self,
2079 cancel_token: CancellationToken,
2080 call_state_ref: ActiveCallStateRef,
2081 track_id: &String,
2082 mut invite_option: InviteOption,
2083 call_option: &CallOption,
2084 moh: Option<String>,
2085 auto_hangup: bool,
2086 ) -> Result<String, rsipstack::Error> {
2087 let ssrc = call_state_ref.read().await.ssrc;
2088 let per_call_srtp = call_option.sip.as_ref().and_then(|s| s.enable_srtp);
2089 let rtp_track = self
2090 .create_rtp_track(track_id.clone(), ssrc, per_call_srtp)
2091 .await
2092 .map_err(|e| rsipstack::Error::Error(e.to_string()))?;
2093
2094 let offer = Some(
2095 rtp_track
2096 .local_description()
2097 .await
2098 .map_err(|e| rsipstack::Error::Error(e.to_string()))?,
2099 );
2100
2101 {
2102 let mut cs = call_state_ref.write().await;
2103 if let Some(o) = cs.option.as_mut() {
2104 o.offer = offer.clone();
2105 }
2106 cs.start_time = Utc::now();
2107 };
2108
2109 invite_option.offer = offer.clone().map(|s| s.into());
2110
2111 let needs_contact = contact_needs_public_resolution(&invite_option.contact);
2114
2115 if needs_contact {
2116 let addrs = self.invitation.dialog_layer.endpoint.get_addrs();
2117 if let Some(addr) = find_local_addr_for_uri(&addrs, &invite_option.callee) {
2118 let contact_username = invite_option
2119 .contact
2120 .auth
2121 .as_ref()
2122 .map(|auth| auth.user.as_str())
2123 .or_else(|| {
2124 invite_option
2125 .caller
2126 .auth
2127 .as_ref()
2128 .map(|auth| auth.user.as_str())
2129 });
2130 invite_option.contact = build_public_contact_uri(
2131 &self.app_state.learned_public_address,
2132 self.app_state.auto_learn_public_address_enabled(),
2133 &addr,
2134 contact_username,
2135 Some(&invite_option.contact),
2136 );
2137 } else {
2138 return Err(rsipstack::Error::Error(format!(
2139 "missing local SIP address for callee transport: {}",
2140 invite_option.callee
2141 )));
2142 }
2143 }
2144
2145 let mut rtp_track_to_setup = Some(Box::new(rtp_track) as Box<dyn Track>);
2146
2147 if let Some(moh) = moh {
2148 let ssrc_and_moh = {
2149 let mut state = call_state_ref.write().await;
2150 state.moh = Some(moh.clone());
2151 if state.current_play_id.is_none() {
2152 let ssrc = rand::random::<u32>();
2153 Some((ssrc, moh.clone()))
2154 } else {
2155 info!(
2156 session_id = self.session_id,
2157 "Something is playing, MOH will start after it ends"
2158 );
2159 None
2160 }
2161 };
2162
2163 if let Some((ssrc, moh_path)) = ssrc_and_moh {
2164 let file_track = FileTrack::new(self.server_side_track_id.clone())
2165 .with_play_id(Some(moh_path.clone()))
2166 .with_ssrc(ssrc)
2167 .with_path(moh_path.clone())
2168 .with_cancel_token(self.cancel_token.child_token());
2169 self.update_track_wrapper(Box::new(file_track), Some(moh_path))
2170 .await;
2171 }
2172 } else {
2173 let track = rtp_track_to_setup.take().unwrap();
2174 self.setup_track_with_stream(&call_option, track)
2175 .await
2176 .map_err(|e| rsipstack::Error::Error(e.to_string()))?;
2177 }
2178
2179 info!(
2180 session_id = self.session_id,
2181 track_id,
2182 contact = %invite_option.contact,
2183 "invite {} -> {} offer: \n{}",
2184 invite_option.caller,
2185 invite_option.callee,
2186 offer.as_ref().map(|s| s.as_str()).unwrap_or("<NO OFFER>")
2187 );
2188
2189 let (dlg_state_sender, dlg_state_receiver) =
2190 self.invitation.dialog_layer.new_dialog_state_channel();
2191
2192 let states = InviteDialogStates {
2193 is_client: true,
2194 session_id: self.session_id.clone(),
2195 track_id: track_id.clone(),
2196 event_sender: self.event_sender.clone(),
2197 media_stream: self.media_stream.clone(),
2198 call_state: call_state_ref.clone(),
2199 cancel_token,
2200 terminated_reason: None,
2201 has_early_media: false,
2202 };
2203
2204 let hangup_headers = call_option
2205 .sip
2206 .as_ref()
2207 .and_then(|s| s.hangup_headers.as_ref())
2208 .map(|headers_map| {
2209 headers_map
2210 .iter()
2211 .map(|(k, v)| rsipstack::rsip::Header::Other(k.clone(), v.clone()))
2212 .collect::<Vec<rsipstack::rsip::Header>>()
2213 });
2214
2215 let mut client_dialog_handler = DialogStateReceiverGuard::new(
2216 self.invitation.dialog_layer.clone(),
2217 dlg_state_receiver,
2218 hangup_headers,
2219 );
2220
2221 crate::spawn(async move {
2222 client_dialog_handler.process_dialog(states).await;
2223 });
2224
2225 let (dialog_id, answer) = self
2226 .invitation
2227 .invite(invite_option, dlg_state_sender)
2228 .await?;
2229
2230 self.call_state.write().await.moh = None;
2231
2232 if let Some(track) = rtp_track_to_setup {
2233 info!(
2234 session_id = self.session_id,
2235 track_id, "Stopping MOH and setting up RTP track"
2236 );
2237 self.media_stream
2238 .remove_track(&self.server_side_track_id, false)
2239 .await;
2240
2241 self.setup_track_with_stream(&call_option, track)
2242 .await
2243 .map_err(|e| rsipstack::Error::Error(e.to_string()))?;
2244 }
2245
2246 let answer = match answer {
2247 Some(answer) => {
2248 let s = String::from_utf8_lossy(&answer).to_string();
2249 if s.trim().is_empty() {
2250 let cs = call_state_ref.read().await;
2254 match cs.answer.clone() {
2255 Some(early_sdp) if !early_sdp.is_empty() => {
2256 info!(
2257 session_id = self.session_id,
2258 "200 OK has empty body; using early-media SDP from 183"
2259 );
2260 (early_sdp, true )
2261 }
2262 _ => {
2263 warn!(
2264 session_id = self.session_id,
2265 "200 OK has empty body and no early-media SDP available"
2266 );
2267 (s, false)
2268 }
2269 }
2270 } else {
2271 (s, false)
2272 }
2273 }
2274 None => {
2275 let cs = call_state_ref.read().await;
2277 match cs.answer.clone() {
2278 Some(early_sdp) if !early_sdp.is_empty() => {
2279 info!(
2280 session_id = self.session_id,
2281 "200 OK had no answer; using early-media SDP from 183"
2282 );
2283 (early_sdp, true )
2284 }
2285 _ => {
2286 warn!(session_id = self.session_id, "no answer received");
2287 return Err(rsipstack::Error::DialogError(
2288 "No answer received".to_string(),
2289 dialog_id,
2290 rsipstack::rsip::StatusCode::NotAcceptableHere,
2291 ));
2292 }
2293 }
2294 }
2295 };
2296 let (answer, remote_description_already_applied) = answer;
2297
2298 {
2299 let mut cs = call_state_ref.write().await;
2300 if cs.answer.is_none() {
2301 cs.answer = Some(answer.clone());
2302 }
2303 if auto_hangup {
2304 cs.auto_hangup = Some((ssrc, CallRecordHangupReason::ByRefer));
2305 }
2306 }
2307 if !remote_description_already_applied {
2308 self.media_stream
2309 .update_remote_description(&track_id, &answer)
2310 .await
2311 .ok();
2312 }
2313
2314 Ok(answer)
2315 }
2316
2317 pub fn is_webrtc_sdp(sdp: &str) -> bool {
2319 (sdp.contains("a=ice-ufrag:") || sdp.contains("a=ice-pwd:"))
2320 && sdp.contains("a=fingerprint:")
2321 }
2322
2323 pub async fn setup_answer_track(
2324 &self,
2325 ssrc: u32,
2326 option: &CallOption,
2327 offer: String,
2328 ) -> Result<(String, Box<dyn Track>)> {
2329 let offer = match option.enable_ipv6 {
2330 Some(false) | None => strip_ipv6_candidates(&offer),
2331 _ => offer.clone(),
2332 };
2333
2334 let timeout = option.handshake_timeout.map(|t| Duration::from_secs(t));
2335
2336 let mut media_track = if Self::is_webrtc_sdp(&offer) {
2337 let mut rtc_config = RtcTrackConfig::default();
2338 rtc_config.mode = rustrtc::TransportMode::WebRtc;
2339 rtc_config.ice_servers = self.app_state.config.ice_servers.clone();
2340 if let Some(ref external_ip) = self.app_state.config.external_ip {
2341 rtc_config.external_ip = Some(external_ip.clone());
2342 }
2343 if let Some(ref bind_ip) = self.app_state.config.rtp_bind_ip {
2344 rtc_config.bind_ip = Some(bind_ip.clone());
2345 }
2346 rtc_config.enable_latching = self.app_state.config.enable_rtp_latching;
2347 rtc_config.enable_ice_lite = self.app_state.config.enable_ice_lite;
2348
2349 let webrtc_track = RtcTrack::new(
2350 self.cancel_token.child_token(),
2351 self.session_id.clone(),
2352 self.track_config.clone(),
2353 rtc_config,
2354 )
2355 .with_ssrc(ssrc);
2356
2357 Box::new(webrtc_track) as Box<dyn Track>
2358 } else {
2359 let per_call_srtp = option.sip.as_ref().and_then(|s| s.enable_srtp);
2360 let rtp_track = self
2361 .create_rtp_track(self.session_id.clone(), ssrc, per_call_srtp)
2362 .await?;
2363 Box::new(rtp_track) as Box<dyn Track>
2364 };
2365
2366 let answer = match media_track.handshake(offer.clone(), timeout).await {
2367 Ok(answer) => answer,
2368 Err(e) => {
2369 return Err(anyhow::anyhow!("handshake failed: {e}"));
2370 }
2371 };
2372
2373 return Ok((answer, media_track));
2374 }
2375
2376 pub async fn prepare_incoming_sip_track(
2377 &self,
2378 cancel_token: CancellationToken,
2379 call_state_ref: ActiveCallStateRef,
2380 track_id: &String,
2381 pending_dialog: PendingDialog,
2382 hangup_headers: Option<Vec<rsipstack::rsip::Header>>,
2383 ) -> Result<()> {
2384 let state_receiver = pending_dialog.state_receiver;
2385 let states = InviteDialogStates {
2388 is_client: false,
2389 session_id: self.session_id.clone(),
2390 track_id: track_id.clone(),
2391 event_sender: self.event_sender.clone(),
2392 media_stream: self.media_stream.clone(),
2393 call_state: self.call_state.clone(),
2394 cancel_token,
2395 terminated_reason: None,
2396 has_early_media: false,
2397 };
2398
2399 let initial_request = pending_dialog.dialog.initial_request();
2400 let offer = String::from_utf8_lossy(&initial_request.body).to_string();
2401
2402 let (ssrc, option) = {
2403 let call_state = call_state_ref.read().await;
2404 (
2405 call_state.ssrc,
2406 call_state.option.clone().unwrap_or_default(),
2407 )
2408 };
2409
2410 match self.setup_answer_track(ssrc, &option, offer).await {
2411 Ok((offer, track)) => {
2412 self.setup_track_with_stream(&option, track).await?;
2413 {
2414 let mut state = self.call_state.write().await;
2415 state.ready_to_answer = Some((offer, None, pending_dialog.dialog));
2416 }
2417 }
2418 Err(e) => {
2419 return Err(anyhow::anyhow!("error creating track: {}", e));
2420 }
2421 }
2422
2423 let mut client_dialog_handler = DialogStateReceiverGuard::new(
2424 self.invitation.dialog_layer.clone(),
2425 state_receiver,
2426 hangup_headers,
2427 );
2428
2429 crate::spawn(async move {
2430 client_dialog_handler.process_dialog(states).await;
2431 });
2432 Ok(())
2433 }
2434}
2435
2436impl Drop for ActiveCall {
2437 fn drop(&mut self) {
2438 info!(session_id = self.session_id, "dropping active call");
2439 if let Some(sender) = self.app_state.callrecord_sender.as_ref() {
2440 if let Some(record) = self.get_callrecord() {
2441 if let Err(e) = sender.send(record) {
2442 warn!(
2443 session_id = self.session_id,
2444 "failed to send call record: {}", e
2445 );
2446 }
2447 }
2448 }
2449 }
2450}
2451
2452impl ActiveCallState {
2453 pub fn merge_option(&self, mut option: CallOption) -> CallOption {
2454 if let Some(existing) = &self.option {
2455 if option.asr.is_none() {
2456 option.asr = existing.asr.clone();
2457 }
2458 if option.tts.is_none() {
2459 option.tts = existing.tts.clone();
2460 }
2461 if option.vad.is_none() {
2462 option.vad = existing.vad.clone();
2463 }
2464 if option.denoise.is_none() {
2465 option.denoise = existing.denoise;
2466 }
2467 if option.recorder.is_none() {
2468 option.recorder = existing.recorder.clone();
2469 }
2470 if option.eou.is_none() {
2471 option.eou = existing.eou.clone();
2472 }
2473 if option.extra.is_none() {
2474 option.extra = existing.extra.clone();
2475 }
2476 if option.ambiance.is_none() {
2477 option.ambiance = existing.ambiance.clone();
2478 }
2479 }
2480 option
2481 }
2482
2483 pub fn set_hangup_reason(&mut self, reason: CallRecordHangupReason) {
2484 if self.hangup_reason.is_none() {
2485 self.hangup_reason = Some(reason);
2486 }
2487 }
2488
2489 pub fn build_hangup_event(
2490 &self,
2491 track_id: TrackId,
2492 initiator: Option<String>,
2493 ) -> crate::event::SessionEvent {
2494 let from = self.option.as_ref().and_then(|o| o.caller.as_ref());
2495 let to = self.option.as_ref().and_then(|o| o.callee.as_ref());
2496 let extra = self.extras.clone();
2497
2498 crate::event::SessionEvent::Hangup {
2499 track_id,
2500 timestamp: crate::media::get_timestamp(),
2501 reason: Some(format!("{:?}", self.hangup_reason)),
2502 initiator,
2503 start_time: self.start_time.to_rfc3339(),
2504 answer_time: self.answer_time.map(|t| t.to_rfc3339()),
2505 ringing_time: self.ring_time.map(|t| t.to_rfc3339()),
2506 hangup_time: Utc::now().to_rfc3339(),
2507 extra,
2508 from: from.map(|f| f.into()),
2509 to: to.map(|f| f.into()),
2510 refer: Some(self.is_refer),
2511 }
2512 }
2513
2514 pub fn build_callrecord(
2515 &self,
2516 app_state: AppState,
2517 session_id: String,
2518 call_type: ActiveCallType,
2519 ) -> CallRecord {
2520 let option = self.option.clone().unwrap_or_default();
2521 let recorder = if option.recorder.is_some() {
2522 let recorder_file = app_state.get_recorder_file(&session_id);
2523 if std::path::Path::new(&recorder_file).exists() {
2524 let file_size = std::fs::metadata(&recorder_file)
2525 .map(|m| m.len())
2526 .unwrap_or(0);
2527 vec![crate::callrecord::CallRecordMedia {
2528 track_id: session_id.clone(),
2529 path: recorder_file,
2530 size: file_size,
2531 extra: None,
2532 }]
2533 } else {
2534 vec![]
2535 }
2536 } else {
2537 vec![]
2538 };
2539
2540 let dump_event_file = app_state.get_dump_events_file(&session_id);
2541 let dump_event_file = if std::path::Path::new(&dump_event_file).exists() {
2542 Some(dump_event_file)
2543 } else {
2544 None
2545 };
2546
2547 let refer_callrecord = self.refer_callstate.as_ref().and_then(|rc| {
2548 if let Ok(rc) = rc.try_read() {
2549 Some(Box::new(rc.build_callrecord(
2550 app_state.clone(),
2551 rc.session_id.clone(),
2552 ActiveCallType::B2bua,
2553 )))
2554 } else {
2555 None
2556 }
2557 });
2558
2559 let caller = option.caller.clone().unwrap_or_default();
2560 let callee = option.callee.clone().unwrap_or_default();
2561
2562 CallRecord {
2563 option: Some(option),
2564 call_id: session_id,
2565 call_type,
2566 start_time: self.start_time,
2567 ring_time: self.ring_time.clone(),
2568 answer_time: self.answer_time.clone(),
2569 end_time: Utc::now(),
2570 caller,
2571 callee,
2572 hangup_reason: self.hangup_reason.clone(),
2573 hangup_messages: Vec::new(),
2574 status_code: self.last_status_code,
2575 extras: self.extras.clone(),
2576 dump_event_file,
2577 recorder,
2578 refer_callrecord,
2579 }
2580 }
2581}