Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 14 additions & 6 deletions src/dialog/invitation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -187,16 +187,22 @@ pub(super) struct DialogGuardForUnconfirmed<'a> {
pub dialog_layer_inner: &'a DialogLayerInnerRef,
pub id: &'a DialogId,
invite_tx: Option<Transaction>,
/// The INVITE's dialog, `None` once `process_invite` returned.
dialog: Option<InviteDialog>,
}

impl<'a> Drop for DialogGuardForUnconfirmed<'a> {
fn drop(&mut self) {
let Some(dlg) = self.dialog_layer_inner.dialogs.remove(&self.id.to_string()) else {
return;
};

let Dialog::Invite(client_dialog) = dlg else {
return;
let client_dialog = match self.dialog_layer_inner.dialogs.remove(&self.id.to_string()) {
Some(Dialog::Invite(client_dialog)) => client_dialog,
Some(_) => return,
// The application already removed the dialog from the layer
// (`DialogLayer::remove_dialog`): its INVITE still has to end.
None => match self.dialog.take() {
Some(client_dialog) => client_dialog,
// `process_invite` returned: `do_invite` handles the outcome.
None => return,
},
};

match client_dialog.state() {
Expand Down Expand Up @@ -688,6 +694,7 @@ impl DialogLayer {
dialog_layer_inner: &self.inner,
id: &id,
invite_tx: Some(tx),
dialog: Some(dialog.clone()),
};

let tx = guard
Expand All @@ -696,6 +703,7 @@ impl DialogLayer {
.expect("transcation should be avaible");

let r = dialog.process_invite(tx).boxed().await;
guard.dialog = None;
self.inner.dialogs.remove(&id.to_string());

match r {
Expand Down
81 changes: 68 additions & 13 deletions src/dialog/tests/test_cancel_2xx_race.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ use crate::dialog::{
dialog::{DialogState, DialogStateReceiver, TerminatedReason},
dialog_layer::DialogLayer,
invitation::InviteOption,
DialogId,
};
use crate::sip::{prelude::HeadersExt, Method, Request, SipMessage, Uri};
use crate::transport::{udp::UdpConnection, TransportLayer};
Expand Down Expand Up @@ -181,7 +182,19 @@ enum Order {
InviteOkLate,
}

async fn run_crossing_2xx(order: Order) -> crate::Result<()> {
/// The id the `Calling` state reports, the one the dialog is registered under.
async fn calling_id(states: &mut DialogStateReceiver) -> DialogId {
match wait_for_state(states, "Calling", Duration::from_secs(2), |s| {
matches!(s, DialogState::Calling(_))
})
.await
{
DialogState::Calling(id) => id,
_ => unreachable!(),
}
}

async fn run_crossing_2xx(order: Order, provisional: u16, remove: bool) -> crate::Result<()> {
let token = CancellationToken::new();
let Uac {
dialog_layer,
Expand All @@ -191,15 +204,27 @@ async fn run_crossing_2xx(order: Order) -> crate::Result<()> {
let wait = Duration::from_secs(2);

let (state_sender, mut states) = unbounded_channel();
let invite = tokio::spawn(async move { dialog_layer.do_invite(option, state_sender).await });
let layer = dialog_layer.clone();
let invite = tokio::spawn(async move { layer.do_invite(option, state_sender).await });

let (inv, uac) = recv_request(&peer, Method::Invite, wait).await;
reply(&peer, uac, &inv, 180, "Ringing").await;
wait_for_state(&mut states, "Early", wait, |s| {
matches!(s, DialogState::Early(_, _))
let id = calling_id(&mut states).await;
let reason = if provisional == 100 {
"Trying"
} else {
"Ringing"
};
reply(&peer, uac, &inv, provisional, reason).await;
wait_for_state(&mut states, "Trying or Early", wait, |s| {
matches!(s, DialogState::Trying(_) | DialogState::Early(_, _))
})
.await;

if remove {
// The application removes the dialog from the layer first.
dialog_layer.remove_dialog(&id);
assert!(dialog_layer.is_empty());
}
// The application abandons the call: dropping the `do_invite` future
// cancels the INVITE.
invite.abort();
Expand Down Expand Up @@ -296,17 +321,25 @@ async fn run_crossing_2xx(order: Order) -> crate::Result<()> {

#[tokio::test]
async fn test_2xx_before_cancel_response_is_acked_and_byed() -> crate::Result<()> {
run_crossing_2xx(Order::InviteOkFirst).await
run_crossing_2xx(Order::InviteOkFirst, 180, false).await
}

#[tokio::test]
async fn test_2xx_after_cancel_response_is_acked_and_byed() -> crate::Result<()> {
run_crossing_2xx(Order::CancelOkFirst).await
run_crossing_2xx(Order::CancelOkFirst, 180, false).await
}

#[tokio::test]
async fn test_2xx_after_cancel_settle_window_is_acked_and_byed() -> crate::Result<()> {
run_crossing_2xx(Order::InviteOkLate).await
run_crossing_2xx(Order::InviteOkLate, 180, false).await
}

/// The same when the application removed the dialog from the layer before
/// dropping the `do_invite` future, in Early (180) and in Trying (100).
#[tokio::test]
async fn test_removed_dialog_2xx_crossing_the_cancel_is_acked_and_byed() -> crate::Result<()> {
run_crossing_2xx(Order::InviteOkFirst, 180, true).await?;
run_crossing_2xx(Order::InviteOkFirst, 100, true).await
}

/// A CANCEL that wins the race (487 to the INVITE) ends the call as before:
Expand Down Expand Up @@ -355,7 +388,7 @@ async fn test_cancel_answered_487_sends_no_bye() -> crate::Result<()> {
/// Dropped before any response: Terminated(UacCancel) is reported at once and
/// nothing is sent while no provisional has arrived (RFC 3261 §9.1). The first
/// response then gets a CANCEL (180), or an ACK and a BYE (200).
async fn run_dropped_before_provisional(first: u16) -> crate::Result<()> {
async fn run_dropped_before_provisional(first: u16, remove: bool) -> crate::Result<()> {
let token = CancellationToken::new();
let Uac {
dialog_layer,
Expand All @@ -365,8 +398,15 @@ async fn run_dropped_before_provisional(first: u16) -> crate::Result<()> {
let wait = Duration::from_secs(2);

let (state_sender, mut states) = unbounded_channel();
let invite = tokio::spawn(async move { dialog_layer.do_invite(option, state_sender).await });
let layer = dialog_layer.clone();
let invite = tokio::spawn(async move { layer.do_invite(option, state_sender).await });
let (inv, uac) = recv_request(&peer, Method::Invite, wait).await;
let id = calling_id(&mut states).await;
if remove {
// The application removes the dialog from the layer first.
dialog_layer.remove_dialog(&id);
assert!(dialog_layer.is_empty());
}
invite.abort();
let _ = invite.await;
let terminated = wait_for_state(&mut states, "Terminated", Duration::from_millis(200), |s| {
Expand Down Expand Up @@ -418,15 +458,30 @@ async fn run_dropped_before_provisional(first: u16) -> crate::Result<()> {

#[tokio::test]
async fn test_dropped_before_provisional_is_cancelled_after_the_180() -> crate::Result<()> {
run_dropped_before_provisional(180).await
run_dropped_before_provisional(180, false).await
}

#[tokio::test]
async fn test_dropped_before_provisional_2xx_is_acked_and_byed() -> crate::Result<()> {
run_dropped_before_provisional(200).await
run_dropped_before_provisional(200, false).await
}

#[tokio::test]
async fn test_removed_before_provisional_is_cancelled_after_the_180() -> crate::Result<()> {
run_dropped_before_provisional(180, true).await
}

#[tokio::test]
async fn test_removed_before_provisional_2xx_is_acked_and_byed() -> crate::Result<()> {
run_dropped_before_provisional(200, true).await
}

#[tokio::test]
async fn test_dropped_before_provisional_final_failure_is_acked_only() -> crate::Result<()> {
run_dropped_before_provisional(486).await
run_dropped_before_provisional(486, false).await
}

#[tokio::test]
async fn test_removed_before_provisional_final_failure_is_acked_only() -> crate::Result<()> {
run_dropped_before_provisional(486, true).await
}
Loading