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
8 changes: 4 additions & 4 deletions src/dialog/client_dialog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -157,6 +157,9 @@ impl ClientInviteDialog {
/// # Returns
/// * `Ok(())` - BYE was sent successfully or dialog is already terminated.
/// * `Err(Error)` - Failed to build/send BYE request, or dialog is in a state where BYE does not apply.
///
/// Once the BYE is handed to its transaction the dialog is `Terminated`,
/// even when an error is returned (RFC 3261 §15.1.1).
pub async fn bye_with_headers(&self, headers: Option<Vec<crate::sip::Header>>) -> Result<()> {
if !self.inner.is_confirmed() {
if !self.inner.is_terminated() {
Expand All @@ -178,10 +181,7 @@ impl ClientInviteDialog {
self.inner
.make_request(crate::sip::Method::Bye, None, None, None, headers, None)?;

self.inner.do_request(request).await?;
self.inner
.transition(DialogState::Terminated(self.id(), TerminatedReason::UacBye))?;
Ok(())
self.inner.send_bye(request, TerminatedReason::UacBye).await
}

/// Send a BYE request with a SIP `Reason` header.
Expand Down
11 changes: 11 additions & 0 deletions src/dialog/dialog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1312,6 +1312,17 @@ impl DialogInner {
self.send_dialog_request(request).boxed().await
}

/// Send the BYE that ends this dialog. RFC 3261 §15.1.1: the session ends
/// once the BYE is handed to its transaction, and a 481, a 408 or no
/// response ends the dialog. So whatever the transaction returns, the
/// dialog is `Terminated(reason)`; a failure is still returned.
pub(super) async fn send_bye(&self, request: Request, reason: TerminatedReason) -> Result<()> {
let result = self.do_request(request).await;
let id = self.id.lock().clone();
self.transition(DialogState::Terminated(id, reason))?;
result.map(|_| ())
}

/// RFC 3261 §13.3.1.4: the server transaction of an INVITE or re-INVITE
/// retransmitted our 2xx (`answered_2xx`) for 64*T1 and ended without an
/// ACK. The dialog is terminated with [`TerminatedReason::Timeout`] and the
Expand Down
8 changes: 4 additions & 4 deletions src/dialog/invite_dialog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -422,6 +422,9 @@ impl InviteDialog {
/// # Returns
/// * `Ok(())` - BYE was sent successfully or dialog is already terminated.
/// * `Err(Error)` - Failed to build/send BYE request, or dialog is in a state where BYE does not apply.
///
/// Once the BYE is handed to its transaction the dialog is `Terminated`,
/// even when an error is returned (RFC 3261 §15.1.1).
pub async fn bye_with_headers(&self, headers: Option<Vec<Header>>) -> Result<()> {
let confirmed_or_waiting_ack = self.inner.is_confirmed()
|| (self.role() == TransactionRole::Server && self.inner.waiting_ack());
Expand All @@ -445,14 +448,11 @@ impl InviteDialog {
.inner
.make_request(Method::Bye, None, None, None, headers, None)?;

self.inner.do_request(request).await?;
let reason = match self.role() {
TransactionRole::Server => TerminatedReason::UasBye,
TransactionRole::Client => TerminatedReason::UacBye,
};
self.inner
.transition(DialogState::Terminated(self.id(), reason))?;
Ok(())
self.inner.send_bye(request, reason).await
}

/// Send a BYE request with a SIP `Reason` header.
Expand Down
8 changes: 4 additions & 4 deletions src/dialog/server_dialog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -351,6 +351,9 @@ impl ServerInviteDialog {
/// # Returns
/// * `Ok(())` - BYE was sent successfully or dialog is already terminated.
/// * `Err(Error)` - Failed to build/send BYE request, or dialog is in a state where BYE does not apply.
///
/// Once the BYE is handed to its transaction the dialog is `Terminated`,
/// even when an error is returned (RFC 3261 §15.1.1).
pub async fn bye_with_headers(&self, headers: Option<Vec<crate::sip::Header>>) -> Result<()> {
if !self.inner.is_confirmed() && !self.inner.waiting_ack() {
if !self.inner.is_terminated() {
Expand All @@ -372,10 +375,7 @@ impl ServerInviteDialog {
self.inner
.make_request(crate::sip::Method::Bye, None, None, None, headers, None)?;

self.inner.do_request(request).await?;
self.inner
.transition(DialogState::Terminated(self.id(), TerminatedReason::UasBye))?;
Ok(())
self.inner.send_bye(request, TerminatedReason::UasBye).await
}

/// Send a BYE request with a SIP `Reason` header.
Expand Down
162 changes: 160 additions & 2 deletions src/dialog/tests/test_client_dialog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,15 +3,19 @@ use crate::dialog::{
dialog::{DialogInner, DialogState, TerminatedReason},
DialogId,
};
use crate::sip::{headers::*, prelude::HeadersExt, Request, Response, StatusCode, Uri};
use crate::transaction::endpoint::TargetLocator;
use crate::sip::{headers::*, prelude::HeadersExt, Method, Request, Response, SipMessage};
use crate::sip::{StatusCode, Uri};
use crate::transaction::endpoint::{EndpointOption, TargetLocator};
use crate::transaction::key::TransactionRole;
use crate::transport::transport_layer::DomainResolver;
use crate::transport::SipConnection;
use crate::transport::{udp::UdpConnection, SipAddr, TransportLayer};
use crate::EndpointBuilder;
use async_trait::async_trait;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::time::Duration;
use tokio::net::UdpSocket;
use tokio::sync::mpsc::unbounded_channel;
use tokio::sync::oneshot;
use tokio_util::sync::CancellationToken;
Expand Down Expand Up @@ -1695,3 +1699,157 @@ async fn test_cancel_returns_error_on_transaction_timeout() -> crate::Result<()>

Ok(())
}

/// Resolves like the default rule until the flag is set, then fails.
struct FlakyLocator(Arc<AtomicBool>);

#[async_trait]
impl TargetLocator for FlakyLocator {
async fn locate(&self, uri: &Uri) -> crate::Result<SipAddr> {
if self.0.load(Ordering::SeqCst) {
return Err(crate::Error::Error("no route".to_string()));
}
SipAddr::try_from(uri)
}
}

/// Waits for a `method` request at `peer` and answers it with `status` (if any).
async fn answer_next(peer: &UdpSocket, method: Method, status: Option<&str>) -> crate::Result<()> {
let mut buf = vec![0u8; 4096];
loop {
let recv = tokio::time::timeout(Duration::from_secs(2), peer.recv_from(&mut buf));
let (len, from) = recv
.await
.unwrap_or_else(|_| panic!("no {method} request"))?;
let Ok(SipMessage::Request(req)) = SipMessage::try_from(&buf[..len]) else {
continue;
};
if req.method != method {
continue;
}
if let Some(status) = status {
let to = req.to_header()?.value().to_string();
let to = if to.contains(";tag=") {
to
} else {
format!("{to};tag=bob")
};
let resp = format!(
"SIP/2.0 {status}\r\nVia: {}\r\nFrom: {}\r\nTo: {to}\r\nCall-ID: {}\r\n\
CSeq: {}\r\nContact: <sip:bob@{}>\r\nContent-Length: 0\r\n\r\n",
req.via_header()?.value(),
req.from_header()?.value(),
req.call_id_header()?.value(),
req.cseq_header()?.value(),
peer.local_addr()?,
);
peer.send_to(resp.as_bytes(), from).await?;
}
return Ok(());
}
}

/// RFC 3261 §15.1.1: once the BYE is handed to its client transaction the
/// session is over, and a 481, a 408 or no response also ends the dialog.
/// Every outcome of `bye()` leaves the dialog `Terminated`, notified once.
#[tokio::test]
async fn test_bye_terminates_the_dialog_whatever_the_outcome() -> crate::Result<()> {
use crate::dialog::{authenticate::Credential, dialog_layer::DialogLayer};
use crate::dialog::{invitation::InviteOption, server_dialog::ServerInviteDialog};
// (case, answer to the BYE, the BYE has no route, bye() via: the
// deprecated wrappers are driven over the same dialog)
for (case, answer, no_route, via) in [
(
"481",
Some("481 Call/Transaction Does Not Exist"),
false,
"invite",
),
("408", Some("408 Request Timeout"), false, "invite"),
("no response", None, false, "invite"),
(
"401 without a challenge",
Some("401 Unauthorized"),
false,
"invite",
),
("no route", None, true, "invite"),
("no route", None, true, "client"),
("no route", None, true, "server"),
] {
let token = CancellationToken::new();
let peer = UdpSocket::bind("127.0.0.1:0").await?;
let no_route_flag = Arc::new(AtomicBool::new(false));
let tl = TransportLayer::new(token.child_token());
let udp = UdpConnection::create_connection("127.0.0.1:0".parse()?, None, None).await?;
let local = udp.get_addr().get_socketaddr()?;
tl.add_transport(udp.into());
let endpoint = EndpointBuilder::new()
.with_transport_layer(tl)
.with_cancel_token(token.child_token())
.with_target_locator(Box::new(FlakyLocator(no_route_flag.clone())))
.with_option(EndpointOption {
t1: Duration::from_millis(10),
t1x64: Duration::from_millis(640),
..Default::default()
})
.build();
let inner = endpoint.inner.clone();
tokio::spawn(async move { inner.serve().await });
let (state_sender, mut states) = unbounded_channel();
let invite = InviteOption {
caller: Uri::try_from("sip:alice@example.com")?,
callee: Uri::try_from(format!("sip:bob@{}", peer.local_addr()?).as_str())?,
contact: Uri::try_from(format!("sip:alice@{local}").as_str())?,
// Answers a 401 to the BYE, which carries no challenge.
credential: Some(Credential {
username: "alice".into(),
password: "secret".into(),
realm: None,
auth_username: None,
}),
..Default::default()
};
let layer = DialogLayer::new(endpoint.inner.clone());
let invite = tokio::spawn(async move { layer.do_invite(invite, state_sender).await });
answer_next(&peer, Method::Invite, Some("200 OK")).await?;
let (dialog, _) = invite.await.unwrap()?;
while states.try_recv().is_ok() {}

no_route_flag.store(no_route, Ordering::SeqCst);
let inner = dialog.inner.clone();
let state = dialog.inner.clone();
let bye = tokio::spawn(async move {
match via {
"client" => ClientInviteDialog { inner }.bye().await,
"server" => ServerInviteDialog { inner }.bye().await,
_ => dialog.bye().await,
}
});
if !no_route {
answer_next(&peer, Method::Bye, answer).await?;
}
let result = tokio::time::timeout(Duration::from_secs(3), bye)
.await
.unwrap_or_else(|_| panic!("{case} via {via}: bye() hangs"))
.unwrap();
let mut reasons = Vec::new();
while let Ok(state) = states.try_recv() {
if let DialogState::Terminated(_, reason) = state {
reasons.push(reason);
}
}
let expected = if via == "server" {
"[UasBye]"
} else {
"[UacBye]"
};
assert_eq!(format!("{reasons:?}"), expected, "{case} via {via}");
assert!(state.is_terminated(), "{case} via {via}");
// The BYE's failure is still reported.
let failed = no_route || answer == Some("401 Unauthorized");
assert_eq!(result.is_err(), failed, "{case} via {via}: {result:?}");
token.cancel();
}
Ok(())
}
Loading