Compare commits

..

4 Commits

Author SHA1 Message Date
lebaudantoine d465291cf5 🧑‍💻(devx) add a WebRTC stats and network throttling devtool
Introduce an in-app devtool that monitors WebRTC statistics in
real time and lets developers simulate various network scenarios,
including constraining the uplink and downlink bandwidth.

Makes it much easier to reproduce and investigate connectivity or
quality issues locally without depending on external tools.

The code was AI generated, and might contain some smell.
It's only enabled in dev, and not included in the production
build. Feel free to enhance it as needed.
2026-09-01 19:40:48 +02:00
lebaudantoine ea771bff0c 🔧(devx) configure a TURN server on the local LiveKit dev stack
Wire a TURN server into the local LiveKit server used by the dev
stack, so ICE negotiation has more candidate types available during
local testing.

Makes it easier to reproduce connectivity scenarios that would
otherwise only show up on stricter networks in production.
2026-09-01 19:40:43 +02:00
leo 0a0cdae896 (agent) support Voxtral realtime as inference engine
The Kyutai open-source model turned out not to be production-ready:
it caused disruptions in the production environment, especially on
long-running meeting sessions.

Switch to the Voxtral realtime model, which looks like a
credible competitor and behaves much better in our setup.

For now, the code handling the Voxtral realtime API lives directly
in the project. It could be extracted into an open-source package
later.

See PR #1277 for the full details of the implementation, proposed
by @cameldev.
2026-08-28 12:18:52 +02:00
leo c7e3168ba3 🔥(backend) remove the S3 storage-event webhook for recordings
Recordings used to be finalized by an inbound notifications sent by
S3. This tied the recording lifecycle to bucket notifications, adding
complexity and dependency to limited S3 services.

