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