Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -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
Expand Down
35 changes: 33 additions & 2 deletions src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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.
Expand All @@ -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,
Expand All @@ -138,6 +144,7 @@ pub(crate) struct Client {
api_secret: Vec<u8>,
server_token: String,
base_url: Url,
coordinator_ws_url: Url,
http: reqwest::Client,
retry: RetryConfig,
log_bodies: bool,
Expand Down Expand Up @@ -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)
Expand All @@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -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();
Expand Down
4 changes: 3 additions & 1 deletion src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)]
Expand Down
70 changes: 2 additions & 68 deletions src/rtc/coordinator/ws.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<MaybeTlsStream<TcpStream>>;

/// 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.
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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<Url> {
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 {
Expand Down Expand Up @@ -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")
Expand Down
3 changes: 1 addition & 2 deletions src/rtc/join/lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
10 changes: 5 additions & 5 deletions src/rtc/join/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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");
Expand Down Expand Up @@ -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");
Expand Down
Loading