Skip to main content

active_call/call/
active_call.rs

1use super::Command;
2use crate::{
3    CallOption, ReferOption,
4    event::{EventReceiver, EventSender, SessionEvent},
5    media::{
6        TrackId,
7        ambiance::AmbianceProcessor,
8        engine::StreamEngine,
9        negotiate::strip_ipv6_candidates,
10        processor::SubscribeProcessor,
11        recorder::RecorderOption,
12        stream::{MediaStream, MediaStreamBuilder},
13        track::{
14            Track, TrackConfig,
15            file::FileTrack,
16            media_pass::MediaPassTrack,
17            rtc::{RtcTrack, RtcTrackConfig},
18            tts::SynthesisHandle,
19            websocket::{WebsocketBytesReceiver, WebsocketTrack},
20        },
21    },
22    synthesis::{SynthesisCommand, SynthesisOption},
23    transcription::TranscriptionOption,
24};
25use crate::{
26    app::AppState,
27    call::{
28        CommandReceiver, CommandSender,
29        sip::{DialogStateReceiverGuard, Invitation, InviteDialogStates},
30    },
31    callrecord::{CallRecord, CallRecordEvent, CallRecordEventType, CallRecordHangupReason},
32    useragent::{
33        invitation::PendingDialog,
34        public_address::{
35            build_public_contact_uri, contact_needs_public_resolution, find_local_addr_for_uri,
36        },
37    },
38};
39use anyhow::Result;
40use audio_codec::CodecType;
41use chrono::{DateTime, Utc};
42use rsipstack::dialog::{invitation::InviteOption, server_dialog::ServerInviteDialog};
43use serde::{Deserialize, Serialize};
44use std::{collections::HashMap, path::Path, sync::Arc, time::Duration};
45use tokio::{fs::File, select, sync::Mutex, sync::RwLock, time::sleep};
46use tokio_util::sync::CancellationToken;
47use tracing::{debug, info, warn};
48
49#[cfg(test)]
50mod tests {
51    use super::*;
52    use crate::app::AppStateBuilder;
53    use crate::callrecord::CallRecordHangupReason;
54    use crate::config::Config;
55    use crate::media::track::tts::SynthesisHandle;
56    use crate::synthesis::SynthesisCommand;
57    use tokio::sync::mpsc;
58
59    #[tokio::test]
60    async fn test_tts_ssrc_reuse_for_autohangup() -> Result<()> {
61        let mut config = Config::default();
62        config.udp_port = 0; // Use random port
63        config.media_cache_path = "/tmp/mediacache".to_string();
64        let stream_engine = Arc::new(StreamEngine::default());
65        let app_state = AppStateBuilder::new()
66            .with_config(config)
67            .with_stream_engine(stream_engine)
68            .build()
69            .await?;
70
71        let cancel_token = CancellationToken::new();
72        let session_id = "test-session".to_string();
73        let track_config = TrackConfig::default();
74
75        let mut option = crate::CallOption::default();
76        option.tts = Some(crate::synthesis::SynthesisOption::default());
77
78        let active_call = Arc::new(ActiveCall::new(
79            ActiveCallType::Sip,
80            cancel_token.clone(),
81            session_id.clone(),
82            app_state.invitation.clone(),
83            app_state.clone(),
84            track_config,
85            None,
86            false,
87            None,
88            None,
89            None,
90        ));
91
92        {
93            let mut state = active_call.call_state.write().await;
94            state.option = Some(option);
95        }
96
97        let (tx, mut rx) = mpsc::unbounded_channel::<SynthesisCommand>();
98        let initial_ssrc = 12345;
99        let handle = SynthesisHandle::new(tx, Some("play_1".to_string()), initial_ssrc);
100
101        // 1. Set initial TTS handle
102        {
103            let mut state = active_call.call_state.write().await;
104            state.tts_handle = Some(handle);
105            state.current_play_id = Some("play_1".to_string());
106        }
107
108        // 2. Call do_tts with auto_hangup=true and same play_id
109        active_call
110            .do_tts(
111                "hangup now".to_string(),
112                None,
113                Some("play_1".to_string()),
114                Some(true),
115                false,
116                true,
117                None,
118                None,
119                false,
120                None,
121            )
122            .await?;
123
124        // 3. Verify auto_hangup state uses the EXISTING SSRC
125        {
126            let state = active_call.call_state.read().await;
127            assert!(state.auto_hangup.is_some());
128            let (h_ssrc, reason) = state.auto_hangup.clone().unwrap();
129            assert_eq!(
130                h_ssrc, initial_ssrc,
131                "SSRC should be reused from existing handle"
132            );
133            assert_eq!(reason, CallRecordHangupReason::BySystem);
134        }
135
136        // 4. Verify command was sent to existing channel
137        let cmd = rx.try_recv().expect("Should have received tts command");
138        assert_eq!(cmd.text, "hangup now");
139
140        Ok(())
141    }
142
143    #[tokio::test]
144    async fn test_tts_new_ssrc_for_different_play_id() -> Result<()> {
145        let mut config = Config::default();
146        config.udp_port = 0; // Use random port
147        config.media_cache_path = "/tmp/mediacache".to_string();
148        let stream_engine = Arc::new(StreamEngine::default());
149        let app_state = AppStateBuilder::new()
150            .with_config(config)
151            .with_stream_engine(stream_engine)
152            .build()
153            .await?;
154
155        let active_call = Arc::new(ActiveCall::new(
156            ActiveCallType::Sip,
157            CancellationToken::new(),
158            "test-session".to_string(),
159            app_state.invitation.clone(),
160            app_state.clone(),
161            TrackConfig::default(),
162            None,
163            false,
164            None,
165            None,
166            None,
167        ));
168
169        let mut tts_opt = crate::synthesis::SynthesisOption::default();
170        tts_opt.provider = Some(crate::synthesis::SynthesisType::Aliyun);
171        let mut option = crate::CallOption::default();
172        option.tts = Some(tts_opt);
173        {
174            let mut state = active_call.call_state.write().await;
175            state.option = Some(option);
176        }
177
178        let (tx, _rx) = mpsc::unbounded_channel();
179        let initial_ssrc = 111;
180        let handle = SynthesisHandle::new(tx, Some("play_1".to_string()), initial_ssrc);
181
182        {
183            let mut state = active_call.call_state.write().await;
184            state.tts_handle = Some(handle);
185            state.current_play_id = Some("play_1".to_string());
186        }
187
188        // Call do_tts with DIFFERENT play_id
189        active_call
190            .do_tts(
191                "new play".to_string(),
192                None,
193                Some("play_2".to_string()),
194                Some(true),
195                false,
196                true,
197                None,
198                None,
199                false,
200                None,
201            )
202            .await?;
203
204        // Verify auto_hangup uses a NEW SSRC (because it should interrupt and start fresh)
205        {
206            let state = active_call.call_state.read().await;
207            let (h_ssrc, _) = state.auto_hangup.clone().unwrap();
208            assert_ne!(
209                h_ssrc, initial_ssrc,
210                "Should use a new SSRC for different play_id"
211            );
212        }
213
214        Ok(())
215    }
216}
217
218#[derive(Deserialize)]
219#[serde(rename_all = "camelCase")]
220pub struct CallParams {
221    pub id: Option<String>,
222    #[serde(rename = "dump")]
223    pub dump_events: Option<bool>,
224    #[serde(rename = "ping")]
225    pub ping_interval: Option<u32>,
226    pub server_side_track: Option<String>,
227}
228
229#[derive(Debug, Serialize, Deserialize, Clone, Default, PartialEq, Eq)]
230#[serde(rename_all = "camelCase")]
231pub enum ActiveCallType {
232    Webrtc,
233    B2bua,
234    WebSocket,
235    #[default]
236    Sip,
237}
238
239#[derive(Default)]
240pub struct ActiveCallState {
241    pub session_id: String,
242    pub start_time: DateTime<Utc>,
243    pub ring_time: Option<DateTime<Utc>>,
244    pub answer_time: Option<DateTime<Utc>>,
245    pub hangup_reason: Option<CallRecordHangupReason>,
246    pub last_status_code: u16,
247    pub option: Option<CallOption>,
248    pub answer: Option<String>,
249    pub ssrc: u32,
250    pub refer_callstate: Option<ActiveCallStateRef>,
251    pub extras: Option<HashMap<String, serde_json::Value>>,
252    pub is_refer: bool,
253    pub sip_hangup_headers_template: Option<HashMap<String, String>>,
254
255    // Runtime state (migrated from ActiveCall to reduce multiple locks)
256    pub tts_handle: Option<SynthesisHandle>,
257    pub auto_hangup: Option<(u32, CallRecordHangupReason)>,
258    pub wait_input_timeout: Option<u32>,
259    pub moh: Option<String>,
260    pub current_play_id: Option<String>,
261    pub audio_receiver: Option<WebsocketBytesReceiver>,
262    pub ready_to_answer: Option<(String, Option<Box<dyn Track>>, ServerInviteDialog)>,
263    pub pending_asr_resume: Option<(u32, TranscriptionOption)>,
264}
265
266pub type ActiveCallRef = Arc<ActiveCall>;
267pub type ActiveCallStateRef = Arc<RwLock<ActiveCallState>>;
268
269pub struct ActiveCall {
270    pub call_state: ActiveCallStateRef,
271    pub cancel_token: CancellationToken,
272    pub call_type: ActiveCallType,
273    pub session_id: String,
274    pub media_stream: Arc<MediaStream>,
275    pub track_config: TrackConfig,
276    pub event_sender: EventSender,
277    pub app_state: AppState,
278    pub invitation: Invitation,
279    pub cmd_sender: CommandSender,
280    pub dump_events: bool,
281    pub server_side_track_id: TrackId,
282}
283
284pub struct ActiveCallGuard {
285    pub call: ActiveCallRef,
286    pub active_calls: usize,
287}
288
289impl ActiveCallGuard {
290    pub fn new(call: ActiveCallRef) -> Self {
291        let active_calls = {
292            call.app_state
293                .total_calls
294                .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
295            let mut calls = call.app_state.active_calls.lock().unwrap();
296            calls.insert(call.session_id.clone(), call.clone());
297            calls.len()
298        };
299        Self { call, active_calls }
300    }
301}
302
303impl Drop for ActiveCallGuard {
304    fn drop(&mut self) {
305        self.call
306            .app_state
307            .active_calls
308            .lock()
309            .unwrap()
310            .remove(&self.call.session_id);
311    }
312}
313
314pub struct ActiveCallReceiver {
315    pub cmd_receiver: CommandReceiver,
316    pub dump_cmd_receiver: CommandReceiver,
317    pub dump_event_receiver: EventReceiver,
318}
319
320impl ActiveCall {
321    pub fn new(
322        call_type: ActiveCallType,
323        cancel_token: CancellationToken,
324        session_id: String,
325        invitation: Invitation,
326        app_state: AppState,
327        track_config: TrackConfig,
328        audio_receiver: Option<WebsocketBytesReceiver>,
329        dump_events: bool,
330        server_side_track_id: Option<TrackId>,
331        extras: Option<HashMap<String, serde_json::Value>>,
332        sip_hangup_headers_template: Option<HashMap<String, String>>,
333    ) -> Self {
334        let event_sender = crate::event::create_event_sender();
335        let cmd_sender = tokio::sync::broadcast::Sender::<Command>::new(32);
336        let media_stream_builder = MediaStreamBuilder::new(event_sender.clone())
337            .with_id(session_id.clone())
338            .with_cancel_token(cancel_token.child_token());
339        let media_stream = Arc::new(media_stream_builder.build());
340        let start_time = Utc::now();
341        // Inject built-in session variables into extras
342        let call_type_str = match &call_type {
343            ActiveCallType::Sip => "sip",
344            ActiveCallType::WebSocket => "websocket",
345            ActiveCallType::Webrtc => "webrtc",
346            ActiveCallType::B2bua => "b2bua",
347        };
348        let extras = {
349            let mut e = extras.unwrap_or_default();
350            e.entry(crate::playbook::BUILTIN_SESSION_ID.to_string())
351                .or_insert_with(|| serde_json::Value::String(session_id.clone()));
352            e.entry(crate::playbook::BUILTIN_CALL_TYPE.to_string())
353                .or_insert_with(|| serde_json::Value::String(call_type_str.to_string()));
354            e.entry(crate::playbook::BUILTIN_START_TIME.to_string())
355                .or_insert_with(|| serde_json::Value::String(start_time.to_rfc3339()));
356            Some(e)
357        };
358        let call_state = Arc::new(RwLock::new(ActiveCallState {
359            session_id: session_id.clone(),
360            start_time,
361            ssrc: rand::random::<u32>(),
362            extras,
363            audio_receiver,
364            sip_hangup_headers_template,
365            ..Default::default()
366        }));
367        Self {
368            cancel_token,
369            call_type,
370            session_id,
371            call_state,
372            media_stream,
373            track_config,
374            event_sender,
375            app_state,
376            invitation,
377            cmd_sender,
378            dump_events,
379            server_side_track_id: server_side_track_id.unwrap_or("server-side-track".to_string()),
380        }
381    }
382
383    pub async fn enqueue_command(&self, command: Command) -> Result<()> {
384        self.cmd_sender
385            .send(command)
386            .map_err(|e| anyhow::anyhow!("Failed to send command: {}", e))?;
387        Ok(())
388    }
389
390    /// Create a new ActiveCallReceiver for this ActiveCall
391    /// `tokio::sync::broadcast` not cached messages, so need to early create receiver
392    /// before calling `serve()`
393    pub fn new_receiver(&self) -> ActiveCallReceiver {
394        ActiveCallReceiver {
395            cmd_receiver: self.cmd_sender.subscribe(),
396            dump_cmd_receiver: self.cmd_sender.subscribe(),
397            dump_event_receiver: self.event_sender.subscribe(),
398        }
399    }
400
401    pub async fn serve(&self, receiver: ActiveCallReceiver) -> Result<()> {
402        let ActiveCallReceiver {
403            mut cmd_receiver,
404            dump_cmd_receiver,
405            dump_event_receiver,
406        } = receiver;
407
408        let process_command_loop = async move {
409            while let Ok(command) = cmd_receiver.recv().await {
410                match Box::pin(self.dispatch(command)).await {
411                    Ok(_) => (),
412                    Err(e) => {
413                        warn!(session_id = self.session_id, "{}", e);
414                        self.event_sender
415                            .send(SessionEvent::Error {
416                                track_id: self.session_id.clone(),
417                                timestamp: crate::media::get_timestamp(),
418                                sender: "command".to_string(),
419                                error: e.to_string(),
420                                code: None,
421                            })
422                            .ok();
423                    }
424                }
425            }
426        };
427        self.app_state
428            .total_calls
429            .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
430
431        tokio::join!(
432            self.dump_loop(self.dump_events, dump_cmd_receiver, dump_event_receiver),
433            async {
434                select! {
435                    _ = process_command_loop => {
436                        info!(session_id = self.session_id, "command loop done");
437                    }
438                    _ = self.process() => {
439                        info!(session_id = self.session_id, "call serve done");
440                    }
441                    _ = self.cancel_token.cancelled() => {
442                        info!(session_id = self.session_id, "call cancelled - cleaning up resources");
443                    }
444                }
445                self.cancel_token.cancel();
446            }
447        );
448        Ok(())
449    }
450
451    async fn process(&self) -> Result<()> {
452        let mut event_receiver = self.event_sender.subscribe();
453
454        let input_timeout_expire = Arc::new(Mutex::new((0u64, 0u32)));
455        let input_timeout_expire_ref = input_timeout_expire.clone();
456        let event_sender = self.event_sender.clone();
457        let wait_input_timeout_loop = async {
458            loop {
459                let (start_time, expire) = { *input_timeout_expire.lock().await };
460                if expire > 0 && crate::media::get_timestamp() >= start_time + expire as u64 {
461                    info!(session_id = self.session_id, "wait input timeout reached");
462                    *input_timeout_expire.lock().await = (0, 0);
463                    event_sender
464                        .send(SessionEvent::Silence {
465                            track_id: self.server_side_track_id.clone(),
466                            timestamp: crate::media::get_timestamp(),
467                            start_time,
468                            duration: expire as u64,
469                            samples: None,
470                        })
471                        .ok();
472                }
473                sleep(Duration::from_millis(100)).await;
474            }
475        };
476        let server_side_track_id = self.server_side_track_id.clone();
477        let event_hook_loop = async move {
478            while let Ok(event) = event_receiver.recv().await {
479                match event {
480                    SessionEvent::Speaking { .. }
481                    | SessionEvent::Dtmf { .. }
482                    | SessionEvent::AsrDelta { .. }
483                    | SessionEvent::AsrFinal { .. }
484                    | SessionEvent::TrackStart { .. } => {
485                        *input_timeout_expire_ref.lock().await = (0, 0);
486                    }
487                    SessionEvent::TrackEnd {
488                        track_id,
489                        play_id,
490                        ssrc,
491                        ..
492                    } => {
493                        if track_id != server_side_track_id {
494                            continue;
495                        }
496
497                        let (moh_path, auto_hangup, wait_timeout_val) = {
498                            let mut state = self.call_state.write().await;
499                            if play_id != state.current_play_id {
500                                debug!(
501                                    session_id = self.session_id,
502                                    ?play_id,
503                                    current = ?state.current_play_id,
504                                    "ignoring interrupted track end"
505                                );
506                                continue;
507                            }
508                            state.current_play_id = None;
509                            (
510                                state.moh.clone(),
511                                state.auto_hangup.clone(),
512                                state.wait_input_timeout.take(),
513                            )
514                        };
515
516                        if let Some(path) = moh_path {
517                            info!(session_id = self.session_id, "looping moh: {}", path);
518                            let ssrc = rand::random::<u32>();
519                            let file_track = FileTrack::new(self.server_side_track_id.clone())
520                                .with_play_id(Some(path.clone()))
521                                .with_ssrc(ssrc)
522                                .with_path(path.clone())
523                                .with_cancel_token(self.cancel_token.child_token());
524                            self.update_track_wrapper(Box::new(file_track), Some(path))
525                                .await;
526                            continue;
527                        }
528
529                        if let Some((hangup_ssrc, hangup_reason)) = auto_hangup {
530                            if hangup_ssrc == ssrc {
531                                info!(
532                                    session_id = self.session_id,
533                                    ssrc, "auto hangup when track end track_id:{}", track_id
534                                );
535                                self.do_hangup(Some(hangup_reason), None, None).await.ok();
536                            }
537                        }
538
539                        if let Some(timeout) = wait_timeout_val {
540                            let expire = if timeout > 0 {
541                                (crate::media::get_timestamp(), timeout)
542                            } else {
543                                (0, 0)
544                            };
545                            *input_timeout_expire_ref.lock().await = expire;
546                        }
547                    }
548                    SessionEvent::Interrupt { receiver } => {
549                        let track_id =
550                            receiver.unwrap_or_else(|| self.server_side_track_id.clone());
551                        if track_id == self.server_side_track_id {
552                            debug!(
553                                session_id = self.session_id,
554                                "received interrupt event, stopping playback"
555                            );
556                            self.do_interrupt(true).await.ok();
557                        }
558                    }
559                    SessionEvent::Inactivity { track_id, .. } => {
560                        info!(
561                            session_id = self.session_id,
562                            track_id, "inactivity timeout reached, hanging up"
563                        );
564                        self.do_hangup(Some(CallRecordHangupReason::InactivityTimeout), None, None)
565                            .await
566                            .ok();
567                    }
568                    SessionEvent::Hangup { refer, .. } => {
569                        // Check if we need to resume ASR after refer hangup
570                        if refer == Some(true) {
571                            let mut cs = self.call_state.write().await;
572                            if let Some((refer_ssrc, asr_option)) = cs.pending_asr_resume.take() {
573                                // Verify it's the refer call that ended
574                                let is_refer_hangup = cs
575                                    .refer_callstate
576                                    .as_ref()
577                                    .map(|rcs| {
578                                        rcs.try_read()
579                                            .map(|g| g.ssrc == refer_ssrc)
580                                            .unwrap_or(false)
581                                    })
582                                    .unwrap_or(false);
583
584                                if is_refer_hangup {
585                                    drop(cs); // Release lock before async operations
586                                    info!(
587                                        session_id = self.session_id,
588                                        "Refer call ended, resuming parent ASR"
589                                    );
590
591                                    // Resume ASR
592                                    match self
593                                        .app_state
594                                        .stream_engine
595                                        .create_asr_processor(
596                                            self.server_side_track_id.clone(),
597                                            self.cancel_token.child_token(),
598                                            asr_option,
599                                            self.event_sender.clone(),
600                                        )
601                                        .await
602                                    {
603                                        Ok(asr_processor) => {
604                                            if let Err(e) = self
605                                                .media_stream
606                                                .append_processor(
607                                                    &self.server_side_track_id,
608                                                    asr_processor,
609                                                )
610                                                .await
611                                            {
612                                                warn!(
613                                                    session_id = self.session_id,
614                                                    "Failed to resume ASR after refer: {}", e
615                                                );
616                                            }
617                                        }
618                                        Err(e) => {
619                                            warn!(
620                                                session_id = self.session_id,
621                                                "Failed to create ASR processor for resume: {}", e
622                                            );
623                                        }
624                                    }
625                                }
626                            }
627                        }
628                    }
629                    SessionEvent::Error { track_id, .. } => {
630                        if track_id != server_side_track_id {
631                            continue;
632                        }
633
634                        let moh_info = {
635                            let mut state = self.call_state.write().await;
636                            if let Some(path) = state.moh.clone() {
637                                let fallback = "./config/sounds/refer_moh.wav".to_string();
638                                let next_path = if path != fallback
639                                    && std::path::Path::new(&fallback).exists()
640                                {
641                                    info!(
642                                        session_id = self.session_id,
643                                        "moh error, switching to fallback: {}", fallback
644                                    );
645                                    state.moh = Some(fallback.clone());
646                                    fallback
647                                } else {
648                                    info!(
649                                        session_id = self.session_id,
650                                        "looping moh on error: {}", path
651                                    );
652                                    path
653                                };
654                                Some(next_path)
655                            } else {
656                                None
657                            }
658                        };
659
660                        if let Some(next_path) = moh_info {
661                            let ssrc = rand::random::<u32>();
662                            let file_track = FileTrack::new(self.server_side_track_id.clone())
663                                .with_play_id(Some(next_path.clone()))
664                                .with_ssrc(ssrc)
665                                .with_path(next_path.clone())
666                                .with_cancel_token(self.cancel_token.child_token());
667                            self.update_track_wrapper(Box::new(file_track), Some(next_path))
668                                .await;
669                            continue;
670                        }
671                    }
672                    _ => {}
673                }
674            }
675        };
676
677        select! {
678            _ = wait_input_timeout_loop=>{
679                info!(session_id = self.session_id, "wait input timeout loop done");
680            }
681            _ = self.media_stream.serve() => {
682                info!(session_id = self.session_id, "media stream loop done");
683            }
684            _ = event_hook_loop => {
685                info!(session_id = self.session_id, "event loop done");
686            }
687        }
688        Ok(())
689    }
690
691    async fn dispatch(&self, command: Command) -> Result<()> {
692        match command {
693            Command::Invite { option } => self.do_invite(option).await,
694            Command::Accept { option } => self.do_accept(option).await,
695            Command::Reject { reason, code } => {
696                self.do_reject(code.map(|c| (c as u16).into()), Some(reason))
697                    .await
698            }
699            Command::Ringing {
700                ringtone,
701                recorder,
702                early_media,
703            } => self.do_ringing(ringtone, recorder, early_media).await,
704            Command::Tts {
705                text,
706                speaker,
707                play_id,
708                auto_hangup,
709                streaming,
710                end_of_stream,
711                option,
712                wait_input_timeout,
713                base64,
714                cache_key,
715            } => {
716                self.do_tts(
717                    text,
718                    speaker,
719                    play_id,
720                    auto_hangup,
721                    streaming.unwrap_or_default(),
722                    end_of_stream.unwrap_or_default(),
723                    option,
724                    wait_input_timeout,
725                    base64.unwrap_or_default(),
726                    cache_key,
727                )
728                .await
729            }
730            Command::Play {
731                url,
732                play_id,
733                auto_hangup,
734                wait_input_timeout,
735                offset_ms,
736            } => {
737                self.do_play(url, play_id, auto_hangup, wait_input_timeout, offset_ms)
738                    .await
739            }
740            Command::Hangup {
741                reason,
742                initiator,
743                headers,
744            } => {
745                let reason = reason.map(|r| {
746                    r.parse::<CallRecordHangupReason>()
747                        .unwrap_or(CallRecordHangupReason::BySystem)
748                });
749                self.do_hangup(reason, initiator, headers).await
750            }
751            Command::Refer {
752                caller,
753                callee,
754                options,
755            } => self.do_refer(caller, callee, options).await,
756            Command::Mute { track_id } => self.do_mute(track_id).await,
757            Command::Unmute { track_id } => self.do_unmute(track_id).await,
758            Command::Pause {} => self.do_pause().await,
759            Command::Resume {} => self.do_resume().await,
760            Command::Interrupt {
761                graceful: passage,
762                fade_out_ms: _,
763            } => self.do_interrupt(passage.unwrap_or_default()).await,
764            Command::History { speaker, text } => self.do_history(speaker, text).await,
765            Command::Custom { sender, data } => self.do_custom(sender, data),
766        }
767    }
768
769    fn build_record_option(&self, option: &CallOption) -> Option<RecorderOption> {
770        if let Some(recorder_option) = &option.recorder {
771            let mut recorder_file = recorder_option.recorder_file.clone();
772            if recorder_file.contains("{id}") {
773                recorder_file = recorder_file.replace("{id}", &self.session_id);
774            }
775
776            let recorder_file = if recorder_file.is_empty() {
777                self.app_state.get_recorder_file(&self.session_id)
778            } else {
779                let p = Path::new(&recorder_file);
780                p.is_absolute()
781                    .then(|| recorder_file.clone())
782                    .unwrap_or_else(|| self.app_state.get_recorder_file(&recorder_file))
783            };
784            info!(
785                session_id = self.session_id,
786                recorder_file, "created recording file"
787            );
788
789            let track_samplerate = self.track_config.samplerate;
790            let recorder_samplerate = if track_samplerate > 0 {
791                track_samplerate
792            } else {
793                recorder_option.samplerate
794            };
795            let recorder_ptime = if recorder_option.ptime == 0 {
796                200
797            } else {
798                recorder_option.ptime
799            };
800            let requested_format = recorder_option
801                .format
802                .unwrap_or(self.app_state.config.recorder_format());
803            let format = requested_format.effective();
804            if requested_format != format {
805                warn!(
806                    session_id = self.session_id,
807                    requested = requested_format.extension(),
808                    "Recorder format fallback to wav due to unsupported feature"
809                );
810            }
811            let mut recorder_config = RecorderOption {
812                recorder_file,
813                samplerate: recorder_samplerate,
814                ptime: recorder_ptime,
815                format: Some(format),
816            };
817            recorder_config.ensure_path_extension(format);
818            Some(recorder_config)
819        } else {
820            None
821        }
822    }
823
824    async fn invite_or_accept(&self, mut option: CallOption, sender: String) -> Result<CallOption> {
825        // Merge with existing configuration (e.g., from playbook)
826        {
827            let state = self.call_state.read().await;
828            option = state.merge_option(option);
829        }
830
831        option.check_default();
832        if let Some(opt) = self.build_record_option(&option) {
833            self.media_stream.update_recorder_option(opt).await;
834        }
835
836        if let Some(opt) = &option.media_pass {
837            let track_id = self.server_side_track_id.clone();
838            let cancel_token = self.cancel_token.child_token();
839            let ssrc = rand::random::<u32>();
840            let media_pass_track = MediaPassTrack::new(
841                self.session_id.clone(),
842                ssrc,
843                track_id,
844                cancel_token,
845                opt.clone(),
846            );
847            self.update_track_wrapper(Box::new(media_pass_track), None)
848                .await;
849        }
850
851        info!(
852            session_id = self.session_id,
853            call_type = ?self.call_type,
854            sender,
855            ?option,
856            "caller with option"
857        );
858
859        match self.setup_caller_track(&option).await {
860            Ok(_) => return Ok(option),
861            Err(e) => {
862                self.app_state
863                    .total_failed_calls
864                    .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
865                let error_event = crate::event::SessionEvent::Error {
866                    track_id: self.session_id.clone(),
867                    timestamp: crate::media::get_timestamp(),
868                    sender,
869                    error: e.to_string(),
870                    code: None,
871                };
872                self.event_sender.send(error_event).ok();
873                self.do_hangup(Some(CallRecordHangupReason::BySystem), None, None)
874                    .await
875                    .ok();
876                return Err(e);
877            }
878        }
879    }
880
881    async fn do_invite(&self, option: CallOption) -> Result<()> {
882        self.invite_or_accept(option, "invite".to_string())
883            .await
884            .map(|_| ())
885    }
886
887    async fn do_accept(&self, mut option: CallOption) -> Result<()> {
888        let has_pending = self
889            .invitation
890            .find_dialog_id_by_session_id(&self.session_id)
891            .is_some();
892        let ready_to_answer_val = {
893            let state = self.call_state.read().await;
894            state.ready_to_answer.is_none()
895        };
896
897        if ready_to_answer_val {
898            if !has_pending {
899                // emit reject event
900                warn!(session_id = self.session_id, "no pending call to accept");
901                let rejet_event = crate::event::SessionEvent::Reject {
902                    track_id: self.session_id.clone(),
903                    timestamp: crate::media::get_timestamp(),
904                    reason: "no pending call".to_string(),
905                    refer: None,
906                    code: Some(486),
907                };
908                self.event_sender.send(rejet_event).ok();
909                self.do_hangup(Some(CallRecordHangupReason::BySystem), None, None)
910                    .await
911                    .ok();
912                return Err(anyhow::anyhow!("no pending call to accept"));
913            }
914            option = self.invite_or_accept(option, "accept".to_string()).await?;
915        } else {
916            option.check_default();
917            self.call_state.write().await.option = Some(option.clone());
918        }
919        info!(session_id = self.session_id, ?option, "accepting call");
920        let ready = self.call_state.write().await.ready_to_answer.take();
921        if let Some((answer, track, dialog)) = ready {
922            info!(
923                session_id = self.session_id,
924                track_id = track.as_ref().map(|t| t.id()),
925                "ready to answer with track"
926            );
927
928            let headers = vec![rsipstack::rsip::Header::ContentType(
929                "application/sdp".to_string().into(),
930            )];
931
932            match dialog.accept(Some(headers), Some(answer.as_bytes().to_vec())) {
933                Ok(_) => {
934                    {
935                        let mut state = self.call_state.write().await;
936                        state.answer = Some(answer);
937                        state.answer_time = Some(Utc::now());
938                    }
939                    self.finish_caller_stack(&option, track).await?;
940                }
941                Err(e) => {
942                    warn!(session_id = self.session_id, "failed to accept call: {}", e);
943                    return Err(anyhow::anyhow!("failed to accept call"));
944                }
945            }
946        }
947        return Ok(());
948    }
949
950    async fn do_reject(
951        &self,
952        code: Option<rsipstack::rsip::StatusCode>,
953        reason: Option<String>,
954    ) -> Result<()> {
955        match self
956            .invitation
957            .find_dialog_id_by_session_id(&self.session_id)
958        {
959            Some(id) => {
960                info!(
961                    session_id = self.session_id,
962                    ?reason,
963                    ?code,
964                    "rejecting call"
965                );
966                let result = self.invitation.hangup(id, code, reason).await;
967                if result.is_ok() {
968                    self.cancel_token.cancel();
969                }
970                result
971            }
972            None => {
973                let ready = self.call_state.write().await.ready_to_answer.take();
974                if let Some((_, _, dialog)) = ready {
975                    info!(
976                        session_id = self.session_id,
977                        ?reason,
978                        ?code,
979                        "rejecting call from ready_to_answer"
980                    );
981                    let dialog_id = dialog.id();
982                    dialog.reject(code, reason).ok();
983                    self.invitation.dialog_layer.remove_dialog(&dialog_id);
984                    self.cancel_token.cancel();
985                }
986                Ok(())
987            }
988        }
989    }
990
991    async fn do_ringing(
992        &self,
993        ringtone: Option<String>,
994        recorder: Option<RecorderOption>,
995        early_media: Option<bool>,
996    ) -> Result<()> {
997        let ready_to_answer_val = self.call_state.read().await.ready_to_answer.is_none();
998        if ready_to_answer_val {
999            let option = CallOption {
1000                recorder,
1001                ..Default::default()
1002            };
1003            let _ = self.invite_or_accept(option, "ringing".to_string()).await?;
1004        }
1005
1006        let state = self.call_state.read().await;
1007        if let Some((answer, _, dialog)) = state.ready_to_answer.as_ref() {
1008            let (headers, body) = if early_media.unwrap_or_default() || ringtone.is_some() {
1009                let headers = vec![rsipstack::rsip::Header::ContentType(
1010                    "application/sdp".to_string().into(),
1011                )];
1012                (Some(headers), Some(answer.as_bytes().to_vec()))
1013            } else {
1014                (None, None)
1015            };
1016
1017            dialog.ringing(headers, body).ok();
1018            info!(
1019                session_id = self.session_id,
1020                ringtone, early_media, "playing ringtone"
1021            );
1022            if let Some(ringtone_url) = ringtone {
1023                drop(state);
1024                self.do_play(ringtone_url, None, None, None, None).await.ok();
1025            } else {
1026                info!(session_id = self.session_id, "no ringtone to play");
1027            }
1028        }
1029        Ok(())
1030    }
1031
1032    async fn do_tts(
1033        &self,
1034        text: String,
1035        speaker: Option<String>,
1036        play_id: Option<String>,
1037        auto_hangup: Option<bool>,
1038        streaming: bool,
1039        end_of_stream: bool,
1040        option: Option<SynthesisOption>,
1041        wait_input_timeout: Option<u32>,
1042        base64: bool,
1043        cache_key: Option<String>,
1044    ) -> Result<()> {
1045        let tts_option = {
1046            let call_state = self.call_state.read().await;
1047            match call_state.option.clone().unwrap_or_default().tts {
1048                Some(opt) => opt.merge_with(option),
1049                None => {
1050                    if let Some(opt) = option {
1051                        opt
1052                    } else {
1053                        return Err(anyhow::anyhow!("no tts option available"));
1054                    }
1055                }
1056            }
1057        };
1058        let speaker = match speaker {
1059            Some(s) => Some(s),
1060            None => tts_option.speaker.clone(),
1061        };
1062
1063        let mut play_command = SynthesisCommand {
1064            text,
1065            speaker,
1066            play_id: play_id.clone(),
1067            streaming,
1068            end_of_stream: if !streaming { true } else { end_of_stream },
1069            option: tts_option,
1070            base64,
1071            cache_key,
1072        };
1073        info!(
1074            session_id = self.session_id,
1075            provider = ?play_command.option.provider,
1076            text = %play_command.text.chars().take(10).collect::<String>(),
1077            speaker = play_command.speaker.as_deref(),
1078            auto_hangup = auto_hangup.unwrap_or_default(),
1079            play_id = play_command.play_id.as_deref(),
1080            streaming = play_command.streaming,
1081            end_of_stream = play_command.end_of_stream,
1082            wait_input_timeout = wait_input_timeout.unwrap_or_default(),
1083            is_base64 = play_command.base64,
1084            cache_key = play_command.cache_key.as_deref(),
1085            "new synthesis"
1086        );
1087
1088        let ssrc = rand::random::<u32>();
1089        let (should_interrupt, picked_ssrc) = {
1090            let mut state = self.call_state.write().await;
1091
1092            let (target_ssrc, changed) = if let Some(handle) = &state.tts_handle {
1093                if play_id.is_some() && state.current_play_id != play_id {
1094                    (ssrc, true)
1095                } else {
1096                    (handle.ssrc, false)
1097                }
1098            } else {
1099                (ssrc, false)
1100            };
1101
1102            // Defer auto_hangup setting until after potential interrupt.
1103            // auto_hangup will be set below after do_interrupt() to avoid being cleared.
1104            state.wait_input_timeout = wait_input_timeout;
1105
1106            state.current_play_id = play_id.clone();
1107            (changed, target_ssrc)
1108        };
1109
1110        if should_interrupt {
1111            let _ = self.do_interrupt(false).await;
1112        }
1113
1114        // Set auto_hangup AFTER potential interrupt to avoid it being cleared by do_interrupt().
1115        // Only preserve auto_hangup when reusing the same handle (same play_id).
1116        // When starting a new track or interrupting, clear stale auto_hangup.
1117        {
1118            let mut state = self.call_state.write().await;
1119            state.auto_hangup = match auto_hangup {
1120                Some(true) => Some((picked_ssrc, CallRecordHangupReason::BySystem)),
1121                _ => {
1122                    // Only preserve auto_hangup when reusing the same handle (same play_id).
1123                    // When starting a new track (different play_id or no existing handle),
1124                    // clear stale auto_hangup to prevent orphaned hangup.
1125                    if state.tts_handle.is_some() && !should_interrupt {
1126                        state.auto_hangup.clone()
1127                    } else {
1128                        None
1129                    }
1130                }
1131            };
1132        }
1133
1134        let existing_handle = self.call_state.read().await.tts_handle.clone();
1135        if let Some(tts_handle) = existing_handle {
1136            match tts_handle.try_send(play_command) {
1137                Ok(_) => return Ok(()),
1138                Err(e) => {
1139                    play_command = e.0;
1140                }
1141            }
1142        }
1143
1144        let (new_handle, tts_track) = StreamEngine::create_tts_track(
1145            self.app_state.stream_engine.clone(),
1146            self.cancel_token.child_token(),
1147            self.session_id.clone(),
1148            self.server_side_track_id.clone(),
1149            picked_ssrc,
1150            play_id.clone(),
1151            streaming,
1152            &play_command.option,
1153        )
1154        .await?;
1155
1156        new_handle.try_send(play_command)?;
1157        self.call_state.write().await.tts_handle = Some(new_handle);
1158        self.update_track_wrapper(tts_track, play_id).await;
1159        Ok(())
1160    }
1161
1162    async fn do_play(
1163        &self,
1164        url: String,
1165        play_id: Option<String>,
1166        auto_hangup: Option<bool>,
1167        wait_input_timeout: Option<u32>,
1168        offset_ms: Option<u32>,
1169    ) -> Result<()> {
1170        let ssrc = rand::random::<u32>();
1171        info!(
1172            session_id = self.session_id,
1173            ssrc, url, play_id, auto_hangup, "play file track"
1174        );
1175
1176        let play_id = play_id.or(Some(url.clone()));
1177
1178        let mut file_track = FileTrack::new(self.server_side_track_id.clone())
1179            .with_play_id(play_id.clone())
1180            .with_ssrc(ssrc)
1181            .with_path(url)
1182            .with_cancel_token(self.cancel_token.child_token());
1183
1184        if let Some(offset) = offset_ms {
1185            file_track = file_track.with_offset_ms(offset);
1186        }
1187
1188        {
1189            let mut state = self.call_state.write().await;
1190            state.tts_handle = None;
1191            state.auto_hangup = match auto_hangup {
1192                Some(true) => Some((ssrc, CallRecordHangupReason::BySystem)),
1193                _ => None,
1194            };
1195            state.wait_input_timeout = wait_input_timeout;
1196        }
1197
1198        self.update_track_wrapper(Box::new(file_track), play_id)
1199            .await;
1200        Ok(())
1201    }
1202
1203    async fn do_history(&self, speaker: String, text: String) -> Result<()> {
1204        self.event_sender
1205            .send(SessionEvent::AddHistory {
1206                sender: Some(self.session_id.clone()),
1207                timestamp: crate::media::get_timestamp(),
1208                speaker,
1209                text,
1210            })
1211            .map(|_| ())
1212            .map_err(Into::into)
1213    }
1214
1215    fn do_custom(&self, sender: Option<String>, data: serde_json::Value) -> Result<()> {
1216        self.event_sender
1217            .send(SessionEvent::Custom {
1218                timestamp: crate::media::get_timestamp(),
1219                sender,
1220                data,
1221            })
1222            .map(|_| ())
1223            .map_err(Into::into)
1224    }
1225
1226    async fn do_interrupt(&self, graceful: bool) -> Result<()> {
1227        {
1228            let mut state = self.call_state.write().await;
1229            state.tts_handle = None;
1230            state.moh = None;
1231            state.auto_hangup = None;
1232        }
1233        self.media_stream
1234            .remove_track(&self.server_side_track_id, graceful)
1235            .await;
1236        Ok(())
1237    }
1238    async fn do_pause(&self) -> Result<()> {
1239        Ok(())
1240    }
1241    async fn do_resume(&self) -> Result<()> {
1242        Ok(())
1243    }
1244    async fn do_hangup(
1245        &self,
1246        reason: Option<CallRecordHangupReason>,
1247        initiator: Option<String>,
1248        headers: Option<HashMap<String, String>>,
1249    ) -> Result<()> {
1250        info!(
1251            session_id = self.session_id,
1252            ?reason,
1253            ?initiator,
1254            ?headers,
1255            "do_hangup"
1256        );
1257
1258        // Store headers in extras if provided
1259        if let Some(headers) = headers {
1260            let h_val = serde_json::to_value(&headers).unwrap_or_default();
1261
1262            {
1263                let mut state = self.call_state.write().await;
1264                let mut extras = state.extras.take().unwrap_or_default();
1265                extras.insert("_hangup_headers".to_string(), h_val.clone());
1266                state.extras = Some(extras);
1267            }
1268        } else {
1269        }
1270
1271        let hangup_reason = match initiator.as_deref() {
1272            Some("caller") => CallRecordHangupReason::ByCaller,
1273            Some("callee") => CallRecordHangupReason::ByCallee,
1274            Some("system") => CallRecordHangupReason::Autohangup,
1275            _ => reason.unwrap_or(CallRecordHangupReason::BySystem),
1276        };
1277
1278        self.media_stream
1279            .stop(Some(hangup_reason.to_string()), initiator);
1280
1281        self.call_state
1282            .write()
1283            .await
1284            .set_hangup_reason(hangup_reason);
1285        Ok(())
1286    }
1287
1288    async fn do_refer(
1289        &self,
1290        caller: String,
1291        callee: String,
1292        refer_option: Option<ReferOption>,
1293    ) -> Result<()> {
1294        self.do_interrupt(false).await.ok();
1295
1296        // Check if we should pause parent ASR
1297        let pause_parent_asr = refer_option
1298            .as_ref()
1299            .and_then(|o| o.pause_parent_asr)
1300            .unwrap_or(false);
1301
1302        // Save original ASR option for later resume
1303        let original_asr_option = if pause_parent_asr {
1304            let cs = self.call_state.read().await;
1305            cs.option.as_ref().and_then(|o| o.asr.clone())
1306        } else {
1307            None
1308        };
1309
1310        // Pause parent ASR if requested
1311        if pause_parent_asr {
1312            info!(
1313                session_id = self.session_id,
1314                "Pausing parent call ASR during refer"
1315            );
1316            self.media_stream
1317                .remove_processor::<crate::media::asr_processor::AsrProcessor>(
1318                    &self.server_side_track_id,
1319                )
1320                .await
1321                .ok();
1322        }
1323
1324        let mut moh = refer_option.as_ref().and_then(|o| o.moh.clone());
1325        if let Some(ref path) = moh {
1326            if !path.starts_with("http") && !std::path::Path::new(path).exists() {
1327                let fallback = "./config/sounds/refer_moh.wav";
1328                if std::path::Path::new(fallback).exists() {
1329                    info!(
1330                        session_id = self.session_id,
1331                        "moh {} not found, using fallback {}", path, fallback
1332                    );
1333                    moh = Some(fallback.to_string());
1334                }
1335            }
1336        }
1337        let ref_call_id = refer_option
1338            .as_ref()
1339            .and_then(|o| o.call_id.clone())
1340            .unwrap_or_else(|| format!("ref-{}-{}", rand::random::<u32>(), self.session_id));
1341
1342        let session_id = self.session_id.clone();
1343        let track_id = self.server_side_track_id.clone();
1344
1345        let (recorder, parent_caller) = {
1346            let cs = self.call_state.read().await;
1347            let option = cs.option.as_ref();
1348            (
1349                option.map(|o| o.recorder.clone()).unwrap_or_default(),
1350                option.and_then(|o| o.caller.clone()),
1351            )
1352        };
1353        let caller = if caller.trim().is_empty() {
1354            parent_caller.unwrap_or_default()
1355        } else {
1356            caller
1357        };
1358
1359        let call_option = CallOption {
1360            caller: Some(caller),
1361            callee: Some(callee.clone()),
1362            sip: refer_option.as_ref().and_then(|o| o.sip.clone()),
1363            asr: refer_option.as_ref().and_then(|o| o.asr.clone()),
1364            denoise: refer_option.as_ref().and_then(|o| o.denoise.clone()),
1365            recorder,
1366            ..Default::default()
1367        };
1368
1369        let mut invite_option = call_option.build_invite_option()?;
1370        invite_option.call_id = Some(ref_call_id);
1371
1372        let headers = invite_option.headers.get_or_insert_with(|| Vec::new());
1373
1374        {
1375            let cs = self.call_state.read().await;
1376            if let Some(opt) = cs.option.as_ref() {
1377                if let Some(callee) = opt.callee.as_ref() {
1378                    headers.push(rsipstack::rsip::Header::Other(
1379                        "X-Referred-To".to_string(),
1380                        callee.clone(),
1381                    ));
1382                }
1383                if let Some(caller) = opt.caller.as_ref() {
1384                    headers.push(rsipstack::rsip::Header::Other(
1385                        "X-Referred-From".to_string(),
1386                        caller.clone(),
1387                    ));
1388                }
1389            }
1390        }
1391
1392        headers.push(rsipstack::rsip::Header::Other(
1393            "X-Referred-Id".to_string(),
1394            self.session_id.clone(),
1395        ));
1396
1397        let ssrc = rand::random::<u32>();
1398        let refer_call_state = Arc::new(RwLock::new(ActiveCallState {
1399            start_time: Utc::now(),
1400            ssrc,
1401            option: Some(call_option.clone()),
1402            is_refer: true,
1403            ..Default::default()
1404        }));
1405
1406        {
1407            let mut cs = self.call_state.write().await;
1408            cs.refer_callstate.replace(refer_call_state.clone());
1409        }
1410
1411        let auto_hangup_requested = refer_option
1412            .as_ref()
1413            .and_then(|o| o.auto_hangup)
1414            .unwrap_or(true);
1415
1416        if auto_hangup_requested {
1417            self.call_state.write().await.auto_hangup =
1418                Some((ssrc, CallRecordHangupReason::ByRefer));
1419        } else {
1420            self.call_state.write().await.auto_hangup = None;
1421        }
1422
1423        // Setup ASR resume after refer ends (if not auto_hangup and ASR was paused)
1424        if !auto_hangup_requested && pause_parent_asr && original_asr_option.is_some() {
1425            let asr_option = original_asr_option.unwrap();
1426            self.call_state.write().await.pending_asr_resume = Some((ssrc, asr_option));
1427        }
1428
1429        let timeout_secs = refer_option.as_ref().and_then(|o| o.timeout).unwrap_or(30);
1430
1431        info!(
1432            session_id = self.session_id,
1433            ssrc,
1434            auto_hangup = auto_hangup_requested,
1435            callee,
1436            timeout_secs,
1437            "do_refer"
1438        );
1439
1440        let r = tokio::time::timeout(
1441            Duration::from_secs(timeout_secs as u64),
1442            self.create_outgoing_sip_track(
1443                self.cancel_token.child_token(),
1444                refer_call_state.clone(),
1445                &track_id,
1446                invite_option,
1447                &call_option,
1448                moh,
1449                auto_hangup_requested,
1450            ),
1451        )
1452        .await;
1453
1454        {
1455            self.call_state.write().await.moh = None;
1456        }
1457
1458        let result = match r {
1459            Ok(res) => res,
1460            Err(_) => {
1461                warn!(
1462                    session_id = session_id,
1463                    "refer sip track creation timed out after {} seconds", timeout_secs
1464                );
1465                self.event_sender
1466                    .send(SessionEvent::Reject {
1467                        track_id,
1468                        timestamp: crate::media::get_timestamp(),
1469                        reason: "Timeout when refer".into(),
1470                        code: Some(408),
1471                        refer: Some(true),
1472                    })
1473                    .ok();
1474                return Err(anyhow::anyhow!("refer sip track creation timed out").into());
1475            }
1476        };
1477
1478        match result {
1479            Ok(answer) => {
1480                self.event_sender
1481                    .send(SessionEvent::Answer {
1482                        timestamp: crate::media::get_timestamp(),
1483                        track_id,
1484                        sdp: answer,
1485                        refer: Some(true),
1486                    })
1487                    .ok();
1488            }
1489            Err(e) => {
1490                warn!(
1491                    session_id = session_id,
1492                    "failed to create refer sip track: {}", e
1493                );
1494                match &e {
1495                    rsipstack::Error::DialogError(reason, _, code) => {
1496                        self.event_sender
1497                            .send(SessionEvent::Reject {
1498                                track_id,
1499                                timestamp: crate::media::get_timestamp(),
1500                                reason: reason.clone(),
1501                                code: Some(code.code() as u32),
1502                                refer: Some(true),
1503                            })
1504                            .ok();
1505                    }
1506                    _ => {}
1507                }
1508                return Err(e.into());
1509            }
1510        }
1511        Ok(())
1512    }
1513
1514    async fn do_mute(&self, track_id: Option<String>) -> Result<()> {
1515        self.media_stream.mute_track(track_id).await;
1516        Ok(())
1517    }
1518
1519    async fn do_unmute(&self, track_id: Option<String>) -> Result<()> {
1520        self.media_stream.unmute_track(track_id).await;
1521        Ok(())
1522    }
1523
1524    pub async fn cleanup(&self) -> Result<()> {
1525        if matches!(self.call_type, ActiveCallType::Sip | ActiveCallType::B2bua) {
1526            self.do_reject(
1527                Some(rsipstack::rsip::StatusCode::Decline),
1528                Some("handler disconnected".to_string()),
1529            )
1530            .await
1531            .ok();
1532        }
1533        self.call_state.write().await.tts_handle = None;
1534        self.media_stream.cleanup().await.ok();
1535        Ok(())
1536    }
1537
1538    pub fn get_callrecord(&self) -> Option<CallRecord> {
1539        self.call_state.try_read().ok().map(|call_state| {
1540            call_state.build_callrecord(
1541                self.app_state.clone(),
1542                self.session_id.clone(),
1543                self.call_type.clone(),
1544            )
1545        })
1546    }
1547
1548    async fn dump_to_file(
1549        &self,
1550        dump_file: &mut File,
1551        cmd_receiver: &mut CommandReceiver,
1552        event_receiver: &mut EventReceiver,
1553    ) {
1554        loop {
1555            select! {
1556                _ = self.cancel_token.cancelled() => {
1557                    break;
1558                }
1559                Ok(cmd) = cmd_receiver.recv() => {
1560                    CallRecordEvent::write(CallRecordEventType::Command, cmd, dump_file)
1561                        .await;
1562                }
1563                Ok(event) = event_receiver.recv() => {
1564                    if matches!(event, SessionEvent::Binary{..}) {
1565                        continue;
1566                    }
1567                    CallRecordEvent::write(CallRecordEventType::Event, event, dump_file)
1568                        .await;
1569                }
1570            };
1571        }
1572    }
1573
1574    async fn dump_loop(
1575        &self,
1576        dump_events: bool,
1577        mut dump_cmd_receiver: CommandReceiver,
1578        mut dump_event_receiver: EventReceiver,
1579    ) {
1580        if !dump_events {
1581            return;
1582        }
1583
1584        let file_name = self.app_state.get_dump_events_file(&self.session_id);
1585        let mut dump_file = match File::options()
1586            .create(true)
1587            .append(true)
1588            .open(&file_name)
1589            .await
1590        {
1591            Ok(file) => file,
1592            Err(e) => {
1593                warn!(
1594                    session_id = self.session_id,
1595                    file_name, "failed to open dump events file: {}", e
1596                );
1597                return;
1598            }
1599        };
1600        self.dump_to_file(
1601            &mut dump_file,
1602            &mut dump_cmd_receiver,
1603            &mut dump_event_receiver,
1604        )
1605        .await;
1606
1607        while let Ok(event) = dump_event_receiver.try_recv() {
1608            if matches!(event, SessionEvent::Binary { .. }) {
1609                continue;
1610            }
1611            CallRecordEvent::write(CallRecordEventType::Event, event, &mut dump_file).await;
1612        }
1613    }
1614
1615    pub async fn create_rtp_track(
1616        &self,
1617        track_id: TrackId,
1618        ssrc: u32,
1619        enable_srtp: Option<bool>,
1620    ) -> Result<RtcTrack> {
1621        let mut rtc_config = RtcTrackConfig::default();
1622        // Per-call flag takes precedence over global config.
1623        let use_srtp = enable_srtp
1624            .or(self.app_state.config.enable_srtp)
1625            .unwrap_or(false);
1626        rtc_config.mode = if use_srtp {
1627            rustrtc::TransportMode::Srtp
1628        } else {
1629            rustrtc::TransportMode::Rtp
1630        };
1631
1632        if let Some(codecs) = &self.app_state.config.codecs {
1633            let mut codec_types = Vec::new();
1634            for c in codecs {
1635                match c.to_lowercase().as_str() {
1636                    "pcmu" => codec_types.push(CodecType::PCMU),
1637                    "pcma" => codec_types.push(CodecType::PCMA),
1638                    "g722" => codec_types.push(CodecType::G722),
1639                    "g729" => codec_types.push(CodecType::G729),
1640                    #[cfg(feature = "opus")]
1641                    "opus" => codec_types.push(CodecType::Opus),
1642                    "dtmf" | "2833" | "telephone_event" => {
1643                        codec_types.push(CodecType::TelephoneEvent)
1644                    }
1645                    _ => {}
1646                }
1647            }
1648            if !codec_types.is_empty() {
1649                rtc_config.preferred_codec = Some(codec_types[0].clone());
1650                rtc_config.codecs = codec_types;
1651            }
1652        }
1653
1654        if rtc_config.preferred_codec.is_none() {
1655            rtc_config.preferred_codec = Some(self.track_config.codec.clone());
1656        }
1657
1658        rtc_config.rtp_port_range = self
1659            .app_state
1660            .config
1661            .rtp_start_port
1662            .zip(self.app_state.config.rtp_end_port);
1663
1664        if let Some(ref external_ip) = self.app_state.config.external_ip {
1665            rtc_config.external_ip = Some(external_ip.clone());
1666        }
1667        if let Some(ref bind_ip) = self.app_state.config.rtp_bind_ip {
1668            rtc_config.bind_ip = Some(bind_ip.clone());
1669        }
1670
1671        rtc_config.enable_latching = self.app_state.config.enable_rtp_latching;
1672        rtc_config.enable_ice_lite = self.app_state.config.enable_ice_lite;
1673
1674        let mut track = RtcTrack::new(
1675            self.cancel_token.child_token(),
1676            track_id,
1677            self.track_config.clone(),
1678            rtc_config,
1679        )
1680        .with_ssrc(ssrc);
1681
1682        track.create().await?;
1683
1684        Ok(track)
1685    }
1686
1687    async fn setup_caller_track(&self, option: &CallOption) -> Result<()> {
1688        let hangup_headers = option
1689            .sip
1690            .as_ref()
1691            .and_then(|s| s.hangup_headers.as_ref())
1692            .map(|headers_map| {
1693                headers_map
1694                    .iter()
1695                    .map(|(k, v)| rsipstack::rsip::Header::Other(k.clone(), v.clone()))
1696                    .collect::<Vec<rsipstack::rsip::Header>>()
1697            });
1698        self.call_state.write().await.option = Some(option.clone());
1699        info!(
1700            session_id = self.session_id,
1701            call_type = ?self.call_type,
1702            "setup caller track"
1703        );
1704
1705        let track = match self.call_type {
1706            ActiveCallType::Webrtc => Some(self.create_webrtc_track().await?),
1707            ActiveCallType::WebSocket => {
1708                let audio_receiver = self.call_state.write().await.audio_receiver.take();
1709                if let Some(receiver) = audio_receiver {
1710                    Some(self.create_websocket_track(receiver).await?)
1711                } else {
1712                    None
1713                }
1714            }
1715            ActiveCallType::Sip => {
1716                if let Some(dialog_id) = self
1717                    .invitation
1718                    .find_dialog_id_by_session_id(&self.session_id)
1719                {
1720                    if let Some(pending_dialog) = self.invitation.get_pending_call(&dialog_id) {
1721                        return self
1722                            .prepare_incoming_sip_track(
1723                                self.cancel_token.clone(),
1724                                self.call_state.clone(),
1725                                &self.session_id,
1726                                pending_dialog,
1727                                hangup_headers,
1728                            )
1729                            .await;
1730                    }
1731                }
1732
1733                // Auto-inject credentials from registered users if not already provided
1734                let mut option = option.clone();
1735                if option.sip.is_none()
1736                    || option
1737                        .sip
1738                        .as_ref()
1739                        .and_then(|s| s.username.as_ref())
1740                        .is_none()
1741                {
1742                    if let Some(callee) = &option.callee {
1743                        if let Some(cred) = self.app_state.find_credentials_for_callee(callee) {
1744                            if option.sip.is_none() {
1745                                option.sip = Some(crate::SipOption {
1746                                    username: Some(cred.username.clone()),
1747                                    password: Some(cred.password.clone()),
1748                                    realm: cred.realm.clone(),
1749                                    ..Default::default()
1750                                });
1751                            }
1752                        }
1753                    }
1754                }
1755
1756                let mut invite_option = option.build_invite_option()?;
1757                invite_option.call_id = Some(self.session_id.clone());
1758
1759                match self
1760                    .create_outgoing_sip_track(
1761                        self.cancel_token.clone(),
1762                        self.call_state.clone(),
1763                        &self.session_id,
1764                        invite_option,
1765                        &option,
1766                        None,
1767                        false,
1768                    )
1769                    .await
1770                {
1771                    Ok(answer) => {
1772                        self.event_sender
1773                            .send(SessionEvent::Answer {
1774                                timestamp: crate::media::get_timestamp(),
1775                                track_id: self.session_id.clone(),
1776                                sdp: answer,
1777                                refer: Some(false),
1778                            })
1779                            .ok();
1780                        return Ok(());
1781                    }
1782                    Err(e) => {
1783                        warn!(
1784                            session_id = self.session_id,
1785                            "failed to create sip track: {}", e
1786                        );
1787                        match &e {
1788                            rsipstack::Error::DialogError(reason, _, code) => {
1789                                self.event_sender
1790                                    .send(SessionEvent::Reject {
1791                                        track_id: self.session_id.clone(),
1792                                        timestamp: crate::media::get_timestamp(),
1793                                        reason: reason.clone(),
1794                                        code: Some(code.code() as u32),
1795                                        refer: Some(false),
1796                                    })
1797                                    .ok();
1798                            }
1799                            _ => {}
1800                        }
1801                        return Err(e.into());
1802                    }
1803                }
1804            }
1805            ActiveCallType::B2bua => {
1806                if let Some(dialog_id) = self
1807                    .invitation
1808                    .find_dialog_id_by_session_id(&self.session_id)
1809                {
1810                    if let Some(pending_dialog) = self.invitation.get_pending_call(&dialog_id) {
1811                        return self
1812                            .prepare_incoming_sip_track(
1813                                self.cancel_token.clone(),
1814                                self.call_state.clone(),
1815                                &self.session_id,
1816                                pending_dialog,
1817                                hangup_headers,
1818                            )
1819                            .await;
1820                    }
1821                }
1822
1823                warn!(
1824                    session_id = self.session_id,
1825                    "no pending dialog found for B2BUA call"
1826                );
1827                return Err(anyhow::anyhow!(
1828                    "no pending dialog found for session_id: {}",
1829                    self.session_id
1830                ));
1831            }
1832        };
1833        match track {
1834            Some(track) => {
1835                self.finish_caller_stack(&option, Some(track)).await?;
1836            }
1837            None => {
1838                warn!(session_id = self.session_id, "no track created for caller");
1839                return Err(anyhow::anyhow!("no track created for caller"));
1840            }
1841        }
1842        Ok(())
1843    }
1844
1845    async fn finish_caller_stack(
1846        &self,
1847        option: &CallOption,
1848        track: Option<Box<dyn Track>>,
1849    ) -> Result<()> {
1850        if let Some(track) = track {
1851            self.setup_track_with_stream(&option, track).await?;
1852        }
1853
1854        {
1855            let call_state = self.call_state.read().await;
1856            if let Some(ref answer) = call_state.answer {
1857                info!(
1858                    session_id = self.session_id,
1859                    "sending answer event: {}", answer,
1860                );
1861                self.event_sender
1862                    .send(SessionEvent::Answer {
1863                        timestamp: crate::media::get_timestamp(),
1864                        track_id: self.session_id.clone(),
1865                        sdp: answer.clone(),
1866                        refer: Some(false),
1867                    })
1868                    .ok();
1869            } else {
1870                warn!(
1871                    session_id = self.session_id,
1872                    "no answer in state to send event"
1873                );
1874            }
1875        }
1876        Ok(())
1877    }
1878
1879    pub async fn setup_track_with_stream(
1880        &self,
1881        option: &CallOption,
1882        mut track: Box<dyn Track>,
1883    ) -> Result<()> {
1884        let processors = match StreamEngine::create_processors(
1885            self.app_state.stream_engine.clone(),
1886            track.as_ref(),
1887            self.cancel_token.child_token(),
1888            self.event_sender.clone(),
1889            self.media_stream.packet_sender.clone(),
1890            option,
1891        )
1892        .await
1893        {
1894            Ok(processors) => processors,
1895            Err(e) => {
1896                warn!(
1897                    session_id = self.session_id,
1898                    "failed to prepare stream processors: {}", e
1899                );
1900                vec![]
1901            }
1902        };
1903
1904        // Add all processors from the hook
1905        for processor in processors {
1906            track.append_processor(processor);
1907        }
1908
1909        self.update_track_wrapper(track, None).await;
1910        Ok(())
1911    }
1912
1913    pub async fn update_track_wrapper(&self, mut track: Box<dyn Track>, play_id: Option<String>) {
1914        let (ambiance_opt, subscribe) = {
1915            let state = self.call_state.read().await;
1916            let mut opt = state
1917                .option
1918                .as_ref()
1919                .and_then(|o| o.ambiance.clone())
1920                .unwrap_or_default();
1921
1922            if let Some(global) = &self.app_state.config.ambiance {
1923                opt.merge(global);
1924            }
1925
1926            let subscribe = state
1927                .option
1928                .as_ref()
1929                .and_then(|o| o.subscribe)
1930                .unwrap_or_default();
1931
1932            (opt, subscribe)
1933        };
1934        if track.id() == &self.server_side_track_id && ambiance_opt.path.is_some() {
1935            match AmbianceProcessor::new(ambiance_opt).await {
1936                Ok(ambiance) => {
1937                    info!(session_id = self.session_id, "loaded ambiance processor");
1938                    track.append_processor(Box::new(ambiance));
1939                }
1940                Err(e) => {
1941                    tracing::error!("failed to load ambiance wav {}", e);
1942                }
1943            }
1944        }
1945
1946        if subscribe && self.call_type != ActiveCallType::WebSocket {
1947            let (track_index, sub_track_id) = if track.id() == &self.server_side_track_id {
1948                (0, self.server_side_track_id.clone())
1949            } else {
1950                (1, self.session_id.clone())
1951            };
1952            let sub_processor =
1953                SubscribeProcessor::new(self.event_sender.clone(), sub_track_id, track_index);
1954            track.append_processor(Box::new(sub_processor));
1955        }
1956
1957        self.call_state.write().await.current_play_id = play_id.clone();
1958        self.media_stream.update_track(track, play_id).await;
1959    }
1960
1961    pub async fn create_websocket_track(
1962        &self,
1963        audio_receiver: WebsocketBytesReceiver,
1964    ) -> Result<Box<dyn Track>> {
1965        let (ssrc, codec) = {
1966            let call_state = self.call_state.read().await;
1967            (
1968                call_state.ssrc,
1969                call_state
1970                    .option
1971                    .as_ref()
1972                    .map(|o| o.codec.clone())
1973                    .unwrap_or_default(),
1974            )
1975        };
1976
1977        let ws_track = WebsocketTrack::new(
1978            self.cancel_token.child_token(),
1979            self.session_id.clone(),
1980            self.track_config.clone(),
1981            self.event_sender.clone(),
1982            audio_receiver,
1983            codec,
1984            ssrc,
1985        );
1986
1987        {
1988            let mut call_state = self.call_state.write().await;
1989            call_state.answer_time = Some(Utc::now());
1990            call_state.answer = Some("".to_string());
1991            call_state.last_status_code = 200;
1992        }
1993
1994        Ok(Box::new(ws_track))
1995    }
1996
1997    pub(super) async fn create_webrtc_track(&self) -> Result<Box<dyn Track>> {
1998        let (ssrc, option) = {
1999            let call_state = self.call_state.read().await;
2000            (
2001                call_state.ssrc,
2002                call_state.option.clone().unwrap_or_default(),
2003            )
2004        };
2005
2006        let mut rtc_config = RtcTrackConfig::default();
2007        rtc_config.mode = rustrtc::TransportMode::WebRtc; // WebRTC
2008        rtc_config.ice_servers = self.app_state.config.ice_servers.clone();
2009
2010        if let Some(codecs) = &self.app_state.config.codecs {
2011            let mut codec_types = Vec::new();
2012            for c in codecs {
2013                match c.to_lowercase().as_str() {
2014                    "pcmu" => codec_types.push(CodecType::PCMU),
2015                    "pcma" => codec_types.push(CodecType::PCMA),
2016                    "g722" => codec_types.push(CodecType::G722),
2017                    "g729" => codec_types.push(CodecType::G729),
2018                    #[cfg(feature = "opus")]
2019                    "opus" => codec_types.push(CodecType::Opus),
2020                    "dtmf" | "2833" | "telephone_event" => {
2021                        codec_types.push(CodecType::TelephoneEvent)
2022                    }
2023                    _ => {}
2024                }
2025            }
2026            if !codec_types.is_empty() {
2027                rtc_config.preferred_codec = Some(codec_types[0].clone());
2028                rtc_config.codecs = codec_types;
2029            }
2030        }
2031
2032        if let Some(ref external_ip) = self.app_state.config.external_ip {
2033            rtc_config.external_ip = Some(external_ip.clone());
2034        }
2035        if let Some(ref bind_ip) = self.app_state.config.rtp_bind_ip {
2036            rtc_config.bind_ip = Some(bind_ip.clone());
2037        }
2038
2039        let mut webrtc_track = RtcTrack::new(
2040            self.cancel_token.child_token(),
2041            self.session_id.clone(),
2042            self.track_config.clone(),
2043            rtc_config,
2044        )
2045        .with_ssrc(ssrc);
2046
2047        let timeout = option.handshake_timeout.map(|t| Duration::from_secs(t));
2048        let offer = match option.enable_ipv6 {
2049            Some(false) | None => {
2050                strip_ipv6_candidates(option.offer.as_ref().unwrap_or(&"".to_string()))
2051            }
2052            _ => option.offer.clone().unwrap_or("".to_string()),
2053        };
2054        let answer: Option<String>;
2055        match webrtc_track.handshake(offer, timeout).await {
2056            Ok(answer_sdp) => {
2057                answer = match option.enable_ipv6 {
2058                    Some(false) | None => Some(strip_ipv6_candidates(&answer_sdp)),
2059                    Some(true) => Some(answer_sdp.to_string()),
2060                };
2061            }
2062            Err(e) => {
2063                warn!(session_id = self.session_id, "failed to setup track: {}", e);
2064                return Err(anyhow::anyhow!("Failed to setup track: {}", e));
2065            }
2066        }
2067
2068        {
2069            let mut call_state = self.call_state.write().await;
2070            call_state.answer_time = Some(Utc::now());
2071            call_state.answer = answer;
2072            call_state.last_status_code = 200;
2073        }
2074        Ok(Box::new(webrtc_track))
2075    }
2076
2077    async fn create_outgoing_sip_track(
2078        &self,
2079        cancel_token: CancellationToken,
2080        call_state_ref: ActiveCallStateRef,
2081        track_id: &String,
2082        mut invite_option: InviteOption,
2083        call_option: &CallOption,
2084        moh: Option<String>,
2085        auto_hangup: bool,
2086    ) -> Result<String, rsipstack::Error> {
2087        let ssrc = call_state_ref.read().await.ssrc;
2088        let per_call_srtp = call_option.sip.as_ref().and_then(|s| s.enable_srtp);
2089        let rtp_track = self
2090            .create_rtp_track(track_id.clone(), ssrc, per_call_srtp)
2091            .await
2092            .map_err(|e| rsipstack::Error::Error(e.to_string()))?;
2093
2094        let offer = Some(
2095            rtp_track
2096                .local_description()
2097                .await
2098                .map_err(|e| rsipstack::Error::Error(e.to_string()))?,
2099        );
2100
2101        {
2102            let mut cs = call_state_ref.write().await;
2103            if let Some(o) = cs.option.as_mut() {
2104                o.offer = offer.clone();
2105            }
2106            cs.start_time = Utc::now();
2107        };
2108
2109        invite_option.offer = offer.clone().map(|s| s.into());
2110
2111        // Set contact to local SIP endpoint address if not already set explicitly
2112        // Check if contact is still default (no scheme set) or if host is localhost-like
2113        let needs_contact = contact_needs_public_resolution(&invite_option.contact);
2114
2115        if needs_contact {
2116            let addrs = self.invitation.dialog_layer.endpoint.get_addrs();
2117            if let Some(addr) = find_local_addr_for_uri(&addrs, &invite_option.callee) {
2118                let contact_username = invite_option
2119                    .contact
2120                    .auth
2121                    .as_ref()
2122                    .map(|auth| auth.user.as_str())
2123                    .or_else(|| {
2124                        invite_option
2125                            .caller
2126                            .auth
2127                            .as_ref()
2128                            .map(|auth| auth.user.as_str())
2129                    });
2130                invite_option.contact = build_public_contact_uri(
2131                    &self.app_state.learned_public_address,
2132                    self.app_state.auto_learn_public_address_enabled(),
2133                    &addr,
2134                    contact_username,
2135                    Some(&invite_option.contact),
2136                );
2137            } else {
2138                return Err(rsipstack::Error::Error(format!(
2139                    "missing local SIP address for callee transport: {}",
2140                    invite_option.callee
2141                )));
2142            }
2143        }
2144
2145        let mut rtp_track_to_setup = Some(Box::new(rtp_track) as Box<dyn Track>);
2146
2147        if let Some(moh) = moh {
2148            let ssrc_and_moh = {
2149                let mut state = call_state_ref.write().await;
2150                state.moh = Some(moh.clone());
2151                if state.current_play_id.is_none() {
2152                    let ssrc = rand::random::<u32>();
2153                    Some((ssrc, moh.clone()))
2154                } else {
2155                    info!(
2156                        session_id = self.session_id,
2157                        "Something is playing, MOH will start after it ends"
2158                    );
2159                    None
2160                }
2161            };
2162
2163            if let Some((ssrc, moh_path)) = ssrc_and_moh {
2164                let file_track = FileTrack::new(self.server_side_track_id.clone())
2165                    .with_play_id(Some(moh_path.clone()))
2166                    .with_ssrc(ssrc)
2167                    .with_path(moh_path.clone())
2168                    .with_cancel_token(self.cancel_token.child_token());
2169                self.update_track_wrapper(Box::new(file_track), Some(moh_path))
2170                    .await;
2171            }
2172        } else {
2173            let track = rtp_track_to_setup.take().unwrap();
2174            self.setup_track_with_stream(&call_option, track)
2175                .await
2176                .map_err(|e| rsipstack::Error::Error(e.to_string()))?;
2177        }
2178
2179        info!(
2180            session_id = self.session_id,
2181            track_id,
2182            contact = %invite_option.contact,
2183            "invite {} -> {} offer: \n{}",
2184            invite_option.caller,
2185            invite_option.callee,
2186            offer.as_ref().map(|s| s.as_str()).unwrap_or("<NO OFFER>")
2187        );
2188
2189        let (dlg_state_sender, dlg_state_receiver) =
2190            self.invitation.dialog_layer.new_dialog_state_channel();
2191
2192        let states = InviteDialogStates {
2193            is_client: true,
2194            session_id: self.session_id.clone(),
2195            track_id: track_id.clone(),
2196            event_sender: self.event_sender.clone(),
2197            media_stream: self.media_stream.clone(),
2198            call_state: call_state_ref.clone(),
2199            cancel_token,
2200            terminated_reason: None,
2201            has_early_media: false,
2202        };
2203
2204        let hangup_headers = call_option
2205            .sip
2206            .as_ref()
2207            .and_then(|s| s.hangup_headers.as_ref())
2208            .map(|headers_map| {
2209                headers_map
2210                    .iter()
2211                    .map(|(k, v)| rsipstack::rsip::Header::Other(k.clone(), v.clone()))
2212                    .collect::<Vec<rsipstack::rsip::Header>>()
2213            });
2214
2215        let mut client_dialog_handler = DialogStateReceiverGuard::new(
2216            self.invitation.dialog_layer.clone(),
2217            dlg_state_receiver,
2218            hangup_headers,
2219        );
2220
2221        crate::spawn(async move {
2222            client_dialog_handler.process_dialog(states).await;
2223        });
2224
2225        let (dialog_id, answer) = self
2226            .invitation
2227            .invite(invite_option, dlg_state_sender)
2228            .await?;
2229
2230        self.call_state.write().await.moh = None;
2231
2232        if let Some(track) = rtp_track_to_setup {
2233            info!(
2234                session_id = self.session_id,
2235                track_id, "Stopping MOH and setting up RTP track"
2236            );
2237            self.media_stream
2238                .remove_track(&self.server_side_track_id, false)
2239                .await;
2240
2241            self.setup_track_with_stream(&call_option, track)
2242                .await
2243                .map_err(|e| rsipstack::Error::Error(e.to_string()))?;
2244        }
2245
2246        let answer = match answer {
2247            Some(answer) => {
2248                let s = String::from_utf8_lossy(&answer).to_string();
2249                if s.trim().is_empty() {
2250                    // 200 OK had no body — this is valid per RFC 3261 when the answer was
2251                    // already negotiated in a 183 Session Progress (early media).
2252                    // Fall back to the early SDP stored by the Early handler.
2253                    let cs = call_state_ref.read().await;
2254                    match cs.answer.clone() {
2255                        Some(early_sdp) if !early_sdp.is_empty() => {
2256                            info!(
2257                                session_id = self.session_id,
2258                                "200 OK has empty body; using early-media SDP from 183"
2259                            );
2260                            (early_sdp, true /* already applied */)
2261                        }
2262                        _ => {
2263                            warn!(
2264                                session_id = self.session_id,
2265                                "200 OK has empty body and no early-media SDP available"
2266                            );
2267                            (s, false)
2268                        }
2269                    }
2270                } else {
2271                    (s, false)
2272                }
2273            }
2274            None => {
2275                // No answer body at all — check if early media SDP is available before failing
2276                let cs = call_state_ref.read().await;
2277                match cs.answer.clone() {
2278                    Some(early_sdp) if !early_sdp.is_empty() => {
2279                        info!(
2280                            session_id = self.session_id,
2281                            "200 OK had no answer; using early-media SDP from 183"
2282                        );
2283                        (early_sdp, true /* already applied */)
2284                    }
2285                    _ => {
2286                        warn!(session_id = self.session_id, "no answer received");
2287                        return Err(rsipstack::Error::DialogError(
2288                            "No answer received".to_string(),
2289                            dialog_id,
2290                            rsipstack::rsip::StatusCode::NotAcceptableHere,
2291                        ));
2292                    }
2293                }
2294            }
2295        };
2296        let (answer, remote_description_already_applied) = answer;
2297
2298        {
2299            let mut cs = call_state_ref.write().await;
2300            if cs.answer.is_none() {
2301                cs.answer = Some(answer.clone());
2302            }
2303            if auto_hangup {
2304                cs.auto_hangup = Some((ssrc, CallRecordHangupReason::ByRefer));
2305            }
2306        }
2307        if !remote_description_already_applied {
2308            self.media_stream
2309                .update_remote_description(&track_id, &answer)
2310                .await
2311                .ok();
2312        }
2313
2314        Ok(answer)
2315    }
2316
2317    /// Detect if SDP is WebRTC format
2318    pub fn is_webrtc_sdp(sdp: &str) -> bool {
2319        (sdp.contains("a=ice-ufrag:") || sdp.contains("a=ice-pwd:"))
2320            && sdp.contains("a=fingerprint:")
2321    }
2322
2323    pub async fn setup_answer_track(
2324        &self,
2325        ssrc: u32,
2326        option: &CallOption,
2327        offer: String,
2328    ) -> Result<(String, Box<dyn Track>)> {
2329        let offer = match option.enable_ipv6 {
2330            Some(false) | None => strip_ipv6_candidates(&offer),
2331            _ => offer.clone(),
2332        };
2333
2334        let timeout = option.handshake_timeout.map(|t| Duration::from_secs(t));
2335
2336        let mut media_track = if Self::is_webrtc_sdp(&offer) {
2337            let mut rtc_config = RtcTrackConfig::default();
2338            rtc_config.mode = rustrtc::TransportMode::WebRtc;
2339            rtc_config.ice_servers = self.app_state.config.ice_servers.clone();
2340            if let Some(ref external_ip) = self.app_state.config.external_ip {
2341                rtc_config.external_ip = Some(external_ip.clone());
2342            }
2343            if let Some(ref bind_ip) = self.app_state.config.rtp_bind_ip {
2344                rtc_config.bind_ip = Some(bind_ip.clone());
2345            }
2346            rtc_config.enable_latching = self.app_state.config.enable_rtp_latching;
2347            rtc_config.enable_ice_lite = self.app_state.config.enable_ice_lite;
2348
2349            let webrtc_track = RtcTrack::new(
2350                self.cancel_token.child_token(),
2351                self.session_id.clone(),
2352                self.track_config.clone(),
2353                rtc_config,
2354            )
2355            .with_ssrc(ssrc);
2356
2357            Box::new(webrtc_track) as Box<dyn Track>
2358        } else {
2359            let per_call_srtp = option.sip.as_ref().and_then(|s| s.enable_srtp);
2360            let rtp_track = self
2361                .create_rtp_track(self.session_id.clone(), ssrc, per_call_srtp)
2362                .await?;
2363            Box::new(rtp_track) as Box<dyn Track>
2364        };
2365
2366        let answer = match media_track.handshake(offer.clone(), timeout).await {
2367            Ok(answer) => answer,
2368            Err(e) => {
2369                return Err(anyhow::anyhow!("handshake failed: {e}"));
2370            }
2371        };
2372
2373        return Ok((answer, media_track));
2374    }
2375
2376    pub async fn prepare_incoming_sip_track(
2377        &self,
2378        cancel_token: CancellationToken,
2379        call_state_ref: ActiveCallStateRef,
2380        track_id: &String,
2381        pending_dialog: PendingDialog,
2382        hangup_headers: Option<Vec<rsipstack::rsip::Header>>,
2383    ) -> Result<()> {
2384        let state_receiver = pending_dialog.state_receiver;
2385        //let pending_token_clone = pending_dialog.token;
2386
2387        let states = InviteDialogStates {
2388            is_client: false,
2389            session_id: self.session_id.clone(),
2390            track_id: track_id.clone(),
2391            event_sender: self.event_sender.clone(),
2392            media_stream: self.media_stream.clone(),
2393            call_state: self.call_state.clone(),
2394            cancel_token,
2395            terminated_reason: None,
2396            has_early_media: false,
2397        };
2398
2399        let initial_request = pending_dialog.dialog.initial_request();
2400        let offer = String::from_utf8_lossy(&initial_request.body).to_string();
2401
2402        let (ssrc, option) = {
2403            let call_state = call_state_ref.read().await;
2404            (
2405                call_state.ssrc,
2406                call_state.option.clone().unwrap_or_default(),
2407            )
2408        };
2409
2410        match self.setup_answer_track(ssrc, &option, offer).await {
2411            Ok((offer, track)) => {
2412                self.setup_track_with_stream(&option, track).await?;
2413                {
2414                    let mut state = self.call_state.write().await;
2415                    state.ready_to_answer = Some((offer, None, pending_dialog.dialog));
2416                }
2417            }
2418            Err(e) => {
2419                return Err(anyhow::anyhow!("error creating track: {}", e));
2420            }
2421        }
2422
2423        let mut client_dialog_handler = DialogStateReceiverGuard::new(
2424            self.invitation.dialog_layer.clone(),
2425            state_receiver,
2426            hangup_headers,
2427        );
2428
2429        crate::spawn(async move {
2430            client_dialog_handler.process_dialog(states).await;
2431        });
2432        Ok(())
2433    }
2434}
2435
2436impl Drop for ActiveCall {
2437    fn drop(&mut self) {
2438        info!(session_id = self.session_id, "dropping active call");
2439        if let Some(sender) = self.app_state.callrecord_sender.as_ref() {
2440            if let Some(record) = self.get_callrecord() {
2441                if let Err(e) = sender.send(record) {
2442                    warn!(
2443                        session_id = self.session_id,
2444                        "failed to send call record: {}", e
2445                    );
2446                }
2447            }
2448        }
2449    }
2450}
2451
2452impl ActiveCallState {
2453    pub fn merge_option(&self, mut option: CallOption) -> CallOption {
2454        if let Some(existing) = &self.option {
2455            if option.asr.is_none() {
2456                option.asr = existing.asr.clone();
2457            }
2458            if option.tts.is_none() {
2459                option.tts = existing.tts.clone();
2460            }
2461            if option.vad.is_none() {
2462                option.vad = existing.vad.clone();
2463            }
2464            if option.denoise.is_none() {
2465                option.denoise = existing.denoise;
2466            }
2467            if option.recorder.is_none() {
2468                option.recorder = existing.recorder.clone();
2469            }
2470            if option.eou.is_none() {
2471                option.eou = existing.eou.clone();
2472            }
2473            if option.extra.is_none() {
2474                option.extra = existing.extra.clone();
2475            }
2476            if option.ambiance.is_none() {
2477                option.ambiance = existing.ambiance.clone();
2478            }
2479        }
2480        option
2481    }
2482
2483    pub fn set_hangup_reason(&mut self, reason: CallRecordHangupReason) {
2484        if self.hangup_reason.is_none() {
2485            self.hangup_reason = Some(reason);
2486        }
2487    }
2488
2489    pub fn build_hangup_event(
2490        &self,
2491        track_id: TrackId,
2492        initiator: Option<String>,
2493    ) -> crate::event::SessionEvent {
2494        let from = self.option.as_ref().and_then(|o| o.caller.as_ref());
2495        let to = self.option.as_ref().and_then(|o| o.callee.as_ref());
2496        let extra = self.extras.clone();
2497
2498        crate::event::SessionEvent::Hangup {
2499            track_id,
2500            timestamp: crate::media::get_timestamp(),
2501            reason: Some(format!("{:?}", self.hangup_reason)),
2502            initiator,
2503            start_time: self.start_time.to_rfc3339(),
2504            answer_time: self.answer_time.map(|t| t.to_rfc3339()),
2505            ringing_time: self.ring_time.map(|t| t.to_rfc3339()),
2506            hangup_time: Utc::now().to_rfc3339(),
2507            extra,
2508            from: from.map(|f| f.into()),
2509            to: to.map(|f| f.into()),
2510            refer: Some(self.is_refer),
2511        }
2512    }
2513
2514    pub fn build_callrecord(
2515        &self,
2516        app_state: AppState,
2517        session_id: String,
2518        call_type: ActiveCallType,
2519    ) -> CallRecord {
2520        let option = self.option.clone().unwrap_or_default();
2521        let recorder = if option.recorder.is_some() {
2522            let recorder_file = app_state.get_recorder_file(&session_id);
2523            if std::path::Path::new(&recorder_file).exists() {
2524                let file_size = std::fs::metadata(&recorder_file)
2525                    .map(|m| m.len())
2526                    .unwrap_or(0);
2527                vec![crate::callrecord::CallRecordMedia {
2528                    track_id: session_id.clone(),
2529                    path: recorder_file,
2530                    size: file_size,
2531                    extra: None,
2532                }]
2533            } else {
2534                vec![]
2535            }
2536        } else {
2537            vec![]
2538        };
2539
2540        let dump_event_file = app_state.get_dump_events_file(&session_id);
2541        let dump_event_file = if std::path::Path::new(&dump_event_file).exists() {
2542            Some(dump_event_file)
2543        } else {
2544            None
2545        };
2546
2547        let refer_callrecord = self.refer_callstate.as_ref().and_then(|rc| {
2548            if let Ok(rc) = rc.try_read() {
2549                Some(Box::new(rc.build_callrecord(
2550                    app_state.clone(),
2551                    rc.session_id.clone(),
2552                    ActiveCallType::B2bua,
2553                )))
2554            } else {
2555                None
2556            }
2557        });
2558
2559        let caller = option.caller.clone().unwrap_or_default();
2560        let callee = option.callee.clone().unwrap_or_default();
2561
2562        CallRecord {
2563            option: Some(option),
2564            call_id: session_id,
2565            call_type,
2566            start_time: self.start_time,
2567            ring_time: self.ring_time.clone(),
2568            answer_time: self.answer_time.clone(),
2569            end_time: Utc::now(),
2570            caller,
2571            callee,
2572            hangup_reason: self.hangup_reason.clone(),
2573            hangup_messages: Vec::new(),
2574            status_code: self.last_status_code,
2575            extras: self.extras.clone(),
2576            dump_event_file,
2577            recorder,
2578            refer_callrecord,
2579        }
2580    }
2581}