Skip to content
Merged
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
9 changes: 9 additions & 0 deletions src/dialog/client_dialog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -668,18 +668,27 @@ impl ClientInviteDialog {
.transition(DialogState::Updated(self.id(), tx.original.clone(), handle))?;

self.inner.process_transaction_handle(tx, rx).await?;
let answered_2xx = tx
.last_response
.as_ref()
.is_some_and(|resp| resp.status_code.kind() == crate::sip::StatusCodeKind::Successful);
let mut acked = false;

// wait for ACK
while let Some(msg) = tx.receive().await {
match msg {
SipMessage::Request(req) if req.method == crate::sip::Method::Ack => {
debug!(id = %self.id(), "received ACK for re-INVITE");
self.inner.remote_ack.lock().replace(req);
acked = true;
break;
}
_ => {}
}
}
self.inner
.end_session_without_ack(tx, answered_2xx && !acked)
.await;
Ok(())
}

Expand Down
115 changes: 115 additions & 0 deletions src/dialog/tests/test_uas_ack_timeout.rs
Original file line number Diff line number Diff line change
Expand Up @@ -958,3 +958,118 @@ async fn test_acked_reinvite_sends_no_bye() -> crate::Result<()> {
token.cancel();
Ok(())
}

/// The deprecated `ClientInviteDialog` has its own re-INVITE handler: a UAC
/// dialog re-INVITEd by the callee must end the same way as `InviteDialog`.
async fn legacy_client_dialog_reinvite(ack: bool) -> crate::Result<()> {
use crate::dialog::{client_dialog::ClientInviteDialog, invitation::InviteOption};
let token = CancellationToken::new();
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 uac: SocketAddr = udp.get_addr().get_socketaddr()?;
transport_layer.add_transport(udp.into());
let endpoint = EndpointBuilder::new()
.with_transport_layer(transport_layer)
.with_cancel_token(token.child_token())
.with_option(short_timers())
.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 { endpoint_inner.serve().await });
let layer = dialog_layer.clone();
tokio::spawn(async move {
while let Some(mut tx) = incoming.recv().await {
if let Some(Dialog::Invite(dialog)) = layer.match_dialog(&tx) {
let mut legacy = ClientInviteDialog::try_from(dialog).expect("a UAC dialog");
tokio::spawn(async move { legacy.handle(&mut tx).await });
}
}
});
let peer = Peer {
socket: UdpSocket::bind("127.0.0.1:0").await?,
uas: uac,
};
let callee = peer.socket.local_addr()?;
let option = InviteOption {
caller: format!("sip:alice@{uac}").as_str().try_into()?,
callee: format!("sip:bob@{callee}").as_str().try_into()?,
contact: format!("sip:alice@{uac}").as_str().try_into()?,
call_id: Some(CALL_ID.to_string()),
..Default::default()
};
let (state_sender, mut states) = unbounded_channel();
let invite = tokio::spawn(async move { dialog_layer.do_invite(option, state_sender).await });
let mut buf = vec![0u8; 4096];
let (len, _) = tokio::time::timeout(Duration::from_secs(2), peer.socket.recv_from(&mut buf))
.await
.expect("timeout waiting for the INVITE")?;
let SipMessage::Request(req) = SipMessage::try_from(std::str::from_utf8(&buf[..len]).unwrap())?
else {
panic!("expected the INVITE");
};
peer.send(format!(
"SIP/2.0 200 OK\r\nVia: {}\r\nFrom: {}\r\nTo: {};tag={FROM_TAG}\r\nCall-ID: {CALL_ID}\r\n\
CSeq: 1 INVITE\r\nContact: <sip:bob@{callee}>\r\nContent-Length: 0\r\n\r\n",
req.via_header()?.value(),
req.from_header()?.value(),
req.to_header()?.value(),
))
.await;
let (dialog, _) = invite.await.unwrap()?;
let local_tag = dialog.id().local_tag;

// The callee re-INVITEs and the application answers it.
peer.send_request(Method::Invite, 7, Some(&local_tag)).await;
let handle = loop {
match tokio::time::timeout(Duration::from_secs(2), states.recv()).await {
Ok(Some(DialogState::Updated(_, _, handle))) => break handle,
Ok(Some(_)) => {}
_ => panic!("timeout waiting for the re-INVITE"),
}
};
let answered = Instant::now();
handle.reply(crate::sip::StatusCode::OK).await.ok();
let mut messages = Vec::new();
while !messages.iter().any(|(_, m)| is_2xx_to(m, 7)) {
assert!(answered.elapsed() < T1X64, "timeout waiting for the 2xx");
messages.extend(peer.collect(Instant::now() + T1 / 2).await);
}
if ack {
peer.send_request(Method::Ack, 7, Some(&local_tag)).await;
}
messages.extend(peer.collect(answered + T1X64 * 2).await);
let bye = messages
.iter()
.any(|(_, m)| matches!(m, SipMessage::Request(r) if r.method == Method::Bye));
if ack {
assert!(!bye, "an ACKed re-INVITE must not end the session");
assert!(terminated_reason(&mut states).is_none());
assert!(dialog.state().is_confirmed());
} else {
assert!(messages.iter().filter(|(_, m)| is_2xx_to(m, 7)).count() >= 2);
assert!(bye_of_dialog(&messages, &local_tag) - answered >= T1X64 - T1);
assert!(matches!(
terminated_reason(&mut states),
Some(TerminatedReason::Timeout)
));
}
token.cancel();
Ok(())
}

#[tokio::test]
async fn test_unacked_reinvite_2xx_ends_the_session_on_a_legacy_client_dialog() -> crate::Result<()>
{
legacy_client_dialog_reinvite(false).await
}

#[tokio::test]
async fn test_acked_reinvite_keeps_a_legacy_client_dialog() -> crate::Result<()> {
legacy_client_dialog_reinvite(true).await
}
Loading