The LiveKit egress_ended webhook, added as a fallback in aee1847, does
the same job without any object-storage dependency. Make it the only
mechanism: RecordingEventsService.handle_complete is now called on
EGRESS_COMPLETE / EGRESS_LIMIT_REACHED unconditionally, instead of only
when RECORDING_STORAGE_EVENT_ENABLE is False. Remove the storage-event
path entirely.
2026-08-27 23:09:48 +02:00
15 changed files with 2477 additions and 539 deletions
+4
View File
@@ -8,6 +8,10 @@ and this project adheres to
## [Unreleased]
### Added
- ✨(agent) support Voxtral realtime as inference engine
### Changed
- 🔥(backend) remove the S3 storage-event webhook for recordings
+3
View File
@@ -213,6 +213,8 @@ services:
- "7880:7880"
- "7881:7881"
- "7882:7882/udp"
- "3478:3478/udp"
- "30000-30100:30000-30100/udp"
volumes:
- ./docker/livekit/config/livekit-server.yaml:/config.yaml
depends_on:
@@ -251,6 +253,7 @@ services:
build:
context: ./src/agents
target: development
command: ["python", "multi_user_transcriber.py", "dev"]
env_file:
- env.d/development/multi_user_transcriber
volumes:
+13
View File
@@ -8,3 +8,16 @@ webhook:
api_key: devkey
urls:
- http://app-dev:8000/api/v1.0/rooms/webhooks-livekit/
turn:
enabled: true
domain: turn.127.0.0.1.nip.io
udp_port: 3478
tls_port: 0
external_tls: false
relay_range_start: 30000
relay_range_end: 30100
allow_restricted_peer_cidrs:
- 192.168.0.0/16
- 172.16.0.0/12
+13 -4
View File
@@ -1,14 +1,23 @@
AWS_S3_ENDPOINT_URL=minio:9000
AWS_S3_ACCESS_KEY_ID=meet
AWS_S3_SECRET_ACCESS_KEY=password
LIVEKIT_URL=ws://livekit:7880
LIVEKIT_API_KEY=devkey
LIVEKIT_API_SECRET=secret
STT_PROVIDER=kyutai # kyutai, deepgram
STT_PROVIDER=voxtral-vllm # voxtral-vllm, kyutai, deepgram
ENABLE_SILERO_VAD=False
DEEPGRAM_API_KEY=
DEEPGRAM_API_KEY=your-deepgram-api-key
KYUTAI_STT_BASE_URL=
KYUTAI_API_KEY=
KYUTAI_STT_BASE_URL=url
KYUTAI_API_KEY=your-kyutai-api-key
VOXTRAL_VLLM_BASE_URL=wss://<host>/v1/realtime
VOXTRAL_VLLM_MODEL=voxtral-mini-4b-realtime-2602
VOXTRAL_VLLM_API_KEY=your-vllm-api-key
VOXTRAL_VLLM_TARGET_STREAMING_DELAY_MS=480
SENTRY_DSN=
SENTRY_ENVIRONMENT=
+53 -16
View File
@@ -1,6 +1,7 @@
"""Multi user transcription agent."""
import asyncio
import contextlib
import logging
import os
@@ -25,6 +26,7 @@ from livekit.agents import (
)
from livekit.plugins import deepgram, silero
import voxtral_vllm_stt
from observability import configure_sentry, set_job_context
from tasks import done_callback
@@ -36,9 +38,18 @@ TRANSCRIBER_AGENT_NAME = os.getenv("TRANSCRIBER_AGENT_NAME", "multi-user-transcr
STT_PROVIDER = os.getenv("STT_PROVIDER", "deepgram")
ENABLE_SILERO_VAD = os.getenv("ENABLE_SILERO_VAD", "true").lower() == "true"
SESSION_DRAIN_TIMEOUT_S = 15.0
def create_stt_provider():
"""Create STT provider based on environment configuration."""
def create_stt_provider(vad: silero.VAD | None = None):
"""Create STT provider based on environment configuration.
Args:
vad: Shared, prewarmed VAD instance. Required in practice for
voxtral-vllm (no server-side endpointing): if omitted, the plugin
loads its own Silero model synchronously on the event loop, once
per participant, freezing all active sessions for the duration.
"""
if STT_PROVIDER == "deepgram":
# Note: Not all Deepgram API parameters are supported by the LiveKit plugin
# detect_language is NOT supported for real-time streaming
@@ -49,6 +60,9 @@ def create_stt_provider():
)
elif STT_PROVIDER == "kyutai":
_stt_instance = kyutai.STT(base_url=os.getenv("KYUTAI_STT_BASE_URL"))
elif STT_PROVIDER == "voxtral-vllm":
# The plugin resolves base_url / model / api_key from the environment.
_stt_instance = voxtral_vllm_stt.STT(vad=vad)
else:
raise ValueError(f"Unknown STT_PROVIDER: {STT_PROVIDER}")
@@ -58,9 +72,9 @@ def create_stt_provider():
class Transcriber(Agent):
"""Create a transcription agent for a specific participant."""
def __init__(self, *, participant_identity: str):
def __init__(self, *, participant_identity: str, vad: silero.VAD | None = None):
"""Init transcription agent."""
stt = create_stt_provider()
stt = create_stt_provider(vad=vad)
super().__init__(
instructions="not-needed",
@@ -76,6 +90,7 @@ class MultiUserTranscriber:
"""Init multi user transcription agent."""
self.ctx = ctx
self._sessions: dict[str, AgentSession] = {}
self._starting: dict[str, asyncio.Task] = {}
self._tasks: set[asyncio.Task] = set()
def start(self):
@@ -96,22 +111,30 @@ class MultiUserTranscriber:
def on_participant_connected(self, participant: rtc.RemoteParticipant):
"""Handle new participant connection by starting transcription session."""
if participant.identity in self._sessions:
identity = participant.identity
if identity in self._sessions or identity in self._starting:
return
logger.info(f"starting session for {participant.identity}")
logger.info(f"starting session for {identity}")
task = asyncio.create_task(self._start_session(participant))
self._starting[identity] = task
self._tasks.add(task)
task.add_done_callback(lambda t, i=identity: self._starting.pop(i, None))
task.add_done_callback(
done_callback(
logger,
self._tasks,
f"start transcription session for {participant.identity}",
f"start transcription session for {identity}",
)
)
def on_participant_disconnected(self, participant: rtc.RemoteParticipant):
"""Handle participant disconnection by closing transcription session."""
if (start_task := self._starting.pop(participant.identity, None)) is not None:
logger.info(f"cancelling pending session start for {participant.identity}")
start_task.cancel()
return
if (session := self._sessions.pop(participant.identity, None)) is None:
return
@@ -127,10 +150,12 @@ class MultiUserTranscriber:
)
async def _start_session(self, participant: rtc.RemoteParticipant) -> AgentSession:
"""Create and start transcription session for participant."""
if participant.identity in self._sessions:
return self._sessions[participant.identity]
"""Create and start transcription session for participant.
Deduplication happens synchronously in on_participant_connected via
self._starting; by the time this coroutine runs, the identity is
already reserved.
"""
vad = self.ctx.proc.userdata.get("vad", None)
session = AgentSession(vad=vad)
room_io = RoomIO(
@@ -141,18 +166,30 @@ class MultiUserTranscriber:
text_input=False, audio_output=False, text_output=True
),
)
await room_io.start()
await session.start(
agent=Transcriber(
participant_identity=participant.identity,
try:
await room_io.start()
await session.start(
agent=Transcriber(
participant_identity=participant.identity,
vad=vad,
)
)
)
except BaseException:
with contextlib.suppress(Exception):
await session.aclose()
raise
self._sessions[participant.identity] = session
return session
async def _close_session(self, sess: AgentSession) -> None:
"""Close and cleanup transcription session."""
await sess.drain()
try:
await asyncio.wait_for(sess.drain(), timeout=SESSION_DRAIN_TIMEOUT_S)
except (TimeoutError, asyncio.TimeoutError):
logger.warning(
"session drain timed out after %.0fs; forcing close",
SESSION_DRAIN_TIMEOUT_S,
)
await sess.aclose()
+2
View File
@@ -12,6 +12,8 @@ dependencies = [
"protobuf==6.33.6",
"minio==7.2.20",
"sentry-sdk==2.66.1",
"websockets==17.1",
"httpx==0.28.1",
]
[project.optional-dependencies]
+770 -510
View File
File diff suppressed because it is too large Load Diff
+476
View File
@@ -0,0 +1,476 @@
"""LiveKit STT plugin for Voxtral Realtime served via vLLM (/v1/realtime).
vLLM exposes Voxtral Realtime over a WebSocket that follows the OpenAI Realtime
API protocol (not Mistral's proprietary realtime protocol).
"""
from __future__ import annotations
import asyncio
import base64
import json
import logging
import os
import weakref
from collections import deque
from dataclasses import dataclass, field
import websockets
from livekit.agents import (
DEFAULT_API_CONNECT_OPTIONS,
APIConnectionError,
APIConnectOptions,
APIStatusError,
stt,
utils,
)
from livekit.agents import (
vad as vad_module,
)
from livekit.agents.types import NOT_GIVEN, NotGivenOr
from livekit.agents.utils import is_given
logger = logging.getLogger("voxtral-vllm-stt")
SAMPLE_RATE = 16000
NUM_CHANNELS = 1
CHUNK_SAMPLES = 1600 # 100 ms @ 16 kHz mono
PREROLL_CHUNKS = 5 # keep 500 ms of audio before start of speech as detected by VAD
# Reconnect policy: exponential backoff capped at MAX, give up after MAX_ATTEMPTS
# consecutive failures (a successful handshake resets the counter).
RECONNECT_BACKOFF_BASE_S = 0.5
RECONNECT_BACKOFF_MAX_S = 8.0
RECONNECT_MAX_ATTEMPTS = 5
@dataclass
class _STTOptions:
base_url: str
model: str
api_key: str | None
target_streaming_delay_ms: int | None
@dataclass
class _PendingUtterance:
"""An utterance in flight on the shared websocket used for reconnect.
`sent_chunks` holds every chunk we have already enqueued for send on this
or a prior connection; on reconnect we replay them before resuming reads
from `queue`. vLLM concatenates `input_audio_buffer.append` events into a
single audio buffer per generation, so duplicates from a partial prior send
are harmless.
"""
queue: asyncio.Queue[bytes | None]
sent_chunks: list[bytes] = field(default_factory=list)
ended: bool = False
class STT(stt.STT):
"""LiveKit STT speaking the OpenAI Realtime protocol served by vLLM."""
def __init__(
self,
*,
base_url: NotGivenOr[str] = NOT_GIVEN,
model: NotGivenOr[str] = NOT_GIVEN,
api_key: NotGivenOr[str] = NOT_GIVEN,
target_streaming_delay_ms: NotGivenOr[int] = NOT_GIVEN,
vad: vad_module.VAD | None = None,
) -> None:
"""Build the STT.
Args:
base_url: WebSocket URL of the vLLM realtime endpoint, e.g.
ws://example:8000/v1/realtime. Falls back to $VOXTRAL_VLLM_BASE_URL.
model: Model name exposed by vLLM, default
mistralai/Voxtral-Mini-4B-Realtime-2602.
api_key: Optional bearer token. Falls back to $VOXTRAL_VLLM_API_KEY.
target_streaming_delay_ms: Target streaming delay in ms forwarded to
vLLM via session.update. Falls back to
$VOXTRAL_VLLM_TARGET_STREAMING_DELAY_MS, else server default.
vad: Voice Activity Detector. If omitted, Silero VAD is loaded.
"""
super().__init__(
capabilities=stt.STTCapabilities(streaming=True, interim_results=True)
)
resolved_url = (
base_url
if is_given(base_url)
else os.environ.get(
"VOXTRAL_VLLM_BASE_URL", "ws://127.0.0.1:8000/v1/realtime"
)
)
resolved_model = (
model
if is_given(model)
else os.environ.get(
"VOXTRAL_VLLM_MODEL", "mistralai/Voxtral-Mini-4B-Realtime-2602"
)
)
resolved_key = (
api_key if is_given(api_key) else os.environ.get("VOXTRAL_VLLM_API_KEY")
)
resolved_delay = (
target_streaming_delay_ms
if is_given(target_streaming_delay_ms)
else (
int(os.environ["VOXTRAL_VLLM_TARGET_STREAMING_DELAY_MS"])
if os.environ.get("VOXTRAL_VLLM_TARGET_STREAMING_DELAY_MS")
else None
)
)
if vad is None:
try:
from livekit.plugins.silero import VAD as SileroVAD # noqa: PLC0415
except ImportError as exc:
raise ImportError(
"livekit-plugins-silero is required for vLLM Voxtral realtime "
"(no server-side endpointing)."
) from exc
vad = SileroVAD.load()
self._vad = vad
self._opts = _STTOptions(
base_url=resolved_url,
model=resolved_model,
api_key=resolved_key,
target_streaming_delay_ms=resolved_delay,
)
self._streams: weakref.WeakSet[SpeechStream] = weakref.WeakSet()
@property
def model(self) -> str:
"""Return the configured vLLM model name."""
return self._opts.model
@property
def provider(self) -> str:
"""Return the provider identifier."""
return "vllm-voxtral-realtime"
async def _recognize_impl(self, *_args, **_kwargs) -> stt.SpeechEvent:
raise NotImplementedError(
"vLLM Voxtral Realtime STT only supports streaming recognition."
)
def stream(
self,
*,
conn_options: APIConnectOptions = DEFAULT_API_CONNECT_OPTIONS,
) -> SpeechStream:
"""Open a new streaming recognition stream."""
s = SpeechStream(
stt=self,
opts=self._opts,
vad_instance=self._vad,
conn_options=conn_options,
)
self._streams.add(s)
return s
class SpeechStream(stt.RecognizeStream):
"""Voxtral realtime handler."""
def __init__(
self,
*,
stt: STT,
opts: _STTOptions,
vad_instance: vad_module.VAD,
conn_options: APIConnectOptions,
) -> None:
"""Init the speech stream."""
super().__init__(stt=stt, conn_options=conn_options, sample_rate=SAMPLE_RATE)
self._opts = opts
self._vad = vad_instance
self._utterance_q: asyncio.Queue[bytes | None] | None = None
self._speaking = False
self._preroll: deque[bytes] = deque(maxlen=PREROLL_CHUNKS)
# Voxtral realtime is strictly sequential: only one generation runs at a
# time, and a new `commit` is ignored while the previous one is still
# producing. We queue per-utterance audio buffers here and let the
# pipeline process them one by one on the shared websocket.
self._utterance_chan: asyncio.Queue[asyncio.Queue[bytes | None] | None] = (
asyncio.Queue()
)
@utils.log_exceptions(logger=logger)
async def _run(self) -> None:
vad_stream = self._vad.stream()
bstream = utils.audio.AudioByteStream(
sample_rate=SAMPLE_RATE,
num_channels=NUM_CHANNELS,
samples_per_channel=CHUNK_SAMPLES,
)
async def input_task() -> None:
async for data in self._input_ch:
if isinstance(data, self._FlushSentinel):
for frame in bstream.flush():
self._handle_chunk(frame.data.tobytes())
continue
vad_stream.push_frame(data)
for frame in bstream.write(data.data.tobytes()):
self._handle_chunk(frame.data.tobytes())
vad_stream.end_input()
async def vad_task() -> None:
async for ev in vad_stream:
if ev.type == vad_module.VADEventType.START_OF_SPEECH:
self._on_start_of_speech()
elif ev.type == vad_module.VADEventType.END_OF_SPEECH:
self._on_end_of_speech()
pipeline_t = asyncio.create_task(self._utterance_pipeline())
try:
await asyncio.gather(input_task(), vad_task())
# signal end-of-stream; pipeline finishes pending utterances first
self._utterance_chan.put_nowait(None)
await pipeline_t
except (APIStatusError, APIConnectionError, asyncio.CancelledError):
raise
except Exception as exc:
logger.exception("vLLM realtime stream failed")
raise APIConnectionError() from exc
finally:
if not pipeline_t.done():
pipeline_t.cancel()
try:
await pipeline_t
except asyncio.CancelledError:
# CancelledError is the expected flow on cancel()
pass
except Exception:
logger.exception("utterance pipeline failed during finalize")
await vad_stream.aclose()
def _handle_chunk(self, chunk: bytes) -> None:
self._preroll.append(chunk)
if self._speaking and self._utterance_q is not None:
self._utterance_q.put_nowait(chunk)
def _on_start_of_speech(self) -> None:
if self._speaking:
return
self._speaking = True
q: asyncio.Queue[bytes | None] = asyncio.Queue()
for chunk in self._preroll:
q.put_nowait(chunk)
self._utterance_q = q
self._utterance_chan.put_nowait(q)
self._event_ch.send_nowait(
stt.SpeechEvent(type=stt.SpeechEventType.START_OF_SPEECH)
)
def _on_end_of_speech(self) -> None:
if not self._speaking:
return
self._speaking = False
if self._utterance_q is not None:
self._utterance_q.put_nowait(None)
self._utterance_q = None
self._event_ch.send_nowait(
stt.SpeechEvent(type=stt.SpeechEventType.END_OF_SPEECH)
)
async def _handshake(self, ws: websockets.ClientConnection) -> str:
created = json.loads(await ws.recv())
if created.get("type") != "session.created":
raise APIStatusError(
f"expected session.created, got {created}",
status_code=500,
body=created,
)
session_update: dict = {"type": "session.update", "model": self._opts.model}
if self._opts.target_streaming_delay_ms is not None:
session_update["target_streaming_delay_ms"] = (
self._opts.target_streaming_delay_ms
)
await ws.send(json.dumps(session_update))
return created.get("id", "")
def _auth_headers(self) -> dict[str, str]:
if self._opts.api_key:
return {"Authorization": f"Bearer {self._opts.api_key}"}
return {}
async def _utterance_pipeline(self) -> None:
# Owns the websocket lifecycle. On drop, reopens and resumes the
# in-flight utterance (if any) by replaying its already-sent chunks.
pending: _PendingUtterance | None = None
attempt = 0
while True:
try:
async with websockets.connect(
self._opts.base_url,
additional_headers=self._auth_headers(),
open_timeout=self._conn_options.timeout,
) as ws:
request_id = await self._handshake(ws)
attempt = 0
while True:
if pending is None:
q = await self._utterance_chan.get()
if q is None:
return
pending = _PendingUtterance(queue=q)
await self._process_utterance(ws, pending, request_id)
pending = None
except (websockets.WebSocketException, OSError, TimeoutError) as exc:
attempt += 1
if attempt > RECONNECT_MAX_ATTEMPTS:
logger.exception(
"vLLM realtime: giving up after %d reconnect attempts",
RECONNECT_MAX_ATTEMPTS,
)
raise APIConnectionError() from exc
backoff = min(
RECONNECT_BACKOFF_BASE_S * (2 ** (attempt - 1)),
RECONNECT_BACKOFF_MAX_S,
)
if pending is None:
logger.warning(
"vLLM WS connection lost between utterances "
"(attempt %d/%d): %s; retrying in %.1fs",
attempt,
RECONNECT_MAX_ATTEMPTS,
exc,
backoff,
)
else:
logger.warning(
"vLLM WS dropped mid-utterance (%d chunks buffered, "
"ended=%s, attempt %d/%d): %s; retrying in %.1fs",
len(pending.sent_chunks),
pending.ended,
attempt,
RECONNECT_MAX_ATTEMPTS,
exc,
backoff,
)
await asyncio.sleep(backoff)
async def _process_utterance(
self,
ws: websockets.ClientConnection,
pending: _PendingUtterance,
request_id: str,
) -> None:
# Start a fresh generation. Safe to send here: the previous utterance's
# transcription.done has already been received (we await it below), so
# the server-side generation_task is done and won't ignore this commit.
await ws.send(json.dumps({"type": "input_audio_buffer.commit"}))
send_t = asyncio.create_task(self._send_audio(ws, pending))
try:
await self._receive_one_transcription(ws, request_id)
finally:
if not send_t.done():
send_t.cancel()
try:
await send_t
except (asyncio.CancelledError, websockets.WebSocketException):
pass
except Exception:
logger.exception("send-audio task failed during finalize")
@staticmethod
async def _send_audio(
ws: websockets.ClientConnection, pending: _PendingUtterance
) -> None:
# Replay anything already sent on a previous (now-dead) connection.
# sent_chunks is appended before send, so a chunk that failed to send
# last time is still present and gets retried here.
for chunk in pending.sent_chunks:
await ws.send(
json.dumps(
{
"type": "input_audio_buffer.append",
"audio": base64.b64encode(chunk).decode("ascii"),
}
)
)
if pending.ended:
await ws.send(
json.dumps({"type": "input_audio_buffer.commit", "final": True})
)
return
while True:
chunk = await pending.queue.get()
if chunk is None:
pending.ended = True
await ws.send(
json.dumps({"type": "input_audio_buffer.commit", "final": True})
)
return
pending.sent_chunks.append(chunk)
await ws.send(
json.dumps(
{
"type": "input_audio_buffer.append",
"audio": base64.b64encode(chunk).decode("ascii"),
}
)
)
async def _receive_one_transcription(
self, ws: websockets.ClientConnection, request_id: str
) -> None:
# Use recv() rather than `async for`: the latter swallows
# ConnectionClosed on close-mid-iteration, which would let a dropped
# WS look like a clean "no transcription" return.
current_text = ""
while True:
raw = await ws.recv()
data = json.loads(raw)
event_type = data.get("type")
if event_type == "transcription.delta":
delta = data.get("delta", "")
if not delta:
continue
current_text += delta
self._event_ch.send_nowait(
stt.SpeechEvent(
type=stt.SpeechEventType.INTERIM_TRANSCRIPT,
request_id=request_id,
alternatives=[stt.SpeechData(text=current_text, language="")],
)
)
elif event_type == "transcription.done":
final_text = data.get("text") or current_text
self._event_ch.send_nowait(
stt.SpeechEvent(
type=stt.SpeechEventType.FINAL_TRANSCRIPT,
request_id=request_id,
alternatives=[stt.SpeechData(text=final_text, language="")],
)
)
usage = data.get("usage") or {}
self._event_ch.send_nowait(
stt.SpeechEvent(
type=stt.SpeechEventType.RECOGNITION_USAGE,
request_id=request_id,
recognition_usage=stt.RecognitionUsage(
audio_duration=float(
usage.get("audio_seconds")
or usage.get("prompt_audio_seconds")
or 0
),
input_tokens=int(usage.get("prompt_tokens") or 0),
output_tokens=int(usage.get("completion_tokens") or 0),
),
)
)
return
elif event_type == "error":
err = data.get("error")
raise APIStatusError(str(err), status_code=500, body=data)
@@ -0,0 +1,595 @@
import { CSSProperties, useState } from 'react'
import { useConnectionState, useRoomContext } from '@livekit/components-react'
import { ConnectionState } from 'livekit-client'
import { StatsSnapshot, TrackRow, useWebRTCStats } from './useWebRTCStats'
import { readRoomConfig } from './roomConfig'
import {
forceTransport,
releaseForcedTransport,
Scenario,
SCENARIOS,
setDownlinkCap,
setUplinkCap,
stepBitrate,
TransportMode,
transportModeFromRoute,
} from './simulation'
const SANS =
"-apple-system, BlinkMacSystemFont, 'Segoe UI', Roboto, Helvetica, Arial, sans-serif"
const MONO = 'ui-monospace, SFMono-Regular, Menlo, Consolas, monospace'
const COLOR = {
text: '#e8e8e8',
muted: '#9a9a9a',
faint: '#6f6f6f',
hairline: '#2c2d31',
border: '#3d3e44',
surface: '#1b1c1e',
chartBg: '#141517',
down: '#f6821f',
up: '#5a9cf8',
ok: '#10b981',
bad: '#f05a4a',
busy: '#e5a13c',
}
const stateColor = (state: ConnectionState): string => {
switch (state) {
case ConnectionState.Connected:
return COLOR.ok
case ConnectionState.Reconnecting:
case ConnectionState.SignalReconnecting:
return COLOR.busy
case ConnectionState.Connecting:
return COLOR.up
case ConnectionState.Disconnected:
return COLOR.bad
}
}
const styles: Record<string, CSSProperties> = {
toggle: {
position: 'fixed',
bottom: 12,
right: 12,
zIndex: 9999,
fontFamily: SANS,
fontSize: 12,
lineHeight: 1,
padding: '8px 12px',
borderRadius: 999,
borderWidth: 1,
borderStyle: 'solid',
borderColor: COLOR.border,
background: COLOR.surface,
color: COLOR.text,
cursor: 'pointer',
boxShadow: '0 2px 10px #0006',
},
panel: {
position: 'fixed',
bottom: 12,
right: 12,
zIndex: 9999,
width: 460,
maxWidth: 'calc(100vw - 24px)',
maxHeight: 'calc(100vh - 24px)',
overflowY: 'auto',
fontFamily: SANS,
fontSize: 11,
lineHeight: 1.5,
color: COLOR.text,
background: COLOR.surface,
borderWidth: 1,
borderStyle: 'solid',
borderColor: COLOR.border,
borderRadius: 10,
padding: 14,
boxShadow: '0 10px 34px #00000080',
},
sectionTitle: {
margin: '14px 0 5px',
color: COLOR.faint,
textTransform: 'uppercase',
letterSpacing: 1.2,
fontSize: 9,
fontWeight: 600,
display: 'flex',
justifyContent: 'space-between',
},
metricsLine: {
display: 'flex',
alignItems: 'center',
gap: 16,
margin: '0 0 10px',
paddingBottom: 10,
borderBottom: `1px solid ${COLOR.hairline}`,
fontVariantNumeric: 'tabular-nums',
},
legend: {
display: 'flex',
gap: 12,
marginTop: 3,
color: COLOR.faint,
fontSize: 10,
},
list: {
display: 'flex',
flexDirection: 'column',
maxHeight: '400px',
overflowY: 'auto',
},
configList: { display: 'flex', flexWrap: 'wrap' },
configRow: {
display: 'flex',
gap: 6,
flex: '0 0 50%',
minWidth: 0,
boxSizing: 'border-box',
padding: '1px 8px 1px 0',
},
configLabel: {
color: COLOR.muted,
flex: '0 0 45%',
whiteSpace: 'nowrap',
overflow: 'hidden',
textOverflow: 'ellipsis',
},
configValue: {
flex: 1,
minWidth: 0,
whiteSpace: 'nowrap',
overflow: 'hidden',
textOverflow: 'ellipsis',
},
trackRow: {
display: 'flex',
gap: 8,
padding: '3px 0',
borderBottom: `1px solid ${COLOR.hairline}`,
alignItems: 'baseline',
},
trackLabel: {
flex: '0 0 34%',
whiteSpace: 'nowrap',
overflow: 'hidden',
textOverflow: 'ellipsis',
},
trackCodec: {
flex: '0 0 11%',
fontFamily: MONO,
fontSize: 10,
color: COLOR.muted,
},
trackKbps: {
flex: '0 0 9%',
textAlign: 'right',
fontVariantNumeric: 'tabular-nums',
},
trackDetail: {
flex: 1,
minWidth: 0,
whiteSpace: 'nowrap',
overflow: 'hidden',
textOverflow: 'ellipsis',
},
mono: { fontFamily: MONO, fontSize: 10, color: COLOR.muted },
button: {
fontFamily: SANS,
fontSize: 11,
padding: '3px 9px',
borderRadius: 5,
borderWidth: 1,
borderStyle: 'solid',
borderColor: COLOR.border,
background: '#232428',
color: COLOR.text,
cursor: 'pointer',
},
buttonActive: { borderColor: COLOR.down, color: COLOR.down },
buttonDisabled: { color: COLOR.faint, cursor: 'not-allowed' },
row: { display: 'flex', alignItems: 'center', gap: 6, flexWrap: 'wrap' },
stepperValue: {
minWidth: 92,
textAlign: 'center',
borderWidth: 1,
borderStyle: 'solid',
borderColor: COLOR.hairline,
borderRadius: 4,
padding: '2px 6px',
background: COLOR.chartBg,
fontVariantNumeric: 'tabular-nums',
},
close: {
background: 'none',
border: 'none',
color: COLOR.faint,
cursor: 'pointer',
fontSize: 14,
},
}
const Sparkline = ({ history }: { history: StatsSnapshot[] }) => {
const width = 430
const height = 54
const max = Math.max(...history.map((s) => Math.max(s.upKbps, s.downKbps)), 1)
const points = (pick: (s: StatsSnapshot) => number) =>
history
.map(
(s, i) =>
`${((i / Math.max(history.length - 1, 1)) * width).toFixed(1)},${(
height -
(pick(s) / max) * (height - 6) -
3
).toFixed(1)}`
)
.join(' ')
return (
<>
<svg
width={width}
height={height}
role="img"
aria-label="Bandwidth over the last minute"
style={{
display: 'block',
background: COLOR.chartBg,
borderRadius: 6,
border: `1px solid ${COLOR.hairline}`,
}}
>
{history.length >= 2 && (
<>
<polyline
points={points((s) => s.downKbps)}
fill="none"
stroke={COLOR.down}
strokeWidth="1.5"
/>
<polyline
points={points((s) => s.upKbps)}
fill="none"
stroke={COLOR.up}
strokeWidth="1.5"
/>
</>
)}
</svg>
{/* Legend outside the plot: text over the lines was unreadable. */}
<div style={styles.legend}>
<span>
<span style={{ color: COLOR.down }}></span> down
</span>
<span>
<span style={{ color: COLOR.up }}></span> up
</span>
<span>last 60s · max {Math.round(max)} kbps</span>
</div>
</>
)
}
const CapStepper = ({
label,
hint,
valueKbps,
onChange,
}: {
label: string
hint: string
valueKbps: number | null
onChange: (kbps: number | null) => void
}) => (
<div style={styles.row}>
<span style={{ flex: '0 0 70px' }}>{label}</span>
<button
type="button"
style={styles.button}
aria-label={`Decrease ${label} bandwidth`}
onClick={() => onChange(stepBitrate(valueKbps, 'decrease'))}
>
</button>
<span
style={{
...styles.stepperValue,
...(valueKbps !== null ? { color: COLOR.down } : {}),
}}
>
{valueKbps === null ? 'unlimited' : `${valueKbps} kbps`}
</span>
<button
type="button"
style={styles.button}
aria-label={`Increase ${label} bandwidth`}
onClick={() => onChange(stepBitrate(valueKbps, 'increase'))}
>
+
</button>
<span style={styles.mono}>{hint}</span>
</div>
)
const TrackList = ({ title, rows }: { title: string; rows: TrackRow[] }) => (
<>
<div style={styles.sectionTitle}>
<span>
{title} ({rows.length})
</span>
<span>
{Math.round(rows.reduce((sum, t) => sum + t.kbps, 0))} kbps media
</span>
</div>
<div style={styles.list}>
{rows.map((track) => (
<div key={track.key} style={styles.trackRow}>
<span style={styles.trackLabel}>{track.label}</span>
<span style={styles.trackCodec}>{track.codec ?? ''}</span>
<span style={styles.trackKbps}>{Math.round(track.kbps)}</span>
<span style={styles.trackDetail}>
{[
track.res,
track.fps !== undefined && `${Math.round(track.fps)}fps`,
]
.filter(Boolean)
.join(' · ')}
</span>
</div>
))}
{rows.length === 0 && <span style={{ color: COLOR.faint }}>none</span>}
</div>
</>
)
const MeetDevtools = () => {
const room = useRoomContext()
const connState = useConnectionState(room)
const [open, setOpen] = useState(false)
const [firedScenario, setFiredScenario] = useState<string>()
const [uplinkCap, setUplinkCapState] = useState<number | null>(null)
const [downlinkCap, setDownlinkCapState] = useState<number | null>(null)
const { snapshot, history } = useWebRTCStats(room, open)
const transportMode: TransportMode = transportModeFromRoute(snapshot?.route)
const fireScenario = (scenario: Scenario) => {
setFiredScenario(scenario.id)
scenario.run(room).catch((e) => {
console.warn('[MeetDevtools] simulateScenario failed', e)
})
window.setTimeout(
() =>
setFiredScenario((current) =>
current === scenario.id ? undefined : current
),
1200
)
}
if (!open) {
return (
<button
type="button"
style={styles.toggle}
onClick={() => setOpen(true)}
aria-label="Open WebRTC devtools"
>
<span style={{ color: stateColor(connState) }}></span> rtc
</button>
)
}
const config = readRoomConfig(room)
const published = snapshot?.tracks.filter((t) => t.dir === 'up') ?? []
const subscribed = snapshot?.tracks.filter((t) => t.dir === 'down') ?? []
const turnProtocols = snapshot?.turnProtocols ?? []
const transports: Array<{
mode: TransportMode
label: string
requires?: string
title: string
}> = [
{
mode: 'auto',
label: 'auto (udp)',
title:
'clears the server-cached transport preference and relay-only policy, then full reconnect',
},
{
mode: 'tcp',
label: 'tcp',
title: 'force-tcp — server prefers TCP candidates (ICE/TCP)',
},
{
mode: 'turn-udp',
label: 'turn:udp',
requires: 'udp',
title:
'client-side: iceTransportPolicy relay + full reconnect (no server hook exists)',
},
{
mode: 'turn-tcp',
label: 'turn:tcp',
requires: 'tcp',
title: 'force-tcp — lands on TURN/TCP when it is the TCP path',
},
{
mode: 'turn-tls',
label: 'turn:tls',
requires: 'tls',
title: 'force-tls — server switches you to TURN over TLS',
},
]
return (
<section style={styles.panel} aria-label="WebRTC devtools">
<div style={styles.metricsLine}>
<span>
<span style={{ color: COLOR.down }}></span> {snapshot?.downKbps ?? 0}{' '}
kbps
</span>
<span>
<span style={{ color: COLOR.up }}></span> {snapshot?.upKbps ?? 0}{' '}
kbps
</span>
<span style={{ color: COLOR.muted }}>
rtt{' '}
<span style={{ color: COLOR.text }}>{snapshot?.rttMs ?? ''}</span> ms
</span>
<span style={{ color: COLOR.muted }}>
jitter{' '}
<span style={{ color: COLOR.text }}>{snapshot?.jitterMs ?? ''}</span>{' '}
ms
</span>
<span
style={{
color: stateColor(connState),
fontSize: 10,
marginLeft: 'auto',
}}
>
{connState}
</span>
<button
type="button"
style={styles.close}
onClick={() => setOpen(false)}
aria-label="Close WebRTC devtools"
>
</button>
</div>
<Sparkline history={history} />
<div style={styles.sectionTitle}>room configuration (live)</div>
<div style={styles.configList}>
{config.map((entry) => (
<div key={entry.label} style={styles.configRow}>
<span
style={{
color:
entry.on === undefined
? COLOR.faint
: entry.on
? COLOR.ok
: COLOR.bad,
}}
>
</span>
<span style={styles.configLabel}>{entry.label}</span>
<span style={styles.configValue} title={entry.value}>
{entry.value}
</span>
</div>
))}
</div>
<TrackList title="published ↑" rows={published} />
<TrackList title="subscribed ↓" rows={subscribed} />
<div style={styles.sectionTitle}>bandwidth</div>
<div style={{ ...styles.list, gap: 4 }}>
<CapStepper
label="uplink"
hint="encoder cap (setParameters)"
valueKbps={uplinkCap}
onChange={(kbps) => {
setUplinkCapState(kbps)
void setUplinkCap(room, kbps)
}}
/>
<CapStepper
label="downlink"
hint="SFU limit (subscriber-bandwidth)"
valueKbps={downlinkCap}
onChange={(kbps) => {
setDownlinkCapState(kbps)
void setDownlinkCap(room, kbps).catch((e) =>
console.warn('[MeetDevtools] downlink cap failed', e)
)
}}
/>
</div>
<div style={styles.sectionTitle}>transport</div>
<div style={styles.row}>
{transports.map(({ mode, label, requires, title }) => {
const available = !requires || turnProtocols.includes(requires)
const active = transportMode === mode
const clickable = available && !active
return (
<button
key={mode}
type="button"
aria-pressed={active}
disabled={!clickable}
title={available ? title : `${title} — not configured`}
style={{
...styles.button,
...(clickable || active ? {} : styles.buttonDisabled),
...(active ? styles.buttonActive : {}),
}}
onClick={() => {
if (!clickable) return
void (
mode === 'auto'
? releaseForcedTransport(room)
: forceTransport(
room,
mode as 'tcp' | 'turn-udp' | 'turn-tcp' | 'turn-tls'
)
).catch((e) =>
console.warn('[MeetDevtools] transport change failed', e)
)
}}
>
{requires && (
<span style={{ color: available ? COLOR.ok : COLOR.faint }}>
{' '}
</span>
)}
{label}
</button>
)
})}
<span style={styles.mono}>
route:{' '}
{snapshot?.route
? [
snapshot.route.protocol,
snapshot.route.type,
snapshot.route.relayProtocol &&
`relay:${snapshot.route.relayProtocol}`,
]
.filter(Boolean)
.join('·')
: ''}
</span>
</div>
<div style={styles.sectionTitle}>connection scenarios</div>
<div style={styles.row}>
{SCENARIOS.map((scenario) => (
<button
key={scenario.id}
type="button"
title={`${scenario.side}-side simulation — momentary, watch the state dot`}
style={{
...styles.button,
...(firedScenario === scenario.id ? styles.buttonActive : {}),
}}
onClick={() => fireScenario(scenario)}
>
{scenario.side === 'server' ? '☁ ' : ''}
{scenario.label}
</button>
))}
</div>
</section>
)
}
export default MeetDevtools
@@ -0,0 +1,14 @@
import { lazy, Suspense } from 'react'
const LazyPanel = import.meta.env.DEV
? lazy(() => import('./MeetDevtools'))
: null
export const MeetDevtools = () => {
if (!LazyPanel) return null
return (
<Suspense fallback={null}>
<LazyPanel />
</Suspense>
)
}
@@ -0,0 +1,88 @@
import { Room } from 'livekit-client'
export type ConfigEntry = {
label: string
value: string
/** true = feature actively on, false = off, undefined = informational */
on?: boolean
}
const formatBackupCodec = (backup: unknown): string => {
if (backup === undefined || backup === true) return 'auto'
if (backup === false) return 'off'
if (typeof backup === 'object' && backup !== null && 'codec' in backup) {
return String((backup as { codec: unknown }).codec)
}
return String(backup)
}
export const readRoomConfig = (room: Room): ConfigEntry[] => {
const options = room.options
const publish = options.publishDefaults
const adaptive = options.adaptiveStream
const entries: ConfigEntry[] = [
{
label: 'adaptiveStream',
value:
typeof adaptive === 'object'
? JSON.stringify(adaptive)
: String(!!adaptive),
on: !!adaptive,
},
{
label: 'dynacast',
value: String(!!options.dynacast),
on: !!options.dynacast,
},
{
label: 'e2ee',
value: String(room.isE2EEEnabled),
on: room.isE2EEEnabled,
},
{
label: 'videoCodec',
value: publish?.videoCodec ?? 'default',
},
{
label: 'backupCodec',
value: formatBackupCodec(publish?.backupCodec),
},
{
label: 'simulcast',
value: String(publish?.simulcast ?? true),
on: publish?.simulcast ?? true,
},
{
label: 'audio dtx',
value: String(publish?.dtx ?? true),
on: publish?.dtx ?? true,
},
{
label: 'audio red',
value: String(publish?.red ?? true),
on: publish?.red ?? true,
},
{
label: 'quality (local)',
value: room.localParticipant.connectionQuality,
},
]
const server = room.serverInfo
if (server) {
entries.push({
label: 'server',
value: [
server.version && `v${server.version}`,
server.region,
server.protocol !== undefined && `proto ${server.protocol}`,
server.edition !== undefined && `edition ${server.edition}`,
]
.filter(Boolean)
.join(' · '),
})
}
return entries
}
@@ -0,0 +1,192 @@
import { Room, Track } from 'livekit-client'
// Transitive dependency of livekit-client (pinned by it); used only to
// build the one signal request room.simulateScenario cannot express.
import { SimulateScenario } from '@livekit/protocol'
/** Shared bandwidth ladder for the /+ steppers; null = unlimited. */
export const BITRATE_LADDER_KBPS: Array<number | null> = [
null,
2000,
1000,
600,
300,
150,
]
export const stepBitrate = (
current: number | null,
direction: 'decrease' | 'increase'
): number | null => {
const index = BITRATE_LADDER_KBPS.indexOf(current)
const safeIndex = index === -1 ? 0 : index
const next =
direction === 'decrease'
? Math.min(safeIndex + 1, BITRATE_LADDER_KBPS.length - 1)
: Math.max(safeIndex - 1, 0)
return BITRATE_LADDER_KBPS[next]
}
const savedEncodings = new WeakMap<RTCRtpSender, Array<number | undefined>>()
export const setUplinkCap = async (
room: Room,
kbps: number | null
): Promise<void> => {
const senders: RTCRtpSender[] = []
room.localParticipant.trackPublications.forEach((pub) => {
const track = pub.track
if (track?.kind === Track.Kind.Video && track.sender) {
senders.push(track.sender)
}
})
for (const sender of senders) {
const params = sender.getParameters()
if (!params.encodings || params.encodings.length === 0) continue
if (kbps === null) {
const original = savedEncodings.get(sender)
params.encodings.forEach((encoding, i) => {
encoding.maxBitrate = original?.[i]
})
savedEncodings.delete(sender)
} else {
if (!savedEncodings.has(sender)) {
savedEncodings.set(
sender,
params.encodings.map((encoding) => encoding.maxBitrate)
)
}
const activeCount =
params.encodings.filter((encoding) => encoding.active !== false)
.length || 1
// Split the budget across active simulcast layers / SVC encoding.
const perEncoding = Math.max(
30_000,
Math.floor((kbps * 1000) / activeCount)
)
params.encodings.forEach((encoding) => {
encoding.maxBitrate = perEncoding
})
}
try {
await sender.setParameters(params)
} catch (e) {
console.warn('[MeetDevtools] setParameters failed', e)
}
}
}
export const setDownlinkCap = (
room: Room,
kbps: number | null
): Promise<void> =>
room.simulateScenario('subscriber-bandwidth', kbps === null ? 0 : kbps * 1000)
export type Scenario = {
id: string
label: string
/** where the simulation happens */
side: 'client' | 'server'
run: (room: Room) => Promise<void>
}
export const SCENARIOS: Scenario[] = [
{
id: 'resume',
label: 'reconnect (resume)',
side: 'client',
// Replays a signaling WebSocket loss; media keeps flowing, client resumes.
run: (room) => room.simulateScenario('signal-reconnect'),
},
{
id: 'resume-fail',
label: 'reconnect (resume fails)',
side: 'client',
// Same, but the next resume attempt fails → exercises the retry ladder.
run: (room) => room.simulateScenario('resume-reconnect'),
},
{
id: 'full-reconnect',
label: 'full reconnect',
side: 'client',
// Complete rejoin with brand-new peer connections.
run: (room) => room.simulateScenario('full-reconnect'),
},
{
id: 'migration',
label: 'server migration',
side: 'server',
run: (room) => room.simulateScenario('migration'),
},
{
id: 'node-failure',
label: 'SFU node failure',
side: 'server',
run: (room) => room.simulateScenario('node-failure'),
},
{
id: 'server-leave',
label: 'server disconnect',
side: 'server',
// Server-initiated leave: the closest thing to "you got kicked".
run: (room) => room.simulateScenario('server-leave'),
},
]
export type TransportMode =
| 'auto'
| 'tcp'
| 'turn-udp'
| 'turn-tcp'
| 'turn-tls'
const clearServerTransportPreference = (room: Room): Promise<void> =>
room.engine.client.sendSimulateScenario(
new SimulateScenario({
scenario: { case: 'switchCandidateProtocol', value: 0 },
})
)
const setRelayOnly = (room: Room, relay: boolean) => {
room.engine.rtcConfig = {
...room.engine.rtcConfig,
iceTransportPolicy: relay ? 'relay' : 'all',
}
}
const settle = () => new Promise((resolve) => setTimeout(resolve, 300))
export const forceTransport = async (
room: Room,
mode: 'tcp' | 'turn-udp' | 'turn-tcp' | 'turn-tls'
): Promise<void> => {
if (mode === 'turn-udp') {
await clearServerTransportPreference(room)
setRelayOnly(room, true)
await settle()
return room.simulateScenario('full-reconnect')
}
setRelayOnly(room, false)
return room.simulateScenario(mode === 'turn-tls' ? 'force-tls' : 'force-tcp')
}
export const releaseForcedTransport = async (room: Room): Promise<void> => {
setRelayOnly(room, false)
await clearServerTransportPreference(room)
await settle()
await room.simulateScenario('full-reconnect')
}
export const transportModeFromRoute = (route?: {
protocol?: string
relayProtocol?: string
}): TransportMode => {
// relayProtocol is the client→TURN leg; protocol alone means no relay.
if (route?.relayProtocol === 'tls') return 'turn-tls'
if (route?.relayProtocol === 'tcp') return 'turn-tcp'
if (route?.relayProtocol === 'udp') return 'turn-udp'
if (route?.protocol === 'tcp') return 'tcp'
return 'auto'
}
@@ -0,0 +1,243 @@
import { useEffect, useRef, useState } from 'react'
import { Room } from 'livekit-client'
export type TrackRow = {
key: string
dir: 'up' | 'down'
label: string
codec?: string
kbps: number
fps?: number
res?: string
}
export type StatsSnapshot = {
ts: number
/** wire totals from transport stats (includes headers, RTCP, FEC) */
upKbps: number
downKbps: number
rttMs?: number
/** worst inbound RTP jitter across subscribed tracks */
jitterMs?: number
availableOutKbps?: number
/** selected ICE route of the publisher transport (measured, not assumed) */
route?: { protocol?: string; type?: string; relayProtocol?: string }
/**
* relayProtocol values of gathered relay local candidates — i.e. which
* client→TURN transports are actually configured (udp/tcp/tls). Empty
* when no TURN server is configured.
*/
turnProtocols: string[]
tracks: TrackRow[]
}
type StatDict = Record<string, unknown>
type Counters = Record<string, number>
const asNumber = (v: unknown): number | undefined =>
typeof v === 'number' && Number.isFinite(v) ? v : undefined
const asString = (v: unknown): string | undefined =>
typeof v === 'string' ? v : undefined
const shortCodec = (mimeType?: string) =>
mimeType ? mimeType.replace(/^(audio|video)\//, '') : undefined
const buildTrackLabels = (room: Room) => {
const labels = new Map<string, string>()
room.localParticipant.trackPublications.forEach((pub) => {
const id = pub.track?.mediaStreamTrack?.id
if (id) labels.set(id, `local ${pub.source}`)
})
room.remoteParticipants.forEach((participant) => {
// Keep rows scannable: participant names capped at 10 chars.
const rawName = participant.name || participant.identity
const name = rawName.length > 10 ? `${rawName.slice(0, 10)}.` : rawName
participant.trackPublications.forEach((pub) => {
const id = pub.track?.mediaStreamTrack?.id
if (id) labels.set(id, `${name} ${pub.source}`)
})
})
return labels
}
export const useWebRTCStats = (
room: Room,
enabled: boolean,
intervalMs = 1000
) => {
// Single source of truth: the snapshot is just the last history entry.
const [history, setHistory] = useState<StatsSnapshot[]>([])
const prevRef = useRef(new Map<string, Counters>())
useEffect(() => {
if (!enabled) return
let cancelled = false
const prev = prevRef.current
/** Per-stat counter deltas; returns 0 on the first sighting. */
const deltas = (key: string, now: Counters): Counters => {
const before = prev.get(key)
prev.set(key, now)
const out: Counters = {}
for (const [name, value] of Object.entries(now)) {
out[name] = before?.[name] !== undefined ? value - before[name] : 0
}
return out
}
const collect = async () => {
const reports: Array<{ pc: 'pub' | 'sub'; report: RTCStatsReport }> = []
try {
// Not public API — see file header.
const manager = room.engine?.pcManager
const pub = await manager?.publisher?.getStats()
if (pub) reports.push({ pc: 'pub', report: pub })
const sub = await manager?.subscriber?.getStats()
if (sub) reports.push({ pc: 'sub', report: sub })
} catch {
// Engine not ready or SDK internals changed; panel shows nothing.
}
if (cancelled || reports.length === 0) return
const labels = buildTrackLabels(room)
const tracks: TrackRow[] = []
let upKbps = 0
let downKbps = 0
let rttMs: number | undefined
let jitterMs: number | undefined
let availableOutKbps: number | undefined
let route: StatsSnapshot['route']
const turnProtocols = new Set<string>()
for (const { pc, report } of reports) {
const byId = new Map<string, StatDict>()
report.forEach((stat) => byId.set(stat.id as string, stat as StatDict))
report.forEach((raw) => {
const stat = raw as StatDict
const type = asString(stat.type)
const ts = asNumber(stat.timestamp) ?? Date.now()
const key = `${pc}:${asString(stat.id) ?? ''}`
if (
type === 'local-candidate' &&
asString(stat.candidateType) === 'relay'
) {
// Relay candidates are only gathered when a TURN server is
// configured and reachable; relayProtocol says how the client
// reaches it (udp/tcp/tls).
turnProtocols.add(
asString(stat.relayProtocol) ?? asString(stat.protocol) ?? 'udp'
)
}
if (type === 'transport') {
const d = deltas(key, {
sent: asNumber(stat.bytesSent) ?? 0,
received: asNumber(stat.bytesReceived) ?? 0,
ts,
})
if (d.ts > 0) {
upKbps += Math.max(0, (d.sent * 8) / d.ts)
downKbps += Math.max(0, (d.received * 8) / d.ts)
}
}
if (type === 'candidate-pair' && stat.nominated === true) {
const rtt = asNumber(stat.currentRoundTripTime)
if (rtt !== undefined) rttMs = Math.round(rtt * 1000)
const available = asNumber(stat.availableOutgoingBitrate)
if (available !== undefined && pc === 'pub') {
availableOutKbps = Math.round(available / 1000)
}
if (pc === 'pub') {
const local = byId.get(asString(stat.localCandidateId) ?? '')
route = {
protocol: asString(local?.protocol),
type: asString(local?.candidateType),
relayProtocol: asString(local?.relayProtocol),
}
}
}
if (type === 'outbound-rtp' || type === 'inbound-rtp') {
const isUp = type === 'outbound-rtp'
if (!isUp) {
const jitter = asNumber(stat.jitter)
if (jitter !== undefined) {
const ms = Math.round(jitter * 1000)
if (jitterMs === undefined || ms > jitterMs) jitterMs = ms
}
}
const d = deltas(key, {
bytes: asNumber(isUp ? stat.bytesSent : stat.bytesReceived) ?? 0,
ts,
})
const kbps = d.ts > 0 ? Math.max(0, (d.bytes * 8) / d.ts) : 0
// Resolve codec + source track.
const codecStat = byId.get(asString(stat.codecId) ?? '')
let msTrackId = asString(stat.trackIdentifier)
if (!msTrackId && isUp) {
const mediaSource = byId.get(asString(stat.mediaSourceId) ?? '')
msTrackId = asString(mediaSource?.trackIdentifier)
}
const rid = asString(stat.rid)
const baseLabel =
(msTrackId && labels.get(msTrackId)) ??
`${asString(stat.kind) ?? 'media'} ssrc ${asNumber(stat.ssrc) ?? '?'}`
const width = asNumber(stat.frameWidth)
const height = asNumber(stat.frameHeight)
tracks.push({
key,
dir: isUp ? 'up' : 'down',
label: rid ? `${baseLabel} [${rid}]` : baseLabel,
codec: shortCodec(asString(codecStat?.mimeType)),
kbps,
fps: asNumber(stat.framesPerSecond),
res: width && height ? `${width}x${height}` : undefined,
})
}
})
}
// Stable order (direction, then label): sorting by bitrate would
// reshuffle rows on every tick as kbps fluctuates.
tracks.sort((a, b) =>
a.dir === b.dir
? a.label.localeCompare(b.label)
: a.dir === 'up'
? -1
: 1
)
const next: StatsSnapshot = {
ts: Date.now(),
upKbps: Math.round(upKbps),
downKbps: Math.round(downKbps),
rttMs,
jitterMs,
availableOutKbps,
route,
turnProtocols: Array.from(turnProtocols).sort(),
tracks,
}
setHistory((h) => [...h.slice(-59), next])
}
void collect()
const id = window.setInterval(() => void collect(), intervalMs)
return () => {
cancelled = true
window.clearInterval(id)
}
}, [room, enabled, intervalMs])
return { snapshot: history[history.length - 1], history }
}
@@ -43,6 +43,7 @@ import { useSnapshot } from 'valtio'
import { userPreferencesStore } from '@/stores/userPreferences'
import { userStore } from '@/stores/user'
import { WatchMediaDeviceErrors } from './WatchMediaDeviceErrors'
import { MeetDevtools } from '@/features/devtools'
import { VOICE_AUDIO_CONSTRAINTS } from '@/features/rooms/livekit/utils/constants'
export const Conference = ({
@@ -298,6 +299,7 @@ export const Conference = ({
<VideoConference />
{!isMobile && <InviteDialog mode={mode} />}
<PictureInPictureConference />
<MeetDevtools />
</LiveKitRoom>
</Screen>
</QueryAware>
@@ -75,15 +75,15 @@ const useTranscriptionState = () => {
const segment = segments[0]
setTranscriptionSegments((prevSegments) => {
const existingSegmentIds = new Set(prevSegments.map((s) => s.id))
if (existingSegmentIds.has(segment.id)) return prevSegments
return [
...prevSegments,
{
participant: participant,
...segment,
},
]
const existingIndex = prevSegments.findIndex(
(s: TranscriptionSegmentWithParticipant) => s.id === segment.id
)
if (existingIndex === -1) {
return [...prevSegments, { participant, ...segment }]
}
const next = prevSegments.slice()
next[existingIndex] = { ...next[existingIndex], ...segment }
return next
})
}