diff --git a/Cargo.toml b/Cargo.toml index cf7af89b..2cf9dd1a 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "rsipstack" -version = "0.7.1" +version = "0.7.3" edition = "2021" description = "SIP Stack Rust library for building SIP applications" license = "MIT" diff --git a/README.md b/README.md index 3807981b..a1eac531 100644 --- a/README.md +++ b/README.md @@ -7,7 +7,7 @@ A RFC 3261/3262 compliant SIP stack written in Rust. The goal of this project is ## Features - **RFC 3261/3262 Compliant**: Full compliance with SIP specification -- **Multiple Transport Support**: UDP, TCP, TLS, WebSocket (TLS/WebSocket require the `rustls` and `websocket` features, enabled by default) +- **Multiple Transport Support**: UDP, TCP, TLS, WebSocket (TLS/WebSocket require the `rustls` and `websocket` features, enabled by default). The TLS client goes through a backend seam (`platform::tls`) — host default is the built-in rustls; embedded stacks can register their own connector - **Transaction Layer**: Complete SIP transaction state machine - **Dialog Layer**: SIP dialog management - **Reliable Provisionals**: PRACK (RFC 3262 / 100rel) support @@ -15,6 +15,36 @@ A RFC 3261/3262 compliant SIP stack written in Rust. The goal of this project is - **High Performance**: Built with Rust for maximum performance - **Easy to Use**: Simple and intuitive API design +## no_std / Embedded Support + +rsipstack compiles without `std` (`alloc`-only) for embassy-based targets — +verified on `xtensa-esp32s3-none-elf` (ESP32-S3): + +```bash +cargo check --no-default-features --features platform-embassy +``` + +- Core SIP codec, transaction, dialog, and **UDP transport** layers are + `no_std` + `alloc`. +- The `platform-embassy` backend maps the runtime seams onto + `embassy-time` / `embassy-sync` (critical-section based); task spawning is + injected once at startup: + + ```rust,ignore + rsipstack::platform::set_spawn_fn(|fut| spawner.spawn(fut).ok()); + ``` + +- Concurrent waits use the backend-agnostic `select2` / `select3` helpers + instead of `tokio::select!`. +- **SIPS (SIP over TLS, client role)** goes through the + `platform::tls::TlsConnector` seam: register a backend with + `platform::tls::set_client_connector` and every + `TlsConnection::connect` rides it. Host builds without a registered + connector keep using the built-in rustls path. +- TCP / WebSocket transports still require `platform-tokio`; on embedded, + use UDP or SIPS today. +- On no_std, register a `DomainResolver` instead of the tokio DNS resolver. + ## TODO - [x] Transport support - [x] UDP @@ -25,6 +55,7 @@ A RFC 3261/3262 compliant SIP stack written in Rust. The goal of this project is - [x] Transaction Layer - [x] Dialog Layer - [ ] WASM target +- [x] no_std embassy-based targets ## Use Cases diff --git a/RECURSIVECX.md b/RECURSIVECX.md index 420adff9..a97c8f55 100644 --- a/RECURSIVECX.md +++ b/RECURSIVECX.md @@ -6,10 +6,14 @@ below is either offered upstream or removed once rcx stops depending on it. The row ids (R2, R3, ...) are those of the convergence ledger (`rsipstack-contrib/convergence-ledger.md`). -Base: upstream **0.7.1** (3286e8c). Everything the fork used to carry that -0.7.1 contains (R1, R4-R11, R13, the R12 teardown) is upstream's version now. -`ReinviteAck` / `take_reinvite_ack()` (R14) is gone: rcx reads -`InviteDialog::last_remote_ack()` and matches its CSeq. +Base: upstream **0.7.3** (2dfa0d1). Everything the fork used to carry that +0.7.3 contains is upstream's version now: R1, R4-R11, R13, the R12 teardown +(0.7.1), and R2, R3, R21, R22b, R23 (client lookup), the "BYE ends the dialog +whatever the outcome" part of R24, and R27 (0.7.3). `ReinviteAck` / +`take_reinvite_ack()` (R14) is gone: rcx reads +`InviteDialog::last_remote_ack()` and matches its CSeq. The R23 server-side +role check on `get_or_create_server_invite` was dropped: rcx routes +in-dialog requests through `match_dialog` / `get_dialog`. The fork has one role-agnostic `InviteDialog` (`Dialog::Invite`), as upstream. The deprecated `ClientInviteDialog` / `ServerInviteDialog` wrappers carry the @@ -17,29 +21,17 @@ same patches where they apply. ## Patches -### R2 + R19: a 1xx to an in-dialog request is not `Early` +### R19: a 1xx to an in-dialog request is not notified as `Early` - **Where:** `DialogInner::send_dialog_request`, `Provisional` arm. - **What:** `Early` is applied and notified only while the dialog can still be - cancelled (Calling / Trying / Early). A 1xx to a re-INVITE or UPDATE on a - confirmed dialog neither regresses it (R2, upstream PR #147) nor notifies - `Early` (R19), which rcx subscribers read as ringing. + cancelled (Calling / Trying / Early). Upstream (#147, R2) no longer + regresses a confirmed dialog but still notifies the 1xx as `Early`; the + fork does not notify it, since rcx subscribers read `Early` as ringing. - **rcx:** the callee-state handlers that match `Early`; the lease-expiry fence. - **Test:** `dialog::tests::test_in_dialog_provisional`. -### R3: one notification per applied dialog transition - -- **Where:** `DialogInner::transition`. -- **What:** the transition is decided and notified under the state lock. - After `Terminated`, nothing more is applied or notified, a second - `Terminated` included; a `WaitAck` ignored after `Confirmed` is not - notified. Event-only states (`Updated`, `Notify`, `Info`, `Options`, - `Refer`) are notified whatever the lifecycle state. Upstream PR #148. -- **rcx:** the callee-state `Terminated` arm (`sip_session/dialog_events.rs`). -- **Tests:** `dialog::tests::test_state_after_terminated`, the notification - tests at the end of `dialog::tests::test_dialog_states`. - ### R15: response provenance (`Response::synthetic`, `Response::received_from`) - **Where:** `sip::message::{Response, ReceivedFrom}`. `synthetic` is `true` on @@ -50,7 +42,7 @@ same patches where they apply. serialized; `None` / `false` after a reparse. - **rcx:** `callrecord/carrier_response.rs`, `callrecord/diagnostics.rs`, `rcx-call/src/sip.rs`. Struct literals must set - `synthetic: false, received_from: None`. + `synthetic: false, received_from: None` (and `wire_reason: None`, R28). - **Tests:** `transaction::tests::test_response_provenance`, `test_stream_reconnect::test_send_failure_on_stream_is_reported_at_once`. @@ -100,43 +92,21 @@ same patches where they apply. `crates/rcx-call/src/sip.rs`. - **Test:** `dialog_layer::test_take_dialog_hands_the_dialog_to_exactly_one_caller`. -### R21: a 2xx crossing a taken dialog's CANCEL is BYE'd - -- **Where:** `DialogGuardForUnconfirmed` (`src/dialog/invitation.rs`): the - `dialog` / `finished` fields and `watch_taken_dialog`. -- **What:** when another owner took the dialog out of the layer - (`take_dialog`) and the `do_invite` future is dropped, the guard keeps the - INVITE transaction: in Trying / Early it watches for a 2xx (up to 64*T1) - and BYEs it, without a second CANCEL; still in Calling (the owner's hangup - could send no CANCEL), it abandons the INVITE as upstream does for a dialog - still in the layer. -- **rcx:** RWI originate `Hangup` arm, `rwi_originate_trunk_e2e_test`. -- **Tests:** `test_cancel_2xx_race::test_taken_dialog_*`. - -### R22a + R22b: raw SIP messages at DEBUG only +### R22a: raw SIP messages at DEBUG only - **What:** UDP, stream and WebSocket raw send/receive logs are DEBUG - (upstream: INFO). The WebSocket parse failure and the dialog layer's - "failed to send request" WARNs carry only the length or the method; the - message itself is logged at DEBUG. -- **Test:** `transport::tests::test_raw_message_log_level` (`bench` feature). - -### R23: role-typed dialog lookups - -- **What:** `DialogLayer::get_client_dialog_by_call_id` returns UAC dialogs - only, and `get_or_create_server_invite` matches existing UAS dialogs only. - With a transparent Call-ID the inbound (UAS) and outbound (UAC) legs of a - proxied call share it. -- **rcx:** `rwi/processor/originate.rs` (Hangup, media timeout, transfer, DTMF). -- **Test:** `dialog_layer::test_lookups_keep_uac_and_uas_dialogs_apart`. - -### R24: BYE lifecycle per role - -- **Where:** `DialogInner::send_bye`, used by `InviteDialog` and both wrappers. -- **What:** a UAS notifies `Terminated(UasBye)` before sending the BYE and - returns the send result; a UAC notifies `Terminated(UacBye)` after the BYE - transaction whatever its outcome. Upstream terminates only after an `Ok` - BYE. + (upstream: INFO). The WARNs without the message (R22b) are upstream's. +- **Tests:** `transport::tests::test_raw_message_log_level` (`bench` feature); + upstream's `test_warn_logs` checks the WebSocket receive log at DEBUG. + +### R24: a UAS notifies `Terminated` before sending its BYE + +- **Where:** `DialogInner::send_bye`. +- **What:** a UAS notifies `Terminated(UasBye)` before the BYE is sent, so + subscribers never wait for the BYE's response, and returns the send + result. Upstream (0.7.3) notifies after the BYE transaction for both roles; + the UAC side is upstream's (`Terminated(UacBye)` whatever the outcome, the + error still returned). - **rcx:** the `Terminated` arms in `sip_session/dialog_events.rs`. - **Test:** `dialog::tests::test_bye_lifecycle`. @@ -146,21 +116,16 @@ same patches where they apply. `@`, then `callid_suffix`. rcx sets it in `proxy/server.rs`. - **Test:** `transaction::tests::tests::test_make_call_id_random22`. -### R27: `Confirmed` carries the 2xx its ACK confirms +### R28: `Response::wire_reason` -- **Where:** `Transaction::cleanup`: a server INVITE transaction keeps - `last_response` (it still hands a copy to `finished_transactions`). -- **What:** since 0.7.1 the matching ACK ends an `Accepted` server INVITE - transaction (RFC 6026 §7.1, upstream #169) before the dialog reads - `tx.last_response` for `DialogState::Confirmed`. Upstream then notifies - `Confirmed` with `Response::default()` (no CSeq, no headers), for the - initial INVITE and every re-INVITE. rcx correlates `Confirmed` by the - response's CSeq (`confirms_initial_invite`, `is_reconfirmation`, the - re-offer settle). -- **Test:** `dialog::tests::test_late_reinvite_ack` (CSeq of each - `Confirmed`). +- **Where:** `sip::message::Response`, `sip::parser`. +- **What:** a known status code whose Status-Line carries a phrase other than + the standard one keeps it in `wire_reason` (e.g. `403 Caller Origination + Number is Invalid`). Set only by the parser; `Display` is unchanged. +- **rcx:** `crates/rcx-callrecord/src/carrier_response.rs` (rcx #1135). +- **Test:** `sip::parser` tests for the custom phrase. -## Upstream 0.7.1 behavior rcx observes (not fork patches) +## Upstream behavior rcx observes (not fork patches) - RFC 6026 `Accepted` state. Server: the 2xx is retransmitted by Timer G (T1 doubling to T2, every transport) until the matching ACK, which ends the @@ -172,3 +137,11 @@ same patches where they apply. - A dropped INVITE sends no CANCEL before a provisional response; it CANCELs on the first provisional, or ACKs and BYEs a 2xx (#162, #171). - A forked 2xx is ACKed with its own To tag and remote target (#172). +- 0.7.3: a BYE ends the dialog whatever its transaction returns (#180); a + forked 2xx's dialog is BYE'd (#181); an in-dialog REFER returns the dialog + to `Confirmed` once answered; Timer F ends a non-INVITE client + transaction in Proceeding (#189); a server in-dialog request with no route + and no dial-back fails at once (#187); the deprecated `ClientInviteDialog` + ends the session on a never-ACKed re-INVITE 2xx (#185); a dropped INVITE + whose dialog was already removed still ends (#183); injectable TLS client + seam (rustls stays the default). diff --git a/src/dialog/client_dialog.rs b/src/dialog/client_dialog.rs index 0ccc1350..e6a23fe5 100644 --- a/src/dialog/client_dialog.rs +++ b/src/dialog/client_dialog.rs @@ -157,16 +157,19 @@ 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() { warn!( dialog_id = %self.id(), - state = ?self.state(), + state = %self.state(), "bye skipped: dialog not confirmed" ); return Err(crate::Error::Error(format!( - "dialog {} cannot send BYE in state {:?}", + "dialog {} cannot send BYE in state {}", self.id(), self.state() ))); @@ -178,7 +181,7 @@ impl ClientInviteDialog { self.inner .make_request(crate::sip::Method::Bye, None, None, None, headers, None)?; - self.inner.send_bye(request).await + self.inner.send_bye(request, TerminatedReason::UacBye).await } /// Send a BYE request with a SIP `Reason` header. @@ -675,6 +678,11 @@ 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 { @@ -682,11 +690,15 @@ impl ClientInviteDialog { 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(()) } @@ -696,7 +708,12 @@ impl ClientInviteDialog { self.inner .transition(DialogState::Refer(self.id(), tx.original.clone(), handle))?; - self.inner.process_transaction_handle(tx, rx).await + // RFC 3515: the REFER was answered (usually 202) and the dialog must + // go back to Confirmed — the implicit subscription's NOTIFYs are + // in-dialog requests that need the confirmed dialog. + let result = self.inner.process_transaction_handle(tx, rx).await; + let confirmed = self.return_to_confirmed(tx); + result.and(confirmed) } async fn handle_message(&mut self, tx: &mut Transaction) -> Result<()> { diff --git a/src/dialog/dialog.rs b/src/dialog/dialog.rs index 88d5106a..1a437e9b 100644 --- a/src/dialog/dialog.rs +++ b/src/dialog/dialog.rs @@ -715,7 +715,7 @@ impl DialogInner { warn!( id = self.id.lock().to_string(), destination = tx.destination.as_ref().map(|d| d.to_string()).as_deref(), - method = %tx.original.method, + method = %method, "failed to send request error: {}", e ); @@ -1187,6 +1187,7 @@ impl DialogInner { } } let need_fallback_retry; + let mut send_error = None; match tx.send().await { Ok(_) => { debug!( @@ -1214,6 +1215,7 @@ impl DialogInner { debug!(id = self.id.lock().to_string(), req = %tx.original, "request that failed to send"); return Err(e); } + send_error = Some(e); } } @@ -1285,6 +1287,11 @@ impl DialogInner { method = %method, "no usable connection and no dial-back target; giving up after first send" ); + // The failed send never started the transaction: no response + // or timer will ever end it, so report the error now. + if let Some(e) = send_error { + return Err(e); + } } } @@ -1309,9 +1316,9 @@ impl DialogInner { // and never back. A 1xx to a mid-dialog request (re-INVITE, // UPDATE, ...) must not regress an established dialog to // Early, or BYE is refused and hangup() tries to CANCEL. - // Nor is it notified: subscribers treat `Early` as the - // dialog's early state (ringing), not as a provisional - // to a later transaction. + // Nor is it notified (R19, fork-only): subscribers + // treat `Early` as the dialog's early state (ringing), + // not as a provisional to a later transaction. if self.can_cancel() { self.transition(DialogState::Early(self.id.lock().clone(), resp))?; } @@ -1382,28 +1389,84 @@ impl DialogInner { self.send_dialog_request(request).boxed().await } - /// Send the BYE that ends this dialog and notify `Terminated`, with the - /// lifecycle subscribers rely on to finish their teardown: + /// 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. /// - /// * UAS: `Terminated(UasBye)` is notified before the BYE is sent, so it - /// never waits for the BYE's response; the send result is returned. - /// * UAC: the BYE transaction runs, then `Terminated(UacBye)` is notified - /// whatever its outcome. A failed send is logged and the dialog still - /// ends locally; `Ok` is returned. - pub(super) async fn send_bye(&self, request: Request) -> Result<()> { + /// Fork (R24): a UAS notifies `Terminated` before the BYE is sent, so its + /// subscribers never wait for the BYE's response. A UAC notifies it after + /// the BYE transaction, as upstream does. + pub(super) async fn send_bye(&self, request: Request, reason: TerminatedReason) -> Result<()> { let id = self.id.lock().clone(); - match self.role { - TransactionRole::Server => { - self.transition(DialogState::Terminated(id, TerminatedReason::UasBye))?; - self.do_request(request).await.map(|_| ()) - } - TransactionRole::Client => { - if let Err(e) = self.do_request(request).await { - info!(%id, error = %e, "bye error, ending the dialog locally"); - } - self.transition(DialogState::Terminated(id, TerminatedReason::UacBye)) - } + if self.role == TransactionRole::Server { + self.transition(DialogState::Terminated(id, reason))?; + return self.do_request(request).await.map(|_| ()); } + let result = self.do_request(request).await; + self.transition(DialogState::Terminated(id, reason))?; + result.map(|_| ()) + } + + /// End a dialog a forked 2xx established (RFC 3261 §13.2.2.4). + /// + /// Every 2xx to the INVITE with a new To tag creates its own dialog; the + /// UAC keeps a single session (the first 2xx's), so the transaction has + /// already ACKed the forked 2xx and this sends the BYE that terminates + /// the extra branch. The BYE is built from that response's own remote + /// target (Contact) and route set (Record-Route); the confirmed dialog's + /// state is not touched and no dialog is registered for the branch. + pub(super) async fn bye_forked_branch(&self, resp: &Response) -> Result<()> { + let contact_uri = resp + .typed_contact_headers()? + .first() + .map(|c| c.uri.clone()) + .ok_or_else(|| crate::Error::Error("missing Contact header".to_string()))?; + + // §12.2.1.1: the forked dialog's route set is the 2xx's Record-Route. + let mut routes: Vec = resp + .record_route_headers() + .into_iter() + .flat_map(|rr| split_rr_values(rr.value())) + .map(Route::from) + .collect(); + routes.reverse(); + + // To carries the forked branch's tag, From keeps ours. + let to = resp.to_header()?.clone(); + let id = self.id.lock().clone(); + let via = self + .endpoint_inner + .get_via(self.via_addr_for_send_transport(), None)?; + let cseq = CSeq { + seq: self.increment_local_seq(), + method: Method::Bye, + }; + + let mut headers: Vec
= vec![ + Header::Via(via.into()), + Header::CallId(id.call_id.clone().into()), + Header::From(self.from.clone().to_string().into()), + Header::To(to), + Header::CSeq(cseq.into()), + Header::UserAgent(self.endpoint_inner.user_agent.clone().into()), + ]; + if let Some(uri) = self.local_contact.as_ref() { + headers.push(Contact::from(uri.clone()).into()); + } + headers.extend(routes.into_iter().map(Header::Route)); + headers.push(Header::MaxForwards(70.into())); + + debug!(id = %id, uri = %contact_uri, "sending BYE to a forked dialog"); + self.do_request(crate::sip::Request { + method: Method::Bye, + uri: contact_uri, + headers: headers.into(), + body: Vec::new(), + version: crate::sip::Version::V2, + }) + .await?; + Ok(()) } /// RFC 3261 §13.3.1.4: the server transaction of an INVITE or re-INVITE @@ -1677,7 +1740,6 @@ impl DialogInner { Ok(()) } - #[cfg_attr(not(feature = "platform-tokio"), allow(unused_mut))] #[cfg_attr(not(feature = "platform-tokio"), allow(unused_mut))] pub async fn process_transaction_handle( &self, diff --git a/src/dialog/dialog_layer.rs b/src/dialog/dialog_layer.rs index 6b0c3591..001fab78 100644 --- a/src/dialog/dialog_layer.rs +++ b/src/dialog/dialog_layer.rs @@ -156,12 +156,7 @@ impl DialogLayer { if !id.local_tag.is_empty() { let dlg = self.inner.dialogs.get(&id.to_string()); match dlg { - // Only a UAS dialog answers an in-dialog request here, as the - // role-typed `Dialog::ServerInvite` did before the unified - // `InviteDialog`. - Some(Dialog::Invite(dlg)) if dlg.role() == TransactionRole::Server => { - return Ok(dlg) - } + Some(Dialog::Invite(dlg)) => return Ok(dlg), _ => { return Err(crate::Error::DialogError( "the dialog not found".to_string(), @@ -458,8 +453,6 @@ impl DialogLayer { self.inner.dialogs.with(|m| { m.values() .filter_map(|d| match d { - // UAC dialogs only: with a transparent Call-ID, the inbound - // (UAS) leg of a proxied call shares it. Dialog::Invite(client_dlg) if client_dlg.role() == TransactionRole::Client && client_dlg.id().call_id == call_id => diff --git a/src/dialog/invitation.rs b/src/dialog/invitation.rs index af8f39da..cc99dc92 100644 --- a/src/dialog/invitation.rs +++ b/src/dialog/invitation.rs @@ -187,33 +187,8 @@ pub(super) struct DialogGuardForUnconfirmed<'a> { pub dialog_layer_inner: &'a DialogLayerInnerRef, pub id: &'a DialogId, invite_tx: Option, - /// The dialog itself, so a 2xx crossing a CANCEL sent by another owner - /// (one that took the dialog out of the layer) can still be BYE'd. - dialog: InviteDialog, - /// `process_invite` ran to completion: nothing left to watch. - finished: bool, -} - -impl DialogGuardForUnconfirmed<'_> { - /// Dropped mid-INVITE after another owner took the dialog out of the - /// layer (`DialogLayer::take_dialog`) to end it. That owner does not hold - /// the INVITE transaction, so a 2xx crossing its CANCEL would be ACKed - /// by the transaction and never BYE'd. Keep the transaction and watch it - /// (up to 64*T1) for a 2xx to BYE. No second CANCEL is sent: the owner - /// sent it. - fn watch_taken_dialog(&mut self, client_dialog: InviteDialog) { - let Some(mut invite_tx) = self.invite_tx.take() else { - return; - }; - debug!(id = %client_dialog.id(), "taken dialog dropped mid-INVITE, watching for a 2xx"); - crate::platform::spawn(async move { - invite_tx.stop_retransmissions(); - let window = invite_tx.endpoint_inner.option.t1x64; - let final_response = wait_response(&mut invite_tx, window, true).await; - drop(invite_tx); - bye_if_2xx(&client_dialog, final_response).await; - }); - } + /// The INVITE's dialog, `None` once `process_invite` returned. + dialog: Option, } impl<'a> Drop for DialogGuardForUnconfirmed<'a> { @@ -221,24 +196,13 @@ impl<'a> Drop for DialogGuardForUnconfirmed<'a> { let client_dialog = match self.dialog_layer_inner.dialogs.remove(&self.id.to_string()) { Some(Dialog::Invite(client_dialog)) => client_dialog, Some(_) => return, - // Another owner took the dialog out of the layer to end it. - None => { - if self.finished { - return; - } - let client_dialog = self.dialog.clone(); - match client_dialog.state() { - // The owner's hangup sent nothing (no CANCEL before a - // provisional response): abandon the INVITE as if the - // dialog were still in the layer. - DialogState::Calling(_) => client_dialog, - DialogState::Trying(_) | DialogState::Early(_, _) => { - self.watch_taken_dialog(client_dialog); - return; - } - _ => return, - } - } + // The application already removed the dialog from the layer + // (`DialogLayer::remove_dialog`): its INVITE still has to end. + None => match self.dialog.take() { + Some(client_dialog) => client_dialog, + // `process_invite` returned: `do_invite` handles the outcome. + None => return, + }, }; match client_dialog.state() { @@ -730,8 +694,7 @@ impl DialogLayer { dialog_layer_inner: &self.inner, id: &id, invite_tx: Some(tx), - dialog: dialog.clone(), - finished: false, + dialog: Some(dialog.clone()), }; let tx = guard @@ -740,7 +703,7 @@ impl DialogLayer { .expect("transcation should be avaible"); let r = dialog.process_invite(tx).boxed().await; - guard.finished = true; + guard.dialog = None; self.inner.dialogs.remove(&id.to_string()); match r { @@ -764,8 +727,14 @@ impl DialogLayer { // here would leave it (and its timers) in the // endpoint's table and silently stop the re-ACKs. if let Some(mut tx) = guard.invite_tx.take() { + let dlg = dialog.clone(); + let confirmed_tag = new_dialog_id.remote_tag.clone(); crate::platform::spawn(async move { - while tx.receive().await.is_some() {} + let mut seen_forks: Vec = Vec::new(); + while let Some(msg) = tx.receive().await { + dlg.end_forked_branch(&msg, &confirmed_tag, &mut seen_forks) + .await; + } debug!(id = %new_dialog_id, "accepted transaction drained (Timer M expired)"); }); } @@ -840,8 +809,15 @@ impl DialogLayer { // observes forked 2xx, and detaches it from the // endpoint's table). See do_invite for the rationale. let confirmed_id = new_id.clone(); + let confirmed_tag = new_id.remote_tag.clone(); + let forked_dlg = dialog_clone.clone(); crate::platform::spawn(async move { - while tx.receive().await.is_some() {} + let mut seen_forks: Vec = Vec::new(); + while let Some(msg) = tx.receive().await { + forked_dlg + .end_forked_branch(&msg, &confirmed_tag, &mut seen_forks) + .await; + } debug!(id = %confirmed_id, "accepted transaction drained (Timer M expired)"); }); } diff --git a/src/dialog/invite_dialog.rs b/src/dialog/invite_dialog.rs index 1c8b0846..023d66f6 100644 --- a/src/dialog/invite_dialog.rs +++ b/src/dialog/invite_dialog.rs @@ -22,7 +22,7 @@ use crate::transaction::key::TransactionRole; use crate::transaction::transaction::{Transaction, TransactionEvent}; use crate::Result; use core::sync::atomic::Ordering; -use tracing::{debug, trace, warn}; +use tracing::{debug, info, trace, warn}; /// Unified INVITE dialog that can act as either a UAS (Server) or UAC (Client). /// @@ -418,6 +418,55 @@ impl InviteDialog { Ok(()) } + /// End a forked dialog a later 2xx established (RFC 3261 §13.2.2.4). + /// + /// Called from the Accepted-window drainer for every message the client + /// INVITE transaction delivers after the dialog confirmed. A 2xx whose To + /// tag differs from the confirmed dialog's is a forked branch: the + /// transaction has already ACKed it, and the UAC — keeping a single + /// session — terminates it with a BYE. Everything else (retransmitted + /// 2xx with the same tag, non-2xx) is ignored. `seen` holds the forked + /// tags already BYE'd, so a retransmitted forked 2xx sends one BYE only. + pub(super) async fn end_forked_branch( + &self, + msg: &SipMessage, + confirmed_remote_tag: &str, + seen: &mut Vec, + ) { + let SipMessage::Response(resp) = msg else { + return; + }; + if resp.status_code.kind() != StatusCodeKind::Successful { + return; + } + // Same tag: a retransmission of the confirmed 2xx, already re-ACKed. + let tag = match resp + .to_header() + .ok() + .and_then(|to| to.tag().ok().flatten()) + .map(|tag| tag.value().to_string()) + { + Some(tag) => tag, + None => return, + }; + if tag == confirmed_remote_tag { + return; + } + if seen.iter().any(|seen| seen == &tag) { + return; + } + seen.push(tag.clone()); + let id = self.id(); + info!( + id = %id, + tag = %tag, + "forked 2xx acknowledged; ending the extra branch with a BYE (RFC 3261 §13.2.2.4)" + ); + if let Err(e) = self.inner.bye_forked_branch(resp).await { + warn!(id = %id, tag = %tag, error = %e, "failed to BYE the forked branch"); + } + } + // ── Shared request semantics ────────────────────────────────────────── /// Send a BYE request to terminate the dialog. @@ -433,14 +482,12 @@ impl InviteDialog { /// dialogs remain a silent no-op. /// /// # Returns - /// * `Ok(())` - BYE was sent (or, for a UAC, attempted) or the dialog is - /// already terminated. - /// * `Err(Error)` - Failed to build the BYE, the dialog is in a state - /// where BYE does not apply, or (UAS) the BYE could not be sent. + /// * `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. /// - /// `Terminated` is notified even when the BYE fails: a UAS notifies it - /// before sending, a UAC once the BYE transaction ends (see - /// `DialogInner::send_bye`). + /// Once the BYE is handed to its transaction the dialog is `Terminated`, + /// even when an error is returned (RFC 3261 §15.1.1). A UAS notifies it + /// before sending the BYE (fork, R24; see `DialogInner::send_bye`). 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()); @@ -448,11 +495,11 @@ impl InviteDialog { if !self.inner.is_terminated() { warn!( dialog_id = %self.id(), - state = ?self.state(), + state = %self.state(), "bye skipped: dialog not confirmed or waiting ack" ); return Err(crate::Error::Error(format!( - "dialog {} cannot send BYE in state {:?}", + "dialog {} cannot send BYE in state {}", self.id(), self.state() ))); @@ -464,7 +511,11 @@ impl InviteDialog { .inner .make_request(Method::Bye, None, None, None, headers, None)?; - self.inner.send_bye(request).await + let reason = match self.role() { + TransactionRole::Server => TerminatedReason::UasBye, + TransactionRole::Client => TerminatedReason::UacBye, + }; + self.inner.send_bye(request, reason).await } /// Send a BYE request with a SIP `Reason` header. @@ -822,7 +873,12 @@ impl InviteDialog { self.inner .transition(DialogState::Refer(self.id(), tx.original.clone(), handle))?; - self.inner.process_transaction_handle(tx, rx).await + // RFC 3515: the REFER was answered (usually 202) and the dialog must + // go back to Confirmed — the implicit subscription's NOTIFYs are + // in-dialog requests that need the confirmed dialog. + let result = self.inner.process_transaction_handle(tx, rx).await; + let confirmed = self.return_to_confirmed(tx); + result.and(confirmed) } async fn handle_message(&mut self, tx: &mut Transaction) -> Result<()> { diff --git a/src/dialog/server_dialog.rs b/src/dialog/server_dialog.rs index 07e41bc1..8087eb94 100644 --- a/src/dialog/server_dialog.rs +++ b/src/dialog/server_dialog.rs @@ -351,16 +351,19 @@ 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() { warn!( dialog_id = %self.id(), - state = ?self.state(), + state = %self.state(), "bye skipped: dialog not confirmed or waiting ack" ); return Err(crate::Error::Error(format!( - "dialog {} cannot send BYE in state {:?}", + "dialog {} cannot send BYE in state {}", self.id(), self.state() ))); @@ -372,7 +375,7 @@ impl ServerInviteDialog { self.inner .make_request(crate::sip::Method::Bye, None, None, None, headers, None)?; - self.inner.send_bye(request).await + self.inner.send_bye(request, TerminatedReason::UasBye).await } /// Send a BYE request with a SIP `Reason` header. diff --git a/src/dialog/tests/mod.rs b/src/dialog/tests/mod.rs index 658f0b6c..7535246d 100644 --- a/src/dialog/tests/mod.rs +++ b/src/dialog/tests/mod.rs @@ -7,13 +7,14 @@ mod test_client_dialog; mod test_connection_affinity; mod test_dialog_layer; mod test_dialog_states; +mod test_forked_2xx_bye; mod test_in_dialog_provisional; mod test_in_dialog_via; mod test_invite_auth_challenge; -mod test_late_reinvite_ack; mod test_prack; mod test_proxy_headers; mod test_refer; +mod test_refer_notify; mod test_server_dialog; mod test_session_id; mod test_state_after_terminated; @@ -21,3 +22,4 @@ mod test_sticky_transport; mod test_sub_pub; mod test_uac_ack_body; mod test_uas_ack_timeout; +mod test_warn_logs; diff --git a/src/dialog/tests/test_bye_lifecycle.rs b/src/dialog/tests/test_bye_lifecycle.rs index 828ece5b..5ddf7024 100644 --- a/src/dialog/tests/test_bye_lifecycle.rs +++ b/src/dialog/tests/test_bye_lifecycle.rs @@ -1,37 +1,23 @@ -//! `bye()` notifies `Terminated` with the lifecycle subscribers rely on to -//! finish their teardown: a UAC ends the dialog locally even when its BYE -//! cannot be sent, and a UAS notifies `Terminated` before its BYE's response. +//! Fork (R24): a UAS `bye()` notifies `Terminated` before the BYE is sent, +//! so subscribers never wait for the BYE's response. (A UAC terminates after +//! the BYE transaction whatever its outcome: upstream's behaviour.) use crate::dialog::{ dialog::{DialogInner, DialogState, DialogStateReceiver, TerminatedReason}, invite_dialog::InviteDialog, DialogId, }; use crate::sip::{Method, Request, Response, SipMessage, Uri}; -use crate::transaction::endpoint::{Endpoint, TargetLocator}; +use crate::transaction::endpoint::Endpoint; use crate::transaction::key::TransactionRole; -use crate::transport::{udp::UdpConnection, SipAddr, TransportLayer}; +use crate::transport::{udp::UdpConnection, TransportLayer}; use crate::EndpointBuilder; -use async_trait::async_trait; use std::sync::Arc; use std::time::Duration; use tokio::net::UdpSocket; use tokio::sync::mpsc::unbounded_channel; use tokio_util::sync::CancellationToken; -/// A locator that can resolve nothing: every request send fails. -struct FailingLocator; - -#[async_trait] -impl TargetLocator for FailingLocator { - async fn locate(&self, _uri: &Uri) -> crate::Result { - Err(crate::Error::Error("no route".to_string())) - } -} - -async fn endpoint( - token: &CancellationToken, - locator: Option>, -) -> crate::Result { +async fn endpoint(token: &CancellationToken) -> crate::Result { let tl = TransportLayer::new(token.child_token()); let udp = UdpConnection::create_connection( "127.0.0.1:0".parse().unwrap(), @@ -40,15 +26,11 @@ async fn endpoint( ) .await?; tl.add_transport(udp.into()); - let mut builder = EndpointBuilder::new(); - builder + let endpoint = EndpointBuilder::new() .with_user_agent("rsipstack-test") .with_transport_layer(tl) - .with_cancel_token(token.child_token()); - if let Some(locator) = locator { - builder.with_target_locator(locator); - } - let endpoint = builder.build(); + .with_cancel_token(token.child_token()) + .build(); let inner = endpoint.inner.clone(); tokio::spawn(async move { let _ = inner.serve().await; @@ -107,39 +89,10 @@ fn confirmed_dialog( Ok((InviteDialog::from_inner(Arc::new(inner)), states)) } -fn terminated(states: &mut DialogStateReceiver) -> Option { - let mut reason = None; - while let Ok(state) = states.try_recv() { - if let DialogState::Terminated(_, r) = state { - reason = Some(r); - } - } - reason -} - -#[tokio::test] -async fn test_uac_bye_that_cannot_be_sent_still_terminates() -> crate::Result<()> { - let token = CancellationToken::new(); - let endpoint = endpoint(&token, Some(Box::new(FailingLocator))).await?; - let (dialog, mut states) = - confirmed_dialog(&endpoint, TransactionRole::Client, "192.0.2.1:5060")?; - - tokio::time::timeout(Duration::from_secs(2), dialog.bye()) - .await - .expect("bye must not hang")?; - assert!(matches!( - terminated(&mut states), - Some(TerminatedReason::UacBye) - )); - assert!(dialog.state().is_terminated()); - token.cancel(); - Ok(()) -} - #[tokio::test] async fn test_uas_bye_notifies_terminated_before_its_response() -> crate::Result<()> { let token = CancellationToken::new(); - let endpoint = endpoint(&token, None).await?; + let endpoint = endpoint(&token).await?; // A peer that receives the BYE and never answers it. let peer = UdpSocket::bind("127.0.0.1:0").await?; let (dialog, mut states) = confirmed_dialog( diff --git a/src/dialog/tests/test_cancel_2xx_race.rs b/src/dialog/tests/test_cancel_2xx_race.rs index b52fb233..946a01cc 100644 --- a/src/dialog/tests/test_cancel_2xx_race.rs +++ b/src/dialog/tests/test_cancel_2xx_race.rs @@ -5,11 +5,15 @@ //! //! These tests drop the `do_invite` future after a 180 — the documented way to //! abandon an outgoing call, which cancels it — and answer the INVITE with a -//! 200 from a raw UDP peer in each wire ordering. +//! 200 from a raw UDP peer in each wire ordering. The `removed` variants call +//! `DialogLayer::remove_dialog` before dropping the future: the INVITE still +//! has to end (CANCEL, and ACK + BYE for a 2xx) even though the layer entry is +//! already gone. use crate::dialog::{ dialog::{DialogState, DialogStateReceiver, TerminatedReason}, dialog_layer::DialogLayer, invitation::InviteOption, + DialogId, }; use crate::sip::{prelude::HeadersExt, Method, Request, SipMessage, Uri}; use crate::transport::{udp::UdpConnection, TransportLayer}; @@ -181,7 +185,19 @@ enum Order { InviteOkLate, } -async fn run_crossing_2xx(order: Order) -> crate::Result<()> { +/// The id the `Calling` state reports, the one the dialog is registered under. +async fn calling_id(states: &mut DialogStateReceiver) -> DialogId { + match wait_for_state(states, "Calling", Duration::from_secs(2), |s| { + matches!(s, DialogState::Calling(_)) + }) + .await + { + DialogState::Calling(id) => id, + _ => unreachable!(), + } +} + +async fn run_crossing_2xx(order: Order, provisional: u16, remove: bool) -> crate::Result<()> { let token = CancellationToken::new(); let Uac { dialog_layer, @@ -191,15 +207,27 @@ async fn run_crossing_2xx(order: Order) -> crate::Result<()> { let wait = Duration::from_secs(2); let (state_sender, mut states) = unbounded_channel(); - let invite = tokio::spawn(async move { dialog_layer.do_invite(option, state_sender).await }); + let layer = dialog_layer.clone(); + let invite = tokio::spawn(async move { layer.do_invite(option, state_sender).await }); let (inv, uac) = recv_request(&peer, Method::Invite, wait).await; - reply(&peer, uac, &inv, 180, "Ringing").await; - wait_for_state(&mut states, "Early", wait, |s| { - matches!(s, DialogState::Early(_, _)) + let id = calling_id(&mut states).await; + let reason = if provisional == 100 { + "Trying" + } else { + "Ringing" + }; + reply(&peer, uac, &inv, provisional, reason).await; + wait_for_state(&mut states, "Trying or Early", wait, |s| { + matches!(s, DialogState::Trying(_) | DialogState::Early(_, _)) }) .await; + if remove { + // The application removes the dialog from the layer first. + dialog_layer.remove_dialog(&id); + assert!(dialog_layer.is_empty()); + } // The application abandons the call: dropping the `do_invite` future // cancels the INVITE. invite.abort(); @@ -296,23 +324,39 @@ async fn run_crossing_2xx(order: Order) -> crate::Result<()> { #[tokio::test] async fn test_2xx_before_cancel_response_is_acked_and_byed() -> crate::Result<()> { - run_crossing_2xx(Order::InviteOkFirst).await + run_crossing_2xx(Order::InviteOkFirst, 180, false).await } #[tokio::test] async fn test_2xx_after_cancel_response_is_acked_and_byed() -> crate::Result<()> { - run_crossing_2xx(Order::CancelOkFirst).await + run_crossing_2xx(Order::CancelOkFirst, 180, false).await } #[tokio::test] async fn test_2xx_after_cancel_settle_window_is_acked_and_byed() -> crate::Result<()> { - run_crossing_2xx(Order::InviteOkLate).await + run_crossing_2xx(Order::InviteOkLate, 180, false).await } -/// A CANCEL that wins the race (487 to the INVITE) ends the call as before: -/// the 487 is ACKed and no BYE is sent. +/// The same when the application removed the dialog from the layer before +/// dropping the `do_invite` future, in Early (180) and in Trying (100). #[tokio::test] -async fn test_cancel_answered_487_sends_no_bye() -> crate::Result<()> { +async fn test_removed_dialog_2xx_crossing_the_cancel_is_acked_and_byed() -> crate::Result<()> { + run_crossing_2xx(Order::InviteOkFirst, 180, true).await?; + run_crossing_2xx(Order::InviteOkFirst, 100, true).await +} + +/// The same in the other wire ordering (the 200 to the CANCEL arrives before +/// the 2xx) with the dialog already removed from the layer. +#[tokio::test] +async fn test_removed_dialog_2xx_after_the_cancel_response_is_acked_and_byed() -> crate::Result<()> +{ + run_crossing_2xx(Order::CancelOkFirst, 180, true).await +} + +/// A CANCEL that wins the race (487 to the INVITE) ends the call as before: +/// the 487 is ACKed and no BYE is sent — also when the application removed +/// the dialog from the layer before dropping the `do_invite` future. +async fn run_cancel_answered_487(provisional: u16, remove: bool) -> crate::Result<()> { let token = CancellationToken::new(); let Uac { dialog_layer, @@ -322,13 +366,26 @@ async fn test_cancel_answered_487_sends_no_bye() -> crate::Result<()> { let wait = Duration::from_secs(2); let (state_sender, mut states) = unbounded_channel(); - let invite = tokio::spawn(async move { dialog_layer.do_invite(option, state_sender).await }); + let layer = dialog_layer.clone(); + let invite = tokio::spawn(async move { layer.do_invite(option, state_sender).await }); let (inv, uac) = recv_request(&peer, Method::Invite, wait).await; - reply(&peer, uac, &inv, 180, "Ringing").await; - wait_for_state(&mut states, "Early", wait, |s| { - matches!(s, DialogState::Early(_, _)) + let id = calling_id(&mut states).await; + let reason = if provisional == 100 { + "Trying" + } else { + "Ringing" + }; + reply(&peer, uac, &inv, provisional, reason).await; + wait_for_state(&mut states, "Trying or Early", wait, |s| { + matches!(s, DialogState::Trying(_) | DialogState::Early(_, _)) }) .await; + + if remove { + // The application removes the dialog from the layer first. + dialog_layer.remove_dialog(&id); + assert!(dialog_layer.is_empty()); + } invite.abort(); let _ = invite.await; let (cancel, _) = recv_request(&peer, Method::Cancel, wait).await; @@ -348,14 +405,48 @@ async fn test_cancel_answered_487_sends_no_bye() -> crate::Result<()> { "a cancelled call must not be BYE'd", ) .await; + + // Exactly one Terminated(UacCancel), never Confirmed. + tokio::time::sleep(Duration::from_millis(200)).await; + let seen: Vec<_> = std::iter::from_fn(|| states.try_recv().ok()).collect(); + assert!( + !seen + .iter() + .any(|s| matches!(s, DialogState::Confirmed(_, _))), + "a cancelled call must not report Confirmed, got {seen:?}" + ); + let terminations: Vec<_> = seen + .iter() + .filter(|s| matches!(s, DialogState::Terminated(_, _))) + .collect(); + assert!( + matches!( + terminations.as_slice(), + [DialogState::Terminated(_, TerminatedReason::UacCancel)] + ), + "expected exactly one Terminated(UacCancel), got {seen:?}" + ); token.cancel(); Ok(()) } +#[tokio::test] +async fn test_cancel_answered_487_sends_no_bye() -> crate::Result<()> { + run_cancel_answered_487(180, false).await +} + +/// The same when the application removed the dialog from the layer before +/// dropping the `do_invite` future, in Early (180) and in Trying (100). +#[tokio::test] +async fn test_removed_cancel_answered_487_sends_no_bye() -> crate::Result<()> { + run_cancel_answered_487(180, true).await?; + run_cancel_answered_487(100, true).await +} + /// Dropped before any response: Terminated(UacCancel) is reported at once and /// nothing is sent while no provisional has arrived (RFC 3261 §9.1). The first /// response then gets a CANCEL (180), or an ACK and a BYE (200). -async fn run_dropped_before_provisional(first: u16) -> crate::Result<()> { +async fn run_dropped_before_provisional(first: u16, remove: bool) -> crate::Result<()> { let token = CancellationToken::new(); let Uac { dialog_layer, @@ -365,8 +456,15 @@ async fn run_dropped_before_provisional(first: u16) -> crate::Result<()> { let wait = Duration::from_secs(2); let (state_sender, mut states) = unbounded_channel(); - let invite = tokio::spawn(async move { dialog_layer.do_invite(option, state_sender).await }); + let layer = dialog_layer.clone(); + let invite = tokio::spawn(async move { layer.do_invite(option, state_sender).await }); let (inv, uac) = recv_request(&peer, Method::Invite, wait).await; + let id = calling_id(&mut states).await; + if remove { + // The application removes the dialog from the layer first. + dialog_layer.remove_dialog(&id); + assert!(dialog_layer.is_empty()); + } invite.abort(); let _ = invite.await; let terminated = wait_for_state(&mut states, "Terminated", Duration::from_millis(200), |s| { @@ -418,143 +516,30 @@ async fn run_dropped_before_provisional(first: u16) -> crate::Result<()> { #[tokio::test] async fn test_dropped_before_provisional_is_cancelled_after_the_180() -> crate::Result<()> { - run_dropped_before_provisional(180).await + run_dropped_before_provisional(180, false).await } #[tokio::test] async fn test_dropped_before_provisional_2xx_is_acked_and_byed() -> crate::Result<()> { - run_dropped_before_provisional(200).await + run_dropped_before_provisional(200, false).await } #[tokio::test] -async fn test_dropped_before_provisional_final_failure_is_acked_only() -> crate::Result<()> { - run_dropped_before_provisional(486).await +async fn test_removed_before_provisional_is_cancelled_after_the_180() -> crate::Result<()> { + run_dropped_before_provisional(180, true).await } -/// Another owner takes the dialog out of the layer (`take_dialog`) and -/// CANCELs it, then the `do_invite` future is dropped. The guard no longer -/// holds a layer entry, but still holds the INVITE transaction: a 2xx -/// crossing that CANCEL is ACKed and BYE'd, and no second CANCEL is sent. #[tokio::test] -async fn test_taken_dialog_2xx_crossing_the_owners_cancel_is_acked_and_byed() -> crate::Result<()> { - let token = CancellationToken::new(); - let Uac { - dialog_layer, - option, - peer, - } = setup(&token).await?; - let wait = Duration::from_secs(2); - - let (state_sender, mut states) = unbounded_channel(); - let layer = dialog_layer.clone(); - let invite = tokio::spawn(async move { layer.do_invite(option, state_sender).await }); - let (inv, uac) = recv_request(&peer, Method::Invite, wait).await; - // The layer keys the dialog by its id before any To tag. - let calling = wait_for_state(&mut states, "Calling", wait, |s| { - matches!(s, DialogState::Calling(_)) - }) - .await; - reply(&peer, uac, &inv, 180, "Ringing").await; - wait_for_state(&mut states, "Early", wait, |s| { - matches!(s, DialogState::Early(_, _)) - }) - .await; - - // The owner takes the dialog and hangs it up (a CANCEL while early). - let taken = dialog_layer - .take_dialog(calling.id()) - .expect("the early dialog is in the layer"); - let hangup = tokio::spawn(async move { taken.hangup().await }); - let (cancel, _) = recv_request(&peer, Method::Cancel, wait).await; - invite.abort(); - let _ = invite.await; - - // The callee answered before the CANCEL reached it. - reply(&peer, uac, &inv, 200, "OK").await; - reply(&peer, uac, &cancel, 200, "OK").await; - let _ = tokio::time::timeout(wait, hangup).await; - - // Everything the UAC sends from here on, BYE answered, so a second - // CANCEL cannot hide behind a helper that skips other methods. - let requests = collect_requests(&peer, uac, Duration::from_millis(1500)).await; - let ack = requests - .iter() - .find(|r| r.method == Method::Ack) - .expect("the 2xx must be ACKed"); - assert_in_dialog(ack, &inv, "ACK"); - let bye = requests - .iter() - .find(|r| r.method == Method::Bye) - .expect("the 2xx must be BYE'd"); - assert_in_dialog(bye, &inv, "BYE"); - assert!( - !requests.iter().any(|r| r.method == Method::Cancel), - "the guard must not send a second CANCEL, got {:?}", - requests - .iter() - .map(|r| r.method.to_string()) - .collect::>() - ); - token.cancel(); - Ok(()) +async fn test_removed_before_provisional_2xx_is_acked_and_byed() -> crate::Result<()> { + run_dropped_before_provisional(200, true).await } -/// Every request the peer receives within `window`, answering each BYE 200. -async fn collect_requests(socket: &UdpSocket, uac: SocketAddr, window: Duration) -> Vec { - let mut out = Vec::new(); - let mut buf = vec![0u8; 4096]; - let deadline = tokio::time::Instant::now() + window; - while let Ok(Ok((len, _))) = tokio::time::timeout_at(deadline, socket.recv_from(&mut buf)).await - { - let text = std::str::from_utf8(&buf[..len]).expect("non utf-8 SIP message"); - if let Ok(SipMessage::Request(req)) = SipMessage::try_from(text) { - if req.method == Method::Bye { - reply(socket, uac, &req, 200, "OK").await; - } - out.push(req); - } - } - out +#[tokio::test] +async fn test_dropped_before_provisional_final_failure_is_acked_only() -> crate::Result<()> { + run_dropped_before_provisional(486, false).await } -/// Taken while still `Calling`: the owner's hangup cannot CANCEL before a -/// provisional response, so the dropped guard abandons the INVITE itself: -/// it CANCELs on the first provisional. #[tokio::test] -async fn test_taken_dialog_before_provisional_is_cancelled_after_first_provisional( -) -> crate::Result<()> { - let token = CancellationToken::new(); - let Uac { - dialog_layer, - option, - peer, - } = setup(&token).await?; - let wait = Duration::from_secs(2); - - let (state_sender, mut states) = unbounded_channel(); - let layer = dialog_layer.clone(); - let invite = tokio::spawn(async move { layer.do_invite(option, state_sender).await }); - let (inv, uac) = recv_request(&peer, Method::Invite, wait).await; - let calling = wait_for_state(&mut states, "Calling", wait, |s| { - matches!(s, DialogState::Calling(_)) - }) - .await; - let taken = dialog_layer - .take_dialog(calling.id()) - .expect("the dialog is in the layer"); - taken.hangup().await?; - invite.abort(); - let _ = invite.await; - - reply(&peer, uac, &inv, 180, "Ringing").await; - let (cancel, _) = recv_request(&peer, Method::Cancel, wait).await; - assert_eq!( - cancel.call_id_header().unwrap().value(), - inv.call_id_header().unwrap().value() - ); - reply(&peer, uac, &cancel, 200, "OK").await; - reply(&peer, uac, &inv, 487, "Request Terminated").await; - recv_request(&peer, Method::Ack, wait).await; - token.cancel(); - Ok(()) +async fn test_removed_before_provisional_final_failure_is_acked_only() -> crate::Result<()> { + run_dropped_before_provisional(486, true).await } diff --git a/src/dialog/tests/test_client_dialog.rs b/src/dialog/tests/test_client_dialog.rs index abfcceb6..653609fe 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; @@ -1707,3 +1711,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(()) +} diff --git a/src/dialog/tests/test_connection_affinity.rs b/src/dialog/tests/test_connection_affinity.rs index 594e01df..ab5e0c98 100644 --- a/src/dialog/tests/test_connection_affinity.rs +++ b/src/dialog/tests/test_connection_affinity.rs @@ -530,6 +530,69 @@ async fn test_restored_dialog_falls_back_to_initial_via_dialback() { } } +/// Fails every lookup, as a registrar-backed locator does once the target is gone. +struct NoRoute; + +#[async_trait::async_trait] +impl crate::transaction::endpoint::TargetLocator for NoRoute { + async fn locate(&self, _: &crate::sip::Uri) -> crate::Result { + Err(crate::Error::Error("no route".to_string())) + } +} + +#[tokio::test] +async fn test_server_request_without_route_fails_unless_dialed_back() -> crate::Result<()> { + // The locator error fails the send before the transaction starts. With no + // dial-back address nothing will ever end that transaction, so the request + // must return the error at once; with one, it is dialed back as before. + let tl = TransportLayer::new(CancellationToken::new()); + let udp = UdpConnection::create_connection("127.0.0.1:0".parse()?, None, None).await?; + tl.add_transport(udp.into()); + let endpoint = EndpointBuilder::new() + .with_transport_layer(tl) + .with_target_locator(Box::new(NoRoute)) + .build(); + let confirmed = |call_id: &str, via: &str| { + let invite = plain_udp_invite("alice-tag", call_id, via); + let (ep, role) = (endpoint.inner.clone(), TransactionRole::Server); + let inner = server_dialog(ep, role, invite, "sip:bob@127.0.0.1:5060"); + let id = inner.id.lock().clone(); + inner + .transition(DialogState::Confirmed(id, Response::default())) + .unwrap(); + InviteDialog::from_inner(Arc::new(inner)) + }; + + let dialog = confirmed("no-route", "SIP/2.0/UDP a.invalid;branch=z9hG4bKnoroute"); + for method in ["INFO", "re-INVITE", "BYE"] { + let request = async { + match method { + "INFO" => dialog.info(None, None).await.map(|_| ()), + "re-INVITE" => dialog.reinvite(None, None).await.map(|_| ()), + _ => dialog.bye().await, + } + }; + let result = tokio::time::timeout(Duration::from_secs(5), request).await; + let err = result.unwrap_or_else(|_| panic!("{method} must fail at once, not hang")); + let err = err.expect_err(method).to_string(); + assert!(err.contains("no route"), "{method}: {err}"); + } + + let probe = UdpSocket::bind("127.0.0.1:0").await?; + let via = format!( + "SIP/2.0/UDP alice.invalid:5060;branch=z9hG4bKdialback;received=127.0.0.1;rport={}", + probe.local_addr()?.port() + ); + let dialog = confirmed("dialback", &via); + tokio::spawn(async move { dialog.bye().await }); + let mut buf = [0u8; 2048]; + let (len, _) = tokio::time::timeout(Duration::from_secs(3), probe.recv_from(&mut buf)) + .await + .expect("the BYE must be dialed back to the Via's received/rport")?; + assert!(buf[..len].starts_with(b"BYE ")); + Ok(()) +} + /// A cancelled (dead) WebSocket flow must not cause affinity retransmissions /// into a dead socket, and the dial-back ladder must still deliver the BYE /// through tier 1 — the address captured from the connection at creation — diff --git a/src/dialog/tests/test_dialog_layer.rs b/src/dialog/tests/test_dialog_layer.rs index 78337c08..9acc90ed 100644 --- a/src/dialog/tests/test_dialog_layer.rs +++ b/src/dialog/tests/test_dialog_layer.rs @@ -311,87 +311,6 @@ async fn test_take_dialog_hands_the_dialog_to_exactly_one_caller() -> crate::Res Ok(()) } -/// With a transparent Call-ID the inbound (UAS) and outbound (UAC) legs of a -/// proxied call share it. Role-typed lookups must not hand out the other -/// role's dialog. -#[tokio::test] -async fn test_lookups_keep_uac_and_uas_dialogs_apart() -> crate::Result<()> { - let token = CancellationToken::new(); - let tl = TransportLayer::new(token.child_token()); - tl.add_transport(create_mock_connection().await?); - let endpoint = EndpointBuilder::new() - .with_user_agent("rsipstack-test") - .with_transport_layer(tl) - .build(); - let dialog_layer = DialogLayer::new(endpoint.inner.clone()); - let call_id = "shared-call-id"; - - // The inbound leg: a server dialog. - let invite_req = create_invite_request("caller-tag", "", call_id, "z9hG4bKinbound"); - let key = TransactionKey::from_request(&invite_req, TransactionRole::Server)?; - let tx = Transaction::new_server( - key, - invite_req, - endpoint.inner.clone(), - Some(create_mock_connection().await?), - ); - let (state_sender, _) = unbounded_channel(); - let server = dialog_layer.get_or_create_server_invite( - &tx, - state_sender, - None, - Some(crate::sip::Uri::try_from("sip:bob@bob.example.com:5060")?), - )?; - - // The outbound leg: a client dialog with the same Call-ID. - let (state_sender, _) = unbounded_channel(); - let (client, _client_tx) = dialog_layer.create_client_invite_dialog( - crate::dialog::invitation::InviteOption { - caller: crate::sip::Uri::try_from("sip:alice@example.com")?, - callee: crate::sip::Uri::try_from("sip:carol@example.com")?, - contact: crate::sip::Uri::try_from("sip:alice@127.0.0.1:5060")?, - call_id: Some(call_id.to_string()), - ..Default::default() - }, - state_sender, - )?; - assert_eq!(client.id().call_id, call_id); - dialog_layer.inner.dialogs.insert( - client.id().to_string(), - crate::dialog::dialog::Dialog::Invite(client.clone()), - ); - - let found = dialog_layer.get_client_dialog_by_call_id(call_id); - assert_eq!(found.len(), 1, "only the UAC dialog is a client dialog"); - assert_eq!(found[0].role(), TransactionRole::Client); - assert_eq!(found[0].id(), client.id()); - assert_eq!(server.role(), TransactionRole::Server); - - // An in-dialog request whose id maps to a UAC dialog is not answered by - // it as a server dialog. - let in_dialog = create_invite_request("far-tag", "near-tag", call_id, "z9hG4bKreinvite"); - let key = TransactionKey::from_request(&in_dialog, TransactionRole::Server)?; - let tx = Transaction::new_server( - key, - in_dialog, - endpoint.inner.clone(), - Some(create_mock_connection().await?), - ); - let id = DialogId::try_from(&tx)?; - dialog_layer.inner.dialogs.insert( - id.to_string(), - crate::dialog::dialog::Dialog::Invite(client.clone()), - ); - let (state_sender, _) = unbounded_channel(); - assert!( - dialog_layer - .get_or_create_server_invite(&tx, state_sender, None, None) - .is_err(), - "a UAC dialog must not be returned as the server INVITE dialog" - ); - Ok(()) -} - #[tokio::test] async fn test_dialog_layer_with_swapped_tags() -> crate::Result<()> { let endpoint = create_test_endpoint().await?; @@ -503,6 +422,47 @@ async fn test_multiple_dialogs_management() -> crate::Result<()> { Ok(()) } +/// A B2BUA that keeps the Call-ID has both legs of a call in one dialog layer: +/// the inbound one as UAS, the outbound one as UAC. Only the UAC one is a +/// client dialog. +#[tokio::test] +async fn test_get_client_dialog_by_call_id_returns_only_uac_dialogs() -> crate::Result<()> { + let endpoint = create_test_endpoint().await?; + let conn = create_mock_connection().await?; + endpoint.inner.transport_layer.add_transport(conn.clone()); + let dialog_layer = std::sync::Arc::new(DialogLayer::new(endpoint.inner.clone())); + let call_id = "b2bua-call-id"; + + let invite_req = create_invite_request("caller-tag", "", call_id, "z9hG4bKinbound"); + let key = TransactionKey::from_request(&invite_req, TransactionRole::Server)?; + let tx = Transaction::new_server(key, invite_req, endpoint.inner.clone(), Some(conn)); + let (state_sender, _) = unbounded_channel(); + let uas = dialog_layer.get_or_create_server_invite(&tx, state_sender, None, None)?; + + let peer = tokio::net::UdpSocket::bind("127.0.0.1:0").await?; + let callee = format!("sip:carol@{}", peer.local_addr()?); + let (state_sender, _) = unbounded_channel(); + let (uac, _invite) = dialog_layer.do_invite_async( + crate::dialog::invitation::InviteOption { + caller: crate::sip::Uri::try_from("sip:alice@example.com")?, + callee: crate::sip::Uri::try_from(callee.as_str())?, + contact: crate::sip::Uri::try_from("sip:alice@127.0.0.1:5060")?, + call_id: Some(call_id.to_string()), + ..Default::default() + }, + state_sender, + )?; + + let found: Vec<_> = dialog_layer + .get_client_dialog_by_call_id(call_id) + .iter() + .map(|d| (d.role(), d.id())) + .collect(); + assert_eq!(found, vec![(TransactionRole::Client, uac.id())]); + assert!(dialog_layer.get_dialog(&uas.id()).is_some()); + Ok(()) +} + #[tokio::test] async fn test_dialog_error_cases() -> crate::Result<()> { let endpoint = create_test_endpoint().await?; diff --git a/src/dialog/tests/test_forked_2xx_bye.rs b/src/dialog/tests/test_forked_2xx_bye.rs new file mode 100644 index 00000000..8d50f1fb --- /dev/null +++ b/src/dialog/tests/test_forked_2xx_bye.rs @@ -0,0 +1,238 @@ +//! RFC 3261 §13.2.2.4: a 2xx to the INVITE with a new To tag establishes its +//! own dialog. The UAC keeps a single session (the first 2xx's): the +//! transaction ACKs the forked 2xx with its own tag and remote target, and +//! the dialog layer ends the extra branch with a BYE built from that +//! response's Contact. No dialog is registered for the branch and the +//! confirmed dialog's state is untouched. + +use crate::dialog::dialog_layer::DialogLayer; +use crate::dialog::invitation::InviteOption; +use crate::dialog::DialogId; +use crate::sip::prelude::HeadersExt; +use crate::sip::{Method, Request, SipMessage, Uri}; +use crate::transport::udp::UdpConnection; +use crate::transport::TransportLayer; +use crate::EndpointBuilder; +use std::time::Duration; +use tokio::net::UdpSocket; +use tokio::sync::mpsc::unbounded_channel; +use tokio_util::sync::CancellationToken; + +/// Receive the next request of `method` at `peer`. +async fn next_request( + peer: &UdpSocket, + method: Method, + what: &str, +) -> (Request, std::net::SocketAddr) { + let mut buf = vec![0u8; 4096]; + loop { + let (len, from) = tokio::time::timeout(Duration::from_secs(3), peer.recv_from(&mut buf)) + .await + .unwrap_or_else(|_| panic!("timeout waiting for the {what}")) + .expect("peer socket error"); + let Ok(SipMessage::Request(req)) = SipMessage::try_from(&buf[..len]) else { + continue; + }; + if req.method == method { + return (req, from); + } + } +} + +/// Answer a request the way the peer UA would (top Via honored). +async fn reply_ok( + peer: &UdpSocket, + req: &Request, + from: std::net::SocketAddr, +) -> crate::Result<()> { + let ok = format!( + "SIP/2.0 200 OK\r\nVia: {}\r\nFrom: {}\r\nTo: {}\r\nCall-ID: {}\r\nCSeq: {}\r\nContent-Length: 0\r\n\r\n", + req.via_header()?.value(), + req.from_header()?.value(), + req.to_header()?.value(), + req.call_id_header()?.value(), + req.cseq_header()?.value(), + ); + peer.send_to(ok.as_bytes(), from).await?; + Ok(()) +} + +/// RFC 3261 §13.2.2.4: a forked 2xx (To tag `tag-b`) is ACKed with its own +/// tag and remote target (the transaction's job, pinned here end to end) and +/// then ended with a BYE built from its Contact — without touching the +/// confirmed dialog or registering the forked branch. +#[tokio::test] +async fn test_forked_2xx_is_byed_without_touching_the_confirmed_dialog() -> crate::Result<()> { + let token = CancellationToken::new(); + let peer = UdpSocket::bind("127.0.0.1:0").await?; + let tl = TransportLayer::new(token.child_token()); + let udp = UdpConnection::create_connection("127.0.0.1:0".parse()?, None, None).await?; + let uac_addr = udp.get_addr().get_socketaddr()?; + let contact = Uri::try_from(format!("sip:alice@{uac_addr}").as_str())?; + tl.add_transport(udp.into()); + let endpoint = EndpointBuilder::new() + .with_transport_layer(tl) + .with_option(crate::transaction::endpoint::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 layer = DialogLayer::new(endpoint.inner.clone()); + 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, + ..Default::default() + }; + // A second handle over the same layer state, for the assertions below + // (DialogLayer is not Clone; the registry lives in the shared inner). + let assert_layer = DialogLayer { + endpoint: layer.endpoint.clone(), + inner: layer.inner.clone(), + }; + let do_invite = tokio::spawn(async move { layer.do_invite(invite, state_sender).await }); + + // The INVITE arrives; branch A answers with tag-a. + let mut buf = vec![0u8; 4096]; + let (len, from) = tokio::time::timeout(Duration::from_secs(3), peer.recv_from(&mut buf)) + .await + .expect("timeout waiting for the INVITE")?; + let invite_req = match SipMessage::try_from(&buf[..len])? { + SipMessage::Request(req) if req.method == Method::Invite => req, + other => panic!("expected the INVITE, got {other}"), + }; + let ok_a = format!( + "SIP/2.0 200 OK\r\nVia: {}\r\nFrom: {}\r\nTo: {};tag=tag-a\r\nCall-ID: {}\r\nCSeq: {}\r\nContact: \r\nContent-Length: 0\r\n\r\n", + invite_req.via_header()?.value(), + invite_req.from_header()?.value(), + invite_req.to_header()?.value(), + invite_req.call_id_header()?.value(), + invite_req.cseq_header()?.value(), + peer.local_addr()?, + ); + peer.send_to(ok_a.as_bytes(), from).await?; + + let (dialog, _) = tokio::time::timeout(Duration::from_secs(3), do_invite) + .await + .expect("timeout waiting for do_invite") + .expect("do_invite task panicked")?; + assert!(dialog.inner.state.lock().is_confirmed()); + while states.try_recv().is_ok() {} + + let confirmed = dialog.id(); + assert_eq!(confirmed.remote_tag, "tag-a"); + + // Branch B forks in: same Call-ID, new To tag, own Contact. + let ok_b = format!( + "SIP/2.0 200 OK\r\nVia: {}\r\nFrom: {}\r\nTo: {};tag=tag-b\r\nCall-ID: {}\r\nCSeq: {}\r\nContact: \r\nContent-Length: 0\r\n\r\n", + invite_req.via_header()?.value(), + invite_req.from_header()?.value(), + invite_req.to_header()?.value(), + invite_req.call_id_header()?.value(), + invite_req.cseq_header()?.value(), + peer.local_addr()?, + ); + peer.send_to(ok_b.as_bytes(), uac_addr).await?; + + // The first ACK confirms branch A; the forked 2xx is then ACKed with + // its own tag and Contact as R-URI. + let (mut ack, _) = next_request(&peer, Method::Ack, "forked ACK").await; + while !ack + .to_header() + .ok() + .and_then(|to| to.tag().ok().flatten()) + .is_some_and(|tag| tag.value() == "tag-b") + { + (ack, _) = next_request(&peer, Method::Ack, "forked ACK").await; + } + let ack_to = ack.to_header()?.value().to_string(); + assert!(ack_to.contains(";tag=tag-b"), "ACK To: {ack_to}"); + assert!( + ack.uri.to_string().contains("bob-b"), + "ACK R-URI must be the forked Contact: {}", + ack.uri + ); + + // The forked branch is then ended with a BYE from its own Contact. + let (bye, from) = next_request(&peer, Method::Bye, "forked-branch BYE").await; + let bye_to = bye.to_header()?.value().to_string(); + assert!(bye_to.contains(";tag=tag-b"), "BYE To: {bye_to}"); + assert!( + bye.uri.to_string().contains("bob-b"), + "BYE R-URI must be the forked Contact: {}", + bye.uri + ); + let cseq = bye.cseq_header()?.value().to_string(); + assert_eq!(cseq, "2 BYE", "BYE CSeq continues the forked dialog's"); + assert_eq!( + bye.call_id_header()?.value(), + invite_req.call_id_header()?.value() + ); + // The forked callee answers the BYE like a real UA would. + reply_ok(&peer, &bye, from).await?; + + // A retransmission of the forked 2xx (in flight before its ACK landed) + // must not trigger a second BYE. + peer.send_to(ok_b.as_bytes(), uac_addr).await?; + let mut buf = vec![0u8; 4096]; + let quiet_until = tokio::time::Instant::now() + Duration::from_millis(400); + loop { + let remaining = quiet_until.saturating_duration_since(tokio::time::Instant::now()); + if remaining.is_zero() { + break; + } + match tokio::time::timeout(remaining, peer.recv_from(&mut buf)).await { + Err(_) => break, // the quiet window elapsed: no duplicate BYE + Ok(Err(e)) => return Err(e.into()), + Ok(Ok((len, _))) => { + if let Ok(SipMessage::Request(req)) = SipMessage::try_from(&buf[..len]) { + assert_ne!( + req.method, + Method::Bye, + "duplicate BYE for a retransmitted forked 2xx" + ); + } + } + } + } + + // The confirmed dialog is untouched: no new state notifications, no + // dialog registered for the forked branch, and the dialog still works. + assert!( + states.try_recv().is_err(), + "the forked 2xx must not touch the confirmed dialog's state" + ); + let forked_id = DialogId { + call_id: confirmed.call_id.clone(), + local_tag: confirmed.local_tag.clone(), + remote_tag: "tag-b".to_string(), + }; + assert!( + assert_layer.get_dialog(&forked_id).is_none(), + "the forked branch must not be registered as a dialog" + ); + assert!(dialog.inner.state.lock().is_confirmed()); + + // The confirmed dialog is still usable: a BYE ends it normally. + let dialog2 = dialog.clone(); + let bye_task = tokio::spawn(async move { dialog2.bye().await }); + let (bye, _) = next_request(&peer, Method::Bye, "confirmed-dialog BYE").await; + assert!( + bye.to_header()?.value().to_string().contains(";tag=tag-a"), + "the confirmed dialog's BYE keeps tag-a: {}", + bye.to_header()?.value() + ); + reply_ok(&peer, &bye, uac_addr).await?; + tokio::time::timeout(Duration::from_secs(3), bye_task) + .await + .expect("timeout waiting for bye()") + .expect("bye task panicked")?; + assert!(dialog.inner.state.lock().is_terminated()); + token.cancel(); + Ok(()) +} diff --git a/src/dialog/tests/test_in_dialog_provisional.rs b/src/dialog/tests/test_in_dialog_provisional.rs index 76668eed..18cc6685 100644 --- a/src/dialog/tests/test_in_dialog_provisional.rs +++ b/src/dialog/tests/test_in_dialog_provisional.rs @@ -115,6 +115,22 @@ pub(super) async fn establish_with_credential( }); let dialog_layer = DialogLayer::new(endpoint.inner.clone()); + // Pump inbound in-dialog requests (e.g. the peer's BYE) into the layer. + let mut incoming = endpoint.incoming_transactions()?; + let pump_layer = DialogLayer { + endpoint: dialog_layer.endpoint.clone(), + inner: dialog_layer.inner.clone(), + }; + tokio::spawn(async move { + while let Some(mut tx) = incoming.recv().await { + if let Some(mut dialog) = pump_layer.match_dialog(&tx) { + tokio::spawn(async move { + let _ = dialog.handle(&mut tx).await; + }); + } + } + }); + let (state_sender, mut state_receiver) = unbounded_channel(); let invite_option = InviteOption { caller: Uri::try_from("sip:alice@example.com")?, @@ -190,12 +206,13 @@ async fn assert_provisional_keeps_confirmed(method: Method) -> crate::Result<()> "dialog must be Confirmed after the in-dialog {method} completes, state: {}", dialog.state() ); - // The 183 neither changes the dialog's state nor is notified as `Early`: - // subscribers read `Early` as the dialog's own early state. + // Fork (R19): the 183 is not notified either. Subscribers read `Early` + // as the dialog's early state (ringing), not as a provisional to a later + // transaction. let states = drain_states(&mut states); assert!( !states.iter().any(|s| matches!(s, DialogState::Early(_, _))), - "a 183 to an in-dialog {method} must not be notified as Early, got {states:?}" + "a 1xx to an in-dialog {method} must not be notified as Early, got {states:?}" ); // The call can still be hung up with a BYE. @@ -222,3 +239,71 @@ async fn test_reinvite_provisional_keeps_dialog_confirmed() -> crate::Result<()> async fn test_update_provisional_keeps_dialog_confirmed() -> crate::Result<()> { assert_provisional_keeps_confirmed(Method::Update).await } + +/// A provisional response to an in-dialog request is never notified after +/// the dialog terminated: the peer hangs up while a re-INVITE is pending, +/// and only then answers it with a 183. +#[tokio::test] +async fn test_provisional_after_terminated_is_not_notified() -> crate::Result<()> { + let token = CancellationToken::new(); + let (dialog, mut states, peer) = establish(&token).await?; + while states.try_recv().is_ok() {} + + // A re-INVITE goes out and stays pending. + let requester = dialog.clone(); + let pending = tokio::spawn(async move { requester.reinvite(None, None).await }); + let (req, uac) = recv_request(&peer, Method::Invite).await; + + // The peer hangs up while the re-INVITE is pending. + let id = dialog.id(); + let uac_uri = format!("sip:alice@{uac}"); + let bye = format!( + "BYE {uac_uri} SIP/2.0\r\n\ + Via: SIP/2.0/UDP {peer_addr};branch=z9hG4bK-bye\r\n\ + Max-Forwards: 70\r\n\ + From: ;tag={PEER_TAG}\r\n\ + To: <{uac_uri}>;tag={local}\r\n\ + Call-ID: {call_id}\r\n\ + CSeq: 2 BYE\r\n\ + Content-Length: 0\r\n\r\n", + peer_addr = peer.local_addr()?, + local = id.local_tag, + call_id = id.call_id, + ); + peer.send_to(bye.as_bytes(), uac).await?; + // The dialog terminates and answers the BYE with a 200. + tokio::time::timeout(Duration::from_secs(2), async { + loop { + if matches!(states.recv().await, Some(DialogState::Terminated(..))) { + break; + } + } + }) + .await + .expect("timeout waiting for Terminated"); + assert!(dialog.state().is_terminated()); + + // Only now does the peer answer the re-INVITE — with a provisional. + reply(&peer, uac, &req, 183, "Session Progress").await; + tokio::time::sleep(Duration::from_millis(100)).await; + + let after: Vec = drain_states(&mut states) + .into_iter() + .map(|s| s.to_string()) + .collect(); + assert!( + after.is_empty(), + "nothing may be notified after Terminated, got {after:?}" + ); + + // Finish the pending re-INVITE so the transaction does not linger. + reply(&peer, uac, &req, 200, "OK").await; + tokio::time::timeout(Duration::from_secs(2), pending) + .await + .expect("re-INVITE did not complete") + .expect("re-INVITE task panicked")?; + + assert!(dialog.state().is_terminated()); + token.cancel(); + Ok(()) +} diff --git a/src/dialog/tests/test_late_reinvite_ack.rs b/src/dialog/tests/test_late_reinvite_ack.rs deleted file mode 100644 index 1a14d904..00000000 --- a/src/dialog/tests/test_late_reinvite_ack.rs +++ /dev/null @@ -1,277 +0,0 @@ -//! The ACK for a 2xx carries the CSeq number of the INVITE it acknowledges -//! and stops that 2xx's retransmissions (RFC 3261 §13.2.2.4, §13.3.1.4). -//! Over UDP the ACK of one re-INVITE can arrive after the next re-INVITE on -//! the same dialog: it must still reach its own transaction, and must not be -//! taken for the ACK of the newer one. Each `Confirmed` carries the 2xx its -//! ACK confirms, although that ACK ends the server transaction. -use crate::dialog::{ - dialog::{Dialog, DialogState, DialogStateReceiver}, - dialog_layer::DialogLayer, - invite_dialog::InviteDialog, -}; -use crate::sip::{prelude::HeadersExt, Method, Response, SipMessage, StatusCode}; -use crate::transaction::endpoint::EndpointOption; -use crate::transport::{udp::UdpConnection, TransportLayer}; -use crate::EndpointBuilder; -use std::net::SocketAddr; -use std::sync::Arc; -use std::time::Duration; -use tokio::net::UdpSocket; -use tokio::sync::mpsc::{unbounded_channel, UnboundedReceiver}; -use tokio_util::sync::CancellationToken; - -const CALL_ID: &str = "late-reinvite-ack-test"; -const FROM_TAG: &str = "uac-tag"; - -/// A raw UAC peer. -struct Peer { - socket: UdpSocket, - uas: SocketAddr, -} - -impl Peer { - async fn send_request(&self, method: Method, cseq: u32, branch: &str, to_tag: Option<&str>) { - let addr = self.socket.local_addr().unwrap(); - let to = match to_tag { - Some(tag) => format!(";tag={tag}", self.uas), - None => format!("", self.uas), - }; - let msg = format!( - "{method} sip:bob@{uas} SIP/2.0\r\n\ - Via: SIP/2.0/UDP {addr};branch={branch}\r\n\ - Max-Forwards: 70\r\n\ - From: ;tag={FROM_TAG}\r\n\ - To: {to}\r\n\ - Call-ID: {CALL_ID}\r\n\ - CSeq: {cseq} {method}\r\n\ - Contact: \r\n\ - Content-Length: 0\r\n\r\n", - uas = self.uas, - ); - self.socket.send_to(msg.as_bytes(), self.uas).await.unwrap(); - } - - /// Wait for the 2xx to the request with `cseq`, skipping anything else. - async fn recv_2xx(&self, cseq: u32) -> Response { - let mut buf = vec![0u8; 4096]; - loop { - let (len, _) = - tokio::time::timeout(Duration::from_secs(2), self.socket.recv_from(&mut buf)) - .await - .expect("timeout waiting for a 2xx") - .unwrap(); - let text = std::str::from_utf8(&buf[..len]).unwrap(); - if let Ok(SipMessage::Response(resp)) = SipMessage::try_from(text) { - let seq = resp.cseq_header().unwrap().seq().unwrap(); - if seq == cseq && resp.status_code == StatusCode::OK { - return resp; - } - } - } - } - - /// Count the 2xx (re)transmissions per CSeq received during `window`. - async fn count_2xx(&self, window: Duration) -> std::collections::HashMap { - let mut counts = std::collections::HashMap::new(); - let mut buf = vec![0u8; 4096]; - let deadline = tokio::time::Instant::now() + window; - while let Ok(Ok((len, _))) = - tokio::time::timeout_at(deadline, self.socket.recv_from(&mut buf)).await - { - let text = std::str::from_utf8(&buf[..len]).unwrap(); - if let Ok(SipMessage::Response(resp)) = SipMessage::try_from(text) { - if resp.status_code == StatusCode::OK { - let seq = resp.cseq_header().unwrap().seq().unwrap(); - *counts.entry(seq).or_insert(0) += 1; - } - } - } - counts - } -} - -/// A UAS endpoint (short T1) running the usual incoming-transaction loop. -async fn setup( - token: &CancellationToken, -) -> crate::Result<(DialogStateReceiver, UnboundedReceiver, Peer)> { - 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 uas: SocketAddr = udp.get_addr().get_socketaddr()?; - transport_layer.add_transport(udp.into()); - let endpoint = EndpointBuilder::new() - .with_user_agent("rsipstack-test") - .with_transport_layer(transport_layer) - .with_cancel_token(token.child_token()) - .with_option(EndpointOption { - t1: Duration::from_millis(100), - t1x64: Duration::from_millis(64 * 100), - ..Default::default() - }) - .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 { - let _ = endpoint_inner.serve().await; - }); - - let (state_sender, states) = unbounded_channel(); - let (dialog_sender, dialogs) = unbounded_channel(); - tokio::spawn(async move { - while let Some(mut tx) = incoming.recv().await { - let has_to_tag = tx.original.to_header().unwrap().tag().unwrap().is_some(); - let dialog = if has_to_tag { - dialog_layer.match_dialog(&tx) - } else if tx.original.method == Method::Invite { - let dialog = dialog_layer - .get_or_create_server_invite(&tx, state_sender.clone(), None, None) - .expect("server dialog"); - dialog_sender.send(dialog.clone()).unwrap(); - Some(Dialog::Invite(dialog)) - } else { - None - }; - if let Some(mut dialog) = dialog { - tokio::spawn(async move { - let _ = dialog.handle(&mut tx).await; - }); - } - } - }); - - let peer = Peer { - socket: UdpSocket::bind("127.0.0.1:0").await?, - uas, - }; - Ok((states, dialogs, peer)) -} - -async fn next_state(rx: &mut DialogStateReceiver) -> DialogState { - tokio::time::timeout(Duration::from_secs(2), rx.recv()) - .await - .expect("timeout waiting for dialog state") - .expect("state channel closed") -} - -/// CSeq numbers of the responses carried by the `Confirmed` states in `rx`. -fn confirmed_cseqs(rx: &mut DialogStateReceiver) -> Vec { - let mut seqs = Vec::new(); - while let Ok(state) = rx.try_recv() { - if let DialogState::Confirmed(_, resp) = state { - seqs.push(resp.cseq_header().unwrap().seq().unwrap()); - } - } - seqs -} - -/// Answer the next re-INVITE the application is offered with a 200. -async fn accept_reinvite(states: &mut DialogStateReceiver) { - loop { - if let DialogState::Updated(_, _, handle) = next_state(states).await { - handle.reply(StatusCode::OK).await.unwrap(); - return; - } - } -} - -#[tokio::test] -async fn test_late_ack_of_previous_reinvite_reaches_its_own_transaction() -> crate::Result<()> { - let token = CancellationToken::new(); - let (mut states, mut dialogs, peer) = setup(&token).await?; - - // Establish the call: INVITE (CSeq 1), 200, ACK. - peer.send_request(Method::Invite, 1, "z9hG4bK-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 to_tag = peer - .recv_2xx(1) - .await - .to_header()? - .tag()? - .expect("To tag") - .value() - .to_string(); - peer.send_request(Method::Ack, 1, "z9hG4bK-ack-1", Some(&to_tag)) - .await; - // The matching ACK ends the Accepted server transaction (RFC 6026 §7.1); - // `Confirmed` still carries the 2xx it confirms. - let confirmed = loop { - if let DialogState::Confirmed(_, resp) = next_state(&mut states).await { - break resp; - } - }; - assert_eq!(confirmed.status_code, StatusCode::OK); - assert_eq!(confirmed.cseq_header()?.seq()?, 1); - - // re-INVITE CSeq 2 is answered; its ACK is delayed in the network. - peer.send_request(Method::Invite, 2, "z9hG4bK-invite-2", Some(&to_tag)) - .await; - accept_reinvite(&mut states).await; - peer.recv_2xx(2).await; - - // The UAC, having sent that ACK, sends re-INVITE CSeq 3, which is - // answered too. - peer.send_request(Method::Invite, 3, "z9hG4bK-invite-3", Some(&to_tag)) - .await; - accept_reinvite(&mut states).await; - peer.recv_2xx(3).await; - - // Now the delayed ACK of CSeq 2 arrives. - peer.send_request(Method::Ack, 2, "z9hG4bK-ack-2", Some(&to_tag)) - .await; - let confirmed = loop { - if let DialogState::Confirmed(_, resp) = next_state(&mut states).await { - break resp.cseq_header()?.seq()?; - } - }; - assert_eq!( - confirmed, 2, - "the ACK of CSeq 2 must confirm re-INVITE 2, not re-INVITE 3" - ); - // Let anything already in flight arrive, then watch the retransmissions. - peer.count_2xx(Duration::from_millis(150)).await; - let counts = peer.count_2xx(Duration::from_millis(1200)).await; - assert_eq!( - counts.get(&2).copied().unwrap_or(0), - 0, - "the ACK of CSeq 2 must stop the retransmissions of its 200, got {counts:?}" - ); - assert!( - counts.get(&3).copied().unwrap_or(0) > 0, - "the 200 to CSeq 3 is not acknowledged yet and must still be retransmitted, got {counts:?}" - ); - assert_eq!( - confirmed_cseqs(&mut states), - Vec::::new(), - "the ACK of CSeq 2 must confirm nothing else" - ); - - // The ACK of CSeq 3 still finds its transaction. - peer.send_request(Method::Ack, 3, "z9hG4bK-ack-3", Some(&to_tag)) - .await; - let confirmed = loop { - if let DialogState::Confirmed(_, resp) = next_state(&mut states).await { - break resp.cseq_header()?.seq()?; - } - }; - assert_eq!(confirmed, 3, "the ACK of CSeq 3 must confirm re-INVITE 3"); - peer.count_2xx(Duration::from_millis(150)).await; - let counts = peer.count_2xx(Duration::from_millis(1000)).await; - assert!( - counts.is_empty(), - "no 200 may be retransmitted once both are acknowledged, got {counts:?}" - ); - assert_eq!(confirmed_cseqs(&mut states), Vec::::new()); - - token.cancel(); - Ok(()) -} diff --git a/src/dialog/tests/test_refer_notify.rs b/src/dialog/tests/test_refer_notify.rs new file mode 100644 index 00000000..9ce7f455 --- /dev/null +++ b/src/dialog/tests/test_refer_notify.rs @@ -0,0 +1,294 @@ +//! RFC 3515 regression: after an in-dialog REFER is answered (usually 202), +//! the dialog must return to `Confirmed` — the implicit subscription's +//! NOTIFYs are in-dialog requests that need the confirmed dialog. Leaving +//! the dialog in the `Refer` state starved the referrer of every NOTIFY. +use crate::dialog::{ + dialog::{Dialog, DialogState, DialogStateReceiver}, + dialog_layer::DialogLayer, + invite_dialog::InviteDialog, +}; +use crate::sip::{prelude::HeadersExt, Method, SipMessage, StatusCode}; +use crate::transaction::endpoint::EndpointOption; +use crate::transport::udp::UdpConnection; +use crate::transport::TransportLayer; +use crate::EndpointBuilder; +use std::net::SocketAddr; +use std::sync::Arc; +use std::time::{Duration, Instant}; +use tokio::net::UdpSocket; +use tokio::sync::mpsc::{unbounded_channel, UnboundedReceiver}; +use tokio_util::sync::CancellationToken; + +const CALL_ID: &str = "refer-notify-test"; +const FROM_TAG: &str = "referrer-tag"; + +struct Harness { + states: DialogStateReceiver, + dialogs: UnboundedReceiver, + peer: UdpSocket, + uas: SocketAddr, +} + +async fn setup(token: &CancellationToken) -> crate::Result { + 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 uas: SocketAddr = udp.get_addr().get_socketaddr()?; + transport_layer.add_transport(udp.into()); + let endpoint = EndpointBuilder::new() + .with_user_agent("rsipstack-test") + .with_transport_layer(transport_layer) + .with_cancel_token(token.child_token()) + .with_option(EndpointOption { + t1: Duration::from_millis(20), + t1x64: Duration::from_millis(64 * 20), + ..Default::default() + }) + .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 { + let _ = endpoint_inner.serve().await; + }); + + let (state_sender, states) = unbounded_channel(); + let (dialog_sender, dialogs) = unbounded_channel(); + tokio::spawn(async move { + while let Some(mut tx) = incoming.recv().await { + let has_to_tag = tx.original.to_header().unwrap().tag().unwrap().is_some(); + let dialog = if has_to_tag { + dialog_layer.match_dialog(&tx) + } else if tx.original.method == Method::Invite { + let dialog = dialog_layer + .get_or_create_server_invite(&tx, state_sender.clone(), None, None) + .expect("server dialog"); + dialog_sender.send(dialog.clone()).unwrap(); + Some(Dialog::Invite(dialog)) + } else { + None + }; + if let Some(mut dialog) = dialog { + tokio::spawn(async move { + let _ = dialog.handle(&mut tx).await; + }); + } + } + }); + + let peer = UdpSocket::bind("127.0.0.1:0").await?; + Ok(Harness { + states, + dialogs, + peer, + uas, + }) +} + +async fn recv_message( + socket: &UdpSocket, + buf: &mut [u8], + wait: Duration, +) -> Option<(String, SocketAddr)> { + let deadline = Instant::now() + wait; + loop { + let now = Instant::now(); + if now >= deadline { + return None; + } + let Ok(Ok((len, src))) = tokio::time::timeout(deadline - now, socket.recv_from(buf)).await + else { + return None; + }; + return Some((String::from_utf8_lossy(&buf[..len]).to_string(), src)); + } +} + +/// Confirm a call, then REFER it: the dialog must answer 202, return to +/// Confirmed, and `notify_refer` must deliver the 100 Trying NOTIFY. +#[tokio::test] +async fn test_refer_returns_dialog_to_confirmed_and_notify_works() -> crate::Result<()> { + let token = CancellationToken::new(); + let Harness { + mut states, + mut dialogs, + peer, + uas, + } = setup(&token).await?; + + let peer_addr = peer.local_addr()?; + let invite = format!( + "INVITE sip:bob@{uas} SIP/2.0\r\n\ + Via: SIP/2.0/UDP {peer_addr};branch=z9hG4bK-rn-invite\r\n\ + Max-Forwards: 70\r\n\ + From: ;tag={FROM_TAG}\r\n\ + To: \r\n\ + Call-ID: {CALL_ID}\r\n\ + CSeq: 1 INVITE\r\n\ + Contact: \r\n\ + Content-Length: 0\r\n\r\n", + ); + peer.send_to(invite.as_bytes(), uas).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)?; + + // First 200 OK, remember its To tag. + let mut buf = vec![0u8; 4096]; + let to_tag = loop { + let (text, _) = recv_message(&peer, &mut buf, Duration::from_secs(2)) + .await + .expect("timeout waiting for the first 200"); + if let Ok(SipMessage::Response(resp)) = SipMessage::try_from(text.as_str()) { + if resp.status_code.code() == 200 { + break resp + .to_header()? + .tag()? + .expect("the 200 must carry the local tag") + .value() + .to_string(); + } + } + }; + // Confirm the dialog. + let ack = format!( + "ACK sip:alice@{peer_addr} SIP/2.0\r\n\ + Via: SIP/2.0/UDP {peer_addr};branch=z9hG4bK-rn-ack\r\n\ + Max-Forwards: 70\r\n\ + From: ;tag={FROM_TAG}\r\n\ + To: ;tag={to_tag}\r\n\ + Call-ID: {CALL_ID}\r\n\ + CSeq: 1 ACK\r\n\ + Content-Length: 0\r\n\r\n", + ); + peer.send_to(ack.as_bytes(), uas).await?; + + // In-dialog REFER transferring the call elsewhere. + let refer = format!( + "REFER sip:bob@{uas} SIP/2.0\r\n\ + Via: SIP/2.0/UDP {peer_addr};branch=z9hG4bK-rn-refer\r\n\ + Max-Forwards: 70\r\n\ + From: ;tag={FROM_TAG}\r\n\ + To: ;tag={to_tag}\r\n\ + Call-ID: {CALL_ID}\r\n\ + CSeq: 2 REFER\r\n\ + Refer-To: \r\n\ + Contact: \r\n\ + Content-Length: 0\r\n\r\n", + ); + peer.send_to(refer.as_bytes(), uas).await?; + + // The application answers the REFER with 202 (as an RFC 3515 referee). + let mut refer_handle = None; + let deadline = Instant::now() + Duration::from_secs(2); + while Instant::now() < deadline { + while let Ok(state) = states.try_recv() { + if let DialogState::Refer(_, _, handle) = state { + refer_handle = Some(handle); + } + } + if refer_handle.is_some() { + break; + } + tokio::time::sleep(Duration::from_millis(5)).await; + } + let refer_handle = refer_handle.expect("the dialog never surfaced the REFER"); + refer_handle + .reply(StatusCode::Other(202, "Accepted".into())) + .await + .ok(); + + // The referrer sees the 202 ... + let saw_202 = loop { + let (text, _) = recv_message(&peer, &mut buf, Duration::from_secs(2)) + .await + .expect("timeout waiting for the 202"); + if text.starts_with("SIP/2.0 202") { + break true; + } + }; + assert!(saw_202); + + // ... and the dialog re-surfaces a Confirmed event: in-dialog request + // events do not change the stored state, but TU-side logic (e.g. a + // pending refer NOTIFY) waits on the Confirmed event after answering. + // Without `return_to_confirmed` in handle_refer it never arrives. + let reconfirmed = { + let deadline = Instant::now() + Duration::from_secs(2); + let mut seen = false; + while Instant::now() < deadline { + while let Ok(state) = states.try_recv() { + if matches!(state, DialogState::Confirmed(..)) { + seen = true; + } + } + if seen { + break; + } + tokio::time::sleep(Duration::from_millis(5)).await; + } + seen + }; + assert!( + reconfirmed, + "answering the REFER must re-surface a Confirmed event for the dialog" + ); + assert!(dialog.state().is_confirmed()); + + // The 100 Trying NOTIFY for the implicit subscription goes out. + let notify = dialog.notify_refer(StatusCode::Trying, "active").await; + assert!( + notify.is_ok(), + "notify_refer must work once the dialog is Confirmed again: {:?}", + notify.err() + ); + let saw_notify = loop { + let (text, src) = recv_message(&peer, &mut buf, Duration::from_secs(2)) + .await + .expect("timeout waiting for the NOTIFY"); + if text.starts_with("NOTIFY") { + assert!( + text.contains("Event: refer"), + "the NOTIFY must carry Event: refer, got {text}" + ); + assert!( + text.contains("Subscription-State: active"), + "the NOTIFY must carry Subscription-State: active, got {text}" + ); + assert!( + text.contains("SIP/2.0 100"), + "the NOTIFY body must carry the 100 Trying sipfrag, got {text}" + ); + // Answer the NOTIFY so the transaction completes. + let header = |name: &str| { + text.lines() + .find(|l| l.starts_with(name)) + .unwrap() + .strip_prefix(&format!("{name} ")) + .unwrap() + .to_string() + }; + let resp = format!( + "SIP/2.0 200 OK\r\nVia: {}\r\nFrom: {}\r\nTo: {}\r\nCall-ID: {}\r\nCSeq: {}\r\nContent-Length: 0\r\n\r\n", + header("Via:"), + header("From:"), + header("To:"), + header("Call-ID:"), + header("CSeq:"), + ); + peer.send_to(resp.as_bytes(), src).await?; + break true; + } + }; + assert!(saw_notify); + + token.cancel(); + Ok(()) +} diff --git a/src/dialog/tests/test_uas_ack_timeout.rs b/src/dialog/tests/test_uas_ack_timeout.rs index a3bf069a..e6f73224 100644 --- a/src/dialog/tests/test_uas_ack_timeout.rs +++ b/src/dialog/tests/test_uas_ack_timeout.rs @@ -514,14 +514,18 @@ async fn test_unacked_reinvite_2xx_ends_the_session() -> crate::Result<()> { Ok(()) } -async fn wait_confirmed(states: &mut DialogStateReceiver) { +/// Wait for `Confirmed` and return the CSeq number of the 2xx it carries. +async fn wait_confirmed(states: &mut DialogStateReceiver) -> u32 { loop { let state = tokio::time::timeout(Duration::from_secs(2), states.recv()) .await .expect("timeout waiting for the call to be confirmed") .expect("state channel closed"); - if matches!(state, DialogState::Confirmed(..)) { - return; + if let DialogState::Confirmed(_, resp) = state { + let cseq = resp.cseq_header().expect("Confirmed carries the 2xx"); + assert_eq!(cseq.method().unwrap(), Method::Invite); + assert_eq!(resp.status_code.code(), 200); + return cseq.seq().unwrap(); } } } @@ -652,7 +656,7 @@ async fn test_2xx_over_tcp_is_retransmitted_until_the_ack() -> crate::Result<()> .write_all(request(Method::Ack, 1, Some(&to_tag)).as_bytes()) .await?; let acked = Instant::now(); - wait_confirmed(&mut states).await; + assert_eq!(wait_confirmed(&mut states).await, 1); let after_ack = read_tcp_messages(&mut stream, &mut buf, acked + T1 * 30).await; assert!( !after_ack @@ -790,7 +794,7 @@ async fn test_late_ack_of_an_earlier_reinvite_reaches_its_own_transaction() -> c }; 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; + assert_eq!(wait_confirmed(&mut states).await, 1); // re-INVITE 2 and 3 are answered; the ACK of 2 arrives after re-INVITE 3. let mut answered = None; @@ -854,6 +858,9 @@ async fn test_late_ack_of_an_earlier_reinvite_reaches_its_own_transaction() -> c )), "an ACKed call must not be ended" ); + // Each ACK confirms its own re-INVITE, in the order the ACKs arrived. + assert_eq!(wait_confirmed(&mut states).await, 2); + assert_eq!(wait_confirmed(&mut states).await, 3); assert!(terminated_reason(&mut states).is_none()); assert!(dialog.state().is_confirmed()); token.cancel(); @@ -951,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: \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 +} diff --git a/src/dialog/tests/test_warn_logs.rs b/src/dialog/tests/test_warn_logs.rs new file mode 100644 index 00000000..c3382bbf --- /dev/null +++ b/src/dialog/tests/test_warn_logs.rs @@ -0,0 +1,195 @@ +//! WARN records carry metadata only; the whole SIP message is logged below WARN. +use crate::dialog::dialog::{DialogInner, DialogState}; +use crate::dialog::{client_dialog::ClientInviteDialog, invite_dialog::InviteDialog}; +use crate::dialog::{server_dialog::ServerInviteDialog, DialogId}; +use crate::sip::{Request, SipMessage, Uri}; +use crate::transaction::{endpoint::TargetLocator, key::TransactionRole}; +use crate::transport::{SipAddr, TransportLayer}; +use std::sync::{Arc, Mutex}; +use tokio::sync::mpsc::unbounded_channel; +use tracing::{field::Field, span, Event, Level, Metadata, Subscriber}; + +const MARKER: &str = "private-marker"; + +/// Records every event as (level, " field=value ..."). +#[derive(Clone, Default)] +struct Capture(Arc>>); + +impl Subscriber for Capture { + fn enabled(&self, _: &Metadata<'_>) -> bool { + true + } + fn new_span(&self, _: &span::Attributes<'_>) -> span::Id { + span::Id::from_u64(1) + } + fn record(&self, _: &span::Id, _: &span::Record<'_>) {} + fn record_follows_from(&self, _: &span::Id, _: &span::Id) {} + fn event(&self, event: &Event<'_>) { + let mut line = String::new(); + event.record(&mut |f: &Field, v: &dyn std::fmt::Debug| { + line.push_str(&format!(" {}={:?}", f.name(), v)) + }); + self.0 + .lock() + .unwrap() + .push((*event.metadata().level(), line)); + } + fn enter(&self, _: &span::Id) {} + fn exit(&self, _: &span::Id) {} +} + +impl Capture { + fn lines(&self, level: Level, needle: &str) -> Vec { + let records = self.0.lock().unwrap(); + let hits = records + .iter() + .filter(|(l, s)| *l == level && s.contains(needle)); + hits.map(|(_, s)| s.clone()).collect() + } +} + +/// Fails every lookup, so `Transaction::send` returns an error. +struct NoRoute; + +#[async_trait::async_trait] +impl TargetLocator for NoRoute { + async fn locate(&self, _: &Uri) -> crate::Result { + Err(crate::Error::Error("no route".to_string())) + } +} + +/// A request or response whose From carries `MARKER`. +fn message(start_line: &str, cseq: &str) -> SipMessage { + let text = format!( + "{start_line}\r\nVia: SIP/2.0/UDP 127.0.0.1;branch=z9hG4bKwarn\r\nCSeq: {cseq}\r\n\ + From: \"{MARKER}\" ;tag=a\r\n\ + To: ;tag=b\r\nCall-ID: warn-call\r\n\r\n" + ); + SipMessage::try_from(text.as_str()).unwrap() +} + +fn request(method: &str, cseq: u32) -> Request { + let start_line = format!("{method} sip:bob@127.0.0.1:5999 SIP/2.0"); + match message(&start_line, &format!("{cseq} {method}")) { + SipMessage::Request(req) => req, + other => panic!("{other:?}"), + } +} + +/// A UAC dialog on an endpoint whose every request send fails. +fn dialog() -> crate::Result> { + let endpoint = crate::EndpointBuilder::new() + .with_transport_layer(TransportLayer::new(Default::default())) + .with_target_locator(Box::new(NoRoute)) + .build(); + let id = DialogId { + call_id: "warn-call".to_string(), + local_tag: "a".to_string(), + remote_tag: "b".to_string(), + }; + let (state_tx, _) = unbounded_channel(); + let (tu_tx, _) = unbounded_channel(); + let invite = request("INVITE", 1); + let inner = DialogInner::new( + TransactionRole::Client, + id, + invite, + endpoint.inner.clone(), + state_tx, + None, + None, + tu_tx, + )?; + Ok(Arc::new(inner)) +} + +#[tokio::test] +async fn test_failed_request_send_warns_without_the_request() -> crate::Result<()> { + let capture = Capture::default(); + let _guard = tracing::subscriber::set_default(capture.clone()); + let inner = dialog()?; + + // send_dialog_request, then send_prack_request. + assert!(inner.do_request(request("INFO", 2)).await.is_err()); + assert!(inner.send_prack_request(request("PRACK", 3)).await.is_err()); + + let warns = capture.lines(Level::WARN, "failed to send request"); + assert_eq!(warns.len(), 2, "{warns:?}"); + assert!(warns.iter().all(|w| !w.contains(MARKER)), "{warns:?}"); + assert!(warns[0].contains("method=INFO") && warns[1].contains("method=PRACK")); + let debugs = capture.lines(Level::DEBUG, "request that failed to send"); + assert_eq!(debugs.len(), 2, "{debugs:?}"); + assert!(debugs.iter().all(|d| d.contains(MARKER)), "{debugs:?}"); + Ok(()) +} + +#[tokio::test] +async fn test_bye_in_early_state_warns_without_the_response() -> crate::Result<()> { + let capture = Capture::default(); + let _guard = tracing::subscriber::set_default(capture.clone()); + let inner = dialog()?; + let SipMessage::Response(ringing) = message("SIP/2.0 180 Ringing", "1 INVITE") else { + unreachable!() + }; + inner.transition(DialogState::Early(inner.id.lock().clone(), ringing))?; + + let errors = [ + ClientInviteDialog { + inner: inner.clone(), + } + .bye() + .await, + InviteDialog { + inner: inner.clone(), + } + .bye() + .await, + ServerInviteDialog { inner }.bye().await, + ]; + for e in errors { + let e = e.expect_err("BYE in Early must fail").to_string(); + assert!(e.contains("(Early)") && !e.contains(MARKER), "{e}"); + } + let warns = capture.lines(Level::WARN, "bye skipped"); + assert_eq!(warns.len(), 3, "{warns:?}"); + assert!(warns.iter().all(|w| !w.contains(MARKER)), "{warns:?}"); + Ok(()) +} + +#[cfg(feature = "websocket")] +#[tokio::test] +async fn test_websocket_parse_failure_warns_without_the_message() -> crate::Result<()> { + use crate::transport::{stream::StreamConnection, websocket::WebSocketConnection}; + use futures::SinkExt; + use tokio_tungstenite::tungstenite::{handshake::server::Response, Message}; + + let capture = Capture::default(); + let _guard = tracing::subscriber::set_default(capture.clone()); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?; + let addr = listener.local_addr()?; + tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let sip = |_: &_, mut resp: Response| { + let headers = resp.headers_mut(); + headers.insert("sec-websocket-protocol", "sip".parse().unwrap()); + Ok(resp) + }; + let mut ws = tokio_tungstenite::accept_hdr_async(stream, sip) + .await + .unwrap(); + let text = format!("NOT SIP \r\n\r\n"); + ws.send(Message::Text(text.into())).await.unwrap(); + ws.close(None).await.ok(); + }); + let remote = SipAddr::new(crate::sip::transport::Transport::Ws, addr.into()); + let conn = WebSocketConnection::connect(&remote, None).await?; + conn.serve_loop(unbounded_channel().0).await?; + + let warns = capture.lines(Level::WARN, "Error parsing SIP message"); + assert_eq!(warns.len(), 1, "{warns:?}"); + assert!(!warns[0].contains(MARKER), "{warns:?}"); + assert!(warns[0].contains("len=")); + // Still logged, below WARN, when it is received (fork R22a: at DEBUG). + assert!(!capture.lines(Level::DEBUG, MARKER).is_empty()); + Ok(()) +} diff --git a/src/platform/mod.rs b/src/platform/mod.rs index 2fffe4e6..9d929439 100644 --- a/src/platform/mod.rs +++ b/src/platform/mod.rs @@ -32,6 +32,7 @@ pub use tokio::sync::mpsc; pub mod atomic64; pub mod net; pub mod select; +pub mod tls; pub use select::{select2, select3, Either, Select2, Select3, Which3}; diff --git a/src/platform/tls.rs b/src/platform/tls.rs new file mode 100644 index 00000000..e86a2d0d --- /dev/null +++ b/src/platform/tls.rs @@ -0,0 +1,86 @@ +//! Backend-agnostic TLS client seam (SIPS). +//! +//! `transport::tls::TlsConnection::connect` consults the connector registered +//! here before falling back to the built-in rustls path, so custom TLS stacks +//! — including embedded ones over embassy-net (embedded-tls, esp-mbedtls) — +//! can provide SIPS without tokio or rustls. Host builds are unaffected when +//! nothing is registered. +//! +//! The connector **owns the TCP dial**: each backend brings its own TCP stack +//! (tokio on the host, embassy-net on embedded), so the seam only passes the +//! target address. Streams expose async-fn IO with owned buffers, which is +//! directly implementable over `embedded-io-async` / rustls alike; the host +//! adapter buffers chunks when bridging into tokio's poll-based halves. + +use crate::Result; +use alloc::boxed::Box; +use alloc::string::String; +use alloc::sync::Arc; +use alloc::vec::Vec; +use core::net::SocketAddr; +use core::sync::atomic::{AtomicPtr, Ordering}; + +/// Backend-agnostic TLS client configuration (PEM encodings). +#[derive(Debug, Clone, Default)] +pub struct TlsClientConfig { + /// SNI / certificate verification name (e.g. `pbx.example.com`). + pub server_name: String, + /// Root certificates to trust (PEM). `None`/empty: backend default trust. + pub root_certs: Option>, + /// Client certificate for mTLS (PEM). + pub client_cert: Option>, + /// Client key for mTLS (PEM). + pub client_key: Option>, +} + +/// An established TLS client session (handshake already completed). +#[async_trait::async_trait] +pub trait TlsStream: Send + Sync + 'static { + /// Reads the next chunk of application data (one backend buffer's worth; + /// chunking carries no message framing). + async fn recv(&self) -> Result>; + + /// Sends all of `data`. + async fn send_all(&self, data: Vec) -> Result<()>; + + /// Closes the session (TLS close_notify + TCP shutdown when supported). + async fn shutdown(&self) -> Result<()>; + + /// Local socket address of the underlying TCP connection. + fn local_addr(&self) -> Result; +} + +/// Factory for outbound TLS sessions (SIPS client role). Owns the TCP dial. +#[async_trait::async_trait] +pub trait TlsConnector: Send + Sync + 'static { + async fn connect( + &self, + config: &TlsClientConfig, + addr: SocketAddr, + ) -> Result>; +} + +static CLIENT_CONNECTOR: AtomicPtr> = AtomicPtr::new(core::ptr::null_mut()); + +/// Registers the process-wide TLS client connector (thread-safe, leak-once). +/// Call once during startup, before any SIPS connection. +pub fn set_client_connector(connector: Arc) { + let leaked = alloc::boxed::Box::leak(alloc::boxed::Box::new(connector)); + CLIENT_CONNECTOR.store(leaked as *mut Arc, Ordering::Release); +} + +/// Clears a previously registered connector (mainly for tests). +pub fn clear_client_connector() { + CLIENT_CONNECTOR.store(core::ptr::null_mut(), Ordering::Release); +} + +/// Returns the registered connector, if any. +pub fn client_connector() -> Option> { + let ptr = CLIENT_CONNECTOR.load(Ordering::Acquire); + if ptr.is_null() { + None + } else { + // The Arc was leaked at registration and outlives the process. + unsafe { Some((*ptr).clone()) } + } +} diff --git a/src/transaction/tests/test_client.rs b/src/transaction/tests/test_client.rs index 16e0187f..9358bd24 100644 --- a/src/transaction/tests/test_client.rs +++ b/src/transaction/tests/test_client.rs @@ -485,3 +485,84 @@ async fn test_invite_2xx_upstream_via_delivery_and_ack() -> Result<()> { } Ok(()) } + +/// RFC 3261 §17.1.2.2: Timer F still runs after a provisional response. A +/// BYE answered with one 1xx and then nothing must end with a 408 at 64*T1, +/// on unreliable and reliable transports alike. +#[tokio::test] +async fn test_non_invite_timer_f_after_provisional() -> Result<()> { + use crate::transaction::{endpoint::EndpointOption, EndpointBuilder}; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + use tokio::net::{TcpListener, UdpSocket}; + + // The request with its start line swapped for a provisional status line. + fn provisional(code: u16, request: &[u8]) -> Vec { + let text = String::from_utf8_lossy(request); + let (_, headers) = text.split_once("\r\n").unwrap(); + format!("SIP/2.0 {code} Provisional\r\n{headers}").into_bytes() + } + + for (code, tcp) in [(100, false), (180, false), (183, true)] { + let tl = crate::transport::TransportLayer::new(Default::default()); + let udp = UdpConnection::create_connection("127.0.0.1:0".parse()?, None, None).await?; + tl.add_transport(udp.into()); + let t1 = Duration::from_millis(20); + let option = EndpointOption { + t1, + t1x64: t1 * 64, + ..Default::default() + }; + let endpoint = EndpointBuilder::new() + .with_transport_layer(tl) + .with_option(option) + .build(); + + // The peer answers the BYE with one provisional, then stays silent. + let socket = UdpSocket::bind("127.0.0.1:0").await?; + let listener = TcpListener::bind("127.0.0.1:0").await?; + let uri = match tcp { + true => format!("sip:bob@{};transport=tcp", listener.local_addr()?), + false => format!("sip:bob@{}", socket.local_addr()?), + }; + let peer = tokio::spawn(async move { + let mut buf = vec![0u8; 4096]; + if tcp { + let (mut stream, _) = listener.accept().await?; + let mut len = 0; + while !buf[..len].windows(4).any(|w| w == b"\r\n\r\n") { + len += stream.read(&mut buf[len..]).await?; + } + stream.write_all(&provisional(code, &buf[..len])).await?; + std::future::pending::<()>().await; // keep the connection open + } else { + let (len, src) = socket.recv_from(&mut buf).await?; + socket.send_to(&provisional(code, &buf[..len]), src).await?; + } + std::future::pending::>().await + }); + + let mut bye = make_invite_request(&uri)?; + bye.method = crate::sip::Method::Bye; + bye.headers.unique_push(CSeq::new("2 BYE").into()); + let key = TransactionKey::from_request(&bye, TransactionRole::Client)?; + let mut tx = Transaction::new_client(key, bye, endpoint.inner.clone(), None); + let mut codes = vec![]; + let run = async { + tx.send().await?; + while let Some(SipMessage::Response(resp)) = tx.receive().await { + codes.push(resp.status_code.code()); + } + Ok::<_, crate::Error>(()) + }; + // `receive` returning None means the transaction terminated. + let terminated = select! { + r = run => r.map(|_| true)?, + _ = endpoint.serve() => panic!("endpoint stopped"), + _ = sleep(t1 * 64 * 3) => false, + }; + peer.abort(); + assert_eq!(codes, [code, 408], "{code} over tcp={tcp}"); + assert!(terminated, "{code} over tcp={tcp}: not terminated"); + } + Ok(()) +} diff --git a/src/transaction/transaction.rs b/src/transaction/transaction.rs index fe766c1a..56ff6b49 100644 --- a/src/transaction/transaction.rs +++ b/src/transaction/transaction.rs @@ -1238,7 +1238,12 @@ impl Transaction { } } TransactionState::Proceeding => { - if let TransactionTimer::TimerC(_) = timer { + // Timer C (client INVITE), or Timer F, run as Timer B, for a + // non-INVITE client (RFC 3261 §17.1.2.2). + if matches!(timer, TransactionTimer::TimerC(_)) + || (matches!(timer, TransactionTimer::TimerB(_)) + && self.transaction_type == TransactionType::ClientNonInvite) + { // Inform TU about timeout let mut timeout_response = self.endpoint_inner.make_response( &self.original, @@ -1693,9 +1698,9 @@ impl Transaction { TransactionType::ServerNonInvite => { self.last_response.take().map(SipMessage::Response) } - // Kept on the transaction: the matching ACK ends an Accepted - // server INVITE (RFC 6026 §7.1) before the dialog reads the - // 2xx it confirms (`DialogState::Confirmed`). + // Kept: the matching ACK terminates an Accepted server + // INVITE before the dialog reads the 2xx for + // `DialogState::Confirmed`. TransactionType::ServerInvite => { self.last_response.clone().map(SipMessage::Response) } diff --git a/src/transport/tls.rs b/src/transport/tls.rs index f26462b4..ff294b91 100644 --- a/src/transport/tls.rs +++ b/src/transport/tls.rs @@ -7,6 +7,9 @@ use super::{ use crate::platform::CancellationToken; use crate::sip::SipMessage; use crate::{error::Error, transport::transport_layer::TransportLayerInnerRef, Result}; +use core::future::Future; +use core::pin::Pin; +use core::task::{Context, Poll}; use rustls::client::danger::ServerCertVerifier; use rustls::crypto::CryptoProvider; use rustls::server::{ClientHello, ResolvesServerCert}; @@ -436,6 +439,117 @@ impl fmt::Debug for TlsListenerConnection { type TlsClientStream = tokio_rustls::client::TlsStream; type TlsServerStream = tokio_rustls::server::TlsStream; +/// Client halves are boxed so rustls and platform-seam streams share one +/// `TlsConnectionInner::Client` representation. +type ClientReadHalf = Box; +type ClientWriteHalf = Box; + +type BoxFuture = Pin + Send>>; + +/// Bridges a seam TLS stream (async-fn, owned-data) back into the tokio +/// world. Chunks from `recv()` are buffered; `send_all`/`shutdown` futures +/// hold an owned `Arc` + `Vec`, so nothing borrows across polls. +struct TokioSeamStream { + inner: Arc, + recv_fut: Option>>>, + send_fut: Option>>, + shutdown_fut: Option>>, + chunk: alloc::collections::VecDeque, +} + +impl TokioSeamStream { + fn new(inner: Arc) -> Self { + Self { + inner, + recv_fut: None, + send_fut: None, + shutdown_fut: None, + chunk: alloc::collections::VecDeque::new(), + } + } +} + +impl tokio::io::AsyncRead for TokioSeamStream { + fn poll_read( + self: Pin<&mut Self>, + cx: &mut Context<'_>, + rbuf: &mut tokio::io::ReadBuf<'_>, + ) -> Poll> { + let this = self.get_mut(); + loop { + if !this.chunk.is_empty() { + let n = this.chunk.len().min(rbuf.remaining()); + let data: Vec = this.chunk.drain(..n).collect(); + rbuf.put_slice(&data); + return Poll::Ready(Ok(())); + } + if this.recv_fut.is_none() { + let inner = this.inner.clone(); + this.recv_fut = Some(Box::pin(async move { inner.recv().await })); + } + match this.recv_fut.as_mut().unwrap().as_mut().poll(cx) { + Poll::Ready(Ok(data)) => { + this.recv_fut = None; + if data.is_empty() { + // Peer closed the stream. + return Poll::Ready(Ok(())); + } + this.chunk.extend(data); + } + Poll::Ready(Err(e)) => return Poll::Ready(Err(seam_io_err(e))), + Poll::Pending => return Poll::Pending, + } + } + } +} + +fn seam_io_err(e: crate::Error) -> std::io::Error { + std::io::Error::other(e.to_string()) +} + +impl tokio::io::AsyncWrite for TokioSeamStream { + fn poll_write( + self: Pin<&mut Self>, + cx: &mut Context<'_>, + buf: &[u8], + ) -> Poll> { + let this = self.get_mut(); + if this.send_fut.is_none() { + let inner = this.inner.clone(); + let data = buf.to_vec(); + this.send_fut = Some(Box::pin(async move { inner.send_all(data).await })); + } + match this.send_fut.as_mut().unwrap().as_mut().poll(cx) { + Poll::Ready(Ok(())) => { + this.send_fut = None; + Poll::Ready(Ok(buf.len())) + } + Poll::Ready(Err(e)) => Poll::Ready(Err(seam_io_err(e))), + Poll::Pending => Poll::Pending, + } + } + + fn poll_flush(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) // send_all is already awaited to completion + } + + fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + let this = self.get_mut(); + if this.shutdown_fut.is_none() { + let inner = this.inner.clone(); + this.shutdown_fut = Some(Box::pin(async move { inner.shutdown().await })); + } + match this.shutdown_fut.as_mut().unwrap().as_mut().poll(cx) { + Poll::Ready(Ok(())) => { + this.shutdown_fut = None; + Poll::Ready(Ok(())) + } + Poll::Ready(Err(e)) => Poll::Ready(Err(seam_io_err(e))), + Poll::Pending => Poll::Pending, + } + } +} + // TLS connection - uses enum to handle both client and server streams #[derive(Clone)] pub struct TlsConnection { @@ -445,14 +559,7 @@ pub struct TlsConnection { #[derive(Clone)] enum TlsConnectionInner { - Client( - Arc< - StreamConnectionInner< - tokio::io::ReadHalf, - tokio::io::WriteHalf, - >, - >, - ), + Client(Arc>), Server( Arc< StreamConnectionInner< @@ -479,6 +586,13 @@ impl TlsConnection { custom_verifier: Option>, cancel_token: Option, ) -> Result { + // A registered platform connector (embedded backends) takes + // precedence over the built-in rustls path. + if let Some(connector) = crate::platform::tls::client_connector() { + return Self::connect_via_platform(remote_addr, tls_config, cancel_token, connector) + .await; + } + let mut root_store = RootCertStore::empty(); // Load CA certificates if provided @@ -552,6 +666,10 @@ impl TlsConnection { let tls_stream = connector.connect(server_name, stream).await?; let (read_half, write_half) = tokio::io::split(tls_stream); + let (read_half, write_half) = ( + Box::new(read_half) as ClientReadHalf, + Box::new(write_half) as ClientWriteHalf, + ); let connection = Self { inner: TlsConnectionInner::Client(Arc::new(StreamConnectionInner::new( @@ -570,6 +688,67 @@ impl TlsConnection { Ok(connection) } + // Client connect through the platform TLS seam (embedded backends). + async fn connect_via_platform( + remote_addr: &SipAddr, + tls_config: Option<&TlsConfig>, + cancel_token: Option, + connector: Arc, + ) -> Result { + let config = crate::platform::tls::TlsClientConfig { + server_name: tls_config + .and_then(|c| c.sni_hostname.clone()) + .unwrap_or_else(|| match &remote_addr.addr.host { + crate::sip::Host::Domain(domain) => domain.to_string(), + crate::sip::Host::IpAddr(ip) => ip.to_string(), + }), + root_certs: tls_config.and_then(|c| c.ca_certs.clone()), + client_cert: tls_config.and_then(|c| c.client_cert.clone()), + client_key: tls_config.and_then(|c| c.client_key.clone()), + }; + + let socket_addr = match &remote_addr.addr.host { + crate::sip::Host::Domain(domain) => { + let port = remote_addr.addr.port.as_ref().map_or(5061, |p| p.value()); + format!("{}:{}", domain, port).parse()? + } + crate::sip::Host::IpAddr(ip) => { + let port = remote_addr.addr.port.as_ref().map_or(5061, |p| p.value()); + SocketAddr::new(*ip, port) + } + }; + + // The connector owns the TCP dial (backend-specific TCP stack). + let tls_stream = connector.connect(&config, socket_addr).await?; + let local = tls_stream.local_addr()?; + + let local_addr = SipAddr { + r#type: Some(crate::sip::transport::Transport::Tls), + addr: local.into(), + }; + let (read_half, write_half) = tokio::io::split(TokioSeamStream::new(Arc::from(tls_stream))); + let (read_half, write_half) = ( + Box::new(read_half) as ClientReadHalf, + Box::new(write_half) as ClientWriteHalf, + ); + + let connection = Self { + inner: TlsConnectionInner::Client(Arc::new(StreamConnectionInner::new( + local_addr.clone(), + remote_addr.clone(), + read_half, + write_half, + ))), + cancel_token, + }; + debug!( + "Created TLS client connection (platform seam): {} -> {}", + local_addr, remote_addr + ); + + Ok(connection) + } + // Create TLS connection from existing client TLS stream pub async fn from_client_stream( stream: TlsClientStream, @@ -583,6 +762,10 @@ impl TlsConnection { // Split stream into read and write halves let (read_half, write_half) = tokio::io::split(stream); + let (read_half, write_half) = ( + Box::new(read_half) as ClientReadHalf, + Box::new(write_half) as ClientWriteHalf, + ); // Create TLS connection let connection = Self { @@ -704,3 +887,94 @@ impl fmt::Debug for TlsConnection { fmt::Display::fmt(self, f) } } + +#[cfg(test)] +mod platform_seam_tests { + use super::*; + use crate::platform::tls as seam; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + /// Plaintext pass-through "TLS": proves the registered connector was used + /// instead of rustls (which would fail to handshake with a plain echoer). + /// The connector dials `addr` itself, exactly like an embedded backend + /// (embassy-net + embedded-tls) would. + struct MockConnector; + + #[async_trait::async_trait] + impl seam::TlsConnector for MockConnector { + async fn connect( + &self, + _config: &seam::TlsClientConfig, + addr: SocketAddr, + ) -> Result> { + let tcp = tokio::net::TcpStream::connect(addr).await?; + let local = tcp.local_addr()?; + Ok(Box::new(MockStream { + inner: tcp.into(), + local, + })) + } + } + + struct MockStream { + inner: tokio::sync::Mutex, + local: SocketAddr, + } + + #[async_trait::async_trait] + impl seam::TlsStream for MockStream { + async fn recv(&self) -> Result> { + let mut buf = vec![0u8; 512]; + let mut tcp = self.inner.lock().await; + let n = tcp.read(&mut buf).await?; + buf.truncate(n); + Ok(buf) + } + + async fn send_all(&self, data: Vec) -> Result<()> { + let mut tcp = self.inner.lock().await; + tcp.write_all(&data).await?; + Ok(()) + } + + async fn shutdown(&self) -> Result<()> { + Ok(()) + } + + fn local_addr(&self) -> Result { + Ok(self.local) + } + } + + #[tokio::test] + async fn platform_connector_takes_precedence_over_rustls() { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let echo = tokio::spawn(async move { + let (mut sock, _) = listener.accept().await.unwrap(); + let mut buf = [0u8; 4]; + tokio::io::AsyncReadExt::read_exact(&mut sock, &mut buf) + .await + .unwrap(); + assert_eq!(&buf, b"ping"); + tokio::io::AsyncWriteExt::write_all(&mut sock, b"pong") + .await + .unwrap(); + }); + + seam::set_client_connector(std::sync::Arc::new(MockConnector)); + + let remote = SipAddr { + r#type: Some(crate::sip::transport::Transport::Tls), + addr: addr.into(), + }; + let conn = TlsConnection::connect(&remote, None, None, None) + .await + .expect("seam connector should bypass rustls entirely"); + + conn.send_raw(b"ping").await.unwrap(); + echo.await.unwrap(); + + seam::clear_client_connector(); + } +} diff --git a/src/transport/websocket.rs b/src/transport/websocket.rs index 710288ef..54a040d3 100644 --- a/src/transport/websocket.rs +++ b/src/transport/websocket.rs @@ -342,7 +342,6 @@ impl StreamConnection for WebSocketConnection { } Err(e) => { warn!(error = %e, src = %remote_addr, len = text.len(), "Error parsing SIP message"); - debug!(src = %remote_addr, raw_message = ?text.as_str(), "unparsable SIP message"); } } }