Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
90 changes: 90 additions & 0 deletions src/dialog/tests/test_uas_ack_timeout.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(())
}
37 changes: 35 additions & 2 deletions src/transaction/endpoint.rs
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,10 @@ pub struct EndpointInner {
pub finished_transactions: RwMap<TransactionKey, Option<SipMessage>>,
pub transactions: RwMap<TransactionKey, TransactionEventSender>,
pub waiting_ack: RwMap<DialogId, TransactionKey>,
/// 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<Option<TransactionReceiver>>,
cancel_token: CancellationToken,
Expand Down Expand Up @@ -239,6 +243,7 @@ impl EndpointInner {
transactions: RwMap::<TransactionKey, TransactionEventSender>::new(),
finished_transactions: RwMap::<TransactionKey, Option<SipMessage>>::new(),
waiting_ack: RwMap::<DialogId, TransactionKey>::new(),
waiting_ack_cseq: RwMap::new(),
timer_interval: timer_interval.unwrap_or(Duration::from_millis(20)),
cancel_token,
incoming_sender,
Expand Down Expand Up @@ -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;
}

Expand Down Expand Up @@ -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;
}
}
Expand All @@ -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(());
}
Expand Down Expand Up @@ -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<Vec<TransactionKey>> {
Some(self.transactions.with(|m| m.keys().cloned().collect()))
}
Expand Down
11 changes: 7 additions & 4 deletions src/transaction/tests/test_server_invite_ack.rs
Original file line number Diff line number Diff line change
@@ -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;
Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -173,6 +174,7 @@ async fn test_server_invite_ignores_ack_with_other_cseq() {
assert_eq!(tx.state, TransactionState::Confirmed);
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();
Expand Down Expand Up @@ -244,6 +246,7 @@ async fn test_server_invite_ignores_ack_with_other_cseq_over_tcp() {
}
assert_eq!(tx.state, TransactionState::Confirmed);
assert_eq!(endpoint.inner.waiting_ack.len(), 0);
assert_eq!(endpoint.inner.waiting_ack_cseq.len(), 0);

token.cancel();
serve_handle.abort();
Expand Down
37 changes: 23 additions & 14 deletions src/transaction/transaction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -753,8 +753,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();
Expand Down Expand Up @@ -1069,6 +1069,14 @@ impl Transaction {
debug!(key=%self.key, last = self.last_response.is_none(), "entered confirmed state, waiting for ACK");
if let Some(ref resp) = self.last_response {
let dialog_id = DialogId::try_from((resp, TransactionRole::Server))?;
// Only a 2xx ACK is routed by dialog; a non-2xx ACK
// matches by branch (§17.2.3).
let seq = self.original.cseq_header().and_then(|c| c.seq()).ok();
if let Some(seq) = seq.filter(|_| answered_2xx) {
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());
Expand All @@ -1090,7 +1098,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,
);
}
}
}
Expand Down Expand Up @@ -1161,17 +1173,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 = {
Expand Down