Calls, SFU & Signaling Relay
This page documents the real-time calling subsystem of the Rust backend (streaming-backend): the embedded WebRTC SFU built on the str0m crate, the CallRelayService signal bus that fans signaling messages in and out, and the gRPC bridge (signals.rs) that converts internal CallSignal messages into protobuf SignalEvents streamed to clients.
Purpose and Scope
This page covers the end-to-end signaling path for voice/video calls:
- The embedded Rust SFU (
services/sfu/) — an in-process WebRTC Selective Forwarding Unit built onstr0m, which forwards media between participants without transcoding. - The call relay (
services/call_relay.rs) and its subscribers (services/relay_subscriber.rs) — the in-process message bus that carriesCallSignalmessages between the SFU, the gRPC layer, and long-poll/streaming subscribers. - The gRPC signaling bridge (
grpc/signals.rs) that maps internalCallSignalvalues to the wireSignalEventproto. - The startup wiring in
main.rsthat connects the SFU’s outgoing signal channel to the relay, and how both services are registered in the sharedAppState. - Error mapping from
str0m::error::RtcErrorinto the backend’sAppError.
Related topics intentionally left to sibling pages: authentication/authorization middleware (see the auth gRPC service), call ticket issuance (services/call_ticket.rs), chat history persistence (grpc/history.rs), and general service configuration in config.rs.
Overview
Calls in this backend are built around three cooperating pieces:
-
SfuService — an in-process SFU. Unlike a standalone media server (e.g. LiveKit or Janus), the SFU runs inside the same Tokio runtime as the API. It uses the
str0mWebRTC library to manage peer connections, SDP negotiation, and RTP/RTCP media forwarding. The SFU is created with a signal sender channel (sfu_tx) through which it pushes outgoing signaling messages (answer SDP, ICE candidates, session events) back into the relay. -
CallRelayService — a pub/sub relay for
CallSignalmessages. It is the central hub: gRPC handlers and the SFU both publish signals into it, and relay subscribers receive them. This decouples the WebRTC media plane from the control-plane signaling transport, so the same relay can serve gRPC streaming clients, HTTP subscribers, and internal consumers. -
gRPC signaling bridge —
call_signal_to_event()ingrpc/signals.rsconverts an internalCallSignalinto a protobufSignalEvent(event_type,from_id,group_id,payload,created_at) that is streamed to connected clients by the chat gRPC service.
The design intent is separation of concerns: the SFU speaks pure WebRTC (str0m), the relay speaks pure application signaling (CallSignal), and the gRPC layer translates between CallSignal and the wire proto. No layer knows the transport details of the others.
Architecture
The diagram reflects the wiring verified in main.rs: SfuService::new(sfu_tx) receives an unbounded Tokio channel sender; a spawned task reads sfu_rx and forwards every signal into call_relay.send_signal(sig), closing the loop from the SFU back to subscribers. Both call_relay and sfu are stored as Arc in AppState (see state.rs), making them available to every gRPC handler.
Main Content
Startup wiring: connecting the SFU to the relay
The composition root in main.rs establishes the signal path before any service starts. It creates an unbounded mpsc channel, hands the sender to the SFU, and spawns a relay-forwarder task that consumes the receiver:
// Embedded Rust SFU (str0m) — outgoing signals forward to the relay
let (sfu_tx, mut sfu_rx) = tokio::sync::mpsc::unbounded_channel();
let relay_for_sfu = call_relay.clone();
tokio::spawn(async move {
while let Some(sig) = sfu_rx.recv().await {
let _ = relay_for_sfu.send_signal(sig).await;
}
});
let sfu = crate::services::sfu::SfuService::new(sfu_tx).await?;
Source: main.rs
Key design decisions visible here:
- Unbounded channel: signaling messages are small, latency-sensitive control messages; an unbounded channel avoids back-pressure stalls in the media path while the relay forwarder drains it.
Arc-shared relay:call_relayis cloned (relay_for_sfu) and moved into the spawned task, while the original handle is stored inAppStatefor gRPC handlers — one relay instance, many senders.- Errors are swallowed (
let _ = ...): a failed relay send must not crash the SFU task; the media plane continues even if a subscriber is gone. - Async construction:
SfuService::new(...).await?implies the SFU performs asynchronous setup (e.g. binding UDP candidates or pre-allocating session state) before it is stored inAppStateasArc<SfuService>.
The SFU service module layout
The SFU is not a single file; it is a cohesive module with focused submodules:
| File | Role |
|---|---|
services/sfu/mod.rs |
SfuService definition, new() constructor, public API |
services/sfu/session.rs |
Per-call WebRTC session lifecycle (SDP offer/answer, ICE) |
services/sfu/session_tracks.rs |
Track bookkeeping: which participant’s media tracks are forwarded where |
services/sfu/events.rs |
Outbound signal/event types emitted by the SFU |
services/sfu/forward.rs |
Media forwarding logic (selective forwarding of RTP packets) |
services/sfu/run.rs |
The Tokio task/loop that drives the SFU (str0m Rtc poll loop) |
services/sfu/tests.rs, testutil.rs, scratch_test.rs |
Unit/integration tests and helpers |
The module layout below reflects the production-hardened implementation.
session.rsowns per-call WebRTC state — SDP offer/answer handling, ICE, forward-track bookkeeping, and the re-sync recovery path described under Reliability & Production Hardening.session_tracks.rsadds and negotiates forward m-lines and writes RTP into them.forward.rsfansMediaAdded/MediaData/keyframe events out to peers.run.rsdrives thestr0mpoll loop and the stale-offer recovery timer.
The backend deliberately embeds the SFU rather than talking to an external media server: str0m’s Rtc object runs in-process, so signaling and media share the same Tokio executor, and the AppError type already provides a conversion from str0m’s error type:
impl From<str0m::error::RtcError> for AppError {
fn from(e: str0m::error::RtcError) -> Self {
Self::Internal(format!("WebRTC error: {e}"))
}
}
Source: error.rs
This means any str0m failure (ICE failure, DTLS handshake error, SDP parse error) surfaces as a typed AppError::Internal with the WebRTC error string attached, so gRPC handlers and the SFU can use ? to propagate WebRTC failures uniformly.
The call relay: a signaling bus
CallRelayService (services/call_relay.rs) is the central message hub. Its verified surface, from usage, includes send_signal(sig) which is async and takes a CallSignal. The relay is the single integration point named in the SFU comment — “outgoing signals forward to the relay” — and it is also consumed by relay subscribers:
use crate::services::signal::CallSignal;
Source: relay_subscriber.rs
RelaySubscriber (in services/relay_subscriber.rs, using tracing::{error, info}) represents a long-lived consumer of the relay — for example, the gRPC streaming response stream for one client, or an HTTP long-poll handle. The relay fan-out is what makes it possible for one participant’s SFU-generated answer to reach every other participant’s gRPC stream. subscribe(user_id) additionally replays any signals buffered while the user had no active stream — see Relay-side buffering of undelivered media offers under Reliability & Production Hardening.
The gRPC signaling bridge
The bridge between the internal signaling domain and the wire protocol lives in grpc/signals.rs. It converts a CallSignal (internal) into a SignalEvent (protobuf, defined in chat_service::chat_proto):
use std::time::{SystemTime, UNIX_EPOCH};
use crate::services::signal::CallSignal;
use super::chat_service::chat_proto::SignalEvent;
fn now() -> i64 {
SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_secs() as i64
}
pub fn call_signal_to_event(sig: CallSignal) -> SignalEvent {
SignalEvent {
event_type: sig.signal_type,
from_id: sig.caller_id,
group_id: sig.group_id.unwrap_or(0),
payload: sig.payload,
created_at: now(),
}
}
Source: signals.rs
Observations:
- Lossless field mapping:
signal_type→event_type,caller_id→from_id,payload→payload. The only transformation isgroup_id, whereNone(a direct call, not a group call) becomes0on the wire. - Server-side timestamp:
created_atis stamped at conversion time using wall-clock seconds (SystemTime::now()), not taken from the client, so all clients in a call agree on event ordering basis. - Pure function: no I/O, no errors — the function is total and trivially testable.
Core Flow
Outbound path: SFU → relay → gRPC clients
When the SFU needs to tell participants something (an answer SDP, an ICE candidate, a session state change), it emits a CallSignal through sfu_tx. The forwarder task drains sfu_rx and calls relay.send_signal(), the relay fans out to subscribers, and each subscriber’s CallSignal is converted to a SignalEvent and written to the client’s gRPC stream.
Inbound path: client → gRPC → SFU
Inbound signaling (a client’s offer or ICE candidate) enters through the gRPC chat service, is converted from the wire form into a CallSignal and published to the relay, where the SFU session consumer applies it to the corresponding str0m Rtc object. (The exact inbound conversion function lives in grpc/chat_service.rs, which was not read in full for this page; the relay’s dual direction is implied by the verified outbound loop and by CallSignal being the shared message type on both sides of the bridge.)
Lifecycle of a call
The SFU session (services/sfu/session.rs) drives the str0m state machine: it starts in negotiation when an offer arrives, transitions to Connected once DTLS/ICE complete, may renegotiate when tracks are added or removed, and terminates on hangup or a WebRTC error (which surfaces through the RtcError → AppError conversion above).
Data Model
CallSignal (internal domain type)
Defined in services/signal.rs. Fields verified from field usage in grpc/signals.rs:
| Field | Type | Notes |
|---|---|---|
signal_type |
string-ish (mapped to event_type) |
Kind of signal: offer, answer, ICE candidate, hangup, etc. |
caller_id |
numeric/id (mapped to from_id) |
Sender of the signal |
group_id |
Option<...> |
None for direct calls; mapped to 0 on the wire |
payload |
string/bytes | Opaque signal body (e.g. SDP or ICE candidate JSON) |
SignalEvent (wire proto)
Defined in chat_service::chat_proto. Populated by call_signal_to_event():
| Field | Type | Source |
|---|---|---|
event_type |
string | sig.signal_type |
from_id |
id | sig.caller_id |
group_id |
i64 | sig.group_id.unwrap_or(0) |
payload |
string/bytes | sig.payload |
created_at |
i64 (unix seconds) | now() at conversion time |
The Option handling on group_id is the only semantic transformation in the bridge: the internal model distinguishes “no group” (None) from “group 0”, while the wire protocol uses 0 as the sentinel for direct calls.
API Reference
SfuService::new(signal_tx) -> impl Future<Output = Result<SfuService, AppError>>
Creates the embedded str0m-based SFU. The caller passes the sending half of the unbounded Tokio channel through which the SFU emits outgoing CallSignal messages.
Parameters:
signal_tx—tokio::sync::mpsc::UnboundedSender<CallSignal>(constructed assfu_txinmain.rs)
Returns: a ready-to-use SfuService (stored as Arc<SfuService> in AppState).
Notes: asynchronous (awaited) — startup performs any I/O or allocation required before serving. Verified call site: main.rs.
CallRelayService::send_signal(sig: CallSignal) -> impl Future<Output = Result<...>>
Publishes a CallSignal into the relay for fan-out to all subscribers.
Parameters:
sig(CallSignal) — the signaling message to distribute.
Notes: called both by the SFU forwarder task (main.rs) and by inbound paths. Errors are ignored at the forwarder call site (let _ = ...), indicating the relay tolerates failed subscriber deliveries. Verified call site: main.rs.
call_signal_to_event(sig: CallSignal) -> SignalEvent
Pure conversion from the internal CallSignal domain type to the protobuf SignalEvent used on the gRPC wire. Total function — no error cases.
Parameters:
sig(CallSignal) — internal signal to convert.
Returns: SignalEvent with event_type = sig.signal_type, from_id = sig.caller_id, group_id = sig.group_id.unwrap_or(0), payload = sig.payload, created_at = now() (unix seconds). Verified source: signals.rs.
impl From<str0m::error::RtcError> for AppError
Maps any str0m WebRTC error into AppError::Internal with the message "WebRTC error: {e}". Enables ?-based propagation of media-plane failures through the application error type. Verified source: error.rs.
Configuration Options
The signaling/SFU subsystem’s runtime configuration lives in the backend-wide config.rs and build.rs (compile-time). Specific SFU knobs (e.g. candidate ports, codec preferences) were not read during this page’s source pass; the following table lists the wiring-level facts verified from source:
| Item | Type | Value / Default | Description |
|---|---|---|---|
| SFU signal channel | tokio::sync::mpsc::UnboundedSender<CallSignal> |
created per-boot in main.rs |
Carries outgoing SFU signals to the relay |
Relay handle in AppState |
Arc<CallRelayService> |
call_relay |
Shared by all gRPC handlers |
SFU handle in AppState |
Arc<SfuService> |
sfu |
Shared by all gRPC handlers |
SFU_PORT |
u16 |
ephemeral if unset | Fixed UDP port the SFU binds, so the host firewall can allow it |
SFU_PUBLIC_IP |
std::net::IpAddr |
machine primary IP if unset | Public address advertised as the SFU’s ICE candidate |
SignalEvent.created_at |
i64 unix seconds |
server wall clock | Stamped at conversion time |
For backend-wide settings (database, ports, etc.), see the config page for config.rs.
Reliability & Production Hardening
The SFU went through several rounds of production hardening to make forwarded video reliable across mixed clients (browsers and react-native). The sections below describe each hardening measure, the failure mode it addresses, and the code that implements it.
Codec pinning to OPUS + VP8
Every SFU session is constrained to a single audio codec (Opus) and a single video codec (VP8) via RtcConfig:
let mut rtc = Rtc::builder()
.clear_codecs()
.enable_opus(true)
.enable_vp8(true)
.build(Instant::now());
Source: session.rs — SfuSession::new
Failure mode addressed. Browsers typically send VP8, while react-native (Android) negotiates H.264. With the default codec set, a forward m-line negotiated H.264 while the origin sent VP8, so writer.match_params() could not match the origin’s packets and the video was silently dropped — the forward looked Open but no picture arrived. Audio was unaffected because both sides use Opus. Pinning both sessions to OPUS + VP8 guarantees the origin’s payload parameters always match the receiver’s negotiated forward, so forwarding works in both directions regardless of client type.
One forward m-line per renegotiation offer
Forward-track renegotiations add exactly one new m-line per offer. Both the queued-track flush (flush_forward_tracks) and the stuck-track re-sync (reoffer_stuck_tracks) add a single track and apply() it:
let (origin_user, origin_mid, kind) = self.pending_tracks.remove(0);
let mut change = self.rtc.sdp_api();
let mid = change.add_media(kind, str0m::media::Direction::SendOnly, Some(origin_user.to_string()), None, None);
let (offer, pending) = change.apply()?;
Source: session_tracks.rs — flush_forward_tracks
Failure mode addressed. Batching several new m-lines into a single renegotiation produced client answers with fewer m-lines than the offer (Differing m-line count in offer vs answer: 2 != 1), which str0m rejects. Issuing one m-line per offer keeps the answer 1 with the offer.
Opening only the accepted forward track
accept_answer tracks which mid the current pending offer added (pending_mid) and marks only that track Open on a successful answer:
if let Some(mid) = self.pending_mid.take() {
if let Some(t) = self.tracks_out.iter_mut().find(|t| matches!(t.state, TrackOutState::Negotiating(m) if m == mid)) {
t.state = TrackOutState::Open(mid);
}
}
Source: session.rs — accept_answer
Failure mode addressed. The re-sync path can leave several superseded Negotiating copies of a forward. Marking every Negotiating track Open fabricated tracks that never received a str0m writer — video then dropped with no writer mid=…. Opening only the track in the accepted offer guarantees every Open forward has a writer.
Stale-answer handling
An answer whose m-lines are not present in the current pending offer (i.e. it was generated against a superseded offer) is dropped without consuming the pending:
if let Some(offer) = &self.pending_offer_sdp {
let answer_mids = sdp_mids(&answer.to_sdp_string());
let offer_mids = sdp_mids(offer);
if answer_mids.iter().any(|m| !offer_mids.contains(m))
|| offer_mids.iter().any(|m| !answer_mids.contains(m))
{
tracing::warn!("SFU ignoring stale answer user={} room={} (superseded offer)", self.user_id, self.room_id);
return Ok(None);
}
}
Source: session.rs — accept_answer
Failure mode addressed. A client that answered a superseded forward offer late must not invalidate the current negotiation. The check is bidirectional — it drops answers referencing a mid the current offer lacks or offers missing a mid the answer reflects (e.g. an answer to an offer that was re-synced into a merged re-offer) — so a stale answer is never fed to str0m as an m-line count mismatch that re-triggers the re-sync loop. By ignoring the stale answer and leaving the pending intact, the client’s answer to the current offer is still accepted when it arrives.
ACK-gated re-push of stale forward offers
Clients ACK a forward offer the moment they receive it (media_offer_ack), so the run loop can distinguish a lost offer (no ack) from a lost answer (acked but unanswered). The SFU re-sends the same SDP — same m-lines, same change id — after an ACK-gated stale threshold:
if guard.rtc.is_connected() {
let threshold = if guard.pending_acked() {
Duration::from_secs(6) // acked but unanswered: the answer was lost
} else {
Duration::from_secs(3) // never acked: the offer itself was lost
};
let stale_pending = guard.take_stale_pending_offer(threshold, 5);
let orphaned = !stale_pending && guard.has_orphaned_stuck_track(threshold);
if stale_pending || orphaned {
if let Some(sdp) = guard.reoffer_stuck_tracks(stale_pending) {
reoffers.push((guard.user_id, guard.room_id, sdp));
}
}
}
Source: run.rs — the SFU poll loop
repush_pending_offer re-issues the identical SDP rather than allocating a fresh mid:
fn repush_pending_offer(&mut self) -> Option<String> {
let sdp = self.pending_offer_sdp.clone()?;
tracing::info!("SFU re-pushing stale pending offer user={} mids={:?}", self.user_id, sdp_mids_ordered(&sdp));
self.pending_since = Some(Instant::now());
self.reoffer_count = 0;
self.resync_count += 1;
Some(sdp)
}
Source: session.rs — repush_pending_offer
Failure mode addressed. A forward offer lost on a flaky signaling stream would leave the video forward stuck in Negotiating forever. The recovery has to satisfy three constraints at once:
- Re-pushing the identical SDP is safe. The earlier design re-issued a fresh offer (new change id) after a long stale window, which rejected the peer’s in-flight answer (
str0mChangesOutOfOrder), and the subsequent fresh-mid re-offer reordered the peer’s committed m-lines — the browser rejected it with “The order of m-lines in subsequent offer doesn’t match”. Re-pushing under the existing change id keeps any in-flight answer valid: a peer that never got the offer applies it, a peer that committed it re-applies it (verified bysfu_repush_same_sdp_after_peer_committed). - The threshold is short because it is ACK-gated. A no-ack pending means the offer almost certainly never reached the peer, so 3s is enough; an acked pending means the peer has the offer and its answer was lost, so 6s gives a slow-to-answer client room. A dropped offer is recovered in a few seconds instead of leaving the tile black for the whole stale window.
- Recovery is bounded.
take_stale_pending_offergives up aftermax_reoffers(5) so a session whose peer never answers (e.g. a stale session from an ended call) stops re-offering instead of looping;resync_countcaps the total, and theresyncedset ensures each stuck origin is re-offered at most once per round.
has_orphaned_stuck_track covers the case where a track is left Negotiating with no pending offer at all — a re-sync already flushed the queued track and the answer consumed the only pending (see Resetting the ACK gate when an answer consumes the pending below). Such an offer is provably lost, so the orphan is re-issued as a fresh-mid offer (the only fresh-mid path left).
ICE restart in fresh re-sync offers
Only the fresh-mid recovery path — re-offering an orphaned Negotiating track when nothing is in flight — carries an ICE restart (keeping the existing local candidate):
let mut change = self.rtc.sdp_api();
change.ice_restart(true);
Source: session.rs — reoffer_stuck_tracks
Failure mode addressed. Browsers restart ICE when answering a re-issued offer; if the offer does not itself restart ICE, str0m rejects the answer (Ice restart in answer without one in the preceeding offer) and the web client’s video forward never opens — while react-native (which does not restart ICE this way) works. This was the specific web-side failure that left video working on mobile but not in the browser. The ACK-gated re-push path does not restart ICE — it re-sends the exact bytes already in flight, so connectivity is untouched there.
Preserving peer-committed m-lines across re-syncs
Every recovery path keeps the m-lines the peer has already committed in their existing positions. reoffer_stuck_tracks either re-pushes the current pending (same mid), or — when nothing is in flight — allocates a fresh mid for an orphaned stuck track, never a mix:
// stale UNACKED pending with queued tracks: re-issue the pending ALONE
if self.pending.is_some() && !self.pending_acked && !self.pending_tracks.is_empty() {
return self.repush_pending_offer();
}
// stale ACKED pending (or none): flush ONE queued track, carrying the pending's m-lines forward
if !self.pending_tracks.is_empty() {
return self.flush_forward_tracks();
}
// stale pending: re-push under its EXISTING SDP
if self.pending.is_some() {
return self.repush_pending_offer();
}
Source: session.rs — reoffer_stuck_tracks
Failure mode addressed. Dropping or reordering a peer-committed m-line makes the browser reject every subsequent offer (“The order of m-lines in subsequent offer doesn’t match order from previous offer/answer”). Because acks are lossy, an unacked pending does not prove the peer never applied it — prod showed a peer committing audio MrA, losing the ack, and the flush then dropping MrA while re-offering the queued video alone, which the peer rejected forever. flush_forward_tracks uses SdpApi::merge(pending) to carry the in-flight m-lines forward, and the ambiguous unacked case re-pushes the pending alone before flushing anything behind it.
One forward offer in flight at a time
reoffer_stuck_tracks refuses to issue a new offer while a fresh (non-stale) pending is still being negotiated:
if self.pending.is_some() && !pending_is_stale {
return None;
}
Source: session.rs — reoffer_stuck_tracks
Failure mode addressed. Stacking a second offer on top of an in-flight one either fails to apply on the peer (it is have-remote-offer) or — if the second offer omits the in-flight m-line the peer already committed — corrupts the peer’s negotiated m-line order. Only once the pending goes stale (or resolves) is a new offer issued, and then at most one m-line per offer (see One forward m-line per renegotiation offer).
Resetting the ACK gate when an answer consumes the pending
When an answer is accepted, accept_answer clears pending_acked alongside the consumed pending:
self.pending_offer_sdp = None;
self.pending_since = None;
self.resync_count = 0;
self.resynced.clear();
self.pending_acked = false;
Source: session.rs — accept_answer
Failure mode addressed. If the ack flag survived the answer, the run loop would keep gating orphaned-track re-offers behind the slow acked threshold (6s) when it should use the fast no-ack one (3s) for tracks whose offer was lost, delaying a lost audio/video forward by several extra seconds.
Forward-track de-duplication
add_forward_track refuses to forward a track already forwarded for the same origin:
if self.has_forward(origin_user, origin_mid) {
return None;
}
Source: session_tracks.rs — add_forward_track
Failure mode addressed. A duplicate MediaAdded event (or the fan-out and replay paths both firing for the same track) would add a second m-section for the same origin and derail the negotiation sequence.
Stable SFU UDP port
The SFU binds a fixed UDP port when SFU_PORT is set, and advertises a public candidate via SFU_PUBLIC_IP:
let sfu_port = std::env::var("SFU_PORT").ok().and_then(|s| s.parse::<u16>().ok()).unwrap_or(0);
let socket = Arc::new(UdpSocket::bind(&format!("0.0.0.0:{sfu_port}")).await?);
let ip = std::env::var("SFU_PUBLIC_IP").ok().and_then(|s| s.parse::<std::net::IpAddr>().ok()).unwrap_or_else(local_ip);
Source: mod.rs — SfuService::new
Failure mode addressed. Without a fixed port the SFU binds an ephemeral UDP port that changes on every restart and is usually blocked by the host firewall, so remote clients could never reach the SFU and ICE died at Checking/Disconnected. With SFU_PORT fixed (e.g. 58695), the firewall can open that port once; SFU_PUBLIC_IP ensures the advertised candidate is the machine’s public address rather than the container’s.
Keyframe priming
After a receiver answers a renegotiation offer, prime_new_video_tracks immediately requests a keyframe from the origin so the receiver renders a picture within ~1 RTT instead of waiting for the encoder’s keyframe interval:
Source: forward.rs — prime_new_video_tracks
Relay-side buffering of undelivered media offers
The signal relay never drops an SFU media_offer. When a user has no active stream — a reconnect window, or before their StreamMessages stream is up — the relay buffers the signal in a per-user backlog and flushes it in order the next time the user subscribes:
pub fn subscribe(&self, user_id: i64) -> broadcast::Receiver<CallSignal> {
let tx = self.broadcast_channels.entry(user_id).or_insert_with(|| {
let (t, _) = broadcast::channel(BROADCAST_CAPACITY);
t
}).value().clone();
let rx = tx.subscribe();
// Replay signals buffered while this user had NO active stream.
if let Some(q) = self.pending.get(&user_id) {
let mut q = q.lock().unwrap();
for sig in q.drain(..) { let _ = tx.send(sig); }
}
rx
}
Source: call_relay.rs — subscribe
The backlog is a DashMap<i64, Mutex<VecDeque<CallSignal>>> capped at PENDING_CAP (256) so a never-reconnecting user cannot grow an unbounded queue; the SFU’s re-push is the backstop beyond the cap. Additionally, stream_messages drains the broadcast receiver into an unbounded queue so a busy client stalls the write rather than triggering broadcast::RecvError::Lagged drops (see the gRPC-Web API & Chat Protocol page).
Failure mode addressed. A forward media_offer pushed while the callee’s stream was reconnecting used to be silently lost; only the SFU’s seconds-late re-push recovered it, delaying the callee’s first video. Buffering + replay-on-subscribe removes that delivery gap, and the client-side retry (see the Signaling, Transport & Codec page) covers the reverse direction.
Regression coverage
The hardening measures are covered by the SFU integration suite in services/sfu/tests.rs (21 tests), including:
sfu_converges_under_signal_loss— a deterministic lossy-signaling gate that replays a full two-way call under 15–30% per-message loss at four fixed seeds and asserts both directions’ audio+video open and media flows each time.sfu_repush_accepts_in_flight_answer/sfu_repush_same_sdp_after_peer_committed— pin the same-SDP re-push semantics (an in-flight answer survives a re-push; a committed peer re-applies the identical offer).sfu_resync_preserves_peer_committed_mlines/sfu_does_not_overlap_offers_on_fresh_pending/sfu_reissues_lost_pending_alone_then_flushes_queued— pin the m-line-order and single-pending invariants.- The existing multi-track late-join replay, forwarded audio+video, rejected-answer recovery, and one-m-line-per-offer flush tests.
These tests run in CI on every push to main.
Failure Modes, Edge Cases & Concurrency
WebRTC failures
Any str0m-level failure (ICE connectivity failure, DTLS handshake timeout, malformed SDP, codec mismatch) is converted into AppError::Internal("WebRTC error: ..."). Because the SFU runs in-process, these errors are visible to the surrounding service code and can be logged with tracing or propagated into gRPC status responses.
Relay send failures
The forwarder task deliberately ignores send_signal errors (let _ = relay_for_sfu.send_signal(sig).await). Rationale: the SFU’s media forwarding must never stall because a subscriber is gone or the relay is saturated. A failed delivery drops that one signal; the call continues. This is a fire-and-forget design choice — signaling is treated as best-effort at the relay boundary.
Channel semantics
- Unbounded channel:
sfu_rx.recv().awaitreturnsNoneonly when all senders are dropped and the channel is empty; the forwarder loop then exits. Sincesfu_txis owned bySfuService, the task lives exactly as long as the SFU — no leaked tasks. - Concurrency: the SFU emits signals from the (single) Tokio runtime; the relay fan-out serializes deliveries per subscriber. The
Arc<CallRelayService>sharing means multiple gRPC connections concurrently publish into the same relay; ordering guarantees, if any, come from the relay’s internal design (not read in full).
Edge cases
- Direct calls:
group_id = Noneis encoded as wiregroup_id = 0, so direct (1) and group calls share one event shape. unwrap_or(0): safe becausenow()cannot fail and the conversion has no failure branch — the bridge is total.
Performance & Operational Considerations
- In-process SFU: media forwarding runs on the same Tokio executor as the API, avoiding IPC/network hops between the control plane and the media plane. This lowers latency but means SFU CPU is shared with API work — capacity planning should account for concurrent media sessions per instance.
- Unbounded channel: chosen for low-latency signal delivery; under sustained relay back-pressure, memory could grow. The fire-and-forget forwarder bounds the practical impact.
- Timestamps:
created_atusesSystemTime(nottokio::time::Instant), so event times are wall-clock-consistent across clients and restartable. - Module-level test harness: the presence of
sfu/tests.rs,sfu/testutil.rs, andsfu/scratch_test.rsindicates the SFU is exercised with dedicated test utilities, supporting safe iteration on session and forwarding logic.
Extension Points
- New signal types: extend
CallSignal.signal_typeand the payload contract; the relay and bridge are type-agnostic, so no changes are needed outsideservices/signal.rsand the protoSignalEvent. - New subscriber transports: implement a
RelaySubscriber-style consumer (services/relay_subscriber.rs) to add transports (e.g. HTTP long-poll, WebSocket) without touching the SFU or the relay. - Codec/media policy: lives inside
services/sfu/forward.rs(selective forwarding) andsession_tracks.rs(per-participant track maps); changing forwarding policy is isolated to those modules. - Error handling: extend
error.rswith more specificAppErrorvariants perstr0merror class if finer-grained gRPC status codes are needed.
Related Links
- main.rs — SFU/relay composition root
- grpc/signals.rs — signaling bridge
- services/call_relay.rs — CallRelayService
- services/relay_subscriber.rs — relay subscriber
- services/sfu/ — SFU module (mod, session, session_tracks, events, forward, run)
- error.rs — WebRTC error mapping
- state.rs — AppState registration
- services/signal.rs — CallSignal type
- grpc/chat_service.rs — chat gRPC service and
SignalEventproto
For call ticket issuance, authentication middleware, and backend-wide configuration, see the respective sibling catalog pages.