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
11 changes: 11 additions & 0 deletions src/rtc/coordinator/ws.rs
Original file line number Diff line number Diff line change
Expand Up @@ -442,6 +442,17 @@ mod tests {
assert!(debug.contains("user"));
}

#[test]
fn ws_auth_message_serializes_video_product() {
let auth = WsAuthMessage::video("jwt-token", ConnectUserDetails::new("agent"));
let json = serde_json::to_value(&auth).expect("serialize");
assert_eq!(json["token"], "jwt-token");
assert_eq!(json["user_details"]["id"], "agent");
assert_eq!(json["products"][0], "video");
// Optional user fields are omitted when unset.
assert!(json["user_details"].get("name").is_none());
}

#[test]
fn coordinator_message_limit_accepts_exact_and_rejects_oversized_input() {
ensure_message_size("coordinator test message", 64, 64).expect("exact limit");
Expand Down
28 changes: 25 additions & 3 deletions src/rtc/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,7 @@ pub enum RtcError {
#[error(transparent)]
Join(#[from] SfuJoinError),

/// A client deadline elapsed waiting for the SFU (WS open or `JoinResponse`).
/// A client deadline elapsed, e.g. waiting for the SFU `JoinResponse`.
#[error(transparent)]
Timeout(#[from] SfuTimeoutError),

Expand Down Expand Up @@ -333,9 +333,9 @@ impl SfuJoinError {
}
}

/// A client-side deadline elapsed waiting for the SFU (WS open or `JoinResponse`).
/// A client-side deadline elapsed.
#[derive(Debug, Clone, thiserror::Error)]
#[error("sfu timeout waiting for {what} after {}ms", timeout.as_millis())]
#[error("timeout waiting for {what} after {}ms", timeout.as_millis())]
pub struct SfuTimeoutError {
/// What we were waiting for, e.g. `"join response"`.
pub what: String,
Expand Down Expand Up @@ -421,6 +421,28 @@ pub struct TwirpError {
mod tests {
use super::*;

#[test]
fn from_signal_error_maps_only_real_codes() {
// UNSPECIFIED (and absent) is success.
assert!(RtcError::from_signal_error(None).is_ok());
assert!(
RtcError::from_signal_error(Some(models::Error {
code: models::ErrorCode::Unspecified as i32,
message: String::new(),
should_retry: false,
}))
.is_ok()
);
// A real code becomes an error.
let err = RtcError::from_signal_error(Some(models::Error {
code: models::ErrorCode::ParticipantSignalLost as i32,
message: "boom".to_owned(),
should_retry: true,
}))
.expect_err("should be an error");
assert!(matches!(err, RtcError::Signal { .. }));
}

#[test]
fn join_error_codes_match_sfu() {
assert!(is_join_error_code(ErrorCode::SfuFull as i32));
Expand Down
32 changes: 23 additions & 9 deletions src/rtc/join/connection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -129,10 +129,10 @@ pub(super) fn register_connection_state(
tracer.trace("connectionstatechange", json!(state.to_string()));
if state == RTCPeerConnectionState::Connected {
ever_connected.store(true, Ordering::SeqCst);
core.caps
.lock()
.unwrap_or_else(|e| e.into_inner())
.reset_ice();
let mut lifecycle = core.lifecycle.lock().unwrap_or_else(|e| e.into_inner());
if lifecycle.generation == generation {
lifecycle.failure_limits.reset_ice();
}
}
if state == RTCPeerConnectionState::Failed && reconnect_enabled.load(Ordering::SeqCst) {
let (pub_h, sub_h) = core.pc_health().await;
Expand Down Expand Up @@ -220,6 +220,20 @@ pub(super) async fn handle_event(
return Ok(());
}
use sfu_event::EventPayload as E;
// During a migration the old SFU still sends events. Its WebRTC commands
// are not for the new connection.
if !context.reconnect_enabled.load(Ordering::SeqCst)
&& matches!(
payload,
E::SubscriberOffer(_)
| E::IceTrickle(_)
| E::ChangePublishOptions(_)
| E::ChangePublishQuality(_)
| E::IceRestart(_)
)
{
return Ok(());
}
match payload {
E::SubscriberOffer(offer) => {
negotiate_subscriber(
Expand All @@ -245,30 +259,30 @@ pub(super) async fn handle_event(
}
E::ParticipantJoined(ev) => {
if let Some(p) = ev.participant {
core.roster_upsert(&p);
core.upsert_participant(&p);
core.recompute_subscriptions_for_generation(context.generation)
.await?;
let _ = core.events_tx.send(CallEvent::ParticipantJoined(p));
}
}
E::ParticipantLeft(ev) => {
if let Some(p) = ev.participant {
core.roster_remove(&p.session_id);
core.remove_participant(&p.session_id);
core.recompute_subscriptions_for_generation(context.generation)
.await?;
let _ = core.events_tx.send(CallEvent::ParticipantLeft(p));
}
}
E::ParticipantUpdated(ev) => {
if let Some(p) = ev.participant {
core.roster_upsert(&p);
core.upsert_participant(&p);
core.recompute_subscriptions_for_generation(context.generation)
.await?;
let _ = core.events_tx.send(CallEvent::ParticipantUpdated(p));
}
}
E::TrackPublished(ev) => {
core.roster_add_track(
core.add_published_track(
&ev.user_id,
&ev.session_id,
ev.r#type,
Expand All @@ -283,7 +297,7 @@ pub(super) async fn handle_event(
});
}
E::TrackUnpublished(ev) => {
core.roster_remove_track(&ev.session_id, ev.r#type);
core.remove_published_track(&ev.session_id, ev.r#type);
core.recompute_subscriptions_for_generation(context.generation)
.await?;
let _ = core.events_tx.send(CallEvent::TrackUnpublished {
Expand Down
Loading
Loading