Skip to main content

active_call/call/
active_call.rs

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