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