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