Skip to main content

active_call/call/
tracks.rs

1//! Track construction and wiring for [`ActiveCall`].
2//!
3//! Split out of `active_call.rs` so the call-flow orchestration stays separate
4//! from the (heavily reused) track assembly building blocks.
5
6use super::active_call::{ActiveCall, ActiveCallType, PendingCallerTrack};
7use super::state::LegShared;
8use crate::CallOption;
9use crate::event::SessionEvent;
10use crate::media::TrackId;
11use crate::media::ambiance::SharedAmbianceProcessor;
12use crate::media::engine::StreamEngine;
13use crate::media::negotiate::strip_ipv6_candidates;
14use crate::media::processor::SubscribeProcessor;
15use crate::media::track::Track;
16use crate::media::track::file::FileTrack;
17use crate::media::track::rtc::{RtcTrack, RtcTrackConfig};
18use crate::media::track::websocket::{WebsocketBytesReceiver, WebsocketTrack};
19use crate::useragent::invitation::PendingDialog;
20use crate::useragent::public_address::{
21    build_public_contact_uri, contact_needs_public_resolution, find_local_addr_for_uri,
22};
23use anyhow::Result;
24use audio_codec::CodecType;
25use chrono::Utc;
26use rsipstack::dialog::invitation::InviteOption;
27use rsipstack::rsip::prelude::HeadersExt;
28use std::time::Duration;
29use tokio_util::sync::CancellationToken;
30use tracing::{debug, info, warn};
31
32use super::active_call::PendingSipAnswer;
33use super::sip::{DialogStateReceiverGuard, InviteDialogStates};
34
35/// Track id of the SIP leg of a non-SIP call type (WebSocket/Webrtc) that was
36/// created by an inbound SIP INVITE. Must stay distinct from the caller track
37/// (`session_id`) and the refer leg (`server_side_track_id`) so the three
38/// legs coexist in the media stream's full-mesh forwarding.
39const SIP_LEG_TRACK_ID: &str = "sip-leg-track";
40
41/// Restrict a track's codec capabilities to codecs the remote offer actually
42/// advertises. rustrtc's `create_answer` otherwise emits its default
43/// capability set (e.g. opus-only), which intersects to an *empty* codec list
44/// against a plain PCMU/PCMA SIP offer (`intersect_answer` then produces an
45/// unparseable `m=` line). Configured codecs win when they overlap the offer;
46/// the offer wins when they don't.
47fn restrict_codecs_to_offer(rtc_config: &mut RtcTrackConfig, offer: &str) {
48    let offer_codecs: Vec<CodecType> = offer
49        .lines()
50        .filter_map(|line| {
51            let value = line.trim().strip_prefix("a=rtpmap:")?;
52            crate::media::negotiate::parse_rtpmap(value)
53                .ok()
54                .map(|(_, codec, ..)| codec)
55        })
56        .collect();
57    if offer_codecs.is_empty() {
58        return;
59    }
60    if rtc_config.codecs.is_empty() {
61        rtc_config.codecs = offer_codecs;
62    } else {
63        rtc_config.codecs.retain(|c| offer_codecs.contains(c));
64        if rtc_config.codecs.is_empty() {
65            rtc_config.codecs = offer_codecs;
66        }
67    }
68}
69
70/// Everything needed to place one outgoing INVITE (main leg or refer leg).
71pub(super) struct OutgoingLeg {
72    /// Lifetime token of the leg (call token or the refer-specific child).
73    pub cancel_token: CancellationToken,
74    /// Lock-free shared state of the leg being established.
75    pub leg: LegShared,
76    /// Media track id for the leg.
77    pub track_id: TrackId,
78    /// The outgoing INVITE (offer may be filled in by the caller).
79    pub invite_option: InviteOption,
80    /// Call option used to build media processors / MOH.
81    pub call_option: CallOption,
82    /// Music-on-hold path to play while the INVITE is in flight.
83    pub moh: Option<String>,
84    /// Refer legs with `auto_hangup`: hang up the parent call when this leg ends.
85    pub auto_hangup: bool,
86}
87
88impl ActiveCall {
89    /// Apply the configured codec list to an RTC track config.
90    fn rtc_apply_codecs(&self, rtc_config: &mut RtcTrackConfig) {
91        if let Some(codecs) = &self.app_state.config.codecs {
92            let mut codec_types = Vec::new();
93            for c in codecs {
94                match c.to_lowercase().as_str() {
95                    "pcmu" => codec_types.push(CodecType::PCMU),
96                    "pcma" => codec_types.push(CodecType::PCMA),
97                    "g722" => codec_types.push(CodecType::G722),
98                    "g729" => codec_types.push(CodecType::G729),
99                    "opus" => codec_types.push(CodecType::Opus),
100                    "dtmf" | "2833" | "telephone_event" => {
101                        codec_types.push(CodecType::TelephoneEvent)
102                    }
103                    _ => {}
104                }
105            }
106            if !codec_types.is_empty() {
107                rtc_config.preferred_codec = Some(codec_types[0].clone());
108                rtc_config.codecs = codec_types;
109            }
110        }
111    }
112
113    /// Apply external/bind addresses from the global config.
114    fn rtc_apply_network(&self, rtc_config: &mut RtcTrackConfig) {
115        if let Some(ref external_ip) = self.app_state.config.external_ip {
116            rtc_config.external_ip = Some(external_ip.clone());
117        }
118        if let Some(ref bind_ip) = self.app_state.config.rtp_bind_ip {
119            rtc_config.bind_ip = Some(bind_ip.clone());
120        }
121    }
122
123    /// Apply RTP latching and ICE-lite flags (per-call option overrides config).
124    fn rtc_apply_latching(&self, rtc_config: &mut RtcTrackConfig) {
125        rtc_config.enable_latching = self.app_state.config.enable_rtp_latching;
126        rtc_config.enable_ice_lite = self
127            .progress
128            .load()
129            .option
130            .as_ref()
131            .and_then(|o| o.enable_ice_lite)
132            .or(self.app_state.config.enable_ice_lite);
133    }
134
135    /// Build a looping file track on the server-side track.
136    pub(super) fn make_file_track(&self, path: String, ssrc: u32) -> FileTrack {
137        FileTrack::new(self.server_side_track_id.clone())
138            .with_play_id(Some(path.clone()))
139            .with_ssrc(ssrc)
140            .with_path(path)
141            .with_cancel_token(self.cancel_token.child_token())
142    }
143
144    /// Emit a `Reject` session event from an rsipstack dialog error.
145    pub(super) fn emit_reject_from_rsip_error(
146        &self,
147        track_id: TrackId,
148        refer: bool,
149        e: &rsipstack::Error,
150    ) {
151        if let rsipstack::Error::DialogError(reason, _, code) = e {
152            self.event_sender
153                .send(SessionEvent::Reject {
154                    track_id,
155                    timestamp: crate::media::get_timestamp(),
156                    reason: reason.clone(),
157                    code: Some(code.code() as u32),
158                    refer: Some(refer),
159                })
160                .ok();
161        }
162    }
163
164    /// If a pending incoming dialog exists for this session, start preparing
165    /// the incoming SIP track; shared by the Sip and B2bua call types.
166    pub(super) async fn try_prepare_incoming_sip_track(
167        &self,
168        hangup_headers: Option<Vec<rsipstack::rsip::Header>>,
169    ) -> Option<Result<()>> {
170        let dialog_id = self
171            .invitation
172            .find_dialog_id_by_session_id(&self.session_id)?;
173        let pending_dialog = self.invitation.get_pending_call(&dialog_id)?;
174        Some(
175            self.prepare_incoming_sip_track(pending_dialog, hangup_headers)
176                .await,
177        )
178    }
179
180    /// For non-SIP call types (WebSocket/Webrtc) whose session was created by
181    /// an inbound SIP INVITE, build the SIP leg's answering track and remember
182    /// the prepared 200 OK. `do_accept` sends the answer and registers the
183    /// track, so the carrier leg is actually established instead of sitting in
184    /// `Trying` until the far end times out (production: VOS3000 CANCELs the
185    /// unanswered INVITE after 20s). No-op when there is no pending dialog or
186    /// the INVITE carries no offer.
187    pub(super) async fn try_prepare_pending_sip_answer(
188        &self,
189        option: &CallOption,
190    ) -> Result<()> {
191        let Some(dialog_id) = self
192            .invitation
193            .find_dialog_id_by_session_id(&self.session_id)
194        else {
195            return Ok(());
196        };
197        let Some(pending_dialog) = self.invitation.get_pending_call(&dialog_id) else {
198            return Ok(());
199        };
200
201        let initial_request = pending_dialog.dialog.initial_request();
202        let offer = String::from_utf8_lossy(initial_request.body()).to_string();
203        debug!(
204            session_id = self.session_id,
205            offer = %offer,
206            "preparing sip leg answer"
207        );
208        if offer.trim().is_empty() {
209            warn!(
210                session_id = self.session_id,
211                "inbound SIP dialog has no SDP offer; skipping sip leg preparation"
212            );
213            return Ok(());
214        }
215
216        // The SIP leg needs its own track id and SSRC: `session_id` is the
217        // caller (WebSocket/Webrtc) track and `server_side_track_id` belongs
218        // to the refer leg.
219        let track_id: TrackId = SIP_LEG_TRACK_ID.to_string();
220        let ssrc = rand::random::<u32>();
221
222        let mut rtc_config = RtcTrackConfig::default();
223        let use_srtp = option
224            .sip
225            .as_ref()
226            .and_then(|s| s.enable_srtp)
227            .or(self.app_state.config.enable_srtp)
228            .unwrap_or(false);
229        rtc_config.mode = if use_srtp {
230            rustrtc::TransportMode::Srtp
231        } else {
232            rustrtc::TransportMode::Rtp
233        };
234        self.rtc_apply_codecs(&mut rtc_config);
235        if rtc_config.preferred_codec.is_none() {
236            rtc_config.preferred_codec = Some(self.track_config.codec);
237        }
238        rtc_config.rtp_port_range = self
239            .app_state
240            .config
241            .rtp_start_port
242            .zip(self.app_state.config.rtp_end_port);
243        self.rtc_apply_network(&mut rtc_config);
244        self.rtc_apply_latching(&mut rtc_config);
245        restrict_codecs_to_offer(&mut rtc_config, &offer);
246
247        let mut sip_track = RtcTrack::new(
248            self.cancel_token.child_token(),
249            track_id.clone(),
250            self.track_config.clone(),
251            rtc_config,
252        )
253        .with_ssrc(ssrc);
254        sip_track.create().await?;
255
256        let timeout = option.handshake_timeout.map(|t| Duration::from_secs(t));
257        let answer = sip_track
258            .handshake(offer, timeout)
259            .await
260            .map_err(|e| anyhow::anyhow!("sip leg handshake failed: {e}"))?;
261
262        // Drive the incoming dialog's state machine: hangs the dialog up when
263        // the call ends, and ends cleanly on a far-end BYE/CANCEL.
264        let states = InviteDialogStates::new(
265            false,
266            self.session_id.clone(),
267            track_id.clone(),
268            self.event_sender.clone(),
269            self.media_stream.clone(),
270            self.leg(),
271            self.cancel_token.clone(),
272            None,
273        );
274        let hangup_headers = option
275            .sip
276            .as_ref()
277            .and_then(|s| s.hangup_headers.as_ref())
278            .map(crate::sip_util::sip_headers_from_map);
279        let mut client_dialog_handler = DialogStateReceiverGuard::new(
280            self.invitation.dialog_layer.clone(),
281            pending_dialog.state_receiver,
282            hangup_headers,
283        );
284        crate::spawn(async move {
285            client_dialog_handler.process_dialog(states).await;
286        });
287
288        info!(
289            session_id = self.session_id,
290            track_id, "prepared sip leg answer for non-sip call type"
291        );
292
293        self.set_pending_sip_answer(PendingSipAnswer {
294            answer,
295            dialog: pending_dialog.dialog,
296            track: Box::new(sip_track),
297        });
298        Ok(())
299    }
300
301    /// Per-call ambiance option merged over the global config.
302    fn merged_ambiance(
303        &self,
304        call_ambiance: Option<&crate::media::ambiance::AmbianceOption>,
305    ) -> crate::media::ambiance::AmbianceOption {
306        let mut opt = call_ambiance.cloned().unwrap_or_default();
307        if let Some(global) = &self.app_state.config.ambiance {
308            opt.merge(global);
309        }
310        opt
311    }
312
313    pub(super) async fn create_rtp_track(
314        &self,
315        track_id: TrackId,
316        ssrc: u32,
317        enable_srtp: Option<bool>,
318        offer: Option<&str>,
319    ) -> Result<RtcTrack> {
320        let mut rtc_config = RtcTrackConfig::default();
321        // Per-call flag takes precedence over global config.
322        let use_srtp = enable_srtp
323            .or(self.app_state.config.enable_srtp)
324            .unwrap_or(false);
325        rtc_config.mode = if use_srtp {
326            rustrtc::TransportMode::Srtp
327        } else {
328            rustrtc::TransportMode::Rtp
329        };
330
331        self.rtc_apply_codecs(&mut rtc_config);
332        if let Some(offer) = offer {
333            restrict_codecs_to_offer(&mut rtc_config, offer);
334        }
335
336        if rtc_config.preferred_codec.is_none() {
337            rtc_config.preferred_codec = Some(self.track_config.codec.clone());
338        }
339
340        rtc_config.rtp_port_range = self
341            .app_state
342            .config
343            .rtp_start_port
344            .zip(self.app_state.config.rtp_end_port);
345
346        self.rtc_apply_network(&mut rtc_config);
347        self.rtc_apply_latching(&mut rtc_config);
348
349        let mut track = RtcTrack::new(
350            self.cancel_token.child_token(),
351            track_id,
352            self.track_config.clone(),
353            rtc_config,
354        )
355        .with_ssrc(ssrc);
356
357        track.create().await?;
358
359        Ok(track)
360    }
361
362    pub(super) async fn setup_track_with_stream(
363        &self,
364        option: &CallOption,
365        mut track: Box<dyn Track>,
366    ) -> Result<()> {
367        let processors = match StreamEngine::create_processors(
368            self.app_state.stream_engine.clone(),
369            track.id().clone(),
370            self.cancel_token.child_token(),
371            self.event_sender.clone(),
372            self.media_stream.packet_sender.clone(),
373            option,
374        )
375        .await
376        {
377            Ok(processors) => processors,
378            Err(e) => {
379                warn!(
380                    session_id = self.session_id,
381                    "failed to prepare stream processors: {}", e
382                );
383                vec![]
384            }
385        };
386
387        // Add all processors from the hook
388        for processor in processors {
389            track.append_processor(processor);
390        }
391
392        self.update_track_wrapper(track, None).await;
393        Ok(())
394    }
395
396    pub(super) async fn update_track_wrapper(
397        &self,
398        mut track: Box<dyn Track>,
399        play_id: Option<String>,
400    ) {
401        let (call_ambiance, subscribe) = {
402            let state = self.progress.load_full();
403            (
404                state.option.as_ref().and_then(|o| o.ambiance.clone()),
405                state
406                    .option
407                    .as_ref()
408                    .and_then(|o| o.subscribe)
409                    .unwrap_or_default(),
410            )
411        };
412        let ambiance_opt = self.merged_ambiance(call_ambiance.as_ref());
413
414        let shared_ambiance = match self
415            .media_stream
416            .ensure_ambiance(ambiance_opt, self.server_side_track_id.clone())
417            .await
418        {
419            Ok(shared) => shared,
420            Err(e) => {
421                tracing::error!("failed to load ambiance wav {}", e);
422                None
423            }
424        };
425
426        if track.id() == &self.server_side_track_id {
427            if let Some(shared) = shared_ambiance {
428                info!(session_id = self.session_id, "loaded ambiance processor");
429                track.append_processor(Box::new(SharedAmbianceProcessor::new(shared)));
430            }
431        }
432
433        if subscribe && self.call_type != ActiveCallType::WebSocket {
434            let (track_index, sub_track_id) = if track.id() == &self.server_side_track_id {
435                (0, self.server_side_track_id.clone())
436            } else {
437                (1, self.session_id.clone())
438            };
439            let sub_processor =
440                SubscribeProcessor::new(self.event_sender.clone(), sub_track_id, track_index);
441            track.append_processor(Box::new(sub_processor));
442        }
443
444        self.set_current_play(play_id.clone());
445        self.media_stream.update_track(track, play_id).await;
446    }
447
448    pub(super) async fn setup_caller_track(&self, option: &CallOption) -> Result<()> {
449        let hangup_headers = option
450            .sip
451            .as_ref()
452            .and_then(|s| s.hangup_headers.as_ref())
453            .map(crate::sip_util::sip_headers_from_map);
454        self.set_option(option.clone());
455        info!(
456            session_id = self.session_id,
457            call_type = ?self.call_type,
458            "setup caller track"
459        );
460
461        let track = match self.call_type {
462            ActiveCallType::Webrtc => {
463                let track = self.create_webrtc_track().await?;
464                // The session may be backed by a ringing inbound SIP dialog;
465                // prepare its answer so do_accept establishes the carrier leg.
466                self.try_prepare_pending_sip_answer(option).await?;
467                Some(track)
468            }
469            ActiveCallType::WebSocket => {
470                let audio_receiver = self.audio_receiver.lock().unwrap().take();
471                if let Some(receiver) = audio_receiver {
472                    let track = self.create_websocket_track(receiver).await?;
473                    // The session may be backed by a ringing inbound SIP
474                    // dialog; prepare its answer so do_accept establishes the
475                    // carrier leg (otherwise it stays in Trying until the far
476                    // end CANCELs).
477                    self.try_prepare_pending_sip_answer(option).await?;
478                    Some(track)
479                } else {
480                    None
481                }
482            }
483            ActiveCallType::Sip | ActiveCallType::B2bua => {
484                // Incoming call with a pending dialog: prepare the answering leg.
485                if let Some(result) = self.try_prepare_incoming_sip_track(hangup_headers).await {
486                    return result;
487                }
488
489                if matches!(self.call_type, ActiveCallType::B2bua) {
490                    warn!(
491                        session_id = self.session_id,
492                        "no pending dialog found for B2BUA call"
493                    );
494                    return Err(anyhow::anyhow!(
495                        "no pending dialog found for session_id: {}",
496                        self.session_id
497                    ));
498                }
499
500                // Outgoing SIP call below (Sip only).
501                // Auto-inject credentials from registered users if not already provided
502                let mut option = option.clone();
503                if option.sip.is_none()
504                    || option
505                        .sip
506                        .as_ref()
507                        .and_then(|s| s.username.as_ref())
508                        .is_none()
509                {
510                    if let Some(callee) = &option.callee {
511                        if let Some(cred) = self.app_state.find_credentials_for_callee(callee) {
512                            if option.sip.is_none() {
513                                option.sip = Some(crate::SipOption {
514                                    username: Some(cred.username.clone()),
515                                    password: Some(cred.password.clone()),
516                                    realm: cred.realm.clone(),
517                                    ..Default::default()
518                                });
519                            }
520                        }
521                    }
522                }
523
524                let mut invite_option = option.build_invite_option()?;
525                invite_option.call_id = Some(self.session_id.clone());
526
527                let out = OutgoingLeg {
528                    cancel_token: self.cancel_token.clone(),
529                    leg: self.leg(),
530                    track_id: self.session_id.clone(),
531                    invite_option,
532                    call_option: option.clone(),
533                    moh: None,
534                    auto_hangup: false,
535                };
536                match self.create_outgoing_sip_track(out).await {
537                    Ok(answer) => {
538                        self.event_sender
539                            .send(SessionEvent::Answer {
540                                timestamp: crate::media::get_timestamp(),
541                                track_id: self.session_id.clone(),
542                                sdp: answer,
543                                refer: Some(false),
544                            })
545                            .ok();
546                        return Ok(());
547                    }
548                    Err(e) => {
549                        warn!(
550                            session_id = self.session_id,
551                            "failed to create sip track: {}", e
552                        );
553                        self.emit_reject_from_rsip_error(self.session_id.clone(), false, &e);
554                        return Err(e.into());
555                    }
556                }
557            }
558        };
559        match track {
560            Some(track) => {
561                self.finish_caller_stack(&option, PendingCallerTrack::NotStarted(track))
562                    .await?;
563            }
564            None => {
565                warn!(session_id = self.session_id, "no track created for caller");
566                return Err(anyhow::anyhow!("no track created for caller"));
567            }
568        }
569        Ok(())
570    }
571
572    pub(super) async fn finish_caller_stack(
573        &self,
574        option: &CallOption,
575        pending_track: PendingCallerTrack,
576    ) -> Result<()> {
577        self.ensure_call_ambiance(option).await;
578        match pending_track {
579            PendingCallerTrack::NotStarted(track) => {
580                self.setup_track_with_stream(option, track).await?;
581            }
582            PendingCallerTrack::StartedForEarlyMedia => {
583                // Track is already running in the media stream (started during ringing
584                // for early-media ringtone). Build processors from the accept option
585                // and append them to the running caller track.
586                let track_id = self.session_id.clone();
587                let processors = StreamEngine::create_processors(
588                    self.app_state.stream_engine.clone(),
589                    track_id.clone(),
590                    self.cancel_token.child_token(),
591                    self.event_sender.clone(),
592                    self.media_stream.packet_sender.clone(),
593                    option,
594                )
595                .await
596                .unwrap_or_else(|e| {
597                    warn!(
598                        session_id = self.session_id,
599                        "failed to create processors on accept: {}", e
600                    );
601                    vec![]
602                });
603                for processor in processors {
604                    self.media_stream
605                        .append_processor(&track_id, processor)
606                        .await
607                        .ok();
608                }
609            }
610        }
611
612        {
613            let call_state = self.progress.load_full();
614            if let Some(ref answer) = call_state.answer {
615                info!(
616                    session_id = self.session_id,
617                    "sending answer event: {}", answer,
618                );
619                self.event_sender
620                    .send(SessionEvent::Answer {
621                        timestamp: crate::media::get_timestamp(),
622                        track_id: self.session_id.clone(),
623                        sdp: answer.clone(),
624                        refer: Some(false),
625                    })
626                    .ok();
627            } else {
628                warn!(
629                    session_id = self.session_id,
630                    "no answer in state to send event"
631                );
632            }
633        }
634        Ok(())
635    }
636
637    pub(super) async fn ensure_call_ambiance(&self, option: &CallOption) {
638        let opt = self.merged_ambiance(option.ambiance.as_ref());
639        if let Err(e) = self
640            .media_stream
641            .ensure_ambiance(opt, self.server_side_track_id.clone())
642            .await
643        {
644            tracing::error!(
645                session_id = self.session_id,
646                "failed to load ambiance wav {}",
647                e
648            );
649        }
650    }
651
652    pub(super) async fn create_websocket_track(
653        &self,
654        audio_receiver: WebsocketBytesReceiver,
655    ) -> Result<Box<dyn Track>> {
656        let codec = self
657            .progress
658            .load_full()
659            .option
660            .as_ref()
661            .map(|o| o.codec.clone())
662            .unwrap_or_default();
663
664        let ws_track = WebsocketTrack::new(
665            self.cancel_token.child_token(),
666            self.session_id.clone(),
667            self.track_config.clone(),
668            self.event_sender.clone(),
669            audio_receiver,
670            codec,
671            self.ssrc,
672        );
673
674        self.leg().update_progress(|p| {
675            p.answer = Some("".to_string());
676            p.on_answered();
677        });
678
679        Ok(Box::new(ws_track))
680    }
681
682    pub(super) async fn create_webrtc_track(&self) -> Result<Box<dyn Track>> {
683        let option = self.progress.load_full().option.clone().unwrap_or_default();
684        let ssrc = self.ssrc;
685
686        let mut rtc_config = RtcTrackConfig::default();
687        rtc_config.mode = rustrtc::TransportMode::WebRtc; // WebRTC
688        rtc_config.ice_servers = self.app_state.config.ice_servers.clone();
689
690        self.rtc_apply_codecs(&mut rtc_config);
691        self.rtc_apply_network(&mut rtc_config);
692
693        let mut webrtc_track = RtcTrack::new(
694            self.cancel_token.child_token(),
695            self.session_id.clone(),
696            self.track_config.clone(),
697            rtc_config,
698        )
699        .with_ssrc(ssrc);
700
701        let timeout = option.handshake_timeout.map(|t| Duration::from_secs(t));
702        let offer = match option.enable_ipv6 {
703            Some(false) | None => {
704                strip_ipv6_candidates(option.offer.as_ref().unwrap_or(&"".to_string()))
705            }
706            _ => option.offer.clone().unwrap_or("".to_string()),
707        };
708        let answer: Option<String>;
709        match webrtc_track.handshake(offer, timeout).await {
710            Ok(answer_sdp) => {
711                answer = match option.enable_ipv6 {
712                    Some(false) | None => Some(strip_ipv6_candidates(&answer_sdp)),
713                    Some(true) => Some(answer_sdp.to_string()),
714                };
715            }
716            Err(e) => {
717                warn!(session_id = self.session_id, "failed to setup track: {}", e);
718                return Err(anyhow::anyhow!("Failed to setup track: {}", e));
719            }
720        }
721
722        self.leg().update_progress(|p| {
723            p.answer = answer.clone();
724            p.on_answered();
725        });
726        Ok(Box::new(webrtc_track))
727    }
728
729    pub(super) async fn create_outgoing_sip_track(
730        &self,
731        mut out: OutgoingLeg,
732    ) -> Result<String, rsipstack::Error> {
733        // Apply trunk rules (match + rewrite caller/callee/contact) to the
734        // outgoing INVITE/REFER before it is sent. Covers both normal invite
735        // calls and refer legs since both flow through this function.
736        self.app_state
737            .config
738            .apply_trunk_rules(&mut out.invite_option);
739
740        let track_id = &out.track_id;
741        let ssrc = out.leg.ssrc;
742        let per_call_srtp = out.call_option.sip.as_ref().and_then(|s| s.enable_srtp);
743        let rtp_track = self
744            .create_rtp_track(track_id.clone(), ssrc, per_call_srtp, None)
745            .await
746            .map_err(|e| rsipstack::Error::Error(e.to_string()))?;
747
748        let offer = Some(
749            rtp_track
750                .local_description()
751                .await
752                .map_err(|e| rsipstack::Error::Error(e.to_string()))?,
753        );
754
755        out.leg.update_progress(|p| {
756            if let Some(o) = p.option.as_mut() {
757                o.offer = offer.clone();
758            }
759            p.start_time = Some(Utc::now());
760        });
761
762        let invite_option = &mut out.invite_option;
763        invite_option.offer = offer.clone().map(|s| s.into());
764
765        // Set contact to local SIP endpoint address if not already set explicitly
766        // Check if contact is still default (no scheme set) or if host is localhost-like
767        let needs_contact = contact_needs_public_resolution(&invite_option.contact);
768
769        if needs_contact {
770            let addrs = self.invitation.dialog_layer.endpoint.get_addrs();
771            if let Some(addr) = find_local_addr_for_uri(&addrs, &invite_option.callee) {
772                let contact_username = invite_option
773                    .contact
774                    .auth
775                    .as_ref()
776                    .map(|auth| auth.user.as_str())
777                    .or_else(|| {
778                        invite_option
779                            .caller
780                            .auth
781                            .as_ref()
782                            .map(|auth| auth.user.as_str())
783                    });
784                invite_option.contact = build_public_contact_uri(
785                    &self.app_state.learned_public_address,
786                    self.app_state.auto_learn_public_address_enabled(),
787                    &addr,
788                    contact_username,
789                    Some(&invite_option.contact),
790                );
791            } else {
792                return Err(rsipstack::Error::Error(format!(
793                    "missing local SIP address for callee transport: {}",
794                    invite_option.callee
795                )));
796            }
797        }
798
799        let mut rtp_track_to_setup = Some(Box::new(rtp_track) as Box<dyn Track>);
800
801        if let Some(moh) = out.moh.take() {
802            let ssrc_and_moh = {
803                self.set_moh(Some(moh.clone()));
804                if self.current_play().is_none() {
805                    let ssrc = rand::random::<u32>();
806                    Some((ssrc, moh.clone()))
807                } else {
808                    info!(
809                        session_id = self.session_id,
810                        "Something is playing, MOH will start after it ends"
811                    );
812                    None
813                }
814            };
815
816            if let Some((ssrc, moh_path)) = ssrc_and_moh {
817                let file_track = self.make_file_track(moh_path.clone(), ssrc);
818                self.update_track_wrapper(Box::new(file_track), Some(moh_path))
819                    .await;
820            }
821        } else {
822            let track = rtp_track_to_setup.take().unwrap();
823            self.setup_track_with_stream(&out.call_option, track)
824                .await
825                .map_err(|e| rsipstack::Error::Error(e.to_string()))?;
826        }
827
828        info!(
829            session_id = self.session_id,
830            track_id,
831            contact = %invite_option.contact,
832            "invite {} -> {} offer: \n{}",
833            invite_option.caller,
834            invite_option.callee,
835            offer.as_ref().map(|s| s.as_str()).unwrap_or("<NO OFFER>")
836        );
837
838        let (dlg_state_sender, dlg_state_receiver) =
839            self.invitation.dialog_layer.new_dialog_state_channel();
840
841        let states = InviteDialogStates::new(
842            true,
843            self.session_id.clone(),
844            track_id.clone(),
845            self.event_sender.clone(),
846            self.media_stream.clone(),
847            out.leg.clone(),
848            out.cancel_token.clone(),
849            out.auto_hangup.then_some(crate::callrecord::CallRecordHangupReason::ByRefer),
850        );
851
852        let hangup_headers = out
853            .call_option
854            .sip
855            .as_ref()
856            .and_then(|s| s.hangup_headers.as_ref())
857            .map(crate::sip_util::sip_headers_from_map);
858
859        let mut client_dialog_handler = DialogStateReceiverGuard::new(
860            self.invitation.dialog_layer.clone(),
861            dlg_state_receiver,
862            hangup_headers,
863        );
864
865        crate::spawn(async move {
866            client_dialog_handler.process_dialog(states).await;
867        });
868
869        let (dialog_id, answer) = self
870            .invitation
871            .invite(out.invite_option, dlg_state_sender)
872            .await?;
873
874        if out.cancel_token.is_cancelled() {
875            // CANCEL and the far end's 2xx crossed in flight (RFC 3261 S9.1
876            // glare) - e.g. a slow PSTN gateway had already committed to
877            // alerting the handset by the time our CANCEL arrived, and it
878            // answered anyway. do_invite already ACKed the 2xx; per spec we
879            // must now BYE it instead of treating this as a normal answer -
880            // the ActiveCall session this invite belongs to is already gone.
881            if let Some(dialog) = self.invitation.dialog_layer.get_dialog(&dialog_id) {
882                if let Err(e) = dialog.hangup().await {
883                    warn!(
884                        session_id = self.session_id,
885                        "failed to BYE a late-confirmed cancelled invite: {}", e
886                    );
887                }
888            }
889            return Err(rsipstack::Error::DialogError(
890                "invite was cancelled before this late answer arrived".to_string(),
891                dialog_id,
892                rsipstack::rsip::StatusCode::RequestTerminated,
893            ));
894        }
895
896        self.set_moh(None);
897
898        if let Some(track) = rtp_track_to_setup {
899            info!(
900                session_id = self.session_id,
901                track_id, "Stopping MOH and setting up RTP track"
902            );
903            self.media_stream
904                .remove_track(&self.server_side_track_id, false)
905                .await;
906
907            self.setup_track_with_stream(&out.call_option, track)
908                .await
909                .map_err(|e| rsipstack::Error::Error(e.to_string()))?;
910        }
911
912        // Resolve the final answer SDP, falling back to the early-media (183)
913        // SDP when the 200 OK carries no body.
914        let early_answer = out.leg.progress.load_full().answer.clone();
915        let (answer, remote_description_already_applied) =
916            match crate::call::state::resolve_final_answer(answer, early_answer.as_ref()) {
917                Ok(resolved) => resolved,
918                Err(msg) => {
919                    warn!(session_id = self.session_id, "{}", msg);
920                    return Err(rsipstack::Error::DialogError(
921                        "No answer received".to_string(),
922                        dialog_id,
923                        rsipstack::rsip::StatusCode::NotAcceptableHere,
924                    ));
925                }
926            };
927
928        out.leg.update_progress(|p| p.try_set_answer(&answer));
929
930        if remote_description_already_applied {
931            // The 200 OK carried no body, so the early-media (183) SDP *is*
932            // the answer. It was applied provisionally (Pranswer) while
933            // ringing; re-apply it as a final Answer so the peer connection's
934            // signaling state stabilizes. A plain update would be skipped by
935            // the SDP-equality fast path, so force it.
936            self.media_stream
937                .update_remote_description_force(&track_id, &answer)
938                .await
939                .ok();
940        } else {
941            self.media_stream
942                .update_remote_description(&track_id, &answer)
943                .await
944                .ok();
945        }
946
947        Ok(answer)
948    }
949
950    /// Detect if SDP is WebRTC format
951    pub(super) fn is_webrtc_sdp(sdp: &str) -> bool {
952        (sdp.contains("a=ice-ufrag:") || sdp.contains("a=ice-pwd:"))
953            && sdp.contains("a=fingerprint:")
954    }
955
956    pub(super) async fn setup_answer_track(
957        &self,
958        option: &CallOption,
959        offer: String,
960    ) -> Result<(String, Box<dyn Track>)> {
961        let offer = match option.enable_ipv6 {
962            Some(false) | None => strip_ipv6_candidates(&offer),
963            _ => offer.clone(),
964        };
965
966        let timeout = option.handshake_timeout.map(|t| Duration::from_secs(t));
967
968        let mut media_track = if Self::is_webrtc_sdp(&offer) {
969            let mut rtc_config = RtcTrackConfig::default();
970            rtc_config.mode = rustrtc::TransportMode::WebRtc;
971            rtc_config.ice_servers = self.app_state.config.ice_servers.clone();
972            self.rtc_apply_network(&mut rtc_config);
973            self.rtc_apply_latching(&mut rtc_config);
974            restrict_codecs_to_offer(&mut rtc_config, &offer);
975
976            let webrtc_track = RtcTrack::new(
977                self.cancel_token.child_token(),
978                self.session_id.clone(),
979                self.track_config.clone(),
980                rtc_config,
981            )
982            .with_ssrc(self.ssrc);
983
984            Box::new(webrtc_track) as Box<dyn Track>
985        } else {
986            let per_call_srtp = option.sip.as_ref().and_then(|s| s.enable_srtp);
987            let rtp_track = self
988                .create_rtp_track(
989                    self.session_id.clone(),
990                    self.ssrc,
991                    per_call_srtp,
992                    Some(&offer),
993                )
994                .await?;
995            Box::new(rtp_track) as Box<dyn Track>
996        };
997
998        let answer = match media_track.handshake(offer.clone(), timeout).await {
999            Ok(answer) => answer,
1000            Err(e) => {
1001                return Err(anyhow::anyhow!("handshake failed: {e}"));
1002            }
1003        };
1004
1005        return Ok((answer, media_track));
1006    }
1007
1008    pub(super) async fn prepare_incoming_sip_track(
1009        &self,
1010        pending_dialog: PendingDialog,
1011        hangup_headers: Option<Vec<rsipstack::rsip::Header>>,
1012    ) -> Result<()> {
1013        let state_receiver = pending_dialog.state_receiver;
1014
1015        let states = InviteDialogStates::new(
1016            false,
1017            self.session_id.clone(),
1018            self.session_id.clone(),
1019            self.event_sender.clone(),
1020            self.media_stream.clone(),
1021            self.leg(),
1022            self.cancel_token.clone(),
1023            None,
1024        );
1025
1026        let initial_request = pending_dialog.dialog.initial_request();
1027        let offer = String::from_utf8_lossy(&initial_request.body).to_string();
1028
1029        let caller = initial_request
1030            .from_header()
1031            .ok()
1032            .and_then(|h| h.uri().ok())
1033            .map(|u| u.to_string())
1034            .unwrap_or_default();
1035        let callee = initial_request
1036            .to_header()
1037            .ok()
1038            .and_then(|h| h.uri().ok())
1039            .map(|u| u.to_string())
1040            .unwrap_or_default();
1041        let headers: Option<std::collections::HashMap<String, String>> = {
1042            let mut h = std::collections::HashMap::new();
1043            for header in initial_request.headers.iter() {
1044                if let rsipstack::rsip::Header::Other(name, value) = header {
1045                    h.insert(name.to_string(), value.to_string());
1046                }
1047            }
1048            if h.is_empty() { None } else { Some(h) }
1049        };
1050        self.event_sender
1051            .send(SessionEvent::Incoming {
1052                track_id: self.session_id.clone(),
1053                timestamp: crate::media::get_timestamp(),
1054                caller,
1055                callee,
1056                sdp: offer.clone(),
1057                headers,
1058            })
1059            .ok();
1060
1061        let option = self.progress.load_full().option.clone().unwrap_or_default();
1062
1063        match self.setup_answer_track(&option, offer).await {
1064            Ok((offer, track)) => {
1065                // Start the track in the media stream now — early-media ringtone
1066                // requires the RTP sender loop to be running during ringing.
1067                // Processors are intentionally omitted here; they will be built from
1068                // the accept option (which carries VAD/ASR/AGC config) when Accept
1069                // is issued, via finish_caller_stack(StartedForEarlyMedia).
1070                //
1071                // Do NOT call setup_track_with_stream here: it builds the VAD/ASR/AGC
1072                // processors from the stored option. When Accept arrives without a prior
1073                // Ringing, that option already carries `asr`, so processors would be built
1074                // both here and again in finish_caller_stack(StartedForEarlyMedia),
1075                // resulting in two ASR clients (two WebSocket connections) per session.
1076                // update_track_wrapper only starts the track (plus ambiance/subscribe),
1077                // deferring all VAD/ASR/AGC processors to accept time.
1078                self.update_track_wrapper(track, None).await;
1079                self.set_ready_to_answer(crate::call::active_call::ReadyAnswer {
1080                    answer: offer,
1081                    track: PendingCallerTrack::StartedForEarlyMedia,
1082                    dialog: pending_dialog.dialog,
1083                });
1084            }
1085            Err(e) => {
1086                return Err(anyhow::anyhow!("error creating track: {}", e));
1087            }
1088        }
1089
1090        let mut client_dialog_handler = DialogStateReceiverGuard::new(
1091            self.invitation.dialog_layer.clone(),
1092            state_receiver,
1093            hangup_headers,
1094        );
1095
1096        crate::spawn(async move {
1097            client_dialog_handler.process_dialog(states).await;
1098        });
1099        Ok(())
1100    }
1101}