diff --git a/CHANGELOG.md b/CHANGELOG.md index cd25a8e..b433800 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,17 @@ # Unreleased +## Breaking Changes + +### Coordinator WebSocket URL is configured separately + +`ClientConfig::coordinator_ws_url` sets the coordinator connect WebSocket URL. +It accepts only `ws` and `wss` URLs and defaults to +`DEFAULT_COORDINATOR_WS_URL`. The URL no longer follows `base_url`: a client +for a staging or local environment must set both fields. Code that builds +`ClientConfig` with a struct literal must set the new field or use +`..ClientConfig::default()`. `DEFAULT_COORDINATOR_WS_URL` moved from +`rtc::coordinator::ws` to the crate root. + ## New Features ### Video REST: advanced call statistics and reporting diff --git a/src/client.rs b/src/client.rs index c0c2a05..39ca7ac 100644 --- a/src/client.rs +++ b/src/client.rs @@ -25,6 +25,8 @@ use crate::token; /// Default coordinator base URL (matches getstream-go's `DefaultBaseURL`). pub const DEFAULT_BASE_URL: &str = "https://chat.stream-io-api.com"; +/// Default coordinator connect WebSocket URL (videosdk `defaultOptions.wsURL`). +pub const DEFAULT_COORDINATOR_WS_URL: &str = "wss://video.stream-io-api.com/api/v2/connect"; const DEFAULT_REQUEST_TIMEOUT: Duration = Duration::from_secs(30); const DEFAULT_CONNECT_TIMEOUT: Duration = Duration::from_secs(10); @@ -101,6 +103,9 @@ impl Default for RetryConfig { pub struct ClientConfig { /// Coordinator base URL. Defaults to [`DEFAULT_BASE_URL`]. pub base_url: String, + /// Coordinator connect WebSocket URL, including the path. Only `ws` and `wss` + /// schemes are accepted. Defaults to [`DEFAULT_COORDINATOR_WS_URL`]. + pub coordinator_ws_url: String, /// Per-request timeout. Default 30s. pub request_timeout: Duration, /// TCP + TLS connect timeout. Default 10s. @@ -121,6 +126,7 @@ impl Default for ClientConfig { fn default() -> Self { Self { base_url: DEFAULT_BASE_URL.to_string(), + coordinator_ws_url: DEFAULT_COORDINATOR_WS_URL.to_string(), request_timeout: DEFAULT_REQUEST_TIMEOUT, connect_timeout: DEFAULT_CONNECT_TIMEOUT, idle_timeout: DEFAULT_IDLE_TIMEOUT, @@ -138,6 +144,7 @@ pub(crate) struct Client { api_secret: Vec, server_token: String, base_url: Url, + coordinator_ws_url: Url, http: reqwest::Client, retry: RetryConfig, log_bodies: bool, @@ -187,6 +194,18 @@ impl Client { let base_url = Url::parse(&config.base_url) .map_err(|e| Error::Config(format!("invalid base URL {:?}: {e}", config.base_url)))?; + let coordinator_ws_url = Url::parse(&config.coordinator_ws_url).map_err(|e| { + Error::Config(format!( + "invalid coordinator WebSocket URL {:?}: {e}", + config.coordinator_ws_url + )) + })?; + if !matches!(coordinator_ws_url.scheme(), "ws" | "wss") { + return Err(Error::Config(format!( + "coordinator WebSocket URL {:?} must use ws or wss", + config.coordinator_ws_url + ))); + } let http = reqwest::Client::builder() .pool_max_idle_per_host(config.max_conns_per_host) @@ -204,6 +223,7 @@ impl Client { api_secret: secret_bytes, server_token, base_url, + coordinator_ws_url, http, retry: config.retry, log_bodies: config.log_bodies, @@ -221,8 +241,8 @@ impl Client { &self.api_secret } - pub(crate) fn base_url(&self) -> &Url { - &self.base_url + pub(crate) fn coordinator_ws_url(&self) -> &Url { + &self.coordinator_ws_url } /// The shared `reqwest` client (connection pool). Used by the RTC layer to @@ -621,6 +641,17 @@ mod tests { assert!(!redacted.contains("must-not-leak")); } + #[test] + fn coordinator_ws_url_must_use_a_websocket_scheme() { + let config = ClientConfig { + coordinator_ws_url: "https://video.stream-io-api.com/api/v2/connect".to_owned(), + ..ClientConfig::default() + }; + let error = Client::new("key".to_owned(), "secret".to_owned(), config) + .expect_err("https coordinator WebSocket URL"); + assert!(matches!(error, Error::Config(_))); + } + #[test] fn compatibility_limits_are_conservative_and_configurable() { let limits = NetworkLimits::default(); diff --git a/src/lib.rs b/src/lib.rs index d9e62cd..7c92292 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -15,7 +15,9 @@ use std::sync::Arc; use client::Client; #[doc(inline)] -pub use client::{ClientConfig, DEFAULT_BASE_URL, NetworkLimits, RetryConfig}; +pub use client::{ + ClientConfig, DEFAULT_BASE_URL, DEFAULT_COORDINATOR_WS_URL, NetworkLimits, RetryConfig, +}; #[doc(inline)] pub use error::{ApiError, Error, Result, TokenError, WebhookError}; #[doc(inline)] diff --git a/src/rtc/coordinator/ws.rs b/src/rtc/coordinator/ws.rs index 36efd24..513241b 100644 --- a/src/rtc/coordinator/ws.rs +++ b/src/rtc/coordinator/ws.rs @@ -26,15 +26,13 @@ use tokio_tungstenite::tungstenite::protocol::WebSocketConfig; use tokio_tungstenite::{MaybeTlsStream, WebSocketStream, connect_async_with_config}; use url::Url; -use crate::client::{DEFAULT_BASE_URL, DEFAULT_MAX_WEBSOCKET_MESSAGE_BYTES}; +use crate::client::DEFAULT_MAX_WEBSOCKET_MESSAGE_BYTES; use crate::rtc::error::{Result, RtcError, SfuTimeoutError}; use crate::rtc::identity; type WsStream = WebSocketStream>; -/// Default coordinator connect WebSocket URL (videosdk `defaultOptions.wsURL`). -pub const DEFAULT_COORDINATOR_WS_URL: &str = "wss://video.stream-io-api.com/api/v2/connect"; pub(crate) const DEFAULT_COORDINATOR_HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(15); /// The user identity sent in the coordinator auth message. @@ -139,7 +137,7 @@ impl CoordinatorEvent { /// into a [`CoordinatorWs`] (send path), [`CoordinatorEvents`] (event stream), /// and the [`Connected`] handshake result. /// -/// `ws_url` is usually [`DEFAULT_COORDINATOR_WS_URL`]. The `api_key`/`user_id` +/// `ws_url` is usually [`crate::DEFAULT_COORDINATOR_WS_URL`]. The `api_key`/`user_id` /// are added as query params (`stream-auth-type=jwt`), matching videosdk. pub async fn connect( ws_url: &str, @@ -310,42 +308,6 @@ async fn await_connection_ok( )) } -/// Resolve the coordinator connect WebSocket URL for the configured REST -/// environment. The default production base uses the pinned -/// [`DEFAULT_COORDINATOR_WS_URL`]; any custom base (staging, local) derives the -/// WebSocket URL from it so a reconfigured client never silently talks to -/// production. -pub(crate) fn coordinator_ws_url(base_url: &Url) -> Result { - let default_base = Url::parse(DEFAULT_BASE_URL) - .map_err(|error| RtcError::Url(format!("{DEFAULT_BASE_URL:?}: {error}")))?; - if base_url.scheme() == default_base.scheme() - && base_url.host_str() == default_base.host_str() - && base_url.port_or_known_default() == default_base.port_or_known_default() - { - return Url::parse(DEFAULT_COORDINATOR_WS_URL) - .map_err(|error| RtcError::Url(format!("{DEFAULT_COORDINATOR_WS_URL:?}: {error}"))); - } - - let mut url = base_url.clone(); - let scheme = match url.scheme() { - "http" => "ws", - "https" => "wss", - "ws" => "ws", - "wss" => "wss", - other => { - return Err(RtcError::Url(format!( - "unsupported coordinator base URL scheme {other:?}" - ))); - } - }; - url.set_scheme(scheme) - .map_err(|_| RtcError::Url(format!("cannot set WebSocket scheme on {base_url}")))?; - url.set_path("/api/v2/connect"); - url.set_query(None); - url.set_fragment(None); - Ok(url) -} - /// The send half of an authenticated coordinator WebSocket. #[derive(Debug)] pub struct CoordinatorWs { @@ -466,34 +428,6 @@ mod tests { )); } - #[test] - fn coordinator_url_follows_configured_rest_environment() { - let staging = - Url::parse("https://video-edge-staging.example.com/video?source=test").expect("url"); - assert_eq!( - coordinator_ws_url(&staging) - .expect("staging coordinator URL") - .as_str(), - "wss://video-edge-staging.example.com/api/v2/connect" - ); - - let local = Url::parse("http://127.0.0.1:3030/custom").expect("url"); - assert_eq!( - coordinator_ws_url(&local) - .expect("local coordinator URL") - .as_str(), - "ws://127.0.0.1:3030/api/v2/connect" - ); - - let production = Url::parse(DEFAULT_BASE_URL).expect("url"); - assert_eq!( - coordinator_ws_url(&production) - .expect("production coordinator URL") - .as_str(), - DEFAULT_COORDINATOR_WS_URL - ); - } - #[tokio::test] async fn coordinator_handshake_ignores_other_events_and_times_out() { let listener = TcpListener::bind("127.0.0.1:0") diff --git a/src/rtc/join/lifecycle.rs b/src/rtc/join/lifecycle.rs index 0c07087..0a485c6 100644 --- a/src/rtc/join/lifecycle.rs +++ b/src/rtc/join/lifecycle.rs @@ -701,9 +701,8 @@ impl RtcCore { return Err(join_cancelled()); } let auth = WsAuthMessage::video(user_token, ConnectUserDetails::new(user_id)); - let coordinator_url = coordinator::ws::coordinator_ws_url(self.client.base_url())?; let (mut coordinator, mut events, connected) = coordinator::ws::connect_with_limit( - coordinator_url.as_str(), + self.client.coordinator_ws_url().as_str(), &self.api_key, user_id, &auth, diff --git a/src/rtc/join/tests.rs b/src/rtc/join/tests.rs index 6300f69..b508b4d 100644 --- a/src/rtc/join/tests.rs +++ b/src/rtc/join/tests.rs @@ -189,7 +189,7 @@ async fn fake_coordinator() -> (String, tokio::task::JoinHandle<()>) { .expect("send connection.ok"); while let Some(Ok(_)) = socket.next().await {} }); - (format!("http://{address}"), server) + (format!("ws://{address}"), server) } /// Every request the fake SFU received until the client socket closed. @@ -399,9 +399,9 @@ async fn leave_closes_a_connection_owned_by_a_cancelled_join() { #[tokio::test] async fn generation_change_closes_the_coordinator_socket() { - let (base_url, coordinator) = fake_coordinator().await; + let (coordinator_ws_url, coordinator) = fake_coordinator().await; let core = test_core_with_config(ClientConfig { - base_url, + coordinator_ws_url, ..ClientConfig::default() }); let generation = prepare_joined_core(&core, "alice"); @@ -552,9 +552,9 @@ async fn leave_that_overlaps_a_new_join_keeps_the_new_join_state() { #[tokio::test] async fn stale_coordinator_stop_keeps_the_current_coordinator() { - let (base_url, coordinator) = fake_coordinator().await; + let (coordinator_ws_url, coordinator) = fake_coordinator().await; let core = test_core_with_config(ClientConfig { - base_url, + coordinator_ws_url, ..ClientConfig::default() }); let first = prepare_joined_core(&core, "alice");