From e868b4c8123c7b7618cdcddd3f19f6f0f5457df8 Mon Sep 17 00:00:00 2001 From: jinti Date: Tue, 6 Oct 2026 22:15:32 +0800 Subject: [PATCH] fix(transaction): end the server Accepted state on the matching ACK MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Follow-up review of #164 found that a confirmed server dialog left the Accepted transaction parked in the endpoint's table until Timer L (64*T1): nobody drains it after the dialog's receive loop breaks, and the process_timer fallback detaches with detach(key, None), which does not clean waiting_ack — so every confirmed call also leaked a waiting_ack entry for the process lifetime. - a matching ACK now ends the Accepted state immediately (Terminated): Timer L/G are cancelled, waiting_ack is removed, and the 2xx stays cached in finished_transactions so retransmitted ACKs and late INVITE retransmissions are absorbed below the TU (the pre-#164 post-ACK behavior) - handle_reinvite (InviteDialog + ServerInviteDialog) tracks 'acked' explicitly and passes answered_2xx && !acked to end_session_without_ack, so the timeout path cannot fire on a confirmed re-INVITE now that the ACK ends the transaction - drop the dead 2xx-suppression guard in the Completed Timer G arm (server 2xx never routes to Completed) - new test_accepted_lifecycle suite: server tables are clean right after the ACK; the client Accepted transaction detaches exactly at Timer M with the ACK kept cached --- src/dialog/invite_dialog.rs | 9 +- src/dialog/server_dialog.rs | 9 +- src/transaction/tests/mod.rs | 1 + .../tests/test_accepted_lifecycle.rs | 346 ++++++++++++++++++ .../tests/test_server_invite_ack.rs | 18 +- src/transaction/transaction.rs | 26 +- 6 files changed, 378 insertions(+), 31 deletions(-) create mode 100644 src/transaction/tests/test_accepted_lifecycle.rs diff --git a/src/dialog/invite_dialog.rs b/src/dialog/invite_dialog.rs index ccd9940f..8c15c15b 100644 --- a/src/dialog/invite_dialog.rs +++ b/src/dialog/invite_dialog.rs @@ -841,6 +841,7 @@ impl InviteDialog { .last_response .as_ref() .is_some_and(|resp| resp.status_code.kind() == StatusCodeKind::Successful); + let mut acked = false; while let Some(msg) = tx.receive().await { if let SipMessage::Request(req) = msg { @@ -851,11 +852,17 @@ impl InviteDialog { self.id(), tx.last_response.clone().unwrap_or_default(), ))?; + acked = true; break; } } } - self.inner.end_session_without_ack(tx, answered_2xx).await; + // A matching ACK ends the Accepted transaction; `acked` keeps the + // timeout path (no ACK within 64*T1) from firing on a confirmed + // re-INVITE. + self.inner + .end_session_without_ack(tx, answered_2xx && !acked) + .await; Ok(()) } diff --git a/src/dialog/server_dialog.rs b/src/dialog/server_dialog.rs index 7bc7d4e4..3456600f 100644 --- a/src/dialog/server_dialog.rs +++ b/src/dialog/server_dialog.rs @@ -853,6 +853,7 @@ impl ServerInviteDialog { .last_response .as_ref() .is_some_and(|resp| resp.status_code.kind() == crate::sip::StatusCodeKind::Successful); + let mut acked = false; while let Some(msg) = tx.receive().await { if let SipMessage::Request(req) = msg { @@ -863,11 +864,17 @@ impl ServerInviteDialog { self.id(), tx.last_response.clone().unwrap_or_default(), ))?; + acked = true; break; } } } - self.inner.end_session_without_ack(tx, answered_2xx).await; + // A matching ACK ends the Accepted transaction; `acked` keeps the + // timeout path (no ACK within 64*T1) from firing on a confirmed + // re-INVITE. + self.inner + .end_session_without_ack(tx, answered_2xx && !acked) + .await; Ok(()) } diff --git a/src/transaction/tests/mod.rs b/src/transaction/tests/mod.rs index 1e0ef39c..8a76067b 100644 --- a/src/transaction/tests/mod.rs +++ b/src/transaction/tests/mod.rs @@ -5,6 +5,7 @@ use crate::{ }; use tokio_util::sync::CancellationToken; +mod test_accepted_lifecycle; mod test_ack_flow_reuse; mod test_auto_ack_2xx; mod test_before_send; diff --git a/src/transaction/tests/test_accepted_lifecycle.rs b/src/transaction/tests/test_accepted_lifecycle.rs new file mode 100644 index 00000000..7b096243 --- /dev/null +++ b/src/transaction/tests/test_accepted_lifecycle.rs @@ -0,0 +1,346 @@ +//! Lifecycle guards for the RFC 6026 Accepted state: a confirmed dialog +//! must not leave the transaction parked in the endpoint's tables. +//! +//! - server: a matching ACK ends the Accepted transaction at once +//! (`transactions` and `waiting_ack` are cleaned up immediately, not +//! after Timer L). +//! - client: the Accepted transaction is detached when Timer M expires +//! (the do_invite drainer keeps receiving until then). +use crate::dialog::{ + dialog::{Dialog, DialogStateReceiver}, + dialog_layer::DialogLayer, + invite_dialog::InviteDialog, +}; +use crate::sip::{prelude::HeadersExt, Method, SipMessage}; +use crate::transaction::endpoint::EndpointOption; +use crate::transaction::key::{TransactionKey, TransactionRole}; +use crate::transport::udp::UdpConnection; +use crate::transport::TransportLayer; +use crate::EndpointBuilder; +use std::net::SocketAddr; +use std::sync::Arc; +use std::time::{Duration, Instant}; +use tokio::net::UdpSocket; +use tokio::sync::mpsc::{unbounded_channel, UnboundedReceiver}; +use tokio_util::sync::CancellationToken; + +const CALL_ID: &str = "accepted-lifecycle-test"; +const FROM_TAG: &str = "lifecycle-uac"; + +fn invite_request(peer: SocketAddr, branch: &str) -> crate::sip::Request { + crate::sip::Request { + method: Method::Invite, + uri: crate::sip::Uri::try_from(format!("sip:bob@{peer}").as_str()).unwrap(), + headers: vec![ + crate::sip::headers::Via::new(format!( + "SIP/2.0/UDP {addr};branch={branch}", + addr = "127.0.0.1:5060" + )) + .into(), + crate::sip::headers::CSeq::new("1 INVITE").into(), + crate::sip::headers::From::new(format!(";tag={FROM_TAG}")) + .into(), + crate::sip::headers::To::new("").into(), + crate::sip::headers::CallId::new(CALL_ID).into(), + crate::sip::headers::MaxForwards::new("70").into(), + crate::sip::headers::Contact::new("").into(), + crate::sip::headers::ContentLength::new("0").into(), + ] + .into(), + version: crate::sip::Version::V2, + body: vec![], + } +} + +struct ServerHarness { + endpoint: crate::transaction::endpoint::Endpoint, + dialogs: UnboundedReceiver, + states: DialogStateReceiver, + peer: UdpSocket, + uas: SocketAddr, +} + +async fn server_setup( + token: &CancellationToken, + option: EndpointOption, +) -> crate::Result { + let transport_layer = TransportLayer::new(token.child_token()); + let udp = UdpConnection::create_connection( + "127.0.0.1:0".parse().unwrap(), + None, + Some(token.child_token()), + ) + .await?; + let uas: SocketAddr = udp.get_addr().get_socketaddr()?; + transport_layer.add_transport(udp.into()); + let endpoint = EndpointBuilder::new() + .with_user_agent("rsipstack-test") + .with_transport_layer(transport_layer) + .with_cancel_token(token.child_token()) + .with_option(option) + .build(); + let dialog_layer = Arc::new(DialogLayer::new(endpoint.inner.clone())); + let mut incoming = endpoint.incoming_transactions()?; + let endpoint_inner = endpoint.inner.clone(); + tokio::spawn(async move { + let _ = endpoint_inner.serve().await; + }); + + let (state_sender, states) = unbounded_channel(); + let (dialog_sender, dialogs) = unbounded_channel(); + tokio::spawn(async move { + while let Some(mut tx) = incoming.recv().await { + let has_to_tag = tx.original.to_header().unwrap().tag().unwrap().is_some(); + let dialog = if has_to_tag { + dialog_layer.match_dialog(&tx) + } else if tx.original.method == Method::Invite { + let dialog = dialog_layer + .get_or_create_server_invite(&tx, state_sender.clone(), None, None) + .expect("server dialog"); + dialog_sender.send(dialog.clone()).unwrap(); + Some(Dialog::Invite(dialog)) + } else { + None + }; + if let Some(mut dialog) = dialog { + tokio::spawn(async move { + let _ = dialog.handle(&mut tx).await; + }); + } + } + }); + + let peer = UdpSocket::bind("127.0.0.1:0").await?; + Ok(ServerHarness { + endpoint, + dialogs, + states, + peer, + uas, + }) +} + +/// Server side: once the ACK confirms the dialog, the Accepted transaction +/// must be gone from the endpoint's tables — not parked there until Timer L +/// (64*T1) — and `waiting_ack` must be empty. +#[tokio::test] +async fn test_server_accepted_transaction_ends_on_ack() -> crate::Result<()> { + let token = CancellationToken::new(); + let option = EndpointOption { + t1: Duration::from_millis(20), + t1x64: Duration::from_millis(64 * 20), + ..Default::default() + }; + let harness = server_setup(&token, option).await?; + let ServerHarness { + endpoint, + mut dialogs, + mut states, + peer, + uas, + } = harness; + let inner = endpoint.inner.clone(); + + let invite = invite_request(uas, "z9hG4bK-accepted-lifecycle"); + let key = TransactionKey::from_request(&invite, TransactionRole::Server)?; + peer.send_to(invite.to_string().as_bytes(), uas).await?; + + let dialog = tokio::time::timeout(Duration::from_secs(2), dialogs.recv()) + .await + .expect("timeout waiting for the server dialog") + .unwrap(); + dialog.accept(None, None)?; + + // First 200 OK, remember its To tag. + let mut buf = vec![0u8; 4096]; + let to_tag = loop { + let Ok(Ok((len, _))) = + tokio::time::timeout(Duration::from_secs(2), peer.recv_from(&mut buf)).await + else { + panic!("timeout waiting for the 200"); + }; + let text = std::str::from_utf8(&buf[..len]).unwrap(); + if let Ok(SipMessage::Response(resp)) = SipMessage::try_from(text) { + if resp.status_code.code() == 200 { + break resp + .to_header()? + .tag()? + .expect("the 200 must carry the local tag") + .value() + .to_string(); + } + } + }; + + // The matching ACK confirms the dialog and must end the transaction. + let ack = format!( + "ACK sip:bob@{uas} SIP/2.0\r\n\ + Via: SIP/2.0/UDP 127.0.0.1:5060;branch=z9hG4bK-accepted-lifecycle-ack\r\n\ + Max-Forwards: 70\r\n\ + From: ;tag={FROM_TAG}\r\n\ + To: ;tag={to_tag}\r\n\ + Call-ID: {CALL_ID}\r\n\ + CSeq: 1 ACK\r\n\ + Content-Length: 0\r\n\r\n", + uas = uas, + ); + let acked_at = Instant::now(); + peer.send_to(ack.as_bytes(), uas).await?; + + let mut confirmed = false; + let deadline = Instant::now() + Duration::from_secs(2); + while Instant::now() < deadline { + while let Ok(state) = states.try_recv() { + if matches!(state, crate::dialog::dialog::DialogState::Confirmed(_, _)) { + confirmed = true; + } + } + if confirmed { + break; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + assert!(confirmed, "the ACK must confirm the dialog"); + + // Detached synchronously with the ACK — nothing parked until Timer L. + assert!( + !inner.transactions.contains_key(&key), + "the Accepted transaction must be detached as soon as the ACK confirms the dialog (leaked for {:?} until now)", + acked_at.elapsed() + ); + assert!( + inner.waiting_ack.is_empty(), + "waiting_ack must not retain the confirmed dialog" + ); + // The 2xx stays cached so retransmitted ACKs / late INVITE retransmits + // are absorbed below the TU. + assert!( + inner.finished_transactions.contains_key(&key), + "the 2xx must stay cached for late retransmission absorption" + ); + + token.cancel(); + Ok(()) +} + +/// Client side: the Accepted transaction is detached when Timer M expires. +/// The drainer installed by `DialogLayer::do_invite` must be what keeps it +/// receiving; here the same loop is driven manually with short timers. +#[tokio::test] +async fn test_client_accepted_transaction_detaches_on_timer_m() -> crate::Result<()> { + let token = CancellationToken::new(); + let transport_layer = TransportLayer::new(token.child_token()); + let udp = UdpConnection::create_connection("127.0.0.1:0".parse()?, None, None).await?; + transport_layer.add_transport(udp.into()); + let endpoint = EndpointBuilder::new() + .with_user_agent("rsipstack-test") + .with_transport_layer(transport_layer) + .with_option(EndpointOption { + t1: Duration::from_millis(25), + t1x64: Duration::from_millis(400), + ..Default::default() + }) + .build(); + + let endpoint_inner = endpoint.inner.clone(); + let serve = tokio::spawn(async move { + let _ = endpoint_inner.serve().await; + }); + + let peer = UdpSocket::bind("127.0.0.1:0").await?; + let peer_addr = peer.local_addr()?; + + let invite = invite_request(peer_addr, "z9hG4bK-timer-m"); + let key = TransactionKey::from_request(&invite, TransactionRole::Client)?; + let mut tx = crate::transaction::transaction::Transaction::new_client( + key.clone(), + invite, + endpoint.inner.clone(), + None, + ); + tx.send().await?; + + // Answer the INVITE (and its retransmissions) with a tagged 200 OK. + let mut buf = vec![0u8; 4096]; + let ok = loop { + let Ok(Ok((len, src))) = + tokio::time::timeout(Duration::from_secs(2), peer.recv_from(&mut buf)).await + else { + panic!("timeout waiting for the INVITE"); + }; + let text = String::from_utf8_lossy(&buf[..len]).to_string(); + if text.starts_with("INVITE") { + let resp = format!( + "SIP/2.0 200 OK\r\nVia: {via}\r\nFrom: {from}\r\nTo: {to};tag=peer-tag\r\nCall-ID: {callid}\r\nCSeq: 1 INVITE\r\nContact: \r\nContent-Length: 0\r\n\r\n", + via = text + .lines() + .find(|l| l.starts_with("Via:")) + .unwrap() + .strip_prefix("Via: ") + .unwrap(), + from = text + .lines() + .find(|l| l.starts_with("From:")) + .unwrap() + .strip_prefix("From: ") + .unwrap(), + to = text + .lines() + .find(|l| l.starts_with("To:")) + .unwrap() + .strip_prefix("To: ") + .unwrap(), + callid = CALL_ID, + src = src, + ); + peer.send_to(resp.as_bytes(), src).await?; + break; + } + }; + + // First 2xx → Accepted (auto-ACK fired on the wire). + let first = tokio::time::timeout(Duration::from_secs(2), tx.receive()) + .await + .expect("timeout waiting for the 200") + .expect("transaction ended before the 200"); + assert!(matches!(first, SipMessage::Response(ref r) if r.status_code.code() == 200)); + assert_eq!( + tx.state, + crate::transaction::TransactionState::Accepted, + "the client INVITE 2xx must park the transaction in Accepted (RFC 6026 §7.2)" + ); + assert!( + endpoint.inner.transactions.contains_key(&key), + "the Accepted transaction must stay in the table during the Timer M window" + ); + + // Drive receive() until Timer M ends the transaction (what the do_invite + // drainer does). + let started = Instant::now(); + while tokio::time::timeout(Duration::from_secs(3), tx.receive()) + .await + .expect("drain timed out") + .is_some() + {} + assert!( + started.elapsed() >= Duration::from_millis(300), + "the transaction must stay in Accepted until Timer M, ended after {:?}", + started.elapsed() + ); + assert_eq!( + tx.state, + crate::transaction::TransactionState::Terminated, + "Timer M must terminate the Accepted transaction" + ); + assert!( + !endpoint.inner.transactions.contains_key(&key), + "Timer M must detach the Accepted transaction from the endpoint's table" + ); + // The stored ACK stays cached so late 2xx retransmissions are absorbed. + assert!(endpoint.inner.finished_transactions.contains_key(&key)); + + serve.abort(); + let _ = ok; + token.cancel(); + Ok(()) +} diff --git a/src/transaction/tests/test_server_invite_ack.rs b/src/transaction/tests/test_server_invite_ack.rs index 96752676..cbeca52a 100644 --- a/src/transaction/tests/test_server_invite_ack.rs +++ b/src/transaction/tests/test_server_invite_ack.rs @@ -170,12 +170,12 @@ async fn test_server_invite_ignores_ack_with_other_cseq() { } other => panic!("expected ACK, got {}", other), } - // RFC 6026 §7.1: the ACK does not end the Accepted state (Timer L - // does); retransmissions stop and the waiting_ack entry stays so any - // retransmitted ACK keeps being routed here until then. - assert_eq!(tx.state, TransactionState::Accepted); + // A matching ACK ends the Accepted transaction: retransmissions stop, + // the waiting_ack entry is removed, and late retransmitted ACKs / + // INVITE retransmissions are absorbed via finished_transactions. + assert_eq!(tx.state, TransactionState::Terminated); assert!(tx.timer_g.is_none(), "Timer G must stop once ACKed"); - assert_eq!(endpoint.inner.waiting_ack.len(), 1); + assert_eq!(endpoint.inner.waiting_ack.len(), 0); token.cancel(); serve_handle.abort(); @@ -245,11 +245,9 @@ async fn test_server_invite_ignores_ack_with_other_cseq_over_tcp() { } other => panic!("expected ACK, got {}", other), } - // RFC 6026 §7.1: the ACK does not end the Accepted state (Timer L - // does); the waiting_ack entry stays so retransmitted ACKs keep being - // routed here until then. - assert_eq!(tx.state, TransactionState::Accepted); - assert_eq!(endpoint.inner.waiting_ack.len(), 1); + // A matching ACK ends the Accepted transaction (see the UDP variant). + assert_eq!(tx.state, TransactionState::Terminated); + assert_eq!(endpoint.inner.waiting_ack.len(), 0); token.cancel(); serve_handle.abort(); diff --git a/src/transaction/transaction.rs b/src/transaction/transaction.rs index 5bf2d52e..9427a396 100644 --- a/src/transaction/transaction.rs +++ b/src/transaction/transaction.rs @@ -905,11 +905,11 @@ impl Transaction { return None; } // The peer confirmed the dialog: stop the 2xx - // retransmissions (RFC 3261 §13.3.1.4) and remain in - // Accepted until Timer L ends the transaction. - if let Some(id) = self.timer_g.take() { - self.endpoint_inner.timers.cancel(id); - } + // retransmissions (RFC 3261 §13.3.1.4) and end the + // transaction. Retransmitted ACKs and late INVITE + // retransmissions are absorbed below the TU by + // `finished_transactions`. + self.transition(TransactionState::Terminated).ok(); return Some(req.into()); } } @@ -1136,20 +1136,8 @@ impl Transaction { } TransactionState::Completed => { if let TransactionTimer::TimerG(key, duration) = timer { - // RFC 6026 §7.1 defensive guard: the server transaction - // MUST NOT retransmit 2xx responses on its own. Per the - // RFC 6026 routing in respond() + on_received_response(), - // 2xx finals route to the Accepted state, not Completed — - // so `last_response` here should always be non-2xx. This - // guard catches any legacy / out-of-band code path that - // might land a 2xx in Completed; suppress the retransmit - // and let Timer D / Timer K handle Termination. - if let Some(last_response) = &self.last_response { - if last_response.status_code.kind() == StatusCodeKind::Successful { - return Ok(()); - } - } - // resend the response (non-2xx final — RFC 3261 §17.2.1) + // resend the response (non-2xx final — RFC 3261 §17.2.1; + // 2xx finals route to Accepted, never here) if let Some(last_response) = &self.last_response { if let Some(connection) = &self.connection { let last_response = if let Some(ref inspector) =