diff --git a/src/dialog/client_dialog.rs b/src/dialog/client_dialog.rs index c4b3fd8c..6d859054 100644 --- a/src/dialog/client_dialog.rs +++ b/src/dialog/client_dialog.rs @@ -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>) -> Result<()> { if !self.inner.is_confirmed() { if !self.inner.is_terminated() { @@ -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. diff --git a/src/dialog/dialog.rs b/src/dialog/dialog.rs index 9f62fc2a..46432a4e 100644 --- a/src/dialog/dialog.rs +++ b/src/dialog/dialog.rs @@ -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 diff --git a/src/dialog/invite_dialog.rs b/src/dialog/invite_dialog.rs index 72a3587b..8c185cd1 100644 --- a/src/dialog/invite_dialog.rs +++ b/src/dialog/invite_dialog.rs @@ -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>) -> Result<()> { let confirmed_or_waiting_ack = self.inner.is_confirmed() || (self.role() == TransactionRole::Server && self.inner.waiting_ack()); @@ -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. diff --git a/src/dialog/server_dialog.rs b/src/dialog/server_dialog.rs index 3456600f..c4e09f75 100644 --- a/src/dialog/server_dialog.rs +++ b/src/dialog/server_dialog.rs @@ -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>) -> Result<()> { if !self.inner.is_confirmed() && !self.inner.waiting_ack() { if !self.inner.is_terminated() { @@ -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. diff --git a/src/dialog/tests/test_client_dialog.rs b/src/dialog/tests/test_client_dialog.rs index 9b7838c4..b487a139 100644 --- a/src/dialog/tests/test_client_dialog.rs +++ b/src/dialog/tests/test_client_dialog.rs @@ -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; @@ -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); + +#[async_trait] +impl TargetLocator for FlakyLocator { + async fn locate(&self, uri: &Uri) -> crate::Result { + 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: \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(()) +}