1use crate::call::state::{CallProgress, LegShared};
2use crate::callrecord::CallRecordHangupReason;
3use crate::event::EventSender;
4use crate::media::TrackId;
5use crate::media::stream::MediaStream;
6use crate::useragent::invitation::PendingDialog;
7use anyhow::Result;
8use chrono::Utc;
9use rsipstack::dialog::DialogId;
10use rsipstack::dialog::dialog::{
11 Dialog, DialogState, DialogStateReceiver, DialogStateSender, TerminatedReason,
12};
13use rsipstack::dialog::dialog_layer::DialogLayer;
14use rsipstack::dialog::invitation::InviteOption;
15use std::collections::HashMap;
16use std::sync::Arc;
17use tokio_util::sync::CancellationToken;
18use tracing::{info, warn};
19
20pub(crate) fn remove_dialog(layer: &DialogLayer, id: &DialogId) -> Option<Dialog> {
23 let dialog = layer.get_dialog(id)?;
24 layer.remove_dialog(id);
25 Some(dialog)
26}
27
28pub struct DialogStateReceiverGuard {
29 pub(super) dialog_layer: Arc<DialogLayer>,
30 pub(super) receiver: DialogStateReceiver,
31 pub(super) dialog_id: Option<DialogId>,
32 pub(super) hangup_headers: Option<Vec<rsipstack::rsip::Header>>,
33}
34
35impl DialogStateReceiverGuard {
36 pub fn new(
37 dialog_layer: Arc<DialogLayer>,
38 receiver: DialogStateReceiver,
39 hangup_headers: Option<Vec<rsipstack::rsip::Header>>,
40 ) -> Self {
41 Self {
42 dialog_layer,
43 receiver,
44 dialog_id: None,
45 hangup_headers,
46 }
47 }
48 pub async fn recv(&mut self) -> Option<DialogState> {
49 let state = self.receiver.recv().await;
50 if let Some(ref s) = state {
51 self.dialog_id = Some(s.id().clone());
52 }
53 state
54 }
55
56 fn take_dialog(&mut self) -> Option<Dialog> {
57 let id = self.dialog_id.take()?;
58 info!(%id, "dialog removed on drop");
59 if let Some(dialog) = remove_dialog(&self.dialog_layer, &id) {
60 return Some(dialog);
61 }
62 self.dialog_layer
70 .get_client_dialog_by_call_id(&id.call_id)
71 .into_iter()
72 .next()
73 .map(Dialog::Invite)
74 }
75
76 pub async fn drop_async(&mut self) {
77 if let Some(dialog) = self.take_dialog() {
78 if let Err(e) = dialog.hangup_with_headers(self.hangup_headers.take()).await {
79 warn!(id=%dialog.id(), "error hanging up dialog on drop: {}", e);
80 }
81 }
82 }
83}
84
85impl Drop for DialogStateReceiverGuard {
86 fn drop(&mut self) {
87 if let Some(dialog) = self.take_dialog() {
88 crate::spawn(async move {
89 if let Err(e) = dialog.hangup().await {
90 warn!(id=%dialog.id(), "error hanging up dialog on drop: {}", e);
91 }
92 });
93 }
94 }
95}
96
97pub(super) struct InviteDialogStates {
98 pub is_client: bool,
99 pub session_id: String,
100 pub track_id: TrackId,
101 pub cancel_token: CancellationToken,
102 pub event_sender: EventSender,
103 pub leg: LegShared,
105 pub media_stream: Arc<MediaStream>,
106 pub terminated_reason: Option<TerminatedReason>,
107 pub has_early_media: bool,
108 initial_confirmed: bool,
116 pub hangup_reason: Option<CallRecordHangupReason>,
119}
120
121impl InviteDialogStates {
122 pub(super) fn new(
123 is_client: bool,
124 session_id: String,
125 track_id: TrackId,
126 event_sender: EventSender,
127 media_stream: Arc<MediaStream>,
128 leg: LegShared,
129 cancel_token: CancellationToken,
130 hangup_reason: Option<CallRecordHangupReason>,
131 ) -> Self {
132 Self {
133 is_client,
134 session_id,
135 track_id,
136 cancel_token,
137 event_sender,
138 leg,
139 media_stream,
140 terminated_reason: None,
141 has_early_media: false,
142 initial_confirmed: false,
143 hangup_reason,
144 }
145 }
146}
147
148impl InviteDialogStates {
149 pub(super) fn on_terminated(&mut self) {
153 let term = CallProgress::termination(self.terminated_reason.as_ref());
154 let status_code = term.status_code;
155 let reason = term.hangup_reason;
156 self.leg.update_progress(|p| {
157 p.last_status_code = status_code;
158 p.set_hangup_reason(reason.clone());
159 });
160 let progress = self.leg.progress.load_full();
161
162 self.event_sender
163 .send(crate::event::SessionEvent::TrackEnd {
164 track_id: self.track_id.clone(),
165 timestamp: crate::media::get_timestamp(),
166 duration: progress
167 .answer_time
168 .map(|t| (Utc::now() - t).num_milliseconds())
169 .unwrap_or_default() as u64,
170 ssrc: self.leg.ssrc,
171 play_id: None,
172 auto_hangup: self.hangup_reason.clone(),
173 })
174 .ok();
175 let hangup_event = self
176 .leg
177 .build_hangup_event(self.track_id.clone(), Some(term.initiator.to_string()));
178 self.event_sender.send(hangup_event).ok();
179 }
180}
181
182impl Drop for InviteDialogStates {
183 fn drop(&mut self) {
184 self.on_terminated();
185 self.cancel_token.cancel();
186 }
187}
188
189impl DialogStateReceiverGuard {
190 pub(self) async fn dialog_event_loop(&mut self, states: &mut InviteDialogStates) -> Result<()> {
191 while let Some(event) = self.recv().await {
192 match event {
193 DialogState::Calling(dialog_id) => {
194 info!(session_id=states.session_id, %dialog_id, "dialog calling");
195 states
196 .leg
197 .update_progress(|p| p.session_id = dialog_id.to_string());
198 }
199 DialogState::Trying(_) => {}
200 DialogState::Early(dialog_id, resp) => {
201 let code = resp.status_code.code();
202 let body = resp.body();
203 let answer = String::from_utf8_lossy(body);
204 let has_sdp = !answer.is_empty();
205 info!(session_id=states.session_id, %dialog_id, has_sdp=%has_sdp, "dialog early ({}): \n{}", code, answer);
206
207 states.leg.update_progress(|p| p.on_early(code));
208
209 if !states.is_client {
210 continue;
211 }
212
213 states
214 .event_sender
215 .send(crate::event::SessionEvent::Ringing {
216 track_id: states.track_id.clone(),
217 timestamp: crate::media::get_timestamp(),
218 early_media: has_sdp,
219 refer: Some(states.leg.is_refer),
220 })?;
221
222 if has_sdp {
223 states.has_early_media = true;
224 states.leg.update_progress(|p| p.try_set_answer(&answer));
225 states
226 .media_stream
227 .update_remote_description_provisional(
228 &states.track_id,
229 &answer.to_string(),
230 )
231 .await?;
232 }
233 }
234 DialogState::Confirmed(dialog_id, msg) => {
235 info!(session_id=states.session_id, %dialog_id, has_early_media=%states.has_early_media, "dialog confirmed");
236 states
237 .leg
238 .update_progress(|p| p.on_confirmed(dialog_id.to_string()));
239 if states.is_client && !states.initial_confirmed {
240 states.initial_confirmed = true;
241 let answer = String::from_utf8_lossy(msg.body());
242 let answer = answer.trim();
243 if !answer.is_empty() {
244 if states.has_early_media {
245 info!(
246 session_id = states.session_id,
247 "updating remote description with final answer after early media (force=true)"
248 );
249 if let Err(e) = states
252 .media_stream
253 .update_remote_description_force(
254 &states.track_id,
255 &answer.to_string(),
256 )
257 .await
258 {
259 tracing::warn!(
260 session_id = states.session_id,
261 "failed to force update remote description on confirmed: {}",
262 e
263 );
264 }
265 } else {
266 if let Err(e) = states
267 .media_stream
268 .update_remote_description(
269 &states.track_id,
270 &answer.to_string(),
271 )
272 .await
273 {
274 tracing::warn!(
275 session_id = states.session_id,
276 "failed to update remote description on confirmed: {}",
277 e
278 );
279 }
280 }
281 }
282 }
283 }
284 DialogState::Info(dialog_id, req, tx_handle) => {
285 let body_str = String::from_utf8_lossy(req.body());
286 info!(session_id=states.session_id, %dialog_id, body=%body_str, "dialog info received");
287 if body_str.starts_with("Signal=") {
288 let digit = body_str.trim_start_matches("Signal=").chars().next();
289 if let Some(digit) = digit {
290 states.event_sender.send(crate::event::SessionEvent::Dtmf {
291 track_id: states.track_id.clone(),
292 timestamp: crate::media::get_timestamp(),
293 digit: digit.to_string(),
294 refer: Some(states.leg.is_refer),
295 })?;
296 }
297 }
298 tx_handle.reply(rsipstack::rsip::StatusCode::OK).await.ok();
299 }
300 DialogState::Message(dialog_id, req, tx_handle) => {
301 let body_str = String::from_utf8_lossy(req.body()).to_string();
302 let content_type = req.headers.iter().find_map(|h| {
303 if let rsipstack::rsip::Header::ContentType(content_type) = h {
304 Some(content_type.value().to_string())
305 } else {
306 None
307 }
308 });
309 info!(
310 session_id=states.session_id,
311 %dialog_id,
312 content_type=content_type.as_deref(),
313 body=%body_str,
314 "dialog message received"
315 );
316 states
317 .event_sender
318 .send(crate::event::SessionEvent::Message {
319 track_id: states.track_id.clone(),
320 timestamp: crate::media::get_timestamp(),
321 body: body_str,
322 content_type,
323 refer: Some(states.leg.is_refer),
324 })
325 .ok();
326 tx_handle.reply(rsipstack::rsip::StatusCode::OK).await.ok();
327 }
328 DialogState::Updated(dialog_id, _req, tx_handle) => {
329 info!(session_id = states.session_id, %dialog_id, "dialog update received");
330 let mut answer_sdp = None;
331 if let Some(sdp_body) = _req.body().get(..) {
332 let sdp_str = String::from_utf8_lossy(sdp_body);
333 if !sdp_str.is_empty()
334 && (_req.method == rsipstack::rsip::Method::Invite
335 || _req.method == rsipstack::rsip::Method::Update)
336 {
337 info!(session_id=states.session_id, %dialog_id, method=%_req.method, "handling re-invite/update offer");
338
339 let is_on_hold =
341 crate::media::negotiate::detect_hold_state_from_sdp(&sdp_str);
342 info!(session_id=states.session_id, %dialog_id, is_on_hold=%is_on_hold, "detected hold state from re-invite SDP");
343
344 apply_hold_state(states, is_on_hold).await;
346
347 match states
348 .media_stream
349 .handshake(&states.track_id, sdp_str.to_string(), None)
350 .await
351 {
352 Ok(sdp) => answer_sdp = Some(sdp),
353 Err(e) => {
354 warn!(
355 session_id = states.session_id,
356 "failed to handle re-invite: {}", e
357 );
358 }
359 }
360 } else {
361 info!(session_id=states.session_id, %dialog_id, "updating remote description:\n{}", sdp_str);
362
363 let is_on_hold =
365 crate::media::negotiate::detect_hold_state_from_sdp(&sdp_str);
366 apply_hold_state(states, is_on_hold).await;
367
368 states
369 .media_stream
370 .update_remote_description(&states.track_id, &sdp_str.to_string())
371 .await?;
372 }
373 }
374
375 if let Some(sdp) = answer_sdp {
376 tx_handle
377 .respond(
378 rsipstack::rsip::StatusCode::OK,
379 Some(vec![rsipstack::rsip::Header::ContentType(
380 "application/sdp".to_string().into(),
381 )]),
382 Some(sdp.into()),
383 )
384 .await
385 .ok();
386 } else {
387 tx_handle.reply(rsipstack::rsip::StatusCode::OK).await.ok();
388 }
389 }
390 DialogState::Options(dialog_id, _req, tx_handle) => {
391 info!(session_id = states.session_id, %dialog_id, "dialog options received");
392 tx_handle.reply(rsipstack::rsip::StatusCode::OK).await.ok();
393 }
394 DialogState::Refer(dialog_id, req, tx_handle) => {
395 let refer_to = req
396 .headers
397 .iter()
398 .find_map(|h| {
399 if let rsipstack::rsip::Header::ReferTo(h) = h {
400 return Some(h.value().to_string());
401 }
402 None
403 })
404 .unwrap_or_default();
405 let referred_by = req.headers.iter().find_map(|h| {
406 if let rsipstack::rsip::Header::ReferredBy(h) = h {
407 return Some(h.value().to_string());
408 }
409 None
410 });
411 info!(session_id = states.session_id, %dialog_id, %refer_to, "received REFER");
412 tx_handle
413 .reply(rsipstack::rsip::StatusCode::Other(202, "Accepted".into()))
414 .await
415 .ok();
416 states
417 .event_sender
418 .send(crate::event::SessionEvent::TransferRequest {
419 track_id: states.track_id.clone(),
420 timestamp: crate::media::get_timestamp(),
421 refer_to,
422 referred_by,
423 refer: Some(states.leg.is_refer),
424 })
425 .ok();
426 }
427 DialogState::Terminated(dialog_id, reason) => {
428 info!(
429 session_id = states.session_id,
430 ?dialog_id,
431 ?reason,
432 "dialog terminated"
433 );
434 states.terminated_reason = Some(reason.clone());
435 return Ok(());
436 }
437 other_state => {
438 info!(
439 session_id = states.session_id,
440 %other_state,
441 "dialog received other state"
442 );
443 }
444 }
445 }
446 Ok(())
447 }
448
449 pub(super) async fn process_dialog(&mut self, mut states: InviteDialogStates) {
450 let token = states.cancel_token.clone();
451 tokio::select! {
452 _ = token.cancelled() => {
453 states.terminated_reason = Some(TerminatedReason::UacCancel);
454 }
455 _ = self.dialog_event_loop(&mut states) => {}
456 };
457
458 let extras = states.leg.extras.load_full();
460 if let Some(headers) = crate::sip_util::hangup_headers_from_extras(&extras) {
461 match &mut self.hangup_headers {
462 Some(existing) => existing.extend(headers),
463 None => self.hangup_headers = Some(headers),
464 }
465 }
466
467 self.drop_async().await;
468 }
469}
470
471async fn apply_hold_state(states: &mut InviteDialogStates, is_on_hold: bool) {
473 if is_on_hold {
474 states
475 .media_stream
476 .hold_track(Some(states.track_id.clone()))
477 .await;
478 } else {
479 states
480 .media_stream
481 .resume_track(Some(states.track_id.clone()))
482 .await;
483 }
484 states
485 .event_sender
486 .send(crate::event::SessionEvent::Hold {
487 track_id: states.track_id.clone(),
488 timestamp: crate::media::get_timestamp(),
489 on_hold: is_on_hold,
490 refer: Some(states.leg.is_refer),
491 })
492 .ok();
493}
494
495#[derive(Clone)]
496pub struct Invitation {
497 pub dialog_layer: Arc<DialogLayer>,
498 pub pending_dialogs: Arc<std::sync::Mutex<HashMap<DialogId, PendingDialog>>>,
499 sessions: Arc<std::sync::Mutex<HashMap<String, DialogId>>>,
503}
504
505impl Invitation {
506 pub fn new(dialog_layer: Arc<DialogLayer>) -> Self {
507 Self {
508 dialog_layer,
509 pending_dialogs: Arc::new(std::sync::Mutex::new(HashMap::new())),
510 sessions: Arc::new(std::sync::Mutex::new(HashMap::new())),
511 }
512 }
513
514 pub fn register_session(&self, session_id: &str, dialog_id: &DialogId) {
515 self.sessions
516 .lock()
517 .map(|mut ss| ss.insert(session_id.to_string(), dialog_id.clone()))
518 .ok();
519 }
520
521 pub fn unregister_session(&self, session_id: &str) {
522 self.sessions.lock().map(|mut ss| ss.remove(session_id)).ok();
523 }
524
525 pub fn session_exists(&self, session_id: &str) -> bool {
526 self.sessions
527 .lock()
528 .map(|ss| ss.contains_key(session_id))
529 .unwrap_or(false)
530 }
531
532 pub fn add_pending(&self, dialog_id: DialogId, pending: PendingDialog) {
533 self.pending_dialogs
534 .lock()
535 .map(|mut ps| ps.insert(dialog_id, pending))
536 .ok();
537 }
538
539 pub fn get_pending_call(&self, dialog_id: &DialogId) -> Option<PendingDialog> {
540 self.pending_dialogs
541 .lock()
542 .ok()
543 .and_then(|mut ps| ps.remove(dialog_id))
544 }
545
546 pub fn has_pending_call(&self, dialog_id: &DialogId) -> bool {
547 self.pending_dialogs
548 .lock()
549 .ok()
550 .map(|ps| ps.contains_key(dialog_id))
551 .unwrap_or(false)
552 }
553
554 pub fn find_dialog_id_by_session_id(&self, session_id: &str) -> Option<DialogId> {
555 if let Some(id) = self
556 .sessions
557 .lock()
558 .ok()
559 .and_then(|ss| ss.get(session_id).cloned())
560 {
561 return Some(id);
562 }
563 self.pending_dialogs.lock().ok().and_then(|ps| {
564 ps.iter()
565 .find(|(id, _)| id.to_string() == session_id)
566 .map(|(id, _)| id.clone())
567 })
568 }
569
570 pub async fn hangup(
572 &self,
573 dialog_id: DialogId,
574 code: Option<rsipstack::rsip::StatusCode>,
575 reason: Option<String>,
576 ) -> Result<()> {
577 if let Some(call) = self.get_pending_call(&dialog_id) {
578 call.dialog.reject(code, reason).ok();
579 }
580 if let Some(dialog) = remove_dialog(&self.dialog_layer, &dialog_id) {
581 dialog.hangup().await.ok();
582 }
583 Ok(())
584 }
585
586 pub async fn invite(
587 &self,
588 invite_option: InviteOption,
589 state_sender: DialogStateSender,
590 ) -> Result<(DialogId, Option<Vec<u8>>), rsipstack::Error> {
591 let (dialog, resp) = self
592 .dialog_layer
593 .do_invite(invite_option, state_sender)
594 .await?;
595
596 let offer = match resp {
597 Some(resp) => match resp.status_code.kind() {
598 rsipstack::rsip::StatusCodeKind::Successful => {
599 let offer = resp.body.clone();
600 Some(offer)
601 }
602 _ => {
603 let reason = resp
604 .reason_phrase()
605 .unwrap_or(&resp.status_code.to_string())
606 .to_string();
607 return Err(rsipstack::Error::DialogError(
608 reason,
609 dialog.id(),
610 resp.status_code,
611 ));
612 }
613 },
614 None => {
615 return Err(rsipstack::Error::DialogError(
616 "no response received".to_string(),
617 dialog.id(),
618 rsipstack::rsip::StatusCode::NotAcceptableHere,
619 ));
620 }
621 };
622 Ok((dialog.id(), offer))
623 }
624}
625
626#[cfg(test)]
627mod tests {
628 use super::*;
629 use crate::call::state::CallProgress;
630 use crate::media::stream::MediaStreamBuilder;
631
632 const EARLY_MEDIA_SDP: &str = "v=0\r\n\
634 o=- 1000 1 IN IP4 192.168.1.100\r\n\
635 s=SIP Call\r\n\
636 t=0 0\r\n\
637 m=audio 10000 RTP/AVP 0\r\n\
638 c=IN IP4 192.168.1.100\r\n\
639 a=rtpmap:0 PCMU/8000\r\n\
640 a=sendrecv\r\n";
641
642 fn make_states(has_early_media: bool) -> InviteDialogStates {
643 let (event_tx, _event_rx) = tokio::sync::broadcast::channel(16);
644 let media_stream = Arc::new(
645 MediaStreamBuilder::new(event_tx.clone())
646 .with_id("test-stream".to_string())
647 .build(),
648 );
649 let cancel_token = CancellationToken::new();
650 let leg = LegShared::new(1000, true, CallProgress::default());
651
652 InviteDialogStates {
653 is_client: true,
654 session_id: "test-session".to_string(),
655 track_id: "test-track".to_string(),
656 cancel_token: cancel_token.clone(),
657 event_sender: event_tx.clone(),
658 leg,
659 media_stream,
660 terminated_reason: None,
661 has_early_media,
662 initial_confirmed: false,
663 hangup_reason: None,
664 }
665 }
666
667 #[tokio::test]
671 async fn test_early_sdp_stored_in_leg_progress() {
672 let mut states = make_states(false);
673
674 let answer = EARLY_MEDIA_SDP.to_string();
677 let has_sdp = !answer.is_empty();
678 if states.is_client && has_sdp {
679 states.has_early_media = true;
680 states.leg.update_progress(|p| p.try_set_answer(&answer));
681 }
682
683 let progress = states.leg.progress.load_full();
685 assert!(
686 progress.answer.is_some(),
687 "leg progress answer should be set after 183 with SDP"
688 );
689 assert_eq!(
690 progress.answer.as_deref().unwrap(),
691 EARLY_MEDIA_SDP,
692 "leg progress answer should contain the early SDP"
693 );
694 assert!(states.has_early_media, "has_early_media should be true");
695 }
696
697 #[tokio::test]
704 async fn test_confirmed_empty_body_keeps_early_sdp() {
705 let mut states = make_states(false);
706
707 states.has_early_media = true;
709 states
710 .leg
711 .update_progress(|p| p.try_set_answer(EARLY_MEDIA_SDP));
712
713 states
715 .leg
716 .update_progress(|p| p.on_confirmed("dialog-1".to_string()));
717 let confirmed_answer = String::new();
720 assert!(
721 confirmed_answer.trim().is_empty(),
722 "empty body must not be applied"
723 );
724
725 let progress = states.leg.progress.load_full();
727 assert!(
728 progress.answer.is_some(),
729 "answer must not be None after 200 OK with empty body"
730 );
731 let stored_answer = progress.answer.as_deref().unwrap();
732 assert!(
733 !stored_answer.is_empty(),
734 "answer must not be empty after 200 OK with empty body"
735 );
736 assert_eq!(
737 stored_answer, EARLY_MEDIA_SDP,
738 "answer should still be the early SDP after 200 OK with empty body"
739 );
740 }
741
742 #[tokio::test]
746 async fn test_answer_fallback_to_early_sdp_when_200ok_empty() {
747 let states = make_states(true);
749 states
750 .leg
751 .update_progress(|p| p.try_set_answer(EARLY_MEDIA_SDP));
752
753 let early = states.leg.progress.load_full().answer.clone();
755 let raw_answer: Option<Vec<u8>> = Some(vec![]); let (answer, already_applied) =
758 crate::call::state::resolve_final_answer(raw_answer, early.as_ref()).unwrap();
759
760 assert!(
763 !answer.is_empty(),
764 "Resolved answer must not be empty — should contain the early SDP"
765 );
766 assert_eq!(
767 answer, EARLY_MEDIA_SDP,
768 "Resolved answer should be the early SDP from the 183 handler"
769 );
770 assert!(
771 already_applied,
772 "remote_description_already_applied should be true when using early SDP fallback"
773 );
774 }
775
776 #[tokio::test]
780 async fn test_answer_uses_200ok_sdp_when_present() {
781 const FINAL_SDP: &str = "v=0\r\n\
782 o=- 2000 2 IN IP4 10.0.0.1\r\n\
783 s=SIP Call\r\n\
784 t=0 0\r\n\
785 m=audio 20000 RTP/AVP 0\r\n\
786 c=IN IP4 10.0.0.1\r\n\
787 a=rtpmap:0 PCMU/8000\r\n\
788 a=sendrecv\r\n";
789
790 let states = make_states(true);
791 states
792 .leg
793 .update_progress(|p| p.try_set_answer(EARLY_MEDIA_SDP));
794
795 let early = states.leg.progress.load_full().answer.clone();
796 let (answer, already_applied) = crate::call::state::resolve_final_answer(
797 Some(FINAL_SDP.as_bytes().to_vec()),
798 early.as_ref(),
799 )
800 .unwrap();
801
802 assert_eq!(
803 answer, FINAL_SDP,
804 "When 200 OK has SDP, it should be used (not the early SDP)"
805 );
806 assert!(
807 !already_applied,
808 "remote_description_already_applied should be false when 200 OK has SDP body"
809 );
810 }
811
812 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
819 async fn test_on_terminated_always_emits_events_and_records_state() {
820 use crate::event::SessionEvent;
821 use std::sync::atomic::{AtomicBool, Ordering};
822
823 let mut states = make_states(false);
824 states.terminated_reason = Some(TerminatedReason::UacCancel);
825 let leg = states.leg.clone();
826 let event_sender = states.event_sender.clone();
827 let mut event_receiver = event_sender.subscribe();
828
829 let stop = Arc::new(AtomicBool::new(false));
831 let stop2 = stop.clone();
832 let progress = leg.progress.clone();
833 let writer = crate::spawn(async move {
834 while !stop2.load(Ordering::Relaxed) {
835 progress.rcu(|p| {
838 let mut p = CallProgress::clone(p);
839 p.answer_time.get_or_insert_with(chrono::Utc::now);
840 p
841 });
842 tokio::task::yield_now().await;
843 }
844 });
845
846 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
848
849 drop(states);
851
852 stop.store(true, Ordering::Relaxed);
853 let _ = writer.await;
854
855 let mut saw_track_end = false;
857 let mut saw_hangup = false;
858 while let Ok(event) = event_receiver.try_recv() {
859 match event {
860 SessionEvent::TrackEnd { .. } => saw_track_end = true,
861 SessionEvent::Hangup { refer, .. } => {
862 assert_eq!(refer, Some(true), "refer flag comes from the leg");
863 saw_hangup = true;
864 }
865 _ => {}
866 }
867 }
868 assert!(
869 saw_track_end,
870 "TrackEnd must be emitted from on_terminated even under contention"
871 );
872 assert!(
873 saw_hangup,
874 "Hangup must be emitted from on_terminated even under contention"
875 );
876
877 let progress = leg.progress.load_full();
880 assert_eq!(progress.last_status_code, 487);
881 assert_eq!(
882 progress.hangup_reason,
883 Some(crate::callrecord::CallRecordHangupReason::Canceled)
884 );
885 }
886
887 struct CountingTrack {
890 id: TrackId,
891 config: crate::media::track::TrackConfig,
892 processor_chain: crate::media::processor::ProcessorChain,
893 updates: std::sync::Arc<std::sync::atomic::AtomicUsize>,
894 }
895
896 #[async_trait::async_trait]
897 impl crate::media::track::Track for CountingTrack {
898 fn ssrc(&self) -> u32 {
899 0
900 }
901 fn id(&self) -> &TrackId {
902 &self.id
903 }
904 fn config(&self) -> &crate::media::track::TrackConfig {
905 &self.config
906 }
907 fn processor_chain(&mut self) -> &mut crate::media::processor::ProcessorChain {
908 &mut self.processor_chain
909 }
910 async fn handshake(
911 &mut self,
912 _offer: String,
913 _timeout: Option<tokio::time::Duration>,
914 ) -> Result<String> {
915 Ok(String::new())
916 }
917 async fn update_remote_description(&mut self, _answer: &String) -> Result<()> {
918 self.updates
919 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
920 Ok(())
921 }
922 async fn update_remote_description_force(&mut self, _answer: &String) -> Result<()> {
923 self.updates
924 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
925 Ok(())
926 }
927 async fn start(
928 &mut self,
929 _event_sender: EventSender,
930 _packet_sender: crate::media::track::TrackPacketSender,
931 ) -> Result<()> {
932 Ok(())
933 }
934 async fn stop(&self) -> Result<()> {
935 Ok(())
936 }
937 async fn send_packet(&mut self, _packet: &crate::media::AudioFrame) -> Result<()> {
938 Ok(())
939 }
940 }
941
942 #[tokio::test]
948 async fn test_confirmed_applies_remote_answer_only_once() {
949 let mut states = make_states(false);
950
951 let updates = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
952 let track = CountingTrack {
953 id: states.track_id.clone(),
954 config: crate::media::track::TrackConfig::default(),
955 processor_chain: crate::media::processor::ProcessorChain::new(16000),
956 updates: updates.clone(),
957 };
958 states
959 .media_stream
960 .update_track(Box::new(track), None)
961 .await;
962
963 let endpoint = {
964 let mut builder = rsipstack::EndpointBuilder::new();
965 builder.build()
966 };
967 let dialog_layer = Arc::new(DialogLayer::new(endpoint.inner.clone()));
968
969 let (tx, rx) = tokio::sync::mpsc::unbounded_channel::<DialogState>();
970 let mut guard = DialogStateReceiverGuard::new(dialog_layer, rx, None);
971
972 let dialog_id = DialogId {
973 call_id: "test-call-id".to_string(),
974 local_tag: "test-local-tag".to_string(),
975 remote_tag: "test-remote-tag".to_string(),
976 };
977
978 let mut resp = rsipstack::rsip::Response::default();
979 resp.body = EARLY_MEDIA_SDP.as_bytes().to_vec();
980
981 tx.send(DialogState::Confirmed(dialog_id.clone(), resp.clone()))
983 .unwrap();
984 tx.send(DialogState::Confirmed(dialog_id, resp)).unwrap();
985 drop(tx);
986
987 guard
988 .dialog_event_loop(&mut states)
989 .await
990 .expect("dialog event loop must complete");
991
992 assert_eq!(
993 updates.load(std::sync::atomic::Ordering::SeqCst),
994 1,
995 "remote answer must be applied exactly once, not on the re-INVITE ACK Confirmed event"
996 );
997 }
998
999 #[test]
1000 fn test_session_map_register_find_unregister() {
1001 let endpoint = {
1002 let mut builder = rsipstack::EndpointBuilder::new();
1003 builder.build()
1004 };
1005 let invitation = Invitation::new(Arc::new(DialogLayer::new(endpoint.inner.clone())));
1006
1007 let dialog_id = DialogId {
1008 call_id: "long-call-id-from-carrier".to_string(),
1009 local_tag: String::new(),
1010 remote_tag: "remote-tag".to_string(),
1011 };
1012
1013 assert!(!invitation.session_exists("s.3f9a2b1c4d5e"));
1014 assert!(
1015 invitation
1016 .find_dialog_id_by_session_id("s.3f9a2b1c4d5e")
1017 .is_none()
1018 );
1019
1020 invitation.register_session("s.3f9a2b1c4d5e", &dialog_id);
1021 assert!(invitation.session_exists("s.3f9a2b1c4d5e"));
1022 assert_eq!(
1023 invitation
1024 .find_dialog_id_by_session_id("s.3f9a2b1c4d5e")
1025 .unwrap(),
1026 dialog_id
1027 );
1028
1029 invitation.unregister_session("s.3f9a2b1c4d5e");
1030 assert!(!invitation.session_exists("s.3f9a2b1c4d5e"));
1031 assert!(
1032 invitation
1033 .find_dialog_id_by_session_id("s.3f9a2b1c4d5e")
1034 .is_none()
1035 );
1036 invitation.unregister_session("s.unknown");
1038 }
1039
1040 #[test]
1044 fn test_find_dialog_id_falls_back_to_pending_scan() {
1045 use rsipstack::dialog::dialog::DialogInner;
1046 use rsipstack::dialog::invite_dialog::InviteDialog;
1047 use rsipstack::rsip::typed::{Contact, CSeq, From, To, Via};
1048 use rsipstack::rsip::{Header, Request};
1049 use rsipstack::transaction::key::TransactionRole;
1050
1051 let endpoint = {
1052 let mut builder = rsipstack::EndpointBuilder::new();
1053 builder.build()
1054 };
1055 let invitation = Invitation::new(Arc::new(DialogLayer::new(endpoint.inner.clone())));
1056
1057 let dialog_id = DialogId {
1058 call_id: "legacy-call-id".to_string(),
1059 local_tag: String::new(),
1060 remote_tag: "remote-tag".to_string(),
1061 };
1062
1063 let initial_request = Request {
1064 method: rsipstack::rsip::Method::Invite,
1065 uri: rsipstack::rsip::Uri::try_from("sip:bob@example.com:5060").unwrap(),
1066 headers: vec![
1067 Via::parse("SIP/2.0/UDP alice.example.com:5060;branch=z9hG4bKnashds")
1068 .unwrap()
1069 .into(),
1070 CSeq::parse("1 INVITE").unwrap().into(),
1071 From::parse("Alice <sip:alice@example.com>;tag=remote-tag")
1072 .unwrap()
1073 .into(),
1074 To::parse("Bob <sip:bob@example.com>").unwrap().into(),
1075 Header::CallId("legacy-call-id".into()),
1076 Contact::parse("<sip:alice@alice.example.com:5060>").unwrap().into(),
1077 Header::MaxForwards("70".into()),
1078 ]
1079 .into(),
1080 version: rsipstack::rsip::Version::V2,
1081 body: vec![],
1082 };
1083
1084 let (state_sender, _state_receiver) = tokio::sync::mpsc::unbounded_channel();
1085 let (tu_sender, _tu_receiver) = tokio::sync::mpsc::unbounded_channel();
1086 let inner = std::sync::Arc::new(
1087 DialogInner::new(
1088 TransactionRole::Server,
1089 dialog_id.clone(),
1090 initial_request,
1091 endpoint.inner.clone(),
1092 state_sender,
1093 None,
1094 None,
1095 tu_sender,
1096 )
1097 .expect("failed to create dialog inner"),
1098 );
1099 let dialog = InviteDialog::from_inner(inner);
1100
1101 assert!(
1102 invitation
1103 .find_dialog_id_by_session_id(&dialog_id.to_string())
1104 .is_none()
1105 );
1106
1107 invitation.add_pending(
1108 dialog_id.clone(),
1109 PendingDialog {
1110 token: tokio_util::sync::CancellationToken::new(),
1111 dialog,
1112 state_receiver: {
1113 let (_tx, rx) = tokio::sync::mpsc::unbounded_channel();
1114 rx
1115 },
1116 },
1117 );
1118
1119 assert_eq!(
1120 invitation
1121 .find_dialog_id_by_session_id(&dialog_id.to_string())
1122 .unwrap(),
1123 dialog_id
1124 );
1125 }
1126}