From d89c32c9be61c886210ebf05a7958be1cdb192e6 Mon Sep 17 00:00:00 2001 From: jinti Date: Wed, 7 Oct 2026 20:39:01 +0800 Subject: [PATCH] =?UTF-8?q?fix(dialog):=20BYE=20a=20forked=202xx's=20dialo?= =?UTF-8?q?g=20(RFC=203261=20=C2=A713.2.2.4)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Every 2xx to the INVITE with a new To tag establishes its own dialog. The transaction already re-ACKs a forked 2xx with its own tag and remote target (#172), but the Accepted-window drainer dropped it: no dialog was created for the branch and no BYE was ever sent, so the forked callee kept retransmitting its 2xx until it gave up. The drainer now inspects every message the client INVITE transaction delivers after confirmation: a 2xx whose To tag differs from the confirmed dialog's is ended with a BYE built from that response's own Contact (remote target) and Record-Route (route set), CSeq continuing the forked dialog's. The confirmed dialog's state is untouched, no dialog is registered for the branch, and no state is notified. --- src/dialog/dialog.rs | 61 ++++++++ src/dialog/invitation.rs | 12 +- src/dialog/invite_dialog.rs | 41 ++++- src/dialog/tests/mod.rs | 1 + src/dialog/tests/test_forked_2xx_bye.rs | 197 ++++++++++++++++++++++++ 5 files changed, 309 insertions(+), 3 deletions(-) create mode 100644 src/dialog/tests/test_forked_2xx_bye.rs diff --git a/src/dialog/dialog.rs b/src/dialog/dialog.rs index 992cf98d..c40e9680 100644 --- a/src/dialog/dialog.rs +++ b/src/dialog/dialog.rs @@ -1335,6 +1335,67 @@ impl DialogInner { result.map(|_| ()) } + /// End a dialog a forked 2xx established (RFC 3261 §13.2.2.4). + /// + /// Every 2xx to the INVITE with a new To tag creates its own dialog; the + /// UAC keeps a single session (the first 2xx's), so the transaction has + /// already ACKed the forked 2xx and this sends the BYE that terminates + /// the extra branch. The BYE is built from that response's own remote + /// target (Contact) and route set (Record-Route); the confirmed dialog's + /// state is not touched and no dialog is registered for the branch. + pub(super) async fn bye_forked_branch(&self, resp: &Response) -> Result<()> { + let contact_uri = resp + .typed_contact_headers()? + .first() + .map(|c| c.uri.clone()) + .ok_or_else(|| crate::Error::Error("missing Contact header".to_string()))?; + + // §12.2.1.1: the forked dialog's route set is the 2xx's Record-Route. + let mut routes: Vec = resp + .record_route_headers() + .into_iter() + .flat_map(|rr| split_rr_values(rr.value())) + .map(Route::from) + .collect(); + routes.reverse(); + + // To carries the forked branch's tag, From keeps ours. + let to = resp.to_header()?.clone(); + let id = self.id.lock().clone(); + let via = self + .endpoint_inner + .get_via(self.via_addr_for_send_transport(), None)?; + let cseq = CSeq { + seq: self.increment_local_seq(), + method: Method::Bye, + }; + + let mut headers: Vec
= vec![ + Header::Via(via.into()), + Header::CallId(id.call_id.clone().into()), + Header::From(self.from.clone().to_string().into()), + Header::To(to), + Header::CSeq(cseq.into()), + Header::UserAgent(self.endpoint_inner.user_agent.clone().into()), + ]; + if let Some(uri) = self.local_contact.as_ref() { + headers.push(Contact::from(uri.clone()).into()); + } + headers.extend(routes.into_iter().map(Header::Route)); + headers.push(Header::MaxForwards(70.into())); + + debug!(id = %id, uri = %contact_uri, "sending BYE to a forked dialog"); + self.do_request(crate::sip::Request { + method: Method::Bye, + uri: contact_uri, + headers: headers.into(), + body: Vec::new(), + version: crate::sip::Version::V2, + }) + .await?; + Ok(()) + } + /// RFC 3261 §13.3.1.4: the server transaction of an INVITE or re-INVITE /// retransmitted our 2xx (`answered_2xx`) for 64*T1 and ended without an /// ACK. The dialog is terminated with [`TerminatedReason::Timeout`] and the diff --git a/src/dialog/invitation.rs b/src/dialog/invitation.rs index f39c3269..96da705a 100644 --- a/src/dialog/invitation.rs +++ b/src/dialog/invitation.rs @@ -719,8 +719,12 @@ impl DialogLayer { // here would leave it (and its timers) in the // endpoint's table and silently stop the re-ACKs. if let Some(mut tx) = guard.invite_tx.take() { + let dlg = dialog.clone(); + let confirmed_tag = new_dialog_id.remote_tag.clone(); crate::platform::spawn(async move { - while tx.receive().await.is_some() {} + while let Some(msg) = tx.receive().await { + dlg.end_forked_branch(&msg, &confirmed_tag).await; + } debug!(id = %new_dialog_id, "accepted transaction drained (Timer M expired)"); }); } @@ -795,8 +799,12 @@ impl DialogLayer { // observes forked 2xx, and detaches it from the // endpoint's table). See do_invite for the rationale. let confirmed_id = new_id.clone(); + let confirmed_tag = new_id.remote_tag.clone(); + let forked_dlg = dialog_clone.clone(); crate::platform::spawn(async move { - while tx.receive().await.is_some() {} + while let Some(msg) = tx.receive().await { + forked_dlg.end_forked_branch(&msg, &confirmed_tag).await; + } debug!(id = %confirmed_id, "accepted transaction drained (Timer M expired)"); }); } diff --git a/src/dialog/invite_dialog.rs b/src/dialog/invite_dialog.rs index 3cb0e35a..1bc4a364 100644 --- a/src/dialog/invite_dialog.rs +++ b/src/dialog/invite_dialog.rs @@ -22,7 +22,7 @@ use crate::transaction::key::TransactionRole; use crate::transaction::transaction::{Transaction, TransactionEvent}; use crate::Result; use core::sync::atomic::Ordering; -use tracing::{debug, trace, warn}; +use tracing::{debug, info, trace, warn}; /// Unified INVITE dialog that can act as either a UAS (Server) or UAC (Client). /// @@ -405,6 +405,45 @@ impl InviteDialog { Ok(()) } + /// End a forked dialog a later 2xx established (RFC 3261 §13.2.2.4). + /// + /// Called from the Accepted-window drainer for every message the client + /// INVITE transaction delivers after the dialog confirmed. A 2xx whose To + /// tag differs from the confirmed dialog's is a forked branch: the + /// transaction has already ACKed it, and the UAC — keeping a single + /// session — terminates it with a BYE. Everything else (retransmitted + /// 2xx with the same tag, non-2xx) is ignored. + pub(super) async fn end_forked_branch(&self, msg: &SipMessage, confirmed_remote_tag: &str) { + let SipMessage::Response(resp) = msg else { + return; + }; + if resp.status_code.kind() != StatusCodeKind::Successful { + return; + } + // Same tag: a retransmission of the confirmed 2xx, already re-ACKed. + let tag = match resp + .to_header() + .ok() + .and_then(|to| to.tag().ok().flatten()) + .map(|tag| tag.value().to_string()) + { + Some(tag) => tag, + None => return, + }; + if tag == confirmed_remote_tag { + return; + } + let id = self.id(); + info!( + id = %id, + tag = %tag, + "forked 2xx acknowledged; ending the extra branch with a BYE (RFC 3261 §13.2.2.4)" + ); + if let Err(e) = self.inner.bye_forked_branch(resp).await { + warn!(id = %id, tag = %tag, error = %e, "failed to BYE the forked branch"); + } + } + // ── 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 b3d43cb3..00fe78a2 100644 --- a/src/dialog/tests/mod.rs +++ b/src/dialog/tests/mod.rs @@ -5,6 +5,7 @@ mod test_client_dialog; mod test_connection_affinity; mod test_dialog_layer; mod test_dialog_states; +mod test_forked_2xx_bye; mod test_in_dialog_provisional; mod test_in_dialog_via; mod test_invite_auth_challenge; diff --git a/src/dialog/tests/test_forked_2xx_bye.rs b/src/dialog/tests/test_forked_2xx_bye.rs new file mode 100644 index 00000000..bd601edf --- /dev/null +++ b/src/dialog/tests/test_forked_2xx_bye.rs @@ -0,0 +1,197 @@ +//! RFC 3261 §13.2.2.4: a 2xx to the INVITE with a new To tag establishes its +//! own dialog. The UAC keeps a single session (the first 2xx's): the +//! transaction ACKs the forked 2xx with its own tag and remote target, and +//! the dialog layer ends the extra branch with a BYE built from that +//! response's Contact. No dialog is registered for the branch and the +//! confirmed dialog's state is untouched. + +use crate::dialog::dialog_layer::DialogLayer; +use crate::dialog::invitation::InviteOption; +use crate::dialog::DialogId; +use crate::sip::prelude::HeadersExt; +use crate::sip::{Method, Request, SipMessage, Uri}; +use crate::transport::udp::UdpConnection; +use crate::transport::TransportLayer; +use crate::EndpointBuilder; +use std::time::Duration; +use tokio::net::UdpSocket; +use tokio::sync::mpsc::unbounded_channel; +use tokio_util::sync::CancellationToken; + +/// Receive the next request of `method` at `peer`. +async fn next_request(peer: &UdpSocket, method: Method, what: &str) -> Request { + let mut buf = vec![0u8; 4096]; + loop { + let (len, _) = tokio::time::timeout(Duration::from_secs(3), peer.recv_from(&mut buf)) + .await + .unwrap_or_else(|_| panic!("timeout waiting for the {what}")) + .expect("peer socket error"); + let Ok(SipMessage::Request(req)) = SipMessage::try_from(&buf[..len]) else { + continue; + }; + if req.method == method { + return req; + } + } +} + +/// RFC 3261 §13.2.2.4: a forked 2xx (To tag `tag-b`) is ACKed with its own +/// tag and remote target (the transaction's job, pinned here end to end) and +/// then ended with a BYE built from its Contact — without touching the +/// confirmed dialog or registering the forked branch. +#[tokio::test] +async fn test_forked_2xx_is_byed_without_touching_the_confirmed_dialog() -> crate::Result<()> { + let token = CancellationToken::new(); + let peer = UdpSocket::bind("127.0.0.1:0").await?; + let tl = TransportLayer::new(token.child_token()); + let udp = UdpConnection::create_connection("127.0.0.1:0".parse()?, None, None).await?; + let uac_addr = udp.get_addr().get_socketaddr()?; + let contact = Uri::try_from(format!("sip:alice@{uac_addr}").as_str())?; + tl.add_transport(udp.into()); + let endpoint = EndpointBuilder::new() + .with_transport_layer(tl) + .with_option(crate::transaction::endpoint::EndpointOption { + t1: Duration::from_millis(10), + t1x64: Duration::from_millis(640), + ..Default::default() + }) + .build(); + let inner = endpoint.inner.clone(); + tokio::spawn(async move { inner.serve().await }); + + let layer = DialogLayer::new(endpoint.inner.clone()); + let (state_sender, mut states) = unbounded_channel(); + let invite = InviteOption { + caller: Uri::try_from("sip:alice@example.com")?, + callee: Uri::try_from(format!("sip:bob@{}", peer.local_addr()?).as_str())?, + contact, + ..Default::default() + }; + // A second handle over the same layer state, for the assertions below + // (DialogLayer is not Clone; the registry lives in the shared inner). + let assert_layer = DialogLayer { + endpoint: layer.endpoint.clone(), + inner: layer.inner.clone(), + }; + let do_invite = tokio::spawn(async move { layer.do_invite(invite, state_sender).await }); + + // The INVITE arrives; branch A answers with tag-a. + let mut buf = vec![0u8; 4096]; + let (len, from) = tokio::time::timeout(Duration::from_secs(3), peer.recv_from(&mut buf)) + .await + .expect("timeout waiting for the INVITE")?; + let invite_req = match SipMessage::try_from(&buf[..len])? { + SipMessage::Request(req) if req.method == Method::Invite => req, + other => panic!("expected the INVITE, got {other}"), + }; + let ok_a = format!( + "SIP/2.0 200 OK\r\nVia: {}\r\nFrom: {}\r\nTo: {};tag=tag-a\r\nCall-ID: {}\r\nCSeq: {}\r\nContact: \r\nContent-Length: 0\r\n\r\n", + invite_req.via_header()?.value(), + invite_req.from_header()?.value(), + invite_req.to_header()?.value(), + invite_req.call_id_header()?.value(), + invite_req.cseq_header()?.value(), + peer.local_addr()?, + ); + peer.send_to(ok_a.as_bytes(), from).await?; + + let (dialog, _) = tokio::time::timeout(Duration::from_secs(3), do_invite) + .await + .expect("timeout waiting for do_invite") + .expect("do_invite task panicked")?; + assert!(dialog.inner.state.lock().is_confirmed()); + while states.try_recv().is_ok() {} + + let confirmed = dialog.id(); + assert_eq!(confirmed.remote_tag, "tag-a"); + + // Branch B forks in: same Call-ID, new To tag, own Contact. + let ok_b = format!( + "SIP/2.0 200 OK\r\nVia: {}\r\nFrom: {}\r\nTo: {};tag=tag-b\r\nCall-ID: {}\r\nCSeq: {}\r\nContact: \r\nContent-Length: 0\r\n\r\n", + invite_req.via_header()?.value(), + invite_req.from_header()?.value(), + invite_req.to_header()?.value(), + invite_req.call_id_header()?.value(), + invite_req.cseq_header()?.value(), + peer.local_addr()?, + ); + peer.send_to(ok_b.as_bytes(), uac_addr).await?; + + // The first ACK confirms branch A; the forked 2xx is then ACKed with + // its own tag and Contact as R-URI. + let mut ack = next_request(&peer, Method::Ack, "forked ACK").await; + while !ack + .to_header() + .ok() + .and_then(|to| to.tag().ok().flatten()) + .is_some_and(|tag| tag.value() == "tag-b") + { + ack = next_request(&peer, Method::Ack, "forked ACK").await; + } + let ack_to = ack.to_header()?.value().to_string(); + assert!(ack_to.contains(";tag=tag-b"), "ACK To: {ack_to}"); + assert!( + ack.uri.to_string().contains("bob-b"), + "ACK R-URI must be the forked Contact: {}", + ack.uri + ); + + // The forked branch is then ended with a BYE from its own Contact. + let bye = next_request(&peer, Method::Bye, "forked-branch BYE").await; + let bye_to = bye.to_header()?.value().to_string(); + assert!(bye_to.contains(";tag=tag-b"), "BYE To: {bye_to}"); + assert!( + bye.uri.to_string().contains("bob-b"), + "BYE R-URI must be the forked Contact: {}", + bye.uri + ); + let cseq = bye.cseq_header()?.value().to_string(); + assert_eq!(cseq, "2 BYE", "BYE CSeq continues the forked dialog's"); + assert_eq!( + bye.call_id_header()?.value(), + invite_req.call_id_header()?.value() + ); + + // The confirmed dialog is untouched: no new state notifications, no + // dialog registered for the forked branch, and the dialog still works. + assert!( + states.try_recv().is_err(), + "the forked 2xx must not touch the confirmed dialog's state" + ); + let forked_id = DialogId { + call_id: confirmed.call_id.clone(), + local_tag: confirmed.local_tag.clone(), + remote_tag: "tag-b".to_string(), + }; + assert!( + assert_layer.get_dialog(&forked_id).is_none(), + "the forked branch must not be registered as a dialog" + ); + assert!(dialog.inner.state.lock().is_confirmed()); + + // The confirmed dialog is still usable: a BYE ends it normally. + let dialog2 = dialog.clone(); + let bye_task = tokio::spawn(async move { dialog2.bye().await }); + let bye = next_request(&peer, Method::Bye, "confirmed-dialog BYE").await; + assert!( + bye.to_header()?.value().to_string().contains(";tag=tag-a"), + "the confirmed dialog's BYE keeps tag-a: {}", + bye.to_header()?.value() + ); + let ok_bye = format!( + "SIP/2.0 200 OK\r\nVia: {}\r\nFrom: {}\r\nTo: {}\r\nCall-ID: {}\r\nCSeq: {}\r\nContent-Length: 0\r\n\r\n", + bye.via_header()?.value(), + bye.from_header()?.value(), + bye.to_header()?.value(), + bye.call_id_header()?.value(), + bye.cseq_header()?.value(), + ); + peer.send_to(ok_bye.as_bytes(), uac_addr).await?; + tokio::time::timeout(Duration::from_secs(3), bye_task) + .await + .expect("timeout waiting for bye()") + .expect("bye task panicked")?; + assert!(dialog.inner.state.lock().is_terminated()); + token.cancel(); + Ok(()) +}