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
61 changes: 61 additions & 0 deletions src/dialog/dialog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Route> = 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<Header> = 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
Expand Down
12 changes: 10 additions & 2 deletions src/dialog/invitation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)");
});
}
Expand Down Expand Up @@ -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)");
});
}
Expand Down
41 changes: 40 additions & 1 deletion src/dialog/invite_dialog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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).
///
Expand Down Expand Up @@ -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.
Expand Down
1 change: 1 addition & 0 deletions src/dialog/tests/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
197 changes: 197 additions & 0 deletions src/dialog/tests/test_forked_2xx_bye.rs
Original file line number Diff line number Diff line change
@@ -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: <sip:bob-a@{}>\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: <sip:bob-b@{}>\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(())
}
Loading