Skip to content

An ACK that arrives before the server INVITE transaction has left respond() is dropped silently #192

Description

@beija-wq

Versions affected. 0.7.3 (latest), 0.6.12, 0.6.10. Totals over 3 runs
per setting:

version 1 call at a time (3×300) 10 at a time (3×300) 50 at a time (3×1000)
0.7.3 2/900 10/900 18/3000
0.6.12 0/900 4/900 15/3000
0.6.10 0/900 5/900 7/3000

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::respond sends first and changes state afterwards
    (src/transaction/transaction.rs:499-500):

    connection.send(response, self.destination.as_ref()).await?;
    self.transition(new_state).map(|_| ())
  • The waiting_ack / waiting_ack_cseq entries that route the 2xx ACK (new
    branch) to the INVITE transaction are inserted only inside
    transition(Accepted) (transaction.rs:1313-1373).

  • EndpointInner::on_received_message rewrites the ACK's key only if
    waiting_ack_cseq already 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().await and transition().

How to run. The UAS uses the usual loop: get_or_create_server_invite,
handle() in its own task, then accept(None, None). A raw UDP 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.

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::UdpSocket peer
on 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):

[dependencies]
rsipstack = "=0.7.3"
tokio = { version = "1", features = ["rt-multi-thread", "macros", "net", "time", "sync"] }
tokio-util = "0.7"

src/main.rs:

//! 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"
        )
    }
}

Run:

cargo new repro && cd repro    # then paste the two blocks above
cargo run --release -- 1000 50

Observed.

== issue A, ACK race | rsipstack/0.7.3 | 127.0.0.1 only
calls: 1000, concurrency: 50, ACK lost (not Confirmed within 400 ms): 2
RESULT: REPRODUCED - 2/1000 first ACKs dropped

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 OK at .726843, udp received ... ACK at
.727232, then entered Accepted state ... waiting for ACK at .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 WaitAck on 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_cseq routes) before connection.send, and roll
back 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_ack
routing), #155 and #170 (ACK routed by CSeq). None of them covers this
window.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions