mirror of
https://github.com/suitenumerique/meet.git
synced 2026-08-31 20:58:03 +00:00
Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 432adf84e6 | |||
| 5747381eff | |||
| 0a0cdae896 |
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -8,3 +8,16 @@ webhook:
|
||||
api_key: devkey
|
||||
urls:
|
||||
- http://app-dev:8000/api/v1.0/rooms/webhooks-livekit/
|
||||
|
||||
turn:
|
||||
enabled: true
|
||||
domain: turn.192.168.1.50.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
|
||||
|
||||
|
||||
@@ -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=
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
|
||||
@@ -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]
|
||||
|
||||
Generated
+770
-510
File diff suppressed because it is too large
Load Diff
@@ -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,608 @@
|
||||
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: '200px', overflowY: 'auto' },
|
||||
listRow: { display: 'flex', gap: 8, padding: '1px 0' },
|
||||
listLabel: { flex: '0 0 34%', color: COLOR.muted },
|
||||
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 },
|
||||
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 TrackDetail = ({ track }: { track: TrackRow }) => {
|
||||
const parts: Array<{ text: string; tone?: 'warn' | 'muted' }> = []
|
||||
if (track.res) parts.push({ text: track.res })
|
||||
if (track.fps !== undefined) {
|
||||
parts.push({ text: `${Math.round(track.fps)}fps` })
|
||||
}
|
||||
if (track.lossPct !== undefined && track.lossPct > 0) {
|
||||
parts.push({ text: `loss ${track.lossPct.toFixed(1)}%`, tone: 'warn' })
|
||||
}
|
||||
if (track.froze) parts.push({ text: 'freeze', tone: 'warn' })
|
||||
if (track.limitation) {
|
||||
parts.push({ text: `lim:${track.limitation}`, tone: 'warn' })
|
||||
}
|
||||
const layer = track.layer
|
||||
if (layer) {
|
||||
if (layer.availRes && layer.availRes !== track.res) {
|
||||
parts.push({
|
||||
text: `≤${layer.availRes}${layer.layerCount ? `(${layer.layerCount}L)` : ''}`,
|
||||
tone: 'muted',
|
||||
})
|
||||
}
|
||||
if (layer.reason === 'adaptive' && layer.elementRes) {
|
||||
parts.push({ text: `fit:${layer.elementRes}`, tone: 'muted' })
|
||||
} else if (layer.reason === 'bandwidth') {
|
||||
parts.push({ text: 'bw-limited', tone: 'warn' })
|
||||
} else if (layer.reason === 'paused') {
|
||||
parts.push({ text: 'paused:SFU', tone: 'warn' })
|
||||
} else if (layer.reason === 'off-screen') {
|
||||
parts.push({ text: 'off-screen', tone: 'muted' })
|
||||
}
|
||||
}
|
||||
return (
|
||||
<>
|
||||
{parts.map((part, i) => (
|
||||
<span
|
||||
key={part.text}
|
||||
style={
|
||||
part.tone === 'warn'
|
||||
? { color: COLOR.bad }
|
||||
: part.tone === 'muted'
|
||||
? { color: COLOR.faint }
|
||||
: {}
|
||||
}
|
||||
>
|
||||
{i > 0 && ' · '}
|
||||
{part.text}
|
||||
</span>
|
||||
))}
|
||||
</>
|
||||
)
|
||||
}
|
||||
|
||||
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}>
|
||||
<TrackDetail track={track} />
|
||||
</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.list}>
|
||||
{config.map((entry) => (
|
||||
<div key={entry.label} style={styles.listRow}>
|
||||
<span
|
||||
style={{
|
||||
color:
|
||||
entry.on === undefined
|
||||
? COLOR.faint
|
||||
: entry.on
|
||||
? COLOR.ok
|
||||
: COLOR.bad,
|
||||
}}
|
||||
>
|
||||
●
|
||||
</span>
|
||||
<span style={styles.listLabel}>{entry.label}</span>
|
||||
<span>{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,15 @@
|
||||
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,89 @@
|
||||
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,193 @@
|
||||
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,357 @@
|
||||
import { useEffect, useRef, useState } from 'react'
|
||||
import { Room, Track } from 'livekit-client'
|
||||
|
||||
|
||||
export type LayerInfo = {
|
||||
/** top published layer, e.g. '1280x720' */
|
||||
availRes?: string
|
||||
/** number of published spatial layers (simulcast/SVC), when known */
|
||||
layerCount?: number
|
||||
/** largest attached element size in device px, e.g. '480x270' */
|
||||
elementRes?: string
|
||||
/** why the forwarded layer is below the top one */
|
||||
reason?: 'adaptive' | 'bandwidth' | 'paused' | 'off-screen'
|
||||
}
|
||||
|
||||
export type TrackRow = {
|
||||
key: string
|
||||
dir: 'up' | 'down'
|
||||
kind: string
|
||||
label: string
|
||||
codec?: string
|
||||
kbps: number
|
||||
fps?: number
|
||||
res?: string
|
||||
lossPct?: number
|
||||
limitation?: string
|
||||
/** subscribed video only: SFU layer forwarding context */
|
||||
layer?: LayerInfo
|
||||
/** this track froze during the last tick */
|
||||
froze?: boolean
|
||||
}
|
||||
|
||||
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
|
||||
|
||||
type TrackContext = { label: string; layer?: LayerInfo; publishedH?: number }
|
||||
|
||||
/**
|
||||
* msTrackId -> label + subscribed-layer context, rebuilt on every tick.
|
||||
* The "why is this tile blurry" answer needs SDK state (publication +
|
||||
* element size) joined with getStats (forwarded resolution); this is the
|
||||
* SDK-state half.
|
||||
*/
|
||||
const buildTrackContext = (room: Room) => {
|
||||
const contexts = new Map<string, TrackContext>()
|
||||
room.localParticipant.trackPublications.forEach((pub) => {
|
||||
const id = pub.track?.mediaStreamTrack?.id
|
||||
if (id) contexts.set(id, { label: `local ${pub.source}` })
|
||||
})
|
||||
const dpr = window.devicePixelRatio || 1
|
||||
room.remoteParticipants.forEach((participant) => {
|
||||
// Keep rows scannable: participant names capped at 10 chars.
|
||||
const rawName = participant.name || participant.identity
|
||||
const name = rawName.length > 15 ? `${rawName.slice(0, 15)}.` : rawName
|
||||
participant.trackPublications.forEach((pub) => {
|
||||
const track = pub.track
|
||||
const id = track?.mediaStreamTrack?.id
|
||||
if (!track || !id) return
|
||||
const label = `${name} ${pub.source}`
|
||||
if (track.kind !== Track.Kind.Video) {
|
||||
contexts.set(id, { label })
|
||||
return
|
||||
}
|
||||
let elementRes: string | undefined
|
||||
let elementH: number | undefined
|
||||
for (const element of track.attachedElements) {
|
||||
const w = Math.round(element.clientWidth * dpr)
|
||||
const h = Math.round(element.clientHeight * dpr)
|
||||
if (elementH === undefined || h > elementH) {
|
||||
elementH = h
|
||||
elementRes = `${w}x${h}`
|
||||
}
|
||||
}
|
||||
const dims = pub.dimensions
|
||||
const layers = pub.trackInfo?.layers?.length
|
||||
contexts.set(id, {
|
||||
label,
|
||||
publishedH: dims?.height,
|
||||
layer: {
|
||||
availRes: dims ? `${dims.width}x${dims.height}` : undefined,
|
||||
layerCount: layers && layers > 1 ? layers : undefined,
|
||||
elementRes,
|
||||
reason: !pub.isEnabled
|
||||
? 'off-screen'
|
||||
: track.streamState === Track.StreamState.Paused
|
||||
? 'paused'
|
||||
: undefined,
|
||||
},
|
||||
})
|
||||
})
|
||||
})
|
||||
return contexts
|
||||
}
|
||||
|
||||
/** Decide why a forwarded layer is below the published top layer. */
|
||||
const resolveLayerReason = (
|
||||
context: TrackContext,
|
||||
forwardedH: number | undefined
|
||||
): LayerInfo | undefined => {
|
||||
const layer = context.layer
|
||||
if (!layer) return undefined
|
||||
if (layer.reason) return layer // off-screen / paused already decided
|
||||
const publishedH = context.publishedH
|
||||
if (!publishedH || !forwardedH || forwardedH >= publishedH * 0.9) {
|
||||
return { ...layer, reason: undefined } // full quality, nothing to explain
|
||||
}
|
||||
// Below top layer: if the element only needs about what we get, it's
|
||||
// adaptiveStream fitting the element; otherwise the SFU is holding back
|
||||
// a layer the element could use — congestion.
|
||||
const elementH = layer.elementRes
|
||||
? Number(layer.elementRes.split('x')[1])
|
||||
: undefined
|
||||
const adaptive = elementH !== undefined && forwardedH >= elementH * 0.7
|
||||
return { ...layer, reason: adaptive ? 'adaptive' : 'bandwidth' }
|
||||
}
|
||||
|
||||
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 contexts = buildTrackContext(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) — a turn:…?transport=udp server
|
||||
// must NOT be presented as a TURN/TLS capability.
|
||||
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'
|
||||
const kind = asString(stat.kind)
|
||||
|
||||
if (!isUp) {
|
||||
const jitter = asNumber(stat.jitter)
|
||||
if (jitter !== undefined) {
|
||||
const ms = Math.round(jitter * 1000)
|
||||
if (jitterMs === undefined || ms > jitterMs) jitterMs = ms
|
||||
}
|
||||
}
|
||||
|
||||
// Packet loss: reported directly on inbound-rtp; for outbound it
|
||||
// lives on the matching remote-inbound-rtp (what the SFU got).
|
||||
let packetsLost = asNumber(stat.packetsLost) ?? 0
|
||||
if (isUp) {
|
||||
const remote = byId.get(asString(stat.remoteId) ?? '')
|
||||
packetsLost = asNumber(remote?.packetsLost) ?? 0
|
||||
}
|
||||
|
||||
const d = deltas(key, {
|
||||
bytes: asNumber(isUp ? stat.bytesSent : stat.bytesReceived) ?? 0,
|
||||
packets:
|
||||
asNumber(isUp ? stat.packetsSent : stat.packetsReceived) ?? 0,
|
||||
packetsLost,
|
||||
ts,
|
||||
freezeCount: asNumber(stat.freezeCount) ?? 0,
|
||||
})
|
||||
|
||||
const kbps = d.ts > 0 ? Math.max(0, (d.bytes * 8) / d.ts) : 0
|
||||
const lostTotal = d.packetsLost + d.packets
|
||||
const lossPct =
|
||||
lostTotal > 0
|
||||
? Math.max(0, Math.min(100, (d.packetsLost / lostTotal) * 100))
|
||||
: 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 context = msTrackId ? contexts.get(msTrackId) : undefined
|
||||
const baseLabel =
|
||||
context?.label ??
|
||||
`${kind ?? 'media'} ssrc ${asNumber(stat.ssrc) ?? '?'}`
|
||||
|
||||
const width = asNumber(stat.frameWidth)
|
||||
const height = asNumber(stat.frameHeight)
|
||||
|
||||
tracks.push({
|
||||
key,
|
||||
dir: isUp ? 'up' : 'down',
|
||||
kind: kind ?? 'unknown',
|
||||
label: rid ? `${baseLabel} [${rid}]` : baseLabel,
|
||||
codec: shortCodec(asString(codecStat?.mimeType)),
|
||||
kbps,
|
||||
fps: asNumber(stat.framesPerSecond),
|
||||
res: width && height ? `${width}x${height}` : undefined,
|
||||
lossPct,
|
||||
limitation:
|
||||
asString(stat.qualityLimitationReason) === 'none'
|
||||
? undefined
|
||||
: asString(stat.qualityLimitationReason),
|
||||
layer:
|
||||
!isUp && kind === 'video' && context
|
||||
? resolveLayerReason(context, height)
|
||||
: undefined,
|
||||
froze: !isUp && kind === 'video' && d.freezeCount > 0,
|
||||
})
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// 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
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user