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 = {