From 2e1bd7658435d4010d44a49ff3cbb997424b26f7 Mon Sep 17 00:00:00 2001 From: Tyson George Date: Tue, 6 Oct 2026 09:51:56 -0400 Subject: [PATCH 1/2] fix: route the ACK of a 2xx to the INVITE with the same CSeq MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ACKs for 2xx are routed to their server INVITE transaction by dialog through waiting_ack, which holds one transaction per dialog. Over UDP the ACK of a re-INVITE can arrive after the next re-INVITE on the same dialog: the UAC sends the next re-INVITE once its ACK is out, and the two can be reordered. By then the newer transaction owns the entry, so the late ACK is routed to it and ignored there (its CSeq does not match), while the older transaction never sees its ACK, retransmits its 2xx until 64*T1 and then ends the call with a BYE. An ACK carries the CSeq number of the INVITE it acknowledges (RFC 3261 §13.2.2.4) and stops that 2xx's retransmissions (§13.3.1.4). Keep the waiting transactions by dialog and CSeq number and route an ACK to the one with its CSeq. waiting_ack is kept as before, and a transaction now removes only its own entries, so an older transaction ending no longer drops a newer one's route. Only a 2xx gets a CSeq route: the ACK of a non-2xx is part of the INVITE transaction and matches it by branch (§17.2.3), so it is no longer routed by dialog. --- src/dialog/tests/test_uas_ack_timeout.rs | 90 +++++++++++++++++++ src/transaction/endpoint.rs | 37 +++++++- .../tests/test_server_invite_ack.rs | 11 ++- src/transaction/transaction.rs | 41 ++++++--- 4 files changed, 159 insertions(+), 20 deletions(-) diff --git a/src/dialog/tests/test_uas_ack_timeout.rs b/src/dialog/tests/test_uas_ack_timeout.rs index 80f0f54c..dadfaa48 100644 --- a/src/dialog/tests/test_uas_ack_timeout.rs +++ b/src/dialog/tests/test_uas_ack_timeout.rs @@ -769,3 +769,93 @@ async fn test_stale_ack_is_ignored_and_a_valid_ack_still_confirms() -> crate::Re token.cancel(); Ok(()) } + +/// RFC 3261 §13.2.2.4: the ACK of a 2xx carries the CSeq number of its +/// INVITE. Over UDP the ACK of a re-INVITE can arrive after the next +/// re-INVITE on the dialog: it must still stop its own 2xx, and the newer +/// re-INVITE must still wait for its own ACK. +#[tokio::test] +async fn test_late_ack_of_an_earlier_reinvite_reaches_its_own_transaction() -> crate::Result<()> { + let token = CancellationToken::new(); + let (mut states, mut dialogs, peer) = setup(&token, short_timers()).await?; + peer.send_request(Method::Invite, 1, None).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)?; + let first = peer.collect(Instant::now() + T1 * 2).await; + let Some((_, SipMessage::Response(ok))) = first.iter().find(|(_, m)| is_2xx_to(m, 1)) else { + panic!("the 2xx"); + }; + let local_tag = ok.to_header()?.tag()?.unwrap().value().to_string(); + peer.send_request(Method::Ack, 1, Some(&local_tag)).await; + wait_confirmed(&mut states).await; + + // re-INVITE 2 and 3 are answered; the ACK of 2 arrives after re-INVITE 3. + let mut answered = None; + for cseq in [2, 3] { + peer.send_request(Method::Invite, cseq, Some(&local_tag)) + .await; + let handle = loop { + let state = tokio::time::timeout(Duration::from_secs(2), states.recv()) + .await + .expect("timeout waiting for the re-INVITE") + .expect("state channel closed"); + if let DialogState::Updated(_, _, handle) = state { + break handle; + } + }; + answered.get_or_insert(Instant::now()); + handle.reply(crate::sip::StatusCode::OK).await.ok(); + // The next request follows this 2xx on the wire. + let deadline = Instant::now() + Duration::from_secs(2); + while !peer + .collect(Instant::now() + T1) + .await + .iter() + .any(|(_, m)| is_2xx_to(m, cseq)) + { + assert!( + Instant::now() < deadline, + "timeout waiting for the 2xx to {cseq}" + ); + } + } + peer.send_request(Method::Ack, 2, Some(&local_tag)).await; + let acked_2 = Instant::now(); + let messages = peer.collect(acked_2 + T1 * 16).await; + assert!( + !messages + .iter() + .any(|(at, m)| is_2xx_to(m, 2) && *at > acked_2 + T1 * 4), + "the late ACK of CSeq 2 must stop the retransmissions of its 2xx" + ); + assert!( + messages + .iter() + .any(|(at, m)| is_2xx_to(m, 3) && *at > acked_2), + "the ACK of CSeq 2 must not stop the 2xx to CSeq 3" + ); + + peer.send_request(Method::Ack, 3, Some(&local_tag)).await; + let acked_3 = Instant::now(); + let messages = peer.collect(answered.unwrap() + T1X64 + T1 * 10).await; + assert!( + !messages + .iter() + .any(|(at, m)| is_invite_2xx(m) && *at > acked_3 + T1 * 4), + "no 2xx may be retransmitted once both are ACKed" + ); + assert!( + !messages.iter().any(|(_, m)| matches!( + m, + SipMessage::Request(req) if req.method == Method::Bye + )), + "an ACKed call must not be ended" + ); + assert!(terminated_reason(&mut states).is_none()); + assert!(dialog.state().is_confirmed()); + token.cancel(); + Ok(()) +} diff --git a/src/transaction/endpoint.rs b/src/transaction/endpoint.rs index c6d7b9fc..7a508846 100644 --- a/src/transaction/endpoint.rs +++ b/src/transaction/endpoint.rs @@ -120,6 +120,10 @@ pub struct EndpointInner { pub finished_transactions: RwMap>, pub transactions: RwMap, pub waiting_ack: RwMap, + /// Server INVITE transactions waiting for the ACK of their 2xx, by dialog + /// and INVITE CSeq number: `waiting_ack` holds only a dialog's latest + /// INVITE, and the ACK of an earlier one can arrive after it (§13.2.2.4). + pub(crate) waiting_ack_cseq: RwMap<(DialogId, u32), TransactionKey>, incoming_sender: TransactionSender, incoming_receiver: Mutex>, cancel_token: CancellationToken, @@ -239,6 +243,7 @@ impl EndpointInner { transactions: RwMap::::new(), finished_transactions: RwMap::>::new(), waiting_ack: RwMap::::new(), + waiting_ack_cseq: RwMap::new(), timer_interval: timer_interval.unwrap_or(Duration::from_millis(20)), cancel_token, incoming_sender, @@ -319,6 +324,7 @@ impl EndpointInner { self.transactions.remove(&key); self.finished_transactions.remove(&key); self.waiting_ack.retain(|_, v| v != &key); + self.waiting_ack_cseq.retain(|_, v| v != &key); continue; } @@ -379,7 +385,12 @@ impl EndpointInner { if let Ok(dialog_id) = DialogId::try_from((req, super::key::TransactionRole::Server)) { - if let Some(tx_key) = self.waiting_ack.get(&dialog_id) { + // By dialog AND CSeq number: a late ACK of an earlier + // re-INVITE must reach its own transaction. + let seq = req.cseq_header().and_then(|c| c.seq()).ok(); + if let Some(tx_key) = + seq.and_then(|seq| self.waiting_ack_cseq.get(&(dialog_id, seq))) + { key = tx_key; } } @@ -394,7 +405,7 @@ impl EndpointInner { if let Ok(dialog_id) = DialogId::try_from((req, super::key::TransactionRole::Server)) { - self.waiting_ack.remove(&dialog_id); + self.forget_waiting_ack(dialog_id, req, &key); } return Ok(()); } @@ -604,6 +615,28 @@ impl EndpointInner { Ok(via) } + /// Drop `key`'s ACK routes for `dialog_id`, never another transaction's. + pub(crate) fn forget_waiting_ack( + &self, + dialog_id: DialogId, + req: &crate::sip::Request, + key: &TransactionKey, + ) { + if let Ok(seq) = req.cseq_header().and_then(|c| c.seq()) { + let route = (dialog_id.clone(), seq); + self.waiting_ack_cseq.with_mut(|m| { + if m.get(&route) == Some(key) { + m.remove(&route); + } + }); + } + self.waiting_ack.with_mut(|m| { + if m.get(&dialog_id) == Some(key) { + m.remove(&dialog_id); + } + }); + } + pub fn get_running_transactions(&self) -> Option> { Some(self.transactions.with(|m| m.keys().cloned().collect())) } diff --git a/src/transaction/tests/test_server_invite_ack.rs b/src/transaction/tests/test_server_invite_ack.rs index cbeca52a..827982a4 100644 --- a/src/transaction/tests/test_server_invite_ack.rs +++ b/src/transaction/tests/test_server_invite_ack.rs @@ -1,10 +1,10 @@ //! Tests for ACK matching on a server INVITE transaction (RFC 3261 §17.1.1.3, //! §13.2.2.4): an ACK acknowledges only the INVITE with the same CSeq number. //! -//! ACKs for 2xx are routed per dialog through `waiting_ack`, so a delayed ACK -//! of an earlier (re-)INVITE on the same dialog reaches the transaction of the -//! current re-INVITE. It must not confirm that transaction or stop its 2xx -//! retransmissions (Timer G). +//! ACKs for 2xx are routed by dialog and CSeq (`waiting_ack_cseq`). A delayed +//! ACK of an earlier (re-)INVITE on the same dialog must not confirm the +//! transaction of the current re-INVITE or stop its 2xx retransmissions +//! (Timer G). use crate::sip::headers::*; use crate::sip::prelude::HeadersExt; @@ -145,6 +145,7 @@ async fn test_server_invite_ignores_ack_with_other_cseq() { 1, "the dialog must still route ACKs to the CSeq 2 transaction" ); + assert_eq!(endpoint.inner.waiting_ack_cseq.len(), 1); // The 2xx was retransmitted by Timer G. let (len, _) = timeout(Duration::from_secs(1), client_conn.recv_raw(&mut buf)) @@ -176,6 +177,7 @@ async fn test_server_invite_ignores_ack_with_other_cseq() { 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(), 0); + assert_eq!(endpoint.inner.waiting_ack_cseq.len(), 0); token.cancel(); serve_handle.abort(); @@ -248,6 +250,7 @@ async fn test_server_invite_ignores_ack_with_other_cseq_over_tcp() { // 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); + assert_eq!(endpoint.inner.waiting_ack_cseq.len(), 0); token.cancel(); serve_handle.abort(); diff --git a/src/transaction/transaction.rs b/src/transaction/transaction.rs index 9427a396..37d1deb8 100644 --- a/src/transaction/transaction.rs +++ b/src/transaction/transaction.rs @@ -918,8 +918,8 @@ impl Transaction { { // RFC 3261 §17.1.1.3 / §13.2.2.4: an ACK carries the CSeq // number of the INVITE it acknowledges. ACKs for 2xx are - // routed per dialog (`waiting_ack`), so a delayed ACK of an - // earlier re-INVITE can land here; it must not confirm this + // routed by dialog and CSeq (`waiting_ack_cseq`); one with + // another CSeq that lands here anyway must not confirm this // transaction or stop its 2xx retransmissions. let ack_seq = req.cseq_header().and_then(|c| c.seq()).ok(); let invite_seq = self.original.cseq_header().and_then(|c| c.seq()).ok(); @@ -1358,6 +1358,18 @@ impl Transaction { if let Some(ref resp) = self.last_response { let dialog_id = DialogId::try_from((resp, TransactionRole::Server))?; + // Route the ACK by dialog AND INVITE CSeq: the ACK + // of an earlier re-INVITE can arrive after a newer + // one was answered, and must reach its own + // transaction (RFC 3261 §13.2.2.4). A non-2xx ACK + // matches by branch (§17.2.3) — unreachable here, + // Accepted only carries 2xx. + let seq = self.original.cseq_header().and_then(|c| c.seq()).ok(); + if let Some(seq) = seq { + self.endpoint_inner + .waiting_ack_cseq + .insert((dialog_id.clone(), seq), self.key.clone()); + } self.endpoint_inner .waiting_ack .insert(dialog_id, self.key.clone()); @@ -1470,7 +1482,11 @@ impl Transaction { if self.transaction_type == TransactionType::ServerInvite { if let Some(ref resp) = self.last_response { if let Ok(dialog_id) = DialogId::try_from((resp, self.role())) { - self.endpoint_inner.waiting_ack.remove(&dialog_id); + self.endpoint_inner.forget_waiting_ack( + dialog_id, + &self.original, + &self.key, + ); } } } @@ -1547,17 +1563,14 @@ impl Transaction { matches!(self.transaction_type, TransactionType::ServerInvite) && self.state == TransactionState::Completed; if !is_server_invite_waiting_ack { - match self.last_response { - Some(ref resp) => match DialogId::try_from((resp, self.role())) { - Ok(dialog_id) => self - .endpoint_inner - .waiting_ack - .remove(&dialog_id) - .map(|_| ()), - Err(_) => None, - }, - _ => None, - }; + if let Some(Ok(dialog_id)) = self + .last_response + .as_ref() + .map(|resp| DialogId::try_from((resp, self.role()))) + { + self.endpoint_inner + .forget_waiting_ack(dialog_id, &self.original, &self.key); + } } let last_message = { From a4357d82d9b1fe68cf3548ccb44c6bd8cef11677 Mon Sep 17 00:00:00 2001 From: jinti Date: Tue, 6 Oct 2026 22:33:11 +0800 Subject: [PATCH 2/2] chore: trigger CI after force-push