Skip to main content

active_call/call/
active_call.rs

1use super::Command;
2use crate::{
3    CallOption, ReferOption,
4    call::state::{ActorMsg, CallProgress, CallRuntime, Extras, LegShared, build_callrecord},
5    event::{EventReceiver, EventSender, SessionEvent},
6    media::{
7        TrackId,
8        engine::StreamEngine,
9        recorder::RecorderOption,
10        stream::{MediaStream, MediaStreamBuilder, SERVER_SIDE_TRACK_ID},
11        track::{
12            Track, TrackConfig, forwarding::ForwardingTrack, media_pass::MediaPassTrack,
13            tts::SynthesisHandle, websocket::WebsocketBytesReceiver,
14        },
15    },
16    synthesis::SynthesisCommand,
17    transcription::TranscriptionOption,
18};
19use crate::{
20    app::AppState,
21    call::{CommandReceiver, CommandSender, sip::Invitation},
22    callrecord::{CallRecord, CallRecordEvent, CallRecordEventType, CallRecordHangupReason},
23};
24use anyhow::Result;
25use arc_swap::{ArcSwap, ArcSwapOption};
26use chrono::{DateTime, Utc};
27use rsipstack::dialog::invite_dialog::InviteDialog;
28use serde::{Deserialize, Serialize};
29use std::{
30    collections::HashMap,
31    path::Path,
32    sync::{
33        Arc,
34        atomic::{AtomicBool, Ordering},
35    },
36    time::Duration,
37};
38use tokio::{fs::File, select, sync::mpsc};
39use tokio_util::sync::CancellationToken;
40use tracing::{debug, info, warn};
41
42/// Describes the state of the caller track when an incoming SIP call is waiting to be answered.
43pub enum PendingCallerTrack {
44    /// The track has been started in the media stream during ringing (early media).
45    /// Processors must be built from the accept option and appended to it.
46    StartedForEarlyMedia,
47    /// The track has not been added to the media stream yet.
48    /// setup_track_with_stream will start it and build processors from the accept option.
49    NotStarted(Box<dyn Track>),
50}
51
52#[cfg(test)]
53mod tests {
54    use super::*;
55    use crate::app::AppStateBuilder;
56    use crate::config::Config;
57    use crate::media::track::tts::SynthesisHandle;
58    use crate::synthesis::SynthesisCommand;
59    use tokio::sync::mpsc;
60
61    async fn make_active_call_with_option(option: CallOption) -> Arc<ActiveCall> {
62        let mut config = Config::default();
63        config.udp_port = 0; // Use random port
64        config.media_cache_path = "/tmp/mediacache".to_string();
65        let app_state = AppStateBuilder::new()
66            .with_config(config)
67            .with_stream_engine(Arc::new(StreamEngine::default()))
68            .build()
69            .await
70            .unwrap();
71        let active_call = Arc::new(ActiveCall::new(CallSpec {
72            call_type: ActiveCallType::Sip,
73            cancel_token: CancellationToken::new(),
74            session_id: "test-session".to_string(),
75            invitation: app_state.invitation.clone(),
76            app_state: app_state.clone(),
77            track_config: TrackConfig::default(),
78            audio_receiver: None,
79            dump_events: false,
80            server_side_track_id: None,
81            extras: None,
82        }));
83        active_call.set_option(option);
84        active_call
85    }
86
87    #[tokio::test]
88    async fn test_tts_ssrc_reuse_for_autohangup() -> Result<()> {
89        let mut option = crate::CallOption::default();
90        option.tts = Some(crate::synthesis::SynthesisOption::default());
91        let active_call = make_active_call_with_option(option).await;
92
93        let (tx, mut rx) = mpsc::unbounded_channel::<SynthesisCommand>();
94        let initial_ssrc = 12345;
95        let handle = SynthesisHandle::new(tx, Some("play_1".to_string()), initial_ssrc);
96
97        // 1. Set initial TTS handle
98        active_call.tts_handle.store(Some(Arc::new(handle)));
99        active_call.set_current_play(Some("play_1".to_string()));
100
101        // 2. Call do_tts with auto_hangup=true and same play_id
102        active_call
103            .do_tts(Command::Tts {
104                text: "hangup now".to_string(),
105                speaker: None,
106                play_id: Some("play_1".to_string()),
107                auto_hangup: Some(true),
108                streaming: Some(false),
109                end_of_stream: Some(true),
110                option: None,
111                wait_input_timeout: None,
112                base64: Some(false),
113                cache_key: None,
114            })
115            .await?;
116
117        // 3. Verify the hangup intent rides on the command sent to the existing track
118        let cmd = rx.try_recv().expect("Should have received tts command");
119        assert_eq!(cmd.text, "hangup now");
120        assert_eq!(cmd.auto_hangup, Some(true));
121
122        Ok(())
123    }
124
125    #[tokio::test]
126    async fn test_tts_new_ssrc_for_different_play_id() -> Result<()> {
127        let mut tts_opt = crate::synthesis::SynthesisOption::default();
128        tts_opt.provider = Some(crate::synthesis::SynthesisType::Aliyun);
129        let mut option = crate::CallOption::default();
130        option.tts = Some(tts_opt);
131        let active_call = make_active_call_with_option(option).await;
132
133        let (tx, _rx) = mpsc::unbounded_channel();
134        let initial_ssrc = 111;
135        let handle = SynthesisHandle::new(tx, Some("play_1".to_string()), initial_ssrc);
136
137        active_call.tts_handle.store(Some(Arc::new(handle)));
138        active_call.set_current_play(Some("play_1".to_string()));
139
140        // Call do_tts with DIFFERENT play_id
141        active_call
142            .do_tts(Command::Tts {
143                text: "new play".to_string(),
144                speaker: None,
145                play_id: Some("play_2".to_string()),
146                auto_hangup: Some(true),
147                streaming: Some(false),
148                end_of_stream: Some(true),
149                option: None,
150                wait_input_timeout: None,
151                base64: Some(false),
152                cache_key: None,
153            })
154            .await?;
155
156        // Verify a NEW track was started (new handle with a different ssrc,
157        // because a different play_id interrupts and starts fresh)
158        {
159            let handle = active_call.tts_handle.load_full();
160            assert!(handle.is_some(), "new tts handle should be stored");
161            let handle = handle.unwrap();
162            assert_ne!(
163                handle.ssrc, initial_ssrc,
164                "Should use a new SSRC for different play_id"
165            );
166        }
167
168        Ok(())
169    }
170
171    // refer=Some(true): only the refer call is cancelled, media stream stays alive.
172    #[tokio::test]
173    async fn test_hangup_refer_true_cancels_refer_only() -> Result<()> {
174        let active_call = make_active_call_with_option(crate::CallOption::default()).await;
175
176        let refer_token = active_call.cancel_token.child_token();
177        let refer_leg = LegShared::new(1, true, CallProgress::default());
178        active_call.set_refer_call_token(refer_token.clone());
179        active_call.set_refer_leg(Some(refer_leg.clone()));
180
181        active_call.do_hangup(None, None, None, Some(true)).await?;
182
183        assert!(
184            refer_token.is_cancelled(),
185            "refer token should be cancelled"
186        );
187        assert!(
188            !active_call.media_stream.cancel_token.is_cancelled(),
189            "media stream should NOT stop"
190        );
191        assert!(
192            refer_leg.progress.load_full().hangup_reason.is_some(),
193            "hangup_reason should be set on refer leg"
194        );
195        Ok(())
196    }
197
198    // refer=None: media stream stops and the refer token is also cancelled.
199    #[tokio::test]
200    async fn test_hangup_none_cancels_refer_too() -> Result<()> {
201        let active_call = make_active_call_with_option(crate::CallOption::default()).await;
202
203        let refer_token = active_call.cancel_token.child_token();
204        active_call.set_refer_call_token(refer_token.clone());
205
206        active_call.do_hangup(None, None, None, None).await?;
207
208        assert!(
209            refer_token.is_cancelled(),
210            "refer token should be cancelled"
211        );
212        assert!(
213            active_call.media_stream.cancel_token.is_cancelled(),
214            "media stream should stop"
215        );
216        Ok(())
217    }
218
219    // ---------------------------------------------------------------------------
220    // Regression: ringing-before-accept leaves caller track without processors
221    // ---------------------------------------------------------------------------
222    //
223    // When Ringing is issued before Accept on an incoming SIP call,
224    // prepare_incoming_sip_track starts the caller track in the media stream
225    // (needed for early-media ringtone) with the empty ringing option — no
226    // VAD/ASR/AGC processors.  It stores PendingCallerTrack::StartedForEarlyMedia
227    // in ready_to_answer.  At accept time, finish_caller_stack matches that variant
228    // and calls create_processors + append_processor with the real accept option.
229    //
230    // This test verifies that setup_track_with_stream (same processor-creation path)
231    // fires the ASR builder when given the accept option, proving the fix is sound.
232
233    struct MockCallerTrack {
234        id: TrackId,
235        config: crate::media::track::TrackConfig,
236        processor_chain: crate::media::processor::ProcessorChain,
237    }
238
239    impl MockCallerTrack {
240        fn new(id: TrackId) -> Self {
241            Self {
242                id,
243                config: crate::media::track::TrackConfig::default(),
244                processor_chain: crate::media::processor::ProcessorChain::new(16000),
245            }
246        }
247    }
248
249    #[async_trait::async_trait]
250    impl crate::media::track::Track for MockCallerTrack {
251        fn ssrc(&self) -> u32 {
252            0
253        }
254        fn id(&self) -> &TrackId {
255            &self.id
256        }
257        fn config(&self) -> &crate::media::track::TrackConfig {
258            &self.config
259        }
260        fn processor_chain(&mut self) -> &mut crate::media::processor::ProcessorChain {
261            &mut self.processor_chain
262        }
263        async fn handshake(
264            &mut self,
265            _o: String,
266            _t: Option<tokio::time::Duration>,
267        ) -> Result<String> {
268            Ok(String::new())
269        }
270        async fn update_remote_description(&mut self, _a: &String) -> Result<()> {
271            Ok(())
272        }
273        async fn start(
274            &mut self,
275            _e: crate::event::EventSender,
276            _p: crate::media::track::TrackPacketSender,
277        ) -> Result<()> {
278            Ok(())
279        }
280        async fn stop(&self) -> Result<()> {
281            Ok(())
282        }
283        async fn send_packet(&mut self, _f: &crate::media::AudioFrame) -> Result<()> {
284            Ok(())
285        }
286    }
287
288    struct MockAsrClient;
289
290    #[async_trait::async_trait]
291    impl crate::transcription::TranscriptionClient for MockAsrClient {
292        fn send_audio(
293            &self,
294            _s: &[crate::media::Sample],
295            _src: Option<&crate::media::SourcePacket>,
296        ) -> Result<()> {
297            Ok(())
298        }
299    }
300
301    async fn make_active_call_with_engine_and_option(
302        engine: Arc<StreamEngine>,
303        cache_dir: &str,
304        option: crate::CallOption,
305    ) -> Arc<ActiveCall> {
306        let mut config = Config::default();
307        config.udp_port = 0;
308        config.media_cache_path = cache_dir.to_string();
309        let app_state = AppStateBuilder::new()
310            .with_config(config)
311            .with_stream_engine(engine)
312            .build()
313            .await
314            .unwrap();
315        let session_id = format!("test-{}-{}", cache_dir, uuid::Uuid::new_v4());
316        let active_call = Arc::new(ActiveCall::new(CallSpec {
317            call_type: ActiveCallType::Sip,
318            cancel_token: CancellationToken::new(),
319            session_id: session_id.clone(),
320            invitation: app_state.invitation.clone(),
321            app_state: app_state.clone(),
322            track_config: TrackConfig::default(),
323            audio_receiver: None,
324            dump_events: false,
325            server_side_track_id: None,
326            extras: None,
327        }));
328        active_call.set_option(option);
329        active_call
330    }
331
332    #[tokio::test]
333    async fn test_setup_track_with_stream_builds_processors_from_accept_option() -> Result<()> {
334        let (asr_created_tx, mut asr_created_rx) = mpsc::channel::<()>(1);
335
336        let mock_provider =
337            crate::transcription::TranscriptionType::Other("mock-ringing-asr".to_string());
338
339        let mut engine = StreamEngine::new();
340        engine.register_asr(
341            mock_provider.clone(),
342            Box::new(move |_tid, _tok, _opt, _es| {
343                let tx = asr_created_tx.clone();
344                Box::pin(async move {
345                    let _ = tx.send(()).await;
346                    Ok(Box::new(MockAsrClient)
347                        as Box<dyn crate::transcription::TranscriptionClient>)
348                })
349            }),
350        );
351        let engine = Arc::new(engine);
352
353        let accept_option = crate::CallOption {
354            asr: Some(crate::transcription::TranscriptionOption {
355                provider: Some(mock_provider),
356                ..Default::default()
357            }),
358            ..Default::default()
359        };
360        let active_call = make_active_call_with_engine_and_option(
361            engine,
362            "mediacache_ringing_accept_test",
363            accept_option,
364        )
365        .await;
366        let cancel_token = active_call.cancel_token.clone();
367
368        // Simulate the fixed prepare_incoming_sip_track: track is held in
369        // ready_to_answer, NOT yet added to the media stream.
370        let mock_track = Box::new(MockCallerTrack::new(active_call.session_id.clone()));
371
372        // Simulate finish_caller_stack at accept time: setup_track_with_stream
373        // is called with the full accept option (the code path under test).
374        let accept_option = active_call.progress.load_full().option.clone().unwrap();
375        active_call
376            .setup_track_with_stream(&accept_option, mock_track)
377            .await?;
378
379        // The mock ASR builder must have fired, proving processors were built
380        // from the accept option and attached to the caller track.
381        let received =
382            tokio::time::timeout(std::time::Duration::from_secs(3), asr_created_rx.recv()).await;
383        assert!(
384            received.is_ok() && received.unwrap().is_some(),
385            "ASR processor was NOT created — setup_track_with_stream did not build \
386             processors from the accept option (regression: ringing-before-accept)"
387        );
388
389        cancel_token.cancel();
390        Ok(())
391    }
392
393    // ---------------------------------------------------------------------------
394    // Regression: double-ASR when Accept arrives without a prior Ringing
395    // ---------------------------------------------------------------------------
396    //
397    // prepare_incoming_sip_track starts the caller track during ringing. It must
398    // NOT build VAD/ASR/AGC processors at that point: the stored option already
399    // carries `asr` in the accept-first path, so finish_caller_stack would build a
400    // second set at accept time, producing two ASR clients (two WebSocket
401    // connections) on the same track.
402    //
403    // This test verifies that update_track_wrapper (the prepare-time path) never
404    // fires the ASR builder even when the option carries an asr config.
405
406    #[tokio::test]
407    async fn test_update_track_wrapper_does_not_build_asr_processor() -> Result<()> {
408        let (asr_created_tx, mut asr_created_rx) = mpsc::channel::<()>(1);
409
410        let mock_provider =
411            crate::transcription::TranscriptionType::Other("mock-ringing-asr".to_string());
412
413        let mut engine = StreamEngine::new();
414        engine.register_asr(
415            mock_provider.clone(),
416            Box::new(move |_tid, _tok, _opt, _es| {
417                let tx = asr_created_tx.clone();
418                Box::pin(async move {
419                    let _ = tx.send(()).await;
420                    Ok(Box::new(MockAsrClient)
421                        as Box<dyn crate::transcription::TranscriptionClient>)
422                })
423            }),
424        );
425        let engine = Arc::new(engine);
426
427        // Simulate the accept-first path: setup_caller_track has already stored the
428        // full accept option (with asr) before the track is prepared.
429        let accept_option = crate::CallOption {
430            asr: Some(crate::transcription::TranscriptionOption {
431                provider: Some(mock_provider),
432                ..Default::default()
433            }),
434            ..Default::default()
435        };
436        let active_call = make_active_call_with_engine_and_option(
437            engine,
438            "mediacache_update_track_wrapper_test",
439            accept_option,
440        )
441        .await;
442        let cancel_token = active_call.cancel_token.clone();
443
444        let mock_track = Box::new(MockCallerTrack::new(active_call.session_id.clone()));
445        active_call.update_track_wrapper(mock_track, None).await;
446
447        // update_track_wrapper must NOT fire the ASR builder; the processors are
448        // deferred to finish_caller_stack(StartedForEarlyMedia) at accept time.
449        let received =
450            tokio::time::timeout(std::time::Duration::from_millis(500), asr_created_rx.recv())
451                .await;
452        assert!(
453            received.is_err(),
454            "ASR builder fired during track preparation — double-ASR regression"
455        );
456
457        cancel_token.cancel();
458        Ok(())
459    }
460}
461
462#[derive(Deserialize)]
463#[serde(rename_all = "camelCase")]
464pub struct CallParams {
465    pub id: Option<String>,
466    #[serde(rename = "dump")]
467    pub dump_events: Option<bool>,
468    #[serde(rename = "ping")]
469    pub ping_interval: Option<u32>,
470    pub server_side_track: Option<String>,
471    /// Set when this connection is a one-hop find from another node.
472    /// Only an empty `forward` may be forwarded. A present value is answered
473    /// locally only and must never hop again. `forward=true` must 404 if the
474    /// session is absent (do not create a new call).
475    #[serde(default)]
476    pub forward: Option<bool>,
477    /// Ignored. Kept so older nodes that still send `visited=` can deserialize.
478    #[serde(default)]
479    pub visited: Option<String>,
480}
481
482impl CallParams {
483    /// Build the query string used when forwarding this request to a peer.
484    /// Sets `forward=true` so the peer will not hop (`forward` is no longer empty).
485    pub fn to_forward_query(&self) -> String {
486        let mut parts: Vec<String> = Vec::new();
487        if let Some(id) = &self.id {
488            parts.push(format!("id={}", urlencoding::encode(id)));
489        }
490        if let Some(dump) = self.dump_events {
491            parts.push(format!("dump={}", dump));
492        }
493        if let Some(ping) = self.ping_interval {
494            parts.push(format!("ping={}", ping));
495        }
496        if let Some(track) = &self.server_side_track {
497            parts.push(format!("server_side_track={}", urlencoding::encode(track)));
498        }
499        parts.push("forward=true".to_string());
500        parts.join("&")
501    }
502}
503
504#[derive(Debug, Serialize, Deserialize, Clone, Default, PartialEq, Eq)]
505#[serde(rename_all = "camelCase")]
506pub enum ActiveCallType {
507    Webrtc,
508    B2bua,
509    WebSocket,
510    #[default]
511    Sip,
512}
513
514pub type ActiveCallRef = Arc<ActiveCall>;
515
516/// A call: immutable identity + lock-free shared state (see `call::state`)
517/// plus the infrastructure handles. Actor-owned mutable state lives in
518/// [`CallRuntime`] inside `serve`.
519pub struct ActiveCall {
520    pub cancel_token: CancellationToken,
521    pub call_type: ActiveCallType,
522    pub session_id: String,
523    pub start_time: DateTime<Utc>,
524    pub media_stream: Arc<MediaStream>,
525    pub track_config: TrackConfig,
526    pub event_sender: EventSender,
527    pub app_state: AppState,
528    pub invitation: Invitation,
529    pub cmd_sender: CommandSender,
530    pub dump_events: bool,
531    pub server_side_track_id: TrackId,
532
533    /// Immutable main-leg SSRC.
534    pub ssrc: u32,
535    /// Bridge pause flag shared with the peer call during bridging.
536    pub bridge_paused: Arc<AtomicBool>,
537    /// Main-leg lifecycle progress (dialog id, ring/answer times, answer SDP, option...).
538    pub progress: Arc<ArcSwap<CallProgress>>,
539    /// Main-leg variables (playbook set_var, SIP headers, hangup headers).
540    pub extras: Extras,
541    /// Active music-on-hold path (shared with the spawned refer task).
542    pub moh: ArcSwapOption<String>,
543    /// play id of the current server-side playback (shared for TrackEnd matching).
544    pub current_play_id: ArcSwapOption<String>,
545    /// Live TTS handle (shared so spawned tasks and post-serve cleanup can drop it).
546    pub tts_handle: ArcSwapOption<SynthesisHandle>,
547    /// Shared state of the refer leg, when one is active.
548    pub refer_leg: ArcSwapOption<LegShared>,
549    /// Answer prepared during ringing (SDP + running track + dialog); taken
550    /// by accept/reject. Shared so post-serve `cleanup` can still reject.
551    pub ready_to_answer: ArcSwapOption<ReadyAnswer>,
552    /// Answer prepared for the underlying SIP dialog of non-SIP call types
553    /// (WebSocket/Webrtc); taken by accept, rejected on reject/cleanup.
554    /// Same slot family as `ready_to_answer`: lock-free set/take on the
555    /// actor path, with ownership recoverable from `cleanup`.
556    pub pending_sip_answer: ArcSwapOption<PendingSipAnswer>,
557    /// Cancel this token to hang up only the refer call, leaving the main call alive.
558    pub refer_call_token: ArcSwapOption<CancellationToken>,
559    /// Pending wait-input timeout set by the last Tts/Play command.
560    pub wait_input_timeout: ArcSwapOption<u32>,
561    /// ASR config to resume on the parent leg once the refer leg ends.
562    pub pending_asr_resume: ArcSwapOption<(u32, TranscriptionOption)>,
563    /// WebSocket audio receiver injected at construction, taken once by setup.
564    pub audio_receiver: std::sync::Mutex<Option<WebsocketBytesReceiver>>,
565}
566
567impl ActiveCall {
568    /// Lock-free shared state of the main leg.
569    pub fn leg(&self) -> LegShared {
570        LegShared {
571            ssrc: self.ssrc,
572            is_refer: false,
573            progress: self.progress.clone(),
574            extras: self.extras.clone(),
575        }
576    }
577
578    /// Store/replace the main-leg option in the progress snapshot.
579    pub fn set_option(&self, option: CallOption) {
580        self.progress.rcu(|p| {
581            let mut p = CallProgress::clone(p);
582            p.option = Some(option.clone());
583            p
584        });
585    }
586
587    pub fn moh_path(&self) -> Option<String> {
588        self.moh.load_full().map(|s| s.to_string())
589    }
590
591    pub fn set_moh(&self, v: Option<String>) {
592        self.moh.store(v.map(Arc::new));
593    }
594
595    pub fn current_play(&self) -> Option<String> {
596        self.current_play_id.load_full().map(|s| s.to_string())
597    }
598
599    pub fn set_current_play(&self, v: Option<String>) {
600        self.current_play_id.store(v.map(Arc::new));
601    }
602
603    pub fn refer_leg_value(&self) -> Option<LegShared> {
604        self.refer_leg.load_full().map(|l| l.as_ref().clone())
605    }
606
607    pub fn set_refer_leg(&self, v: Option<LegShared>) {
608        self.refer_leg.store(v.map(Arc::new));
609    }
610
611    /// Prepared answer during ringing; taken once by accept/reject.
612    pub fn set_ready_to_answer(&self, ready: ReadyAnswer) {
613        self.ready_to_answer.store(Some(Arc::new(ready)));
614    }
615
616    pub fn take_ready_to_answer(&self) -> Option<Arc<ReadyAnswer>> {
617        self.ready_to_answer.swap(None)
618    }
619
620    pub fn has_ready_to_answer(&self) -> bool {
621        self.ready_to_answer.load().is_some()
622    }
623
624    /// Prepared 200 OK for the underlying SIP dialog of non-SIP call types.
625    pub fn set_pending_sip_answer(&self, pending: PendingSipAnswer) {
626        self.pending_sip_answer.store(Some(Arc::new(pending)));
627    }
628
629    pub fn take_pending_sip_answer(&self) -> Option<Arc<PendingSipAnswer>> {
630        self.pending_sip_answer.swap(None)
631    }
632
633    /// Cancel token that hangs up only the refer leg, leaving the main call
634    /// alive; set by `do_refer`, taken by hangup.
635    pub fn take_refer_call_token(&self) -> Option<CancellationToken> {
636        self.refer_call_token.swap(None).map(|t| (*t).clone())
637    }
638
639    pub fn set_refer_call_token(&self, token: CancellationToken) {
640        self.refer_call_token.store(Some(Arc::new(token)));
641    }
642
643    /// Pending wait-input timeout set by the last Tts/Play command, consumed
644    /// when the track ends.
645    pub fn take_wait_input_timeout(&self) -> Option<u32> {
646        self.wait_input_timeout.swap(None).map(|t| *t)
647    }
648
649    pub fn set_wait_input_timeout(&self, v: Option<u32>) {
650        self.wait_input_timeout.store(v.map(Arc::new));
651    }
652
653    /// ASR config to resume on the parent leg once the refer leg ends.
654    pub fn set_pending_asr_resume(&self, v: (u32, TranscriptionOption)) {
655        self.pending_asr_resume.store(Some(Arc::new(v)));
656    }
657
658    pub fn take_pending_asr_resume(&self) -> Option<(u32, TranscriptionOption)> {
659        self.pending_asr_resume.swap(None).map(|a| (*a).clone())
660    }
661
662    /// Insert/overwrite one main-leg extras variable.
663    pub fn set_extra(&self, key: &str, value: serde_json::Value) {
664        self.leg().set_extra(key, value);
665    }
666
667    /// Whether a pending (not yet answered) incoming dialog exists for this call.
668    fn has_pending_invite(&self) -> bool {
669        self.invitation
670            .find_dialog_id_by_session_id(&self.session_id)
671            .is_some()
672    }
673
674    /// One-shot hangup for error-cleanup paths outside the actor loop.
675    async fn hangup_now(&self, reason: Option<CallRecordHangupReason>) {
676        self.do_hangup(reason, None, None, None).await.ok();
677    }
678
679    /// One-shot reject for teardown paths outside the actor loop.
680    async fn reject_now(&self, code: Option<rsipstack::rsip::StatusCode>, reason: Option<String>) {
681        self.do_reject(code, reason).await.ok();
682    }
683}
684
685/// Answer prepared during ringing: the SDP to answer with, the already
686/// running caller track (started for early media), and the dialog to accept.
687pub struct ReadyAnswer {
688    pub answer: String,
689    pub track: PendingCallerTrack,
690    pub dialog: InviteDialog,
691}
692
693/// Answer prepared for the underlying inbound SIP dialog of a non-SIP call
694/// type (WebSocket/Webrtc): the 200 OK SDP, the dialog to accept, and the RTP
695/// track that bridges the SIP leg into the media stream.
696pub struct PendingSipAnswer {
697    pub answer: String,
698    pub dialog: InviteDialog,
699    pub track: Box<dyn Track>,
700}
701
702pub struct ActiveCallGuard {
703    pub call: ActiveCallRef,
704    pub active_calls: usize,
705}
706
707impl ActiveCallGuard {
708    pub fn new(call: ActiveCallRef) -> Self {
709        let active_calls = {
710            call.app_state
711                .total_calls
712                .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
713            let mut calls = call.app_state.active_calls.lock().unwrap();
714            calls.insert(call.session_id.clone(), call.clone());
715            calls.len()
716        };
717        Self { call, active_calls }
718    }
719}
720
721impl Drop for ActiveCallGuard {
722    fn drop(&mut self) {
723        self.call
724            .app_state
725            .active_calls
726            .lock()
727            .unwrap()
728            .remove(&self.call.session_id);
729        // Idempotent cleanup of the incoming-call session-id mapping (no-op
730        // for outgoing/websocket sessions that were never registered).
731        self.call
732            .invitation
733            .unregister_session(&self.call.session_id);
734    }
735}
736
737pub struct ActiveCallReceiver {
738    pub cmd_receiver: CommandReceiver,
739    pub dump_cmd_receiver: CommandReceiver,
740    pub dump_event_receiver: EventReceiver,
741}
742
743/// Construction parameters for [`ActiveCall::new`].
744pub struct CallSpec {
745    pub call_type: ActiveCallType,
746    pub cancel_token: CancellationToken,
747    pub session_id: String,
748    pub invitation: Invitation,
749    pub app_state: AppState,
750    pub track_config: TrackConfig,
751    /// WebSocket audio receiver (WebSocket calls only), taken once by setup.
752    pub audio_receiver: Option<WebsocketBytesReceiver>,
753    pub dump_events: bool,
754    /// Overrides the default server-side track id.
755    pub server_side_track_id: Option<TrackId>,
756    /// Initial session variables; built-ins (session id, call type, start
757    /// time) are injected for missing keys.
758    pub extras: Option<HashMap<String, serde_json::Value>>,
759}
760
761impl ActiveCall {
762    pub fn new(spec: CallSpec) -> Self {
763        let CallSpec {
764            call_type,
765            cancel_token,
766            session_id,
767            invitation,
768            app_state,
769            track_config,
770            audio_receiver,
771            dump_events,
772            server_side_track_id,
773            extras,
774        } = spec;
775        let event_sender = crate::event::create_event_sender();
776        let cmd_sender = tokio::sync::broadcast::Sender::<Command>::new(32);
777        let server_side_track_id = server_side_track_id.unwrap_or(SERVER_SIDE_TRACK_ID.to_string());
778        let media_stream_builder = MediaStreamBuilder::new(event_sender.clone())
779            .with_id(session_id.clone())
780            .with_cancel_token(cancel_token.child_token());
781        let media_stream = Arc::new(media_stream_builder.build());
782        let start_time = Utc::now();
783        // Inject built-in session variables into extras
784        let call_type_str = match &call_type {
785            ActiveCallType::Sip => "sip",
786            ActiveCallType::WebSocket => "websocket",
787            ActiveCallType::Webrtc => "webrtc",
788            ActiveCallType::B2bua => "b2bua",
789        };
790        let mut extras = extras.unwrap_or_default();
791        extras
792            .entry(crate::playbook::BUILTIN_SESSION_ID.to_string())
793            .or_insert_with(|| serde_json::Value::String(session_id.clone()));
794        extras
795            .entry(crate::playbook::BUILTIN_CALL_TYPE.to_string())
796            .or_insert_with(|| serde_json::Value::String(call_type_str.to_string()));
797        extras
798            .entry(crate::playbook::BUILTIN_START_TIME.to_string())
799            .or_insert_with(|| serde_json::Value::String(start_time.to_rfc3339()));
800
801        let progress = CallProgress {
802            session_id: session_id.clone(),
803            start_time: Some(start_time),
804            ..Default::default()
805        };
806
807        Self {
808            cancel_token,
809            call_type,
810            session_id,
811            start_time,
812            media_stream,
813            track_config,
814            event_sender,
815            app_state,
816            invitation,
817            cmd_sender,
818            dump_events,
819            server_side_track_id,
820            ssrc: rand::random::<u32>(),
821            bridge_paused: Arc::new(AtomicBool::new(false)),
822            progress: Arc::new(ArcSwap::from_pointee(progress)),
823            extras: Arc::new(ArcSwap::from_pointee(extras)),
824            moh: ArcSwapOption::new(None),
825            current_play_id: ArcSwapOption::new(None),
826            tts_handle: ArcSwapOption::new(None),
827            refer_leg: ArcSwapOption::new(None),
828            ready_to_answer: ArcSwapOption::new(None),
829            pending_sip_answer: ArcSwapOption::new(None),
830            refer_call_token: ArcSwapOption::new(None),
831            wait_input_timeout: ArcSwapOption::new(None),
832            pending_asr_resume: ArcSwapOption::new(None),
833            audio_receiver: std::sync::Mutex::new(audio_receiver),
834        }
835    }
836
837    pub async fn enqueue_command(&self, command: Command) -> Result<()> {
838        self.cmd_sender
839            .send(command)
840            .map_err(|e| anyhow::anyhow!("Failed to send command: {}", e))?;
841        Ok(())
842    }
843
844    /// Create a new ActiveCallReceiver for this ActiveCall
845    /// `tokio::sync::broadcast` not cached messages, so need to early create receiver
846    /// before calling `serve()`
847    pub fn new_receiver(&self) -> ActiveCallReceiver {
848        ActiveCallReceiver {
849            cmd_receiver: self.cmd_sender.subscribe(),
850            dump_cmd_receiver: self.cmd_sender.subscribe(),
851            dump_event_receiver: self.event_sender.subscribe(),
852        }
853    }
854
855    /// The call actor: a single select loop owning the [`CallRuntime`].
856    ///
857    /// Commands, session events, background-completion messages, the
858    /// wait-input timeout tick, the media stream and cancellation all meet in
859    /// one place, so runtime state is plain (lock-free) data owned by this
860    /// task — mirroring the concurrency of the previous separate
861    /// command/event loops without their shared locks.
862    pub async fn serve(self: Arc<Self>, receiver: ActiveCallReceiver) -> Result<()> {
863        let ActiveCallReceiver {
864            mut cmd_receiver,
865            dump_cmd_receiver,
866            dump_event_receiver,
867        } = receiver;
868
869        let mut event_receiver = self.event_sender.subscribe();
870        let (actor_tx, mut actor_rx) = mpsc::channel::<ActorMsg>(16);
871        let mut runtime = CallRuntime::new(actor_tx);
872        runtime.me = Some(self.clone());
873
874        self.app_state
875            .total_calls
876            .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
877
878        let me = self.clone();
879        let actor = async move {
880            // RAII: whatever way the loop exits (break, panic, early drop),
881            // the call's cancel token fires so every child task and the media
882            // stream shut down.
883            let _cancel_on_exit = CancelOnExit(&me.cancel_token);
884            let mut ticker = tokio::time::interval(Duration::from_millis(100));
885            // Keep the media-serve future alive across select iterations.
886            let mut media_serve = Box::pin(me.media_stream.serve());
887            loop {
888                tokio::select! {
889                    cmd = cmd_receiver.recv() => {
890                        match cmd {
891                            Ok(command) => {
892                                // Box::pin keeps the deep do_* future tree off
893                                // the select's stack frame (debug builds overflow
894                                // otherwise).
895                                if let Err(e) = Box::pin(me.dispatch(&mut runtime, command)).await {
896                                    warn!(session_id = me.session_id, "{}", e);
897                                    me.event_sender
898                                        .send(SessionEvent::Error {
899                                            track_id: me.session_id.clone(),
900                                            timestamp: crate::media::get_timestamp(),
901                                            sender: "command".to_string(),
902                                            error: e.to_string(),
903                                            code: None,
904                                        })
905                                        .ok();
906                                }
907                            }
908                            Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue,
909                            Err(_) => {
910                                info!(session_id = me.session_id, "command loop done");
911                                break;
912                            }
913                        }
914                    }
915                    ev = event_receiver.recv() => {
916                        match ev {
917                            Ok(event) => Box::pin(me.handle_event(&mut runtime, event)).await,
918                            Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue,
919                            Err(_) => {
920                                info!(session_id = me.session_id, "event loop done");
921                                break;
922                            }
923                        }
924                    }
925                    Some(msg) = actor_rx.recv() => {
926                        if let Err(e) = Box::pin(me.handle_actor_msg(msg)).await {
927                            warn!(session_id = me.session_id, "{}", e);
928                            me.event_sender
929                                .send(SessionEvent::Error {
930                                    track_id: me.session_id.clone(),
931                                    timestamp: crate::media::get_timestamp(),
932                                    sender: "command".to_string(),
933                                    error: e.to_string(),
934                                    code: None,
935                                })
936                                .ok();
937                        }
938                    }
939                    _ = ticker.tick() => {
940                        Box::pin(me.check_input_timeout(&mut runtime)).await;
941                    }
942                    _ = &mut media_serve => {
943                        info!(session_id = me.session_id, "media stream loop done");
944                        break;
945                    }
946                    _ = me.cancel_token.cancelled() => {
947                        info!(session_id = me.session_id, "call cancelled - cleaning up resources");
948                        break;
949                    }
950                }
951            }
952        };
953
954        tokio::join!(
955            self.dump_loop(self.dump_events, dump_cmd_receiver, dump_event_receiver),
956            actor
957        );
958        Ok(())
959    }
960
961    /// Wait-input silence timeout tick (formerly its own loop + mutex).
962    async fn check_input_timeout(&self, runtime: &mut CallRuntime) {
963        let (start_time, expire) = runtime.input_timeout_expire;
964        if expire > 0 && crate::media::get_timestamp() >= start_time + expire as u64 {
965            info!(session_id = self.session_id, "wait input timeout reached");
966            runtime.input_timeout_expire = (0, 0);
967            self.event_sender
968                .send(SessionEvent::Silence {
969                    track_id: self.server_side_track_id.clone(),
970                    timestamp: crate::media::get_timestamp(),
971                    start_time,
972                    duration: expire as u64,
973                    samples: None,
974                    refer: Some(false),
975                })
976                .ok();
977        }
978    }
979
980    /// Handle a session event (formerly the concurrent event-hook loop).
981    async fn handle_event(&self, runtime: &mut CallRuntime, event: SessionEvent) {
982        match event {
983            SessionEvent::Speaking { .. }
984            | SessionEvent::Dtmf { .. }
985            | SessionEvent::AsrDelta { .. }
986            | SessionEvent::AsrFinal { .. }
987            | SessionEvent::TrackStart { .. } => {
988                runtime.input_timeout_expire = (0, 0);
989            }
990            SessionEvent::TrackEnd {
991                track_id,
992                play_id,
993                ssrc,
994                auto_hangup,
995                ..
996            } => {
997                if track_id != self.server_side_track_id {
998                    return;
999                }
1000
1001                if play_id != self.current_play() {
1002                    debug!(
1003                        session_id = self.session_id,
1004                        ?play_id,
1005                        current = ?self.current_play(),
1006                        "ignoring interrupted track end"
1007                    );
1008                    return;
1009                }
1010                self.set_current_play(None);
1011                let moh_path = self.moh_path();
1012                let wait_timeout_val = self.take_wait_input_timeout();
1013
1014                if let Some(path) = moh_path {
1015                    info!(session_id = self.session_id, "looping moh: {}", path);
1016                    let ssrc = rand::random::<u32>();
1017                    let file_track = self.make_file_track(path.clone(), ssrc);
1018                    self.update_track_wrapper(Box::new(file_track), Some(path))
1019                        .await;
1020                    return;
1021                }
1022
1023                if let Some(hangup_reason) = auto_hangup {
1024                    info!(
1025                        session_id = self.session_id,
1026                        ssrc, "auto hangup when track end track_id:{}", track_id
1027                    );
1028                    self.do_hangup(Some(hangup_reason), None, None, None)
1029                        .await
1030                        .ok();
1031                }
1032
1033                if let Some(timeout) = wait_timeout_val {
1034                    runtime.input_timeout_expire = if timeout > 0 {
1035                        (crate::media::get_timestamp(), timeout)
1036                    } else {
1037                        (0, 0)
1038                    };
1039                }
1040            }
1041            SessionEvent::Interrupt { receiver } => {
1042                let track_id = receiver.unwrap_or_else(|| self.server_side_track_id.clone());
1043                if track_id == self.server_side_track_id {
1044                    debug!(
1045                        session_id = self.session_id,
1046                        "received interrupt event, stopping playback"
1047                    );
1048                    self.do_interrupt(true).await.ok();
1049                }
1050            }
1051            SessionEvent::Inactivity { track_id, .. } => {
1052                info!(
1053                    session_id = self.session_id,
1054                    track_id, "inactivity timeout reached, hanging up"
1055                );
1056                self.do_hangup(
1057                    Some(CallRecordHangupReason::InactivityTimeout),
1058                    None,
1059                    None,
1060                    None,
1061                )
1062                .await
1063                .ok();
1064            }
1065            SessionEvent::Hangup { refer, .. } => {
1066                // Check if we need to resume ASR after refer hangup
1067                if refer == Some(true) {
1068                    if let Some((refer_ssrc, asr_option)) = self.take_pending_asr_resume() {
1069                        // Verify it's the refer call that ended
1070                        let is_refer_hangup = self
1071                            .refer_leg
1072                            .load_full()
1073                            .map(|leg| leg.ssrc == refer_ssrc)
1074                            .unwrap_or(false);
1075
1076                        if is_refer_hangup {
1077                            info!(
1078                                session_id = self.session_id,
1079                                "Refer call ended, resuming parent ASR"
1080                            );
1081
1082                            // Resume ASR
1083                            match self
1084                                .app_state
1085                                .stream_engine
1086                                .create_asr_processor(
1087                                    self.server_side_track_id.clone(),
1088                                    self.cancel_token.child_token(),
1089                                    asr_option,
1090                                    self.event_sender.clone(),
1091                                )
1092                                .await
1093                            {
1094                                Ok(asr_processor) => {
1095                                    if let Err(e) = self
1096                                        .media_stream
1097                                        .append_processor(&self.server_side_track_id, asr_processor)
1098                                        .await
1099                                    {
1100                                        warn!(
1101                                            session_id = self.session_id,
1102                                            "Failed to resume ASR after refer: {}", e
1103                                        );
1104                                    }
1105                                }
1106                                Err(e) => {
1107                                    warn!(
1108                                        session_id = self.session_id,
1109                                        "Failed to create ASR processor for resume: {}", e
1110                                    );
1111                                }
1112                            }
1113                        }
1114                    }
1115                }
1116            }
1117            SessionEvent::Error { track_id, .. } => {
1118                if track_id != self.server_side_track_id {
1119                    return;
1120                }
1121
1122                let moh_info = {
1123                    let path = self.moh_path();
1124                    path.map(|path| {
1125                        let fallback = "./config/sounds/refer_moh.wav".to_string();
1126                        if path != fallback && std::path::Path::new(&fallback).exists() {
1127                            info!(
1128                                session_id = self.session_id,
1129                                "moh error, switching to fallback: {}", fallback
1130                            );
1131                            self.set_moh(Some(fallback.clone()));
1132                            fallback
1133                        } else {
1134                            info!(
1135                                session_id = self.session_id,
1136                                "looping moh on error: {}", path
1137                            );
1138                            path
1139                        }
1140                    })
1141                };
1142
1143                if let Some(next_path) = moh_info {
1144                    let ssrc = rand::random::<u32>();
1145                    let file_track = self.make_file_track(next_path.clone(), ssrc);
1146                    self.update_track_wrapper(Box::new(file_track), Some(next_path))
1147                        .await;
1148                }
1149            }
1150            SessionEvent::Hold { on_hold, .. } => {
1151                self.bridge_paused.store(on_hold, Ordering::Relaxed);
1152            }
1153            _ => {}
1154        }
1155    }
1156
1157    /// Completion of background work spawned by a `do_*` command.
1158    async fn handle_actor_msg(&self, msg: ActorMsg) -> Result<()> {
1159        match msg {
1160            ActorMsg::ReferDone {
1161                track_id,
1162                forward_dtmf,
1163                result,
1164            } => match result {
1165                Ok(answer) => {
1166                    self.media_stream
1167                        .set_track_refer(&track_id, Some(true))
1168                        .await;
1169                    if !forward_dtmf {
1170                        self.media_stream
1171                            .set_track_dtmf_forward(&track_id, false)
1172                            .await;
1173                    }
1174                    self.event_sender
1175                        .send(SessionEvent::Answer {
1176                            timestamp: crate::media::get_timestamp(),
1177                            track_id,
1178                            sdp: answer,
1179                            refer: Some(true),
1180                        })
1181                        .ok();
1182                    Ok(())
1183                }
1184                Err(e) => {
1185                    warn!(
1186                        session_id = self.session_id,
1187                        "failed to create refer sip track: {}", e
1188                    );
1189                    self.emit_reject_from_rsip_error(track_id, true, &e);
1190                    Err(e.into())
1191                }
1192            },
1193        }
1194    }
1195
1196    async fn dispatch(&self, runtime: &mut CallRuntime, command: Command) -> Result<()> {
1197        match command {
1198            Command::Invite { option } => self.do_invite(runtime, option).await,
1199            Command::Accept { option } => self.do_accept(option).await,
1200            Command::Reject { reason, code } => {
1201                self.do_reject(code.map(|c| (c as u16).into()), Some(reason))
1202                    .await
1203            }
1204            Command::Ringing { .. } => self.do_ringing(command).await,
1205            Command::Tts { .. } => self.do_tts(command).await,
1206            Command::Play { .. } => self.do_play(command).await,
1207            Command::Hangup {
1208                reason,
1209                initiator,
1210                headers,
1211                refer,
1212            } => {
1213                let reason = reason.map(|r| {
1214                    r.parse::<CallRecordHangupReason>()
1215                        .unwrap_or(CallRecordHangupReason::BySystem)
1216                });
1217                self.do_hangup(reason, initiator, headers, refer).await
1218            }
1219            Command::Refer {
1220                caller,
1221                callee,
1222                options,
1223            } => self.do_refer(runtime, caller, callee, options).await,
1224            Command::Message {
1225                body,
1226                content_type,
1227                headers,
1228                refer,
1229            } => self.do_message(body, content_type, headers, refer).await,
1230            Command::Bridge { target_session_id } => self.do_bridge(target_session_id).await,
1231            Command::Unbridge { target_session_id } => self.do_unbridge(target_session_id).await,
1232            Command::Mute { track_id } => self.do_mute(track_id).await,
1233            Command::Unmute { track_id } => self.do_unmute(track_id).await,
1234            Command::Pause {} => self.do_pause().await,
1235            Command::Resume {} => self.do_resume().await,
1236            Command::Interrupt {
1237                graceful: passage,
1238                fade_out_ms: _,
1239            } => self.do_interrupt(passage.unwrap_or_default()).await,
1240            Command::History { speaker, text } => self.do_history(speaker, text).await,
1241            Command::Custom { sender, data } => self.do_custom(sender, data),
1242            Command::AddIceCandidate {
1243                candidate,
1244                sdp_mid,
1245                sdp_mline_index,
1246            } => {
1247                self.media_stream
1248                    .add_ice_candidate(&candidate, sdp_mid.as_deref(), sdp_mline_index)
1249                    .await
1250            }
1251        }
1252    }
1253
1254    fn build_record_option(&self, option: &CallOption) -> Option<RecorderOption> {
1255        if let Some(recorder_option) = &option.recorder {
1256            let mut recorder_file = recorder_option.recorder_file.clone();
1257            if recorder_file.contains("{id}") {
1258                recorder_file = recorder_file.replace("{id}", &self.session_id);
1259            }
1260
1261            let recorder_file = if recorder_file.is_empty() {
1262                self.app_state.get_recorder_file(&self.session_id)
1263            } else {
1264                let p = Path::new(&recorder_file);
1265                p.is_absolute()
1266                    .then(|| recorder_file.clone())
1267                    .unwrap_or_else(|| self.app_state.get_recorder_file(&recorder_file))
1268            };
1269            info!(
1270                session_id = self.session_id,
1271                recorder_file, "created recording file"
1272            );
1273
1274            let track_samplerate = self.track_config.samplerate;
1275            let recorder_samplerate = if track_samplerate > 0 {
1276                track_samplerate
1277            } else {
1278                recorder_option.samplerate
1279            };
1280            let recorder_ptime = if recorder_option.ptime == 0 {
1281                200
1282            } else {
1283                recorder_option.ptime
1284            };
1285            let requested_format = recorder_option
1286                .format
1287                .unwrap_or(self.app_state.config.recorder_format());
1288            let format = requested_format.effective();
1289            if requested_format != format {
1290                warn!(
1291                    session_id = self.session_id,
1292                    requested = requested_format.extension(),
1293                    "Recorder format fallback to wav due to unsupported feature"
1294                );
1295            }
1296            let mut recorder_config = RecorderOption {
1297                recorder_file,
1298                samplerate: recorder_samplerate,
1299                ptime: recorder_ptime,
1300                format: Some(format),
1301                native_samplerate: Some(
1302                    recorder_option.native_samplerate.unwrap_or(false)
1303                        || self.app_state.config.recorder_native_samplerate(),
1304                ),
1305            };
1306            recorder_config.ensure_path_extension(format);
1307            Some(recorder_config)
1308        } else {
1309            None
1310        }
1311    }
1312
1313    async fn invite_or_accept(&self, mut option: CallOption, sender: String) -> Result<CallOption> {
1314        // Merge with existing configuration (e.g., from playbook)
1315        {
1316            let state = self.progress.load_full();
1317            option = state.merge_option(option);
1318        }
1319
1320        option.check_default();
1321        if let Some(opt) = self.build_record_option(&option) {
1322            self.media_stream.update_recorder_option(opt).await;
1323        }
1324        self.ensure_call_ambiance(&option).await;
1325
1326        if let Some(opt) = &option.media_pass {
1327            let track_id = self.server_side_track_id.clone();
1328            let cancel_token = self.cancel_token.child_token();
1329            let ssrc = rand::random::<u32>();
1330            let media_pass_track = MediaPassTrack::new(
1331                self.session_id.clone(),
1332                ssrc,
1333                track_id,
1334                cancel_token,
1335                opt.clone(),
1336            );
1337            self.update_track_wrapper(Box::new(media_pass_track), None)
1338                .await;
1339        }
1340
1341        info!(
1342            session_id = self.session_id,
1343            call_type = ?self.call_type,
1344            sender,
1345            ?option,
1346            "caller with option"
1347        );
1348
1349        match self.setup_caller_track(&option).await {
1350            Ok(_) => return Ok(option),
1351            Err(e) => {
1352                self.app_state
1353                    .total_failed_calls
1354                    .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1355                let error_event = crate::event::SessionEvent::Error {
1356                    track_id: self.session_id.clone(),
1357                    timestamp: crate::media::get_timestamp(),
1358                    sender,
1359                    error: e.to_string(),
1360                    code: None,
1361                };
1362                self.event_sender.send(error_event).ok();
1363                self.hangup_now(Some(CallRecordHangupReason::BySystem))
1364                    .await;
1365                return Err(e);
1366            }
1367        }
1368    }
1369
1370    async fn do_invite(&self, runtime: &mut CallRuntime, option: CallOption) -> Result<()> {
1371        // Run the INVITE handshake in the background so the actor keeps
1372        // serving media (and everything else) while the call rings - same
1373        // reasoning as do_refer below. invite_or_accept blocks on the SIP
1374        // transaction for the entire ring duration; handling that inline on
1375        // the actor's own select loop would freeze media_stream.serve() (and
1376        // therefore all live audio bridging for this leg) for just as long.
1377        let me = runtime
1378            .me
1379            .clone()
1380            .ok_or_else(|| anyhow::anyhow!("invite is only supported inside serve()"))?;
1381        crate::spawn(async move {
1382            if let Err(e) = me.invite_or_accept(option, "invite".to_string()).await {
1383                warn!(session_id = me.session_id, "{}", e);
1384                me.event_sender
1385                    .send(SessionEvent::Error {
1386                        track_id: me.session_id.clone(),
1387                        timestamp: crate::media::get_timestamp(),
1388                        sender: "command".to_string(),
1389                        error: e.to_string(),
1390                        code: None,
1391                    })
1392                    .ok();
1393            }
1394        });
1395        Ok(())
1396    }
1397
1398    async fn do_accept(&self, mut option: CallOption) -> Result<()> {
1399        let has_pending = self.has_pending_invite();
1400        let ready_to_answer_val = !self.has_ready_to_answer();
1401
1402        if ready_to_answer_val {
1403            if !has_pending {
1404                // emit reject event
1405                warn!(session_id = self.session_id, "no pending call to accept");
1406                let rejet_event = crate::event::SessionEvent::Reject {
1407                    track_id: self.session_id.clone(),
1408                    timestamp: crate::media::get_timestamp(),
1409                    reason: "no pending call".to_string(),
1410                    refer: None,
1411                    code: Some(486),
1412                };
1413                self.event_sender.send(rejet_event).ok();
1414                self.hangup_now(Some(CallRecordHangupReason::BySystem))
1415                    .await;
1416                return Err(anyhow::anyhow!("no pending call to accept"));
1417            }
1418            option = self.invite_or_accept(option, "accept".to_string()).await?;
1419        } else {
1420            option.check_default();
1421            if let Some(opt) = self.build_record_option(&option) {
1422                self.media_stream.update_recorder_option(opt).await;
1423            }
1424            self.set_option(option.clone());
1425            self.ensure_call_ambiance(&option).await;
1426        }
1427        info!(session_id = self.session_id, ?option, "accepting call");
1428        let ready = self.take_ready_to_answer();
1429        if let Some(ready) = ready {
1430            // Exclusive since `take` swapped it out and the actor serializes
1431            // commands, so the unwrap cannot race another holder.
1432            let ReadyAnswer {
1433                answer,
1434                track: pending_track,
1435                dialog,
1436            } = match Arc::try_unwrap(ready) {
1437                Ok(ready) => ready,
1438                Err(_) => {
1439                    warn!(
1440                        session_id = self.session_id,
1441                        "ready_to_answer held elsewhere; skipping accept"
1442                    );
1443                    return Ok(());
1444                }
1445            };
1446            info!(session_id = self.session_id, "ready to answer with track");
1447
1448            let headers = vec![rsipstack::rsip::Header::ContentType(
1449                "application/sdp".to_string().into(),
1450            )];
1451
1452            match dialog.accept(Some(headers), Some(answer.as_bytes().to_vec())) {
1453                Ok(_) => {
1454                    self.leg().update_progress(|p| {
1455                        p.answer = Some(answer.clone());
1456                        p.answer_time.get_or_insert_with(Utc::now);
1457                    });
1458                    self.finish_caller_stack(&option, pending_track).await?;
1459                }
1460                Err(e) => {
1461                    warn!(session_id = self.session_id, "failed to accept call: {}", e);
1462                    return Err(anyhow::anyhow!("failed to accept call"));
1463                }
1464            }
1465        }
1466
1467        // Non-SIP call types (WebSocket/Webrtc) created by an inbound SIP
1468        // INVITE never go through `ready_to_answer`; the underlying dialog is
1469        // answered here. Without this the carrier leg stays in `Trying` until
1470        // the far end times out and CANCELs (production: VOS3000 20s timeout).
1471        if let Some(pending) = self.take_pending_sip_answer() {
1472            let Ok(pending) = Arc::try_unwrap(pending) else {
1473                warn!(
1474                    session_id = self.session_id,
1475                    "pending sip answer held elsewhere; skipping sip accept"
1476                );
1477                return Ok(());
1478            };
1479            let headers = vec![rsipstack::rsip::Header::ContentType(
1480                "application/sdp".to_string().into(),
1481            )];
1482            match pending
1483                .dialog
1484                .accept(Some(headers), Some(pending.answer.as_bytes().to_vec()))
1485            {
1486                Ok(_) => {
1487                    info!(
1488                        session_id = self.session_id,
1489                        "answered underlying sip dialog"
1490                    );
1491                    self.leg().update_progress(|p| {
1492                        p.answer = Some(pending.answer.clone());
1493                        p.answer_time.get_or_insert_with(Utc::now);
1494                    });
1495                    // Register the SIP leg so customer audio is bridged with
1496                    // the caller track (and later the refer leg).
1497                    self.media_stream.update_track(pending.track, None).await;
1498                }
1499                Err(e) => {
1500                    warn!(
1501                        session_id = self.session_id,
1502                        "failed to accept underlying sip dialog: {}", e
1503                    );
1504                }
1505            }
1506        }
1507        Ok(())
1508    }
1509
1510    async fn do_reject(
1511        &self,
1512        code: Option<rsipstack::rsip::StatusCode>,
1513        reason: Option<String>,
1514    ) -> Result<()> {
1515        if let Some(pending) = self.take_pending_sip_answer() {
1516            info!(
1517                session_id = self.session_id,
1518                ?reason,
1519                ?code,
1520                "rejecting underlying sip dialog"
1521            );
1522            if let Ok(pending) = Arc::try_unwrap(pending) {
1523                pending.dialog.reject(code.clone(), reason.clone()).ok();
1524                self.invitation.dialog_layer.remove_dialog(&pending.dialog.id());
1525            }
1526        }
1527        match self
1528            .invitation
1529            .find_dialog_id_by_session_id(&self.session_id)
1530        {
1531            Some(id) => {
1532                info!(
1533                    session_id = self.session_id,
1534                    ?reason,
1535                    ?code,
1536                    "rejecting call"
1537                );
1538                let result = self.invitation.hangup(id, code, reason).await;
1539                if result.is_ok() {
1540                    self.cancel_token.cancel();
1541                }
1542                result
1543            }
1544            None => {
1545                if let Some(ready) = self.take_ready_to_answer() {
1546                    info!(
1547                        session_id = self.session_id,
1548                        ?reason,
1549                        ?code,
1550                        "rejecting call from ready_to_answer"
1551                    );
1552                    let dialog = &ready.dialog;
1553                    let dialog_id = dialog.id();
1554                    dialog.reject(code, reason).ok();
1555                    self.invitation.dialog_layer.remove_dialog(&dialog_id);
1556                    self.cancel_token.cancel();
1557                }
1558                Ok(())
1559            }
1560        }
1561    }
1562
1563    async fn do_ringing(&self, command: Command) -> Result<()> {
1564        let Command::Ringing {
1565            ringtone,
1566            recorder,
1567            early_media,
1568        } = command
1569        else {
1570            unreachable!("do_ringing called with non-Ringing command");
1571        };
1572
1573        if !self.has_ready_to_answer() {
1574            let option = CallOption {
1575                recorder,
1576                ..Default::default()
1577            };
1578            let _ = self.invite_or_accept(option, "ringing".to_string()).await?;
1579        }
1580
1581        if let Some(ready) = self.ready_to_answer.load_full() {
1582            let (headers, body) = if early_media.unwrap_or_default() || ringtone.is_some() {
1583                let headers = vec![rsipstack::rsip::Header::ContentType(
1584                    "application/sdp".to_string().into(),
1585                )];
1586                (Some(headers), Some(ready.answer.as_bytes().to_vec()))
1587            } else {
1588                (None, None)
1589            };
1590
1591            ready.dialog.ringing(headers, body).ok();
1592            info!(
1593                session_id = self.session_id,
1594                ringtone, early_media, "playing ringtone"
1595            );
1596            if let Some(ringtone_url) = ringtone {
1597                self.do_play(Command::Play {
1598                    url: ringtone_url,
1599                    play_id: None,
1600                    auto_hangup: None,
1601                    wait_input_timeout: None,
1602                    offset_ms: None,
1603                })
1604                .await
1605                .ok();
1606            } else {
1607                info!(session_id = self.session_id, "no ringtone to play");
1608            }
1609        }
1610        Ok(())
1611    }
1612
1613    async fn do_tts(&self, command: Command) -> Result<()> {
1614        let Command::Tts {
1615            text,
1616            speaker,
1617            play_id,
1618            auto_hangup,
1619            streaming,
1620            end_of_stream,
1621            option,
1622            wait_input_timeout,
1623            base64,
1624            cache_key,
1625        } = command
1626        else {
1627            unreachable!("do_tts called with non-Tts command");
1628        };
1629        let streaming = streaming.unwrap_or_default();
1630        let end_of_stream = end_of_stream.unwrap_or_default();
1631        let base64 = base64.unwrap_or_default();
1632
1633        let tts_option = {
1634            let call_state = self.progress.load_full();
1635            match call_state.option.clone().unwrap_or_default().tts {
1636                Some(opt) => opt.merge_with(option),
1637                None => {
1638                    if let Some(opt) = option {
1639                        opt
1640                    } else {
1641                        return Err(anyhow::anyhow!("no tts option available"));
1642                    }
1643                }
1644            }
1645        };
1646        let speaker = match speaker {
1647            Some(s) => Some(s),
1648            None => tts_option.speaker.clone(),
1649        };
1650
1651        let mut play_command = SynthesisCommand {
1652            text,
1653            speaker,
1654            play_id: play_id.clone(),
1655            streaming,
1656            end_of_stream: if !streaming { true } else { end_of_stream },
1657            option: tts_option,
1658            base64,
1659            cache_key,
1660            auto_hangup,
1661        };
1662        info!(
1663            session_id = self.session_id,
1664            provider = ?play_command.option.provider,
1665            text = %play_command.text.chars().take(10).collect::<String>(),
1666            speaker = play_command.speaker.as_deref(),
1667            auto_hangup = auto_hangup.unwrap_or_default(),
1668            play_id = play_command.play_id.as_deref(),
1669            streaming = play_command.streaming,
1670            end_of_stream = play_command.end_of_stream,
1671            wait_input_timeout = wait_input_timeout.unwrap_or_default(),
1672            is_base64 = play_command.base64,
1673            cache_key = play_command.cache_key.as_deref(),
1674            "new synthesis"
1675        );
1676
1677        let ssrc = rand::random::<u32>();
1678        let (should_interrupt, picked_ssrc) = {
1679            let existing_handle = self.tts_handle.load_full();
1680            let current_play_id = self.current_play();
1681
1682            let (target_ssrc, changed) = if let Some(handle) = &existing_handle {
1683                if play_id.is_some() && current_play_id != play_id {
1684                    (ssrc, true)
1685                } else {
1686                    (handle.ssrc, false)
1687                }
1688            } else {
1689                (ssrc, false)
1690            };
1691
1692            // Defer auto_hangup setting until after potential interrupt.
1693            // auto_hangup will be set below after do_interrupt() to avoid being cleared.
1694            self.set_wait_input_timeout(wait_input_timeout);
1695
1696            self.set_current_play(play_id.clone());
1697            (changed, target_ssrc)
1698        };
1699
1700        if should_interrupt {
1701            let _ = self.do_interrupt(false).await;
1702        }
1703
1704        // auto_hangup rides on the track: armed via `with_auto_hangup` when a
1705        // new track is created, or by the command itself for an existing track.
1706
1707        let existing_handle = self.tts_handle.load_full();
1708        if let Some(tts_handle) = existing_handle {
1709            match tts_handle.try_send(play_command) {
1710                Ok(_) => return Ok(()),
1711                Err(e) => {
1712                    play_command = e.0;
1713                }
1714            }
1715        }
1716
1717        let (new_handle, tts_track) = StreamEngine::create_tts_track(
1718            self.app_state.stream_engine.clone(),
1719            self.cancel_token.child_token(),
1720            self.session_id.clone(),
1721            self.server_side_track_id.clone(),
1722            picked_ssrc,
1723            play_id.clone(),
1724            streaming,
1725            &play_command.option,
1726            play_command.auto_hangup,
1727        )
1728        .await?;
1729
1730        new_handle.try_send(play_command)?;
1731        self.tts_handle.store(Some(Arc::new(new_handle)));
1732        self.update_track_wrapper(tts_track, play_id).await;
1733        Ok(())
1734    }
1735
1736    async fn do_play(&self, command: Command) -> Result<()> {
1737        let Command::Play {
1738            url,
1739            play_id,
1740            auto_hangup,
1741            wait_input_timeout,
1742            offset_ms,
1743        } = command
1744        else {
1745            unreachable!("do_play called with non-Play command");
1746        };
1747        let ssrc = rand::random::<u32>();
1748        info!(
1749            session_id = self.session_id,
1750            ssrc, url, play_id, auto_hangup, "play file track"
1751        );
1752
1753        let play_id = play_id.or(Some(url.clone()));
1754
1755        // make_file_track uses the path as play_id; honor an explicit play_id here.
1756        let mut file_track = self
1757            .make_file_track(url, ssrc)
1758            .with_play_id(play_id.clone())
1759            .with_auto_hangup(auto_hangup);
1760
1761        if let Some(offset) = offset_ms {
1762            file_track = file_track.with_offset_ms(offset);
1763        }
1764
1765        {
1766            self.tts_handle.store(None);
1767            self.set_wait_input_timeout(wait_input_timeout);
1768        }
1769
1770        self.update_track_wrapper(Box::new(file_track), play_id)
1771            .await;
1772        Ok(())
1773    }
1774
1775    async fn do_history(&self, speaker: String, text: String) -> Result<()> {
1776        self.event_sender
1777            .send(SessionEvent::AddHistory {
1778                sender: Some(self.session_id.clone()),
1779                timestamp: crate::media::get_timestamp(),
1780                speaker,
1781                text,
1782            })
1783            .map(|_| ())
1784            .map_err(Into::into)
1785    }
1786
1787    fn do_custom(&self, sender: Option<String>, data: serde_json::Value) -> Result<()> {
1788        self.event_sender
1789            .send(SessionEvent::Custom {
1790                timestamp: crate::media::get_timestamp(),
1791                sender,
1792                data,
1793            })
1794            .map(|_| ())
1795            .map_err(Into::into)
1796    }
1797
1798    async fn do_interrupt(&self, graceful: bool) -> Result<()> {
1799        {
1800            self.tts_handle.store(None);
1801            self.set_moh(None);
1802        }
1803        self.media_stream
1804            .remove_track(&self.server_side_track_id, graceful)
1805            .await;
1806        Ok(())
1807    }
1808    async fn do_pause(&self) -> Result<()> {
1809        self.media_stream
1810            .pause_playback(self.server_side_track_id.clone())
1811            .await?;
1812        Ok(())
1813    }
1814    async fn do_resume(&self) -> Result<()> {
1815        self.media_stream
1816            .resume_playback(self.server_side_track_id.clone())
1817            .await?;
1818        Ok(())
1819    }
1820    async fn do_hangup(
1821        &self,
1822        reason: Option<CallRecordHangupReason>,
1823        initiator: Option<String>,
1824        headers: Option<HashMap<String, String>>,
1825        refer: Option<bool>,
1826    ) -> Result<()> {
1827        info!(
1828            session_id = self.session_id,
1829            ?reason,
1830            ?initiator,
1831            ?headers,
1832            ?refer,
1833            "do_hangup"
1834        );
1835
1836        let hangup_reason = match initiator.as_deref() {
1837            Some("caller") => CallRecordHangupReason::ByCaller,
1838            Some("callee") => CallRecordHangupReason::ByCallee,
1839            Some("system") => CallRecordHangupReason::Autohangup,
1840            _ => reason.unwrap_or(CallRecordHangupReason::BySystem),
1841        };
1842
1843        match refer {
1844            Some(true) => {
1845                // Hang up only the refer call, leaving the main call alive.
1846                let refer_token = self.take_refer_call_token();
1847                let refer_leg = self.refer_leg_value();
1848                let has_refer_leg = refer_leg.is_some();
1849                if let Some(leg) = refer_leg {
1850                    if let Some(headers) = headers {
1851                        let h_val = serde_json::to_value(&headers).unwrap_or_default();
1852                        leg.set_extra("_hangup_headers", h_val);
1853                    }
1854                    // Set reason before cancelling so on_terminated() sees it.
1855                    let reason = hangup_reason.clone();
1856                    leg.update_progress(|p| p.set_hangup_reason(reason.clone()));
1857                }
1858                if let Some(token) = refer_token {
1859                    token.cancel();
1860                }
1861                if has_refer_leg {
1862                    self.media_stream
1863                        .remove_track(&self.server_side_track_id, false)
1864                        .await;
1865                }
1866            }
1867            _ => {
1868                if let Some(headers) = headers {
1869                    let h_val = serde_json::to_value(&headers).unwrap_or_default();
1870                    self.leg().set_extra("_hangup_headers", h_val);
1871                }
1872                self.leg()
1873                    .update_progress(|p| p.set_hangup_reason(hangup_reason.clone()));
1874                let refer_token = self.take_refer_call_token();
1875                self.media_stream
1876                    .stop(Some(hangup_reason.to_string()), initiator);
1877                if let Some(token) = refer_token {
1878                    token.cancel();
1879                }
1880            }
1881        }
1882        tokio::task::yield_now().await;
1883        Ok(())
1884    }
1885
1886    /// Initiate a refer (attended transfer) leg.
1887    ///
1888    /// The INVITE handshake can take up to `timeout` seconds, so it runs in a
1889    /// spawned task and reports back through [`ActorMsg::ReferDone`]; the
1890    /// actor loop keeps serving events (e.g. MOH looping) meanwhile.
1891    async fn do_refer(
1892        &self,
1893        runtime: &mut CallRuntime,
1894        caller: String,
1895        callee: String,
1896        refer_option: Option<ReferOption>,
1897    ) -> Result<()> {
1898        self.do_interrupt(false).await.ok();
1899
1900        // Check if we should pause parent ASR
1901        let pause_parent_asr = refer_option
1902            .as_ref()
1903            .and_then(|o| o.pause_parent_asr)
1904            .unwrap_or(false);
1905
1906        // Save original ASR option for later resume
1907        let original_asr_option = if pause_parent_asr {
1908            self.progress
1909                .load_full()
1910                .option
1911                .as_ref()
1912                .and_then(|o| o.asr.clone())
1913        } else {
1914            None
1915        };
1916
1917        // Pause parent ASR if requested
1918        if pause_parent_asr {
1919            info!(
1920                session_id = self.session_id,
1921                "Pausing parent call ASR during refer"
1922            );
1923            self.media_stream
1924                .remove_processor::<crate::media::asr_processor::AsrProcessor>(
1925                    &self.server_side_track_id,
1926                )
1927                .await
1928                .ok();
1929        }
1930
1931        let mut moh = refer_option.as_ref().and_then(|o| o.moh.clone());
1932        if let Some(ref path) = moh {
1933            if !path.starts_with("http") && !std::path::Path::new(path).exists() {
1934                let fallback = "./config/sounds/refer_moh.wav";
1935                if std::path::Path::new(fallback).exists() {
1936                    info!(
1937                        session_id = self.session_id,
1938                        "moh {} not found, using fallback {}", path, fallback
1939                    );
1940                    moh = Some(fallback.to_string());
1941                }
1942            }
1943        }
1944        let ref_call_id = refer_option
1945            .as_ref()
1946            .and_then(|o| o.call_id.clone())
1947            .unwrap_or_else(|| format!("ref-{}-{}", rand::random::<u32>(), self.session_id));
1948
1949        let session_id = self.session_id.clone();
1950        let track_id = self.server_side_track_id.clone();
1951
1952        let (recorder, parent_caller) = {
1953            let progress = self.progress.load_full();
1954            let option = progress.option.as_ref();
1955            (
1956                option.map(|o| o.recorder.clone()).unwrap_or_default(),
1957                option.and_then(|o| o.caller.clone()),
1958            )
1959        };
1960        let caller = if caller.trim().is_empty() {
1961            parent_caller.unwrap_or_default()
1962        } else {
1963            caller
1964        };
1965
1966        let mut call_option = CallOption {
1967            caller: Some(caller),
1968            callee: Some(callee.clone()),
1969            sip: refer_option.as_ref().and_then(|o| o.sip.clone()),
1970            vad: refer_option
1971                .as_ref()
1972                .and_then(|o| o.vad.clone())
1973                .map(|mut opts| {
1974                    opts.refer = Some(true);
1975                    opts
1976                }),
1977            asr: refer_option
1978                .as_ref()
1979                .and_then(|o| o.asr.clone())
1980                .map(|mut opts| {
1981                    opts.refer = Some(true);
1982                    opts
1983                }),
1984            denoise: refer_option.as_ref().and_then(|o| o.denoise.clone()),
1985            agc: refer_option.as_ref().and_then(|o| o.agc.clone()),
1986            recorder,
1987            ..Default::default()
1988        };
1989        call_option.check_default();
1990
1991        let mut invite_option = call_option.build_invite_option()?;
1992        invite_option.call_id = Some(ref_call_id.clone());
1993
1994        let headers = invite_option.headers.get_or_insert_with(|| Vec::new());
1995
1996        {
1997            let progress = self.progress.load_full();
1998            if let Some(opt) = progress.option.as_ref() {
1999                if let Some(callee) = opt.callee.as_ref() {
2000                    headers.push(rsipstack::rsip::Header::Other(
2001                        "X-Referred-To".to_string(),
2002                        callee.clone(),
2003                    ));
2004                }
2005                if let Some(caller) = opt.caller.as_ref() {
2006                    headers.push(rsipstack::rsip::Header::Other(
2007                        "X-Referred-From".to_string(),
2008                        caller.clone(),
2009                    ));
2010                }
2011            }
2012        }
2013
2014        headers.push(rsipstack::rsip::Header::Other(
2015            "X-Referred-Id".to_string(),
2016            self.session_id.clone(),
2017        ));
2018
2019        let ssrc = rand::random::<u32>();
2020        let refer_leg = LegShared::new(
2021            ssrc,
2022            true,
2023            CallProgress {
2024                session_id: ref_call_id.clone(),
2025                start_time: Some(Utc::now()),
2026                option: Some(call_option.clone()),
2027                ..Default::default()
2028            },
2029        );
2030        self.set_refer_leg(Some(refer_leg.clone()));
2031
2032        let auto_hangup_requested = refer_option
2033            .as_ref()
2034            .and_then(|o| o.auto_hangup)
2035            .unwrap_or(true);
2036
2037        // auto_hangup rides on the refer leg's TrackEnd (InviteDialogStates).
2038
2039        // Setup ASR resume after refer ends (if not auto_hangup and ASR was paused)
2040        if !auto_hangup_requested && pause_parent_asr && original_asr_option.is_some() {
2041            let asr_option = original_asr_option.unwrap();
2042            self.set_pending_asr_resume((ssrc, asr_option));
2043        }
2044
2045        let timeout_secs = refer_option.as_ref().and_then(|o| o.timeout).unwrap_or(30);
2046        let forward_dtmf = refer_option
2047            .as_ref()
2048            .and_then(|o| o.forward_dtmf)
2049            .unwrap_or(true);
2050
2051        info!(
2052            session_id = self.session_id,
2053            ssrc,
2054            auto_hangup = auto_hangup_requested,
2055            callee,
2056            timeout_secs,
2057            "do_refer"
2058        );
2059
2060        let refer_cancel_token = self.cancel_token.child_token();
2061        self.set_refer_call_token(refer_cancel_token.clone());
2062
2063        // Run the INVITE handshake in the background so the actor keeps
2064        // serving events (MOH looping, auto-hangup...) while it is in flight.
2065        let me = runtime
2066            .me
2067            .clone()
2068            .ok_or_else(|| anyhow::anyhow!("refer is only supported inside serve()"))?;
2069        let actor_tx = runtime.actor_tx.clone();
2070        let event_sender = self.event_sender.clone();
2071        let log_session_id = session_id.clone();
2072        let reject_track_id = track_id.clone();
2073        crate::spawn(async move {
2074            let out = crate::call::tracks::OutgoingLeg {
2075                cancel_token: refer_cancel_token,
2076                leg: refer_leg,
2077                track_id: track_id.clone(),
2078                invite_option,
2079                call_option,
2080                moh,
2081                auto_hangup: auto_hangup_requested,
2082            };
2083            let result = match tokio::time::timeout(
2084                Duration::from_secs(timeout_secs as u64),
2085                me.create_outgoing_sip_track(out),
2086            )
2087            .await
2088            {
2089                Ok(res) => res,
2090                Err(_) => {
2091                    warn!(
2092                        session_id = log_session_id,
2093                        "refer sip track creation timed out after {} seconds", timeout_secs
2094                    );
2095                    event_sender
2096                        .send(SessionEvent::Reject {
2097                            track_id: reject_track_id,
2098                            timestamp: crate::media::get_timestamp(),
2099                            reason: "Timeout when refer".into(),
2100                            code: Some(408),
2101                            refer: Some(true),
2102                        })
2103                        .ok();
2104                    Err(rsipstack::Error::Error(
2105                        "refer sip track creation timed out".to_string(),
2106                    ))
2107                }
2108            };
2109            me.set_moh(None);
2110            actor_tx
2111                .send(ActorMsg::ReferDone {
2112                    track_id,
2113                    forward_dtmf,
2114                    result,
2115                })
2116                .await
2117                .ok();
2118        });
2119
2120        Ok(())
2121    }
2122
2123    async fn do_message(
2124        &self,
2125        body: String,
2126        content_type: Option<String>,
2127        headers: Option<HashMap<String, String>>,
2128        refer: Option<bool>,
2129    ) -> Result<()> {
2130        if !matches!(self.call_type, ActiveCallType::Sip | ActiveCallType::B2bua) {
2131            return Err(anyhow::anyhow!(
2132                "message command is only supported for SIP calls"
2133            ));
2134        }
2135
2136        let dialog_key = if refer == Some(true) {
2137            self.refer_leg_value()
2138                .map(|leg| leg.progress.load_full().session_id.clone())
2139        } else {
2140            Some(self.progress.load_full().session_id.clone())
2141        };
2142
2143        let mut dialog = dialog_key
2144            .as_ref()
2145            .filter(|id| !id.is_empty())
2146            .and_then(|id| self.invitation.dialog_layer.get_dialog_with(id));
2147
2148        // Incoming sessions use short session ids, so resolve them through the
2149        // session-id mapping before falling back to a dialog-id string scan.
2150        if dialog.is_none() {
2151            if let Some(target_id) = dialog_key.as_ref().filter(|id| !id.is_empty()) {
2152                dialog = self
2153                    .invitation
2154                    .find_dialog_id_by_session_id(target_id)
2155                    .and_then(|dialog_id| self.invitation.dialog_layer.get_dialog(&dialog_id));
2156            }
2157        }
2158
2159        if dialog.is_none() {
2160            if let Some(target_id) = dialog_key.as_ref().filter(|id| !id.is_empty()) {
2161                dialog = self
2162                    .invitation
2163                    .dialog_layer
2164                    .all_dialog_ids()
2165                    .into_iter()
2166                    .filter_map(|id| self.invitation.dialog_layer.get_dialog_with(&id))
2167                    .find(|dialog| dialog.id().to_string() == *target_id);
2168            }
2169        }
2170
2171        // Last resort: look up a confirmed client dialog by call id (the
2172        // dialog id for refer legs, the session id otherwise).
2173        if dialog.is_none() {
2174            let call_id = match (refer == Some(true), dialog_key.as_deref()) {
2175                (true, Some(id)) if !id.is_empty() => Some(id),
2176                (false, _) => Some(self.session_id.as_str()),
2177                _ => None,
2178            };
2179            if let Some(call_id) = call_id {
2180                dialog = self
2181                    .invitation
2182                    .dialog_layer
2183                    .get_client_dialog_by_call_id(call_id)
2184                    .into_iter()
2185                    .find(|d| {
2186                        matches!(
2187                            d.state(),
2188                            rsipstack::dialog::dialog::DialogState::Confirmed(_, _)
2189                        )
2190                    })
2191                    .map(rsipstack::dialog::dialog::Dialog::Invite);
2192            }
2193        }
2194
2195        let dialog = dialog.ok_or_else(|| {
2196            anyhow::anyhow!(
2197                "no established SIP dialog found for message command, refer={}",
2198                refer.unwrap_or_default()
2199            )
2200        })?;
2201
2202        let mut sip_headers = vec![rsipstack::rsip::Header::ContentType(
2203            content_type
2204                .clone()
2205                .unwrap_or_else(|| "text/plain;charset=utf-8".to_string())
2206                .into(),
2207        )];
2208        if let Some(headers) = &headers {
2209            sip_headers.extend(crate::sip_util::sip_headers_from_map(headers));
2210        }
2211
2212        info!(
2213            session_id = self.session_id,
2214            dialog_id = %dialog.id(),
2215            content_type = content_type.as_deref().unwrap_or("text/plain;charset=utf-8"),
2216            refer = refer.unwrap_or_default(),
2217            body = %body.chars().take(64).collect::<String>(),
2218            "sending SIP MESSAGE"
2219        );
2220
2221        let response = dialog
2222            .message(Some(sip_headers), Some(body.into_bytes()))
2223            .await?;
2224        match response {
2225            Some(resp)
2226                if resp.status_code.kind() == rsipstack::rsip::StatusCodeKind::Successful =>
2227            {
2228                Ok(())
2229            }
2230            Some(resp) => Err(anyhow::anyhow!(
2231                "SIP MESSAGE rejected with status {}",
2232                resp.status_code
2233            )),
2234            None => Err(anyhow::anyhow!(
2235                "SIP MESSAGE was not sent because dialog is not confirmed"
2236            )),
2237        }
2238    }
2239
2240    fn bridge_track_id(source_session_id: &str, target_session_id: &str) -> TrackId {
2241        format!("bridge:{}:to:{}", source_session_id, target_session_id)
2242    }
2243
2244    async fn do_bridge(&self, target_session_id: String) -> Result<()> {
2245        let target = {
2246            let calls = self.app_state.active_calls.lock().unwrap();
2247            calls.get(&target_session_id).cloned()
2248        };
2249        let target = target.ok_or_else(|| {
2250            anyhow::anyhow!("bridge target session not found: {}", target_session_id)
2251        })?;
2252
2253        if target.session_id == self.session_id {
2254            return Err(anyhow::anyhow!("cannot bridge a call to itself").into());
2255        }
2256
2257        let self_bridge_track_id = Self::bridge_track_id(&self.session_id, &target.session_id);
2258        let target_bridge_track_id = Self::bridge_track_id(&target.session_id, &self.session_id);
2259
2260        self.media_stream
2261            .remove_track(&self_bridge_track_id, false)
2262            .await;
2263        target
2264            .media_stream
2265            .remove_track(&target_bridge_track_id, false)
2266            .await;
2267
2268        let (self_bridge_sender, self_bridge_receiver) = mpsc::channel(25);
2269        let (target_bridge_sender, target_bridge_receiver) = mpsc::channel(25);
2270
2271        let self_paused = self.bridge_paused.clone();
2272        let target_paused = target.bridge_paused.clone();
2273
2274        let self_forwarding_track = ForwardingTrack::new(
2275            self_bridge_track_id.clone(),
2276            self.session_id.clone(),
2277            target_bridge_sender,
2278            self_bridge_receiver,
2279            self.track_config.clone(),
2280            self.cancel_token.child_token(),
2281            rand::random::<u32>(),
2282            self_paused,
2283        );
2284
2285        let target_forwarding_track = ForwardingTrack::new(
2286            target_bridge_track_id.clone(),
2287            target.session_id.clone(),
2288            self_bridge_sender,
2289            target_bridge_receiver,
2290            target.track_config.clone(),
2291            target.cancel_token.child_token(),
2292            rand::random::<u32>(),
2293            target_paused,
2294        );
2295
2296        self.media_stream
2297            .update_track(Box::new(self_forwarding_track), None)
2298            .await;
2299        target
2300            .media_stream
2301            .update_track(Box::new(target_forwarding_track), None)
2302            .await;
2303
2304        info!(
2305            session_id = self.session_id,
2306            target = target_session_id,
2307            self_bridge_track_id,
2308            target_bridge_track_id,
2309            "audio bridge established"
2310        );
2311        Ok(())
2312    }
2313
2314    async fn do_unbridge(&self, target_session_id: String) -> Result<()> {
2315        let target = {
2316            let calls = self.app_state.active_calls.lock().unwrap();
2317            calls.get(&target_session_id).cloned()
2318        };
2319
2320        let self_bridge_track_id = Self::bridge_track_id(&self.session_id, &target_session_id);
2321        self.media_stream
2322            .remove_track(&self_bridge_track_id, false)
2323            .await;
2324
2325        if let Some(target) = target {
2326            let target_bridge_track_id =
2327                Self::bridge_track_id(&target.session_id, &self.session_id);
2328            target
2329                .media_stream
2330                .remove_track(&target_bridge_track_id, false)
2331                .await;
2332            info!(
2333                session_id = self.session_id,
2334                target = target.session_id,
2335                self_bridge_track_id,
2336                target_bridge_track_id,
2337                "audio bridge removed"
2338            );
2339        } else {
2340            info!(
2341                session_id = self.session_id,
2342                target = target_session_id,
2343                self_bridge_track_id,
2344                "audio bridge removed locally; target session not active"
2345            );
2346        }
2347
2348        Ok(())
2349    }
2350
2351    async fn do_mute(&self, track_id: Option<String>) -> Result<()> {
2352        self.media_stream.mute_track(track_id).await;
2353        Ok(())
2354    }
2355
2356    async fn do_unmute(&self, track_id: Option<String>) -> Result<()> {
2357        self.media_stream.unmute_track(track_id).await;
2358        Ok(())
2359    }
2360
2361    pub async fn cleanup(&self) -> Result<()> {
2362        if matches!(self.call_type, ActiveCallType::Sip | ActiveCallType::B2bua) {
2363            self.reject_now(
2364                Some(rsipstack::rsip::StatusCode::Decline),
2365                Some("handler disconnected".to_string()),
2366            )
2367            .await;
2368        }
2369        // A prepared-but-unaccepted SIP answer (non-SIP call types) must not
2370        // leave the carrier leg ringing forever.
2371        if let Some(pending) = self.take_pending_sip_answer() {
2372            if let Ok(pending) = Arc::try_unwrap(pending) {
2373                pending
2374                    .dialog
2375                    .reject(
2376                        Some(rsipstack::rsip::StatusCode::Decline),
2377                        Some("handler disconnected".to_string()),
2378                    )
2379                    .ok();
2380                self.invitation
2381                    .dialog_layer
2382                    .remove_dialog(&pending.dialog.id());
2383            }
2384        }
2385        self.tts_handle.store(None);
2386        self.media_stream.cleanup().await.ok();
2387        Ok(())
2388    }
2389
2390    /// Build the call record from lock-free snapshots; never blocks, so it is
2391    /// safe (and lossless) from synchronous `Drop`.
2392    pub fn get_callrecord(&self) -> Option<CallRecord> {
2393        let progress = self.progress.load_full();
2394        let extras = self.extras.load_full();
2395        let refer_leg = self.refer_leg_value();
2396        Some(build_callrecord(
2397            &progress,
2398            &extras,
2399            refer_leg.as_ref(),
2400            &self.app_state,
2401            self.session_id.clone(),
2402            self.call_type.clone(),
2403        ))
2404    }
2405
2406    async fn dump_to_file(
2407        &self,
2408        dump_file: &mut File,
2409        cmd_receiver: &mut CommandReceiver,
2410        event_receiver: &mut EventReceiver,
2411    ) {
2412        loop {
2413            select! {
2414                _ = self.cancel_token.cancelled() => {
2415                    break;
2416                }
2417                Ok(cmd) = cmd_receiver.recv() => {
2418                    CallRecordEvent::write(CallRecordEventType::Command, cmd, dump_file)
2419                        .await;
2420                }
2421                Ok(event) = event_receiver.recv() => {
2422                    if matches!(event, SessionEvent::Binary{..}) {
2423                        continue;
2424                    }
2425                    CallRecordEvent::write(CallRecordEventType::Event, event, dump_file)
2426                        .await;
2427                }
2428            };
2429        }
2430    }
2431
2432    async fn dump_loop(
2433        &self,
2434        dump_events: bool,
2435        mut dump_cmd_receiver: CommandReceiver,
2436        mut dump_event_receiver: EventReceiver,
2437    ) {
2438        if !dump_events {
2439            return;
2440        }
2441
2442        let file_name = self.app_state.get_dump_events_file(&self.session_id);
2443        let mut dump_file = match File::options()
2444            .create(true)
2445            .append(true)
2446            .open(&file_name)
2447            .await
2448        {
2449            Ok(file) => file,
2450            Err(e) => {
2451                warn!(
2452                    session_id = self.session_id,
2453                    file_name, "failed to open dump events file: {}", e
2454                );
2455                return;
2456            }
2457        };
2458        self.dump_to_file(
2459            &mut dump_file,
2460            &mut dump_cmd_receiver,
2461            &mut dump_event_receiver,
2462        )
2463        .await;
2464
2465        while let Ok(event) = dump_event_receiver.try_recv() {
2466            if matches!(event, SessionEvent::Binary { .. }) {
2467                continue;
2468            }
2469            CallRecordEvent::write(CallRecordEventType::Event, event, &mut dump_file).await;
2470        }
2471    }
2472}
2473
2474/// Cancels the token on drop, guaranteeing shutdown on every exit path of the
2475/// actor loop (normal break, panic, or task abort).
2476struct CancelOnExit<'a>(&'a CancellationToken);
2477
2478impl Drop for CancelOnExit<'_> {
2479    fn drop(&mut self) {
2480        self.0.cancel();
2481    }
2482}
2483
2484impl Drop for ActiveCall {
2485    fn drop(&mut self) {
2486        info!(session_id = self.session_id, "dropping active call");
2487        if let Some(sender) = self.app_state.callrecord_sender.as_ref() {
2488            if let Some(record) = self.get_callrecord() {
2489                if let Err(e) = sender.send(record) {
2490                    warn!(
2491                        session_id = self.session_id,
2492                        "failed to send call record: {}", e
2493                    );
2494                }
2495            }
2496        }
2497    }
2498}