//! Issue A: an ACK that reaches the stack between the 2xx leaving the socket
//! and the server INVITE transaction entering Accepted is dropped.
//!
//! The peer ACKs each 200 the moment it arrives, and only once (it ignores
//! retransmissions, so a lost ACK stays visible). A dialog that is not
//! `Confirmed` within 400 ms (before the first 2xx retransmission at T1)
//! lost its ACK.
//!
//! Usage: cargo run --release -- [calls] [concurrency] (default 1000 50)
#![allow(dead_code)] // the harness at the bottom is shared by several issues
use rsipstack::dialog::dialog::DialogState;
use std::net::SocketAddr;
use std::sync::atomic::{AtomicUsize, Ordering::SeqCst};
use std::sync::Arc;
use std::time::{Duration, Instant};
use tokio::sync::mpsc::{unbounded_channel, UnboundedReceiver};
#[tokio::main(flavor = "multi_thread")]
async fn main() -> Result<(), BoxError> {
println!(
"== issue A, ACK race | {} | 127.0.0.1 only",
rsipstack::VERSION
);
let args: Vec<usize> = std::env::args()
.skip(1)
.filter_map(|a| a.parse().ok())
.collect();
let calls_n = args.first().copied().unwrap_or(1000);
let conc = args.get(1).copied().unwrap_or(50).max(1);
let stack = Stack::start().await?;
let mut calls = stack.uas_loop()?;
let uas = stack.addr;
// UAS: spawn handle(), accept() at once, then wait 400 ms for Confirmed.
let lost = Arc::new(AtomicUsize::new(0));
let done = Arc::new(AtomicUsize::new(0));
let (lost2, done2) = (lost.clone(), done.clone());
tokio::spawn(async move {
while let Some(mut call) = calls.recv().await {
let (lost3, done3) = (lost2.clone(), done2.clone());
tokio::spawn(async move {
call.spawn_handle();
let _ = call.dialog.accept(None, None);
let deadline = tokio::time::Instant::now() + Duration::from_millis(400);
let mut confirmed = false;
while let Ok(Some(s)) = tokio::time::timeout_at(deadline, call.states.recv()).await
{
if matches!(s, DialogState::Confirmed(..)) {
confirmed = true;
break;
}
}
if !confirmed {
lost3.fetch_add(1, SeqCst);
}
done3.fetch_add(1, SeqCst);
});
}
});
// UACs: `conc` raw peers, one call at a time each.
let next = Arc::new(AtomicUsize::new(0));
let mut workers = Vec::new();
for _ in 0..conc {
let next = next.clone();
workers.push(tokio::spawn(async move {
let peer = Peer::bind().await?;
loop {
let i = next.fetch_add(1, SeqCst);
if i >= calls_n {
return Ok::<(), BoxError>(());
}
let call_id = format!("a-{i}");
peer.send(uas, &peer.invite(uas, &call_id)).await?;
let ok = peer
.recv_until(Duration::from_secs(2), |r| {
r.is_response(200, "INVITE")
&& header(&r.text, "Call-ID").as_deref() == Some(call_id.as_str())
})
.await;
if let Some(ok) = ok {
peer.send(uas, &peer.ack_for(uas, &ok)).await?;
}
// Let the UAS finish its 400 ms window, drop retransmissions.
tokio::time::sleep(Duration::from_millis(450)).await;
while peer.recv(Duration::from_millis(1)).await.is_some() {}
}
}));
}
for w in workers {
w.await??;
}
let started = Instant::now();
while done.load(SeqCst) < calls_n && started.elapsed() < Duration::from_secs(10) {
tokio::time::sleep(Duration::from_millis(50)).await;
}
let lost = lost.load(SeqCst);
println!(
"calls: {calls_n}, concurrency: {conc}, ACK lost (not Confirmed within 400 ms): {lost}"
);
verdict(lost > 0, &format!("{lost}/{calls_n} first ACKs dropped"))
}
// ---- harness: the stack under test and a raw UDP peer, 127.0.0.1 only ----
type BoxError = Box<dyn std::error::Error + Send + Sync>;
fn verdict(reproduced: bool, what: &str) -> ! {
let tag = if reproduced {
"REPRODUCED"
} else {
"not reproduced"
};
println!("RESULT: {tag} - {what}");
std::process::exit(i32::from(reproduced))
}
/// The stack under test: one UDP transport on 127.0.0.1:0, default options.
struct Stack {
endpoint: rsipstack::transaction::Endpoint,
dialogs: std::sync::Arc<rsipstack::dialog::dialog_layer::DialogLayer>,
addr: SocketAddr,
cancel: tokio_util::sync::CancellationToken,
}
impl Stack {
async fn start() -> Result<Stack, BoxError> {
let cancel = tokio_util::sync::CancellationToken::new();
let tl = rsipstack::transport::TransportLayer::new(cancel.child_token());
let udp = rsipstack::transport::udp::UdpConnection::create_connection(
"127.0.0.1:0".parse()?,
None,
Some(cancel.child_token()),
)
.await?;
let addr = udp.get_addr().get_socketaddr()?;
tl.add_transport(udp.into());
let endpoint = rsipstack::EndpointBuilder::new()
.with_user_agent("repro")
.with_transport_layer(tl)
.with_cancel_token(cancel.child_token())
.build();
let inner = endpoint.inner.clone();
tokio::spawn(async move { inner.serve().await });
let dialogs = std::sync::Arc::new(rsipstack::dialog::dialog_layer::DialogLayer::new(
endpoint.inner.clone(),
));
Ok(Stack {
endpoint,
dialogs,
addr,
cancel,
})
}
}
impl Drop for Stack {
fn drop(&mut self) {
self.cancel.cancel();
}
}
/// A new server INVITE dialog, its state receiver and its (not yet handled)
/// transaction.
struct UasCall {
dialog: rsipstack::dialog::invite_dialog::InviteDialog,
states: UnboundedReceiver<DialogState>,
tx: Option<rsipstack::transaction::Transaction>,
}
impl UasCall {
fn spawn_handle(&mut self) {
if let Some(mut tx) = self.tx.take() {
let mut d = self.dialog.clone();
tokio::spawn(async move { d.handle(&mut tx).await });
}
}
}
impl Stack {
/// The usual UAS loop: in-dialog requests go to `match_dialog` + `handle()`;
/// each new INVITE gets `get_or_create_server_invite` and is handed over.
fn uas_loop(&self) -> Result<UnboundedReceiver<UasCall>, BoxError> {
use rsipstack::sip::prelude::HeadersExt;
let mut incoming = self.endpoint.incoming_transactions()?;
let dl = self.dialogs.clone();
let (call_tx, call_rx) = unbounded_channel();
tokio::spawn(async move {
while let Some(mut tx) = incoming.recv().await {
let to_tag = tx
.original
.to_header()
.ok()
.and_then(|h| h.tag().ok().flatten());
if to_tag.is_some() {
if let Some(mut d) = dl.match_dialog(&tx) {
tokio::spawn(async move { d.handle(&mut tx).await });
}
} else if tx.original.method == rsipstack::sip::Method::Invite {
let (state_tx, states) = unbounded_channel();
if let Ok(dialog) = dl.get_or_create_server_invite(&tx, state_tx, None, None) {
let _ = call_tx.send(UasCall {
dialog,
states,
tx: Some(tx),
});
}
}
}
});
Ok(call_rx)
}
}
/// The far end: a raw UDP socket that writes and reads SIP text by hand.
struct Peer {
sock: tokio::net::UdpSocket,
addr: SocketAddr,
}
/// One datagram received by the peer.
struct Rx {
text: String,
from: SocketAddr,
}
impl Rx {
fn first_line(&self) -> &str {
self.text.lines().next().unwrap_or("")
}
fn is_request(&self, method: &str) -> bool {
self.first_line().starts_with(&format!("{method} "))
}
fn is_response(&self, code: u16, cseq_method: &str) -> bool {
self.first_line().starts_with(&format!("SIP/2.0 {code} "))
&& header(&self.text, "CSeq").is_some_and(|c| c.ends_with(cseq_method))
}
}
/// The value of the first header `name` (full name, case-insensitive).
fn header(text: &str, name: &str) -> Option<String> {
let prefix = format!("{}:", name.to_ascii_lowercase());
text.split("\r\n")
.find(|l| l.to_ascii_lowercase().starts_with(&prefix))
.map(|l| l[prefix.len()..].trim().to_string())
}
impl Peer {
async fn bind() -> Result<Peer, BoxError> {
let sock = tokio::net::UdpSocket::bind("127.0.0.1:0").await?;
let addr = sock.local_addr()?;
Ok(Peer { sock, addr })
}
async fn send(&self, to: SocketAddr, text: &str) -> Result<(), BoxError> {
self.sock.send_to(text.as_bytes(), to).await?;
Ok(())
}
async fn recv(&self, wait: Duration) -> Option<Rx> {
let mut buf = vec![0u8; 65536];
let (n, from) = tokio::time::timeout(wait, self.sock.recv_from(&mut buf))
.await
.ok()?
.ok()?;
Some(Rx {
text: String::from_utf8_lossy(&buf[..n]).into_owned(),
from,
})
}
/// The first datagram matching `pred` within `wait`; others are dropped.
async fn recv_until(&self, wait: Duration, pred: impl Fn(&Rx) -> bool) -> Option<Rx> {
let deadline = Instant::now() + wait;
loop {
let rx = self
.recv(deadline.saturating_duration_since(Instant::now()))
.await?;
if pred(&rx) {
return Some(rx);
}
}
}
}
impl Peer {
/// An out-of-dialog INVITE from this peer to `uas`.
fn invite(&self, uas: SocketAddr, call_id: &str) -> String {
let me = self.addr;
format!(
"INVITE sip:bob@{uas} SIP/2.0\r\n\
Via: SIP/2.0/UDP {me};branch=z9hG4bK-{call_id}-inv;rport\r\n\
Max-Forwards: 70\r\n\
From: <sip:alice@{me}>;tag=a1\r\n\
To: <sip:bob@{uas}>\r\n\
Call-ID: {call_id}\r\n\
CSeq: 1 INVITE\r\n\
Contact: <sip:alice@{me}>\r\n\
Content-Length: 0\r\n\r\n"
)
}
/// The ACK for a 2xx (new branch, To copied from the 2xx).
fn ack_for(&self, uas: SocketAddr, ok: &Rx) -> String {
let me = self.addr;
let to = header(&ok.text, "To").unwrap_or_default();
let call_id = header(&ok.text, "Call-ID").unwrap_or_default();
format!(
"ACK sip:bob@{uas} SIP/2.0\r\n\
Via: SIP/2.0/UDP {me};branch=z9hG4bK-{call_id}-ack;rport\r\n\
Max-Forwards: 70\r\n\
From: <sip:alice@{me}>;tag=a1\r\n\
To: {to}\r\n\
Call-ID: {call_id}\r\n\
CSeq: 1 ACK\r\n\
Content-Length: 0\r\n\r\n"
)
}
}
Versions affected. 0.7.3 (latest), 0.6.12, 0.6.10. Totals over 3 runs
per setting:
Environment. rsipstack 0.7.3, the latest version on crates.io, default features. Rust 1.98.1 stable, tokio 1.53 (multi-thread runtime). Linux 6.12 x86_64, 4 CPUs. UDP over loopback only: everything on 127.0.0.1, every port assigned by the OS. Default
EndpointOption(T1 = 500 ms, T2 = 4 s, T4 = 5 s).Where. Line numbers are from the 0.7.3 sources.
Transaction::respondsends first and changes state afterwards(
src/transaction/transaction.rs:499-500):The
waiting_ack/waiting_ack_cseqentries that route the 2xx ACK (newbranch) to the INVITE transaction are inserted only inside
transition(Accepted)(transaction.rs:1313-1373).EndpointInner::on_received_messagerewrites the ACK's key only ifwaiting_ack_cseqalready has the dialog (src/transaction/endpoint.rs:384-396).Otherwise no transaction matches, and
Method::Ack => return Ok(())(
endpoint.rs:528) drops the ACK without logging anything.A UAC on the same host or LAN can send its ACK inside the gap between
send().awaitandtransition().How to run. The UAS uses the usual loop:
get_or_create_server_invite,handle()in its own task, thenaccept(None, None). A raw UDP peer ACKseach 200 the moment it arrives, and only once. It ignores retransmissions, so
a lost ACK stays visible. A dialog that is not
Confirmedwithin 400 ms(before the first 2xx retransmission at T1) lost its ACK.
Reproduction program (one self-contained file, 319 lines) and how to run it
It uses only rsipstack's public API plus a raw
tokio::net::UdpSocketpeeron 127.0.0.1. It prints what it observes and exits with 1 when the problem
reproduces, 0 when it does not.
Cargo.toml(replace the[dependencies]section):src/main.rs:Run:
Observed.
In a separate instrumented run of the same logic (tracing subscriber at
rsipstack=debug), a lost case on 0.7.3 logs this order:udp send ... 200 OKat .726843,udp received ... ACKat.727232, then
entered Accepted state ... waiting for ACKat .727596.Nothing is logged for the ACK itself.
Between 0 and about 1.5 % of first ACKs are dropped per run, more under
concurrent load; one run on 0.7.3 lost 2 of 300 even with one call at a
time. A conformant UAC re-ACKs the 2xx retransmission, so the dialog is confirmed
about 500 ms late. A UAC that sends one ACK only has its call torn down at
64*T1 on 0.7.x, and left in
WaitAckon 0.6.x.Expected. RFC 3261 §13.3.1.4 and §17.2.1; RFC 6026 §8.5 (the server
transaction passes the 2xx to the transport and MUST then transition to
Accepted) and §7.1 (ACKs received in Accepted MUST be passed directly to the
TU and not absorbed): once the 2xx is sent, a matching ACK reaches the TU. The state that a reply depends on must exist before the
request that triggers the reply leaves the host.
Suggested fix. Move to the new state (or at least insert the
waiting_ack/waiting_ack_cseqroutes) beforeconnection.send, and rollback if the send fails. More generally: make every state change that a peer's
reply can observe before the bytes go out.
Related. #149, #164 and #169 (the Accepted state and
waiting_ackrouting), #155 and #170 (ACK routed by CSeq). None of them covers this
window.