From 9cf1a0dad32e88ac0d125ed46c2a1b98cb54d83a Mon Sep 17 00:00:00 2001 From: Tyson George Date: Fri, 25 Sep 2026 14:10:14 -0400 Subject: [PATCH] fix: ACK and BYE a 2xx that crosses the CANCEL of a dropped INVITE MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Dropping the do_invite future while the call is ringing cancels the INVITE. The guard waited up to 2 s for the INVITE's final response and then dropped the transaction, but a 2xx found in that window was only logged, and a 2xx arriving later was ACKed by the detached transaction without anything ending the session. Either way the callee kept a confirmed dialog on dead air until it hung up. A CANCEL that crosses a 2xx has no effect on the INVITE; the UAC must ACK the 2xx and send a BYE (RFC 3261 §9.1, §15; RFC 5407 §3.1.2). Keep the INVITE transaction after the settle window until its final response (up to 64*T1) so a 2xx is ACKed, and send a BYE built from the 2xx's To tag, route set and Contact. Terminated(UacCancel) is still reported at the same point, and no Confirmed state is emitted for the abandoned call. --- src/dialog/invitation.rs | 54 +++- src/dialog/invite_dialog.rs | 45 ++- src/dialog/tests/mod.rs | 1 + src/dialog/tests/test_cancel_2xx_race.rs | 353 +++++++++++++++++++++++ 4 files changed, 437 insertions(+), 16 deletions(-) create mode 100644 src/dialog/tests/test_cancel_2xx_race.rs diff --git a/src/dialog/invitation.rs b/src/dialog/invitation.rs index c8366f65..8017680c 100644 --- a/src/dialog/invitation.rs +++ b/src/dialog/invitation.rs @@ -221,11 +221,15 @@ impl<'a> Drop for DialogGuardForUnconfirmed<'a> { debug!(%self.id, "unconfirmed dialog dropped, cancelling it"); let _handle = crate::platform::spawn(async move { + let started = crate::platform::Instant::now(); + let final_window = invite_tx.endpoint_inner.option.t1x64; let mut timeout = core::pin::pin!(crate::platform::sleep(core::time::Duration::from_secs(2))); invite_tx.stop_retransmissions(); let mut cancel_done = false; - let mut cancel = core::pin::pin!(client_dialog.cancel()); + let mut final_response = None; + // Boxed so it can be dropped before the wait below. + let mut cancel = Box::pin(client_dialog.cancel()); use crate::platform::select::Which3; loop { @@ -265,6 +269,7 @@ impl<'a> Drop for DialogGuardForUnconfirmed<'a> { status = %resp.status_code, "received final response" ); + final_response = Some(resp); break; } Some(_) => {} @@ -273,14 +278,37 @@ impl<'a> Drop for DialogGuardForUnconfirmed<'a> { } } - // `cancel` (pinned above) ends its borrow here; closing the - // command channel terminates the dialog loop. - drop(invite_tx); + drop(cancel); let _ = client_dialog.inner.transition(DialogState::Terminated( client_dialog.id(), TerminatedReason::UacCancel, )); debug!(id = %client_dialog.id(), "dialog terminated"); + + // The callee may still answer after the CANCEL: keep the + // INVITE transaction until its final response (up to + // 64*T1) so a 2xx is ACKed, then end that session with a + // BYE (RFC 3261 §9.1, §15). + if final_response.is_none() { + let elapsed = crate::platform::Instant::now() + .checked_duration_since(started) + .unwrap_or_default(); + final_response = wait_final_response( + &mut invite_tx, + final_window.saturating_sub(elapsed), + ) + .await; + } + // Closing the command channel terminates the dialog loop. + drop(invite_tx); + if let Some(resp) = final_response { + if resp.status_code.kind() == StatusCodeKind::Successful { + info!(id = %client_dialog.id(), "2xx after CANCEL, sending BYE"); + if let Err(e) = client_dialog.bye_2xx_after_cancel(&resp).await { + warn!(id = %client_dialog.id(), error = %e, "BYE after CANCEL failed"); + } + } + } }); } DialogState::Confirmed(_, _) => { @@ -295,6 +323,24 @@ impl<'a> Drop for DialogGuardForUnconfirmed<'a> { } } +/// Wait up to `limit` for the INVITE transaction's final response. +async fn wait_final_response( + invite_tx: &mut Transaction, + limit: core::time::Duration, +) -> Option { + let wait = async { + while let Some(msg) = invite_tx.receive().await { + if let SipMessage::Response(resp) = msg { + if resp.status_code.kind() != StatusCodeKind::Provisional { + return Some(resp); + } + } + } + None + }; + crate::platform::timeout(limit, wait).await.ok().flatten() +} + pub type InviteAsyncResult = Result<(DialogId, Option)>; impl DialogLayer { diff --git a/src/dialog/invite_dialog.rs b/src/dialog/invite_dialog.rs index 36b2eb19..78cf4b12 100644 --- a/src/dialog/invite_dialog.rs +++ b/src/dialog/invite_dialog.rs @@ -331,18 +331,7 @@ impl InviteDialog { self.inner.update_route_set_from_response(&resp); } StatusCode::OK => { - self.inner.update_route_set_from_response(&resp); - let contact = resp.contact_header()?; - self.inner.remote_contact.lock().replace(contact.clone()); - - let contact_uri = resp - .typed_contact_headers()? - .first() - .map(|c| c.uri.clone()) - .ok_or_else(|| { - crate::Error::Error("missing Contact header".to_string()) - })?; - *self.inner.remote_uri.lock() = contact_uri; + self.update_remote_target_from_2xx(&resp)?; self.inner .transition(DialogState::Confirmed(dialog_id.clone(), resp))?; } @@ -360,6 +349,38 @@ impl InviteDialog { Ok((dialog_id, final_response)) } + /// Take the route set, Contact and remote target from a 2xx to the INVITE. + fn update_remote_target_from_2xx(&self, resp: &Response) -> Result<()> { + self.inner.update_route_set_from_response(resp); + let contact = resp.contact_header()?; + self.inner.remote_contact.lock().replace(contact.clone()); + + let contact_uri = resp + .typed_contact_headers()? + .first() + .map(|c| c.uri.clone()) + .ok_or_else(|| crate::Error::Error("missing Contact header".to_string()))?; + *self.inner.remote_uri.lock() = contact_uri; + Ok(()) + } + + /// End the session a 2xx established after we cancelled the INVITE. + /// + /// A CANCEL that crosses a 2xx has no effect on the INVITE (RFC 3261 + /// §9.1, §15), so the UAC has to send a BYE once the 2xx is ACKed. The + /// dialog was already abandoned, so no `Confirmed` state is reported. + pub(super) async fn bye_2xx_after_cancel(&self, resp: &Response) -> Result<()> { + if let Some(tag) = resp.to_header()?.tag()? { + self.inner.update_remote_tag(tag.value())?; + } + self.update_remote_target_from_2xx(resp)?; + let request = self + .inner + .make_request(Method::Bye, None, None, None, None, None)?; + self.inner.do_request(request).await?; + Ok(()) + } + // ── Shared request semantics ────────────────────────────────────────── /// Send a BYE request to terminate the dialog. diff --git a/src/dialog/tests/mod.rs b/src/dialog/tests/mod.rs index 3659bb23..012772dc 100644 --- a/src/dialog/tests/mod.rs +++ b/src/dialog/tests/mod.rs @@ -1,4 +1,5 @@ mod test_authenticate; +mod test_cancel_2xx_race; mod test_client_dialog; mod test_connection_affinity; mod test_dialog_layer; diff --git a/src/dialog/tests/test_cancel_2xx_race.rs b/src/dialog/tests/test_cancel_2xx_race.rs new file mode 100644 index 00000000..b5a939bc --- /dev/null +++ b/src/dialog/tests/test_cancel_2xx_race.rs @@ -0,0 +1,353 @@ +//! A 2xx that answers the INVITE after we sent CANCEL (the callee picked up +//! while the CANCEL was in flight) establishes a dialog at the far end. The +//! CANCEL has no effect on it, so the UAC must ACK the 2xx and then end the +//! session with a BYE (RFC 3261 §9.1, §15; RFC 5407 §3.1.2). +//! +//! These tests drop the `do_invite` future after a 180 — the documented way to +//! abandon an outgoing call, which cancels it — and answer the INVITE with a +//! 200 from a raw UDP peer in each wire ordering. +use crate::dialog::{ + dialog::{DialogState, DialogStateReceiver, TerminatedReason}, + dialog_layer::DialogLayer, + invitation::InviteOption, +}; +use crate::sip::{prelude::HeadersExt, Method, Request, SipMessage, Uri}; +use crate::transport::{udp::UdpConnection, TransportLayer}; +use crate::EndpointBuilder; +use std::net::SocketAddr; +use std::sync::Arc; +use std::time::Duration; +use tokio::net::UdpSocket; +use tokio::sync::mpsc::unbounded_channel; +use tokio_util::sync::CancellationToken; + +const PEER_TAG: &str = "peer-tag"; + +/// Receive the next request with `method` on the raw peer socket within +/// `wait`, skipping anything else. +async fn recv_request(socket: &UdpSocket, method: Method, wait: Duration) -> (Request, SocketAddr) { + let mut buf = vec![0u8; 4096]; + let deadline = tokio::time::Instant::now() + wait; + loop { + let (len, from) = tokio::time::timeout_at(deadline, socket.recv_from(&mut buf)) + .await + .unwrap_or_else(|_| panic!("timeout waiting for {method}")) + .expect("recv_from failed"); + let text = std::str::from_utf8(&buf[..len]).expect("non utf-8 SIP message"); + if let Ok(SipMessage::Request(req)) = SipMessage::try_from(text) { + if req.method == method { + return (req, from); + } + } + } +} + +/// Fail if the peer receives a BYE within `wait`. +async fn assert_no_bye(socket: &UdpSocket, wait: Duration, what: &str) { + let mut buf = vec![0u8; 4096]; + let deadline = tokio::time::Instant::now() + wait; + while let Ok(Ok((len, _))) = tokio::time::timeout_at(deadline, socket.recv_from(&mut buf)).await + { + let text = std::str::from_utf8(&buf[..len]).unwrap(); + if let Ok(SipMessage::Request(req)) = SipMessage::try_from(text) { + assert_ne!(req.method, Method::Bye, "{what}"); + } + } +} + +/// Build a response to `req` as the peer UA, adding the peer's To tag. +fn response(req: &Request, code: u16, reason: &str, contact: &str) -> String { + let to = req.to_header().unwrap().value().to_string(); + let to = if to.contains(";tag=") { + to + } else { + format!("{to};tag={PEER_TAG}") + }; + format!( + "SIP/2.0 {code} {reason}\r\n\ + Via: {}\r\n\ + From: {}\r\n\ + To: {to}\r\n\ + Call-ID: {}\r\n\ + CSeq: {}\r\n\ + Contact: <{contact}>\r\n\ + Content-Length: 0\r\n\r\n", + req.via_header().unwrap().value(), + req.from_header().unwrap().value(), + req.call_id_header().unwrap().value(), + req.cseq_header().unwrap().value(), + ) +} + +async fn reply(socket: &UdpSocket, to: SocketAddr, req: &Request, code: u16, reason: &str) { + let contact = format!("sip:bob@{}", socket.local_addr().unwrap()); + socket + .send_to(response(req, code, reason, &contact).as_bytes(), to) + .await + .expect("send_to failed"); +} + +struct Uac { + dialog_layer: Arc, + option: InviteOption, + peer: UdpSocket, +} + +async fn setup(token: &CancellationToken) -> crate::Result { + let peer = UdpSocket::bind("127.0.0.1:0").await?; + let peer_port = peer.local_addr()?.port(); + + let transport_layer = TransportLayer::new(token.child_token()); + let udp = UdpConnection::create_connection( + "127.0.0.1:0".parse().unwrap(), + None, + Some(token.child_token()), + ) + .await?; + let uac_addr = udp.get_addr().addr.clone(); + transport_layer.add_transport(udp.into()); + let endpoint = EndpointBuilder::new() + .with_user_agent("rsipstack-test") + .with_transport_layer(transport_layer) + .with_cancel_token(token.child_token()) + .build(); + let endpoint_inner = endpoint.inner.clone(); + tokio::spawn(async move { + let _ = endpoint_inner.serve().await; + }); + let dialog_layer = Arc::new(DialogLayer::new(endpoint.inner.clone())); + let option = InviteOption { + caller: Uri::try_from("sip:alice@example.com")?, + callee: Uri::try_from(format!("sip:bob@127.0.0.1:{peer_port};transport=udp").as_str())?, + contact: Uri::try_from(format!("sip:alice@{uac_addr}").as_str())?, + ..Default::default() + }; + Ok(Uac { + dialog_layer, + option, + peer, + }) +} + +async fn wait_for_state( + rx: &mut DialogStateReceiver, + what: &str, + wait: Duration, + pred: impl Fn(&DialogState) -> bool, +) -> DialogState { + tokio::time::timeout(wait, async { + loop { + let state = rx.recv().await.expect("state channel closed"); + if pred(&state) { + return state; + } + } + }) + .await + .unwrap_or_else(|_| panic!("timeout waiting for {what}")) +} + +/// The ACK and the BYE must belong to the dialog the 2xx established: +/// same Call-ID, our From tag, and the peer's To tag. +fn assert_in_dialog(msg: &Request, invite: &Request, what: &str) { + assert_eq!( + msg.call_id_header().unwrap().value(), + invite.call_id_header().unwrap().value(), + "{what} must carry the INVITE's Call-ID" + ); + assert_eq!( + msg.from_header().unwrap().tag().unwrap(), + invite.from_header().unwrap().tag().unwrap(), + "{what} must carry our From tag" + ); + assert_eq!( + msg.to_header() + .unwrap() + .tag() + .unwrap() + .map(|t| t.value().to_string()), + Some(PEER_TAG.to_string()), + "{what} must carry the 2xx's To tag" + ); +} + +#[derive(Clone, Copy, Debug)] +enum Order { + /// 200 to the INVITE, then 200 to the CANCEL. + InviteOkFirst, + /// 200 to the CANCEL, then 200 to the INVITE right away. + CancelOkFirst, + /// 200 to the CANCEL, then 200 to the INVITE after the 2 s settle window. + InviteOkLate, +} + +async fn run_crossing_2xx(order: Order) -> crate::Result<()> { + let token = CancellationToken::new(); + let Uac { + dialog_layer, + option, + peer, + } = setup(&token).await?; + 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 (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(_, _)) + }) + .await; + + // The application abandons the call: dropping the `do_invite` future + // cancels the INVITE. + invite.abort(); + let _ = invite.await; + let (cancel, _) = recv_request(&peer, Method::Cancel, wait).await; + + let mut terminated = None; + match order { + Order::InviteOkFirst => { + reply(&peer, uac, &inv, 200, "OK").await; + reply(&peer, uac, &cancel, 200, "OK").await; + } + Order::CancelOkFirst => { + reply(&peer, uac, &cancel, 200, "OK").await; + reply(&peer, uac, &inv, 200, "OK").await; + } + Order::InviteOkLate => { + reply(&peer, uac, &cancel, 200, "OK").await; + // The dialog reports Terminated(UacCancel) once the 2 s settle + // window passes without a final response. + terminated = Some( + wait_for_state(&mut states, "Terminated", Duration::from_secs(4), |s| { + matches!(s, DialogState::Terminated(_, _)) + }) + .await, + ); + tokio::time::sleep(Duration::from_millis(2500)).await; + reply(&peer, uac, &inv, 200, "OK").await; + } + } + + let (ack, _) = recv_request(&peer, Method::Ack, wait).await; + assert_in_dialog(&ack, &inv, "ACK"); + assert_eq!( + ack.cseq_header().unwrap().seq().unwrap(), + inv.cseq_header().unwrap().seq().unwrap(), + "the ACK must acknowledge the INVITE" + ); + + let (bye, _) = recv_request(&peer, Method::Bye, wait).await; + assert_in_dialog(&bye, &inv, "BYE"); + reply(&peer, uac, &bye, 200, "OK").await; + + // A retransmitted 2xx is ACKed again but starts no second BYE. + reply(&peer, uac, &inv, 200, "OK").await; + let (ack, _) = recv_request(&peer, Method::Ack, wait).await; + assert_in_dialog(&ack, &inv, "ACK of the retransmitted 2xx"); + assert_no_bye( + &peer, + Duration::from_millis(500), + "one BYE per abandoned call", + ) + .await; + + // The abandoned call reports Terminated(UacCancel) once and is never + // reported as Confirmed. + tokio::time::sleep(Duration::from_millis(200)).await; + let mut seen = Vec::new(); + while let Ok(state) = states.try_recv() { + seen.push(state); + } + if let Some(t) = terminated { + seen.insert(0, t); + } + assert!( + !seen + .iter() + .any(|s| matches!(s, DialogState::Confirmed(_, _))), + "an abandoned call must not report Confirmed, got {seen:?}" + ); + let terminations: Vec<_> = seen + .iter() + .filter(|s| matches!(s, DialogState::Terminated(_, _))) + .collect(); + assert!( + matches!( + terminations.as_slice(), + [DialogState::Terminated(_, TerminatedReason::UacCancel)] + ), + "expected exactly one Terminated(UacCancel), got {seen:?}" + ); + assert!( + bye.cseq_header().unwrap().seq().unwrap() > inv.cseq_header().unwrap().seq().unwrap(), + "the BYE must use a new CSeq" + ); + assert_eq!( + bye.uri.to_string(), + format!("sip:bob@{}", peer.local_addr()?), + "the BYE must target the 2xx's Contact" + ); + token.cancel(); + Ok(()) +} + +#[tokio::test] +async fn test_2xx_before_cancel_response_is_acked_and_byed() -> crate::Result<()> { + run_crossing_2xx(Order::InviteOkFirst).await +} + +#[tokio::test] +async fn test_2xx_after_cancel_response_is_acked_and_byed() -> crate::Result<()> { + run_crossing_2xx(Order::CancelOkFirst).await +} + +#[tokio::test] +async fn test_2xx_after_cancel_settle_window_is_acked_and_byed() -> crate::Result<()> { + run_crossing_2xx(Order::InviteOkLate).await +} + +/// A CANCEL that wins the race (487 to the INVITE) ends the call as before: +/// the 487 is ACKed and no BYE is sent. +#[tokio::test] +async fn test_cancel_answered_487_sends_no_bye() -> crate::Result<()> { + let token = CancellationToken::new(); + let Uac { + dialog_layer, + option, + peer, + } = setup(&token).await?; + 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 (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(_, _)) + }) + .await; + invite.abort(); + let _ = invite.await; + let (cancel, _) = recv_request(&peer, Method::Cancel, wait).await; + reply(&peer, uac, &cancel, 200, "OK").await; + reply(&peer, uac, &inv, 487, "Request Terminated").await; + let (ack, _) = recv_request(&peer, Method::Ack, wait).await; + assert_in_dialog(&ack, &inv, "ACK"); + assert_eq!( + ack.cseq_header().unwrap().seq().unwrap(), + inv.cseq_header().unwrap().seq().unwrap(), + "the ACK must acknowledge the INVITE" + ); + + assert_no_bye( + &peer, + Duration::from_millis(1500), + "a cancelled call must not be BYE'd", + ) + .await; + token.cancel(); + Ok(()) +}