1use 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
35const SIP_LEG_TRACK_ID: &str = "sip-leg-track";
40
41fn 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
70pub(super) struct OutgoingLeg {
72 pub cancel_token: CancellationToken,
74 pub leg: LegShared,
76 pub track_id: TrackId,
78 pub invite_option: InviteOption,
80 pub call_option: CallOption,
82 pub moh: Option<String>,
84 pub auto_hangup: bool,
86}
87
88impl ActiveCall {
89 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 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 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 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 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 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 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 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 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 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 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 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 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 self.try_prepare_pending_sip_answer(option).await?;
478 Some(track)
479 } else {
480 None
481 }
482 }
483 ActiveCallType::Sip | ActiveCallType::B2bua => {
484 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 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 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; 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 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 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 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 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 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 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 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}