Skip to main content

active_call/call/
active_call.rs

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