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(()) +}