mirror of
https://github.com/suitenumerique/meet.git
synced 2026-07-27 12:19:10 +00:00
Compare commits
13 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 8ee3cdeab8 | |||
| e7c95f2806 | |||
| 2cb24dacbe | |||
| c7831adc46 | |||
| 4b655fee30 | |||
| 498676270e | |||
| b58fbd9eef | |||
| 2bc60effc8 | |||
| 64bf50116b | |||
| e0f9460606 | |||
| de3764391d | |||
| ec67a12fe4 | |||
| 05f32d008a |
@@ -150,13 +150,14 @@ jobs:
|
||||
uses: actions/setup-python@v6
|
||||
with:
|
||||
python-version: "3.13"
|
||||
cache: "pip"
|
||||
- name: Install development dependencies
|
||||
run: pip install --user .[dev]
|
||||
- name: Install uv
|
||||
uses: astral-sh/setup-uv@v7
|
||||
- name: Install the project
|
||||
run: uv sync --locked --all-extras
|
||||
- name: Check code formatting with ruff
|
||||
run: ~/.local/bin/ruff format . --diff
|
||||
run: uv run ruff format . --diff
|
||||
- name: Lint code with ruff
|
||||
run: ~/.local/bin/ruff check .
|
||||
run: uv run ruff check .
|
||||
|
||||
lint-summary:
|
||||
runs-on: ubuntu-latest
|
||||
|
||||
@@ -18,6 +18,11 @@ and this project adheres to
|
||||
### Changed
|
||||
|
||||
- ♻️(summary) change tasks endpoint signature
|
||||
- ⬆️(dependencies) update urllib3 to v2.7.0 [SECURITY]
|
||||
|
||||
### Changed
|
||||
|
||||
- 🧑💻(agents) use `uv` for package management
|
||||
|
||||
### Fixed
|
||||
|
||||
|
||||
@@ -249,6 +249,7 @@ services:
|
||||
metadata-collector-dev:
|
||||
build:
|
||||
context: ./src/agents
|
||||
target: development
|
||||
command: ["python", "metadata_collector.py", "dev"]
|
||||
environment:
|
||||
- LIVEKIT_URL=ws://livekit:7880
|
||||
@@ -261,6 +262,7 @@ services:
|
||||
- AWS_S3_SECURE_ACCESS=False
|
||||
volumes:
|
||||
- ./src/agents:/app
|
||||
- /app/.venv
|
||||
depends_on:
|
||||
- livekit
|
||||
- minio
|
||||
@@ -333,10 +335,13 @@ services:
|
||||
multi-user-transcriber:
|
||||
build:
|
||||
context: ./src/agents
|
||||
target: development
|
||||
command: ["python", "multi_user_transcriber.py", "dev"]
|
||||
env_file:
|
||||
- env.d/development/multi_user_transcriber
|
||||
volumes:
|
||||
- ./src/agents:/app
|
||||
- /app/.venv
|
||||
|
||||
networks:
|
||||
default:
|
||||
|
||||
@@ -1,9 +1,21 @@
|
||||
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
|
||||
STT_PROVIDER=voxtral-vllm # voxtral-vllm, kyutai, deepgram
|
||||
ENABLE_SILERO_VAD=False
|
||||
|
||||
KYUTAI_STT_BASE_URL=
|
||||
KYUTAI_API_KEY=
|
||||
DEEPGRAM_API_KEY=your-deepgram-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
|
||||
|
||||
|
||||
@@ -0,0 +1,38 @@
|
||||
# Python
|
||||
__pycache__
|
||||
*.pyc
|
||||
**/__pycache__
|
||||
**/*.pyc
|
||||
venv
|
||||
**/.venv
|
||||
|
||||
# System-specific files
|
||||
.DS_Store
|
||||
**/.DS_Store
|
||||
|
||||
# Docker
|
||||
compose.*
|
||||
env.d
|
||||
|
||||
# Docs
|
||||
docs
|
||||
*.md
|
||||
*.log
|
||||
|
||||
# Development/test cache & configurations
|
||||
data
|
||||
.cache
|
||||
.circleci
|
||||
.git
|
||||
.iml
|
||||
db.sqlite3
|
||||
.pylint.d
|
||||
|
||||
**/.idea
|
||||
**/.vscode
|
||||
**/.pytest_cache
|
||||
**/.mypy_cache
|
||||
**/.ruff_cache
|
||||
|
||||
# Env
|
||||
.env
|
||||
+41
-13
@@ -6,31 +6,61 @@ RUN apt-get update && apt-get install -y \
|
||||
libgobject-2.0-0 \
|
||||
&& rm -rf /var/lib/apt/lists/*
|
||||
|
||||
|
||||
# ---- Builder image ----
|
||||
FROM base AS builder
|
||||
|
||||
WORKDIR /builder
|
||||
ENV UV_COMPILE_BYTECODE=1 \
|
||||
UV_LINK_MODE=copy \
|
||||
UV_PYTHON_DOWNLOADS=0
|
||||
|
||||
COPY pyproject.toml .
|
||||
|
||||
RUN mkdir /install && \
|
||||
pip install --prefix=/install .
|
||||
|
||||
FROM base AS development
|
||||
# Install uv
|
||||
COPY --from=ghcr.io/astral-sh/uv:0.10.9 /uv /uvx /bin/
|
||||
|
||||
WORKDIR /app
|
||||
|
||||
COPY pyproject.toml .
|
||||
RUN pip install --no-cache-dir ".[dev]"
|
||||
# Install production dependencies without the project itself (cacheable layer)
|
||||
RUN --mount=type=cache,target=/root/.cache/uv \
|
||||
--mount=type=bind,source=uv.lock,target=uv.lock \
|
||||
--mount=type=bind,source=pyproject.toml,target=pyproject.toml \
|
||||
uv sync --locked --no-install-project --no-dev
|
||||
|
||||
COPY . .
|
||||
# Install the project
|
||||
COPY . /app
|
||||
RUN --mount=type=cache,target=/root/.cache/uv \
|
||||
uv sync --locked --no-dev
|
||||
|
||||
|
||||
# ---- Development image ----
|
||||
FROM base AS development
|
||||
|
||||
ENV UV_COMPILE_BYTECODE=1 \
|
||||
UV_LINK_MODE=copy \
|
||||
UV_PYTHON_DOWNLOADS=0
|
||||
|
||||
COPY --from=ghcr.io/astral-sh/uv:0.10.9 /uv /uvx /bin/
|
||||
|
||||
WORKDIR /app
|
||||
|
||||
COPY . /app
|
||||
|
||||
RUN --mount=type=cache,target=/root/.cache/uv \
|
||||
uv sync --locked --all-extras
|
||||
|
||||
ENV PATH="/app/.venv/bin:$PATH"
|
||||
|
||||
CMD ["python", "metadata_collector.py", "dev"]
|
||||
|
||||
|
||||
# ---- Production image ----
|
||||
FROM base AS production
|
||||
|
||||
WORKDIR /app
|
||||
|
||||
COPY --from=builder /install /usr/local
|
||||
# Copy the pre-built virtualenv and application source
|
||||
COPY --from=builder /app /app
|
||||
|
||||
ENV PATH="/app/.venv/bin:$PATH"
|
||||
|
||||
# Remove pip to reduce attack surface in production
|
||||
RUN pip uninstall -y pip
|
||||
@@ -39,6 +69,4 @@ RUN pip uninstall -y pip
|
||||
ARG DOCKER_USER
|
||||
USER ${DOCKER_USER}
|
||||
|
||||
COPY ./*.py /app/
|
||||
|
||||
CMD ["python", "multi_user_transcriber.py", "start"]
|
||||
|
||||
@@ -25,6 +25,8 @@ from livekit.agents import (
|
||||
)
|
||||
from livekit.plugins import deepgram, silero
|
||||
|
||||
import voxtral_vllm_stt
|
||||
|
||||
load_dotenv()
|
||||
|
||||
logger = logging.getLogger("transcriber")
|
||||
@@ -46,6 +48,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()
|
||||
else:
|
||||
raise ValueError(f"Unknown STT_PROVIDER: {STT_PROVIDER}")
|
||||
|
||||
|
||||
@@ -18,12 +18,8 @@ dev = [
|
||||
"ruff==0.15.6",
|
||||
]
|
||||
|
||||
[tool.setuptools]
|
||||
py-modules = ["multi_user_transcriber", "metadata_collector", "exceptions"]
|
||||
|
||||
[build-system]
|
||||
requires = ["setuptools>=61.0"]
|
||||
build-backend = "setuptools.build_meta"
|
||||
[tool.uv]
|
||||
package = false
|
||||
|
||||
[tool.ruff]
|
||||
target-version = "py313"
|
||||
|
||||
Generated
+1963
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,474 @@
|
||||
"""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):
|
||||
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)
|
||||
@@ -61,6 +61,7 @@ dependencies = [
|
||||
"mozilla-django-oidc==5.0.2",
|
||||
"livekit-api==1.1.0",
|
||||
"aiohttp==3.13.4",
|
||||
"urllib3==2.7.0",
|
||||
]
|
||||
|
||||
[project.urls]
|
||||
|
||||
Generated
+5
-3
@@ -1212,6 +1212,7 @@ dependencies = [
|
||||
{ name = "redis" },
|
||||
{ name = "requests" },
|
||||
{ name = "sentry-sdk" },
|
||||
{ name = "urllib3" },
|
||||
{ name = "whitenoise" },
|
||||
]
|
||||
|
||||
@@ -1273,6 +1274,7 @@ requires-dist = [
|
||||
{ name = "redis", specifier = "==5.2.1" },
|
||||
{ name = "requests", specifier = "==2.33.0" },
|
||||
{ name = "sentry-sdk", specifier = "==2.54.0" },
|
||||
{ name = "urllib3", specifier = "==2.7.0" },
|
||||
{ name = "whitenoise", specifier = "==6.12.0" },
|
||||
]
|
||||
|
||||
@@ -2260,11 +2262,11 @@ wheels = [
|
||||
|
||||
[[package]]
|
||||
name = "urllib3"
|
||||
version = "2.6.3"
|
||||
version = "2.7.0"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/c7/24/5f1b3bdffd70275f6661c76461e25f024d5a38a46f04aaca912426a2b1d3/urllib3-2.6.3.tar.gz", hash = "sha256:1b62b6884944a57dbe321509ab94fd4d3b307075e0c2eae991ac71ee15ad38ed", size = 435556, upload-time = "2026-01-07T16:24:43.925Z" }
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/53/0c/06f8b233b8fd13b9e5ee11424ef85419ba0d8ba0b3138bf360be2ff56953/urllib3-2.7.0.tar.gz", hash = "sha256:231e0ec3b63ceb14667c67be60f2f2c40a518cb38b03af60abc813da26505f4c", size = 433602, upload-time = "2026-05-07T16:13:18.596Z" }
|
||||
wheels = [
|
||||
{ url = "https://files.pythonhosted.org/packages/39/08/aaaad47bc4e9dc8c725e68f9d04865dbcb2052843ff09c97b08904852d84/urllib3-2.6.3-py3-none-any.whl", hash = "sha256:bf272323e553dfb2e87d9bfd225ca7b0f467b919d7bbd355436d3fd37cb0acd4", size = 131584, upload-time = "2026-01-07T16:24:42.685Z" },
|
||||
{ url = "https://files.pythonhosted.org/packages/7f/3e/5db95bcf282c52709639744ca2a8b149baccf648e39c8cc87553df9eae0c/urllib3-2.7.0-py3-none-any.whl", hash = "sha256:9fb4c81ebbb1ce9531cce37674bbc6f1360472bc18ca9a553ede278ef7276897", size = 131087, upload-time = "2026-05-07T16:13:17.151Z" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
|
||||
@@ -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