| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118 |
- """
- LiveKit Worker Dispatcher — per-room ASR worker management.
- App calls POST /connect → gets tokens for a unique room → connects.
- Worker is spawned for that room automatically.
- """
- from __future__ import annotations
- import asyncio
- import logging
- import os
- import sys
- import time
- import jwt
- from aiohttp import web
- sys.path.insert(0, os.path.join(os.path.dirname(os.path.abspath(__file__)), "whisper_asr"))
- logger = logging.getLogger("dispatcher")
- LIVEKIT_URL = os.environ.get("LIVEKIT_URL", "ws://localhost:7880")
- PUBLIC_LIVEKIT_URL = os.environ.get("PUBLIC_LIVEKIT_URL", LIVEKIT_URL)
- API_KEY = "devkey"
- API_SECRET = "secretsecretsecretsecretsecret12"
- def _room_token(room_name: str, identity: str, *, can_create: bool = False) -> str:
- n = int(time.time())
- return jwt.encode({
- "iss": API_KEY, "sub": identity,
- "nbf": n - 60, "exp": n + 24 * 3600,
- "video": {
- "roomJoin": True, "roomCreate": can_create,
- "room": room_name,
- "canPublish": True, "canSubscribe": True, "canPublishData": True,
- },
- }, API_SECRET, algorithm="HS256")
- _workers: dict[str, asyncio.Task] = {}
- async def _spawn_worker(room_name: str):
- """Run a worker for one room. Exits when room is empty or cancelled."""
- from worker import Worker
- try:
- logger.info("Worker starting room=%s", room_name)
- w = Worker(room_name, identity="asr-bot")
- await w.run()
- except asyncio.CancelledError:
- pass
- except Exception:
- logger.exception("Worker %s failed", room_name)
- finally:
- logger.info("Worker stopped room=%s", room_name)
- async def handle_connect(request: web.Request) -> web.Response:
- """POST /connect — return tokens for a room and spawn worker.
- Body: {"identity": "user-xxx", "room": "room-xxx" (optional)}
- If room is provided and a worker already exists, reuse it.
- """
- try:
- body = await request.json()
- app_identity = body.get("identity", f"user-{int(time.time()*1000)%100000}")
- resume_room = body.get("room")
- except Exception:
- app_identity = f"user-{int(time.time()*1000)%100000}"
- resume_room = None
- if resume_room:
- room_name = resume_room
- if room_name in _workers and not _workers[room_name].done():
- logger.info("Session resume: %s (%s)", room_name, app_identity)
- else:
- # Recreate worker for the same room so App can reuse the room name
- # even after the previous worker exited.
- task = asyncio.create_task(_spawn_worker(room_name))
- _workers[room_name] = task
- logger.info("Room %s recreated for %s (total workers: %d)", room_name, app_identity, len(_workers))
- else:
- room_name = f"room-{int(time.time_ns())}"
- task = asyncio.create_task(_spawn_worker(room_name))
- _workers[room_name] = task
- logger.info("Room %s created for %s (total workers: %d)", room_name, app_identity, len(_workers))
- return web.json_response({
- "room": room_name,
- "url": PUBLIC_LIVEKIT_URL,
- "token": _room_token(room_name, app_identity, can_create=True),
- })
- async def handle_health(request: web.Request) -> web.Response:
- # Clean up finished workers
- for name, task in list(_workers.items()):
- if task.done():
- _workers.pop(name, None)
- return web.json_response({"status": "ok", "active_workers": len(_workers)})
- async def main():
- logging.basicConfig(level=logging.INFO)
- app = web.Application()
- app.router.add_post("/connect", handle_connect)
- app.router.add_get("/health", handle_health)
- runner = web.AppRunner(app)
- await runner.setup()
- site = web.TCPSite(runner, "0.0.0.0", 9100)
- await site.start()
- logger.info("Dispatcher listening on :9100")
- await asyncio.Future()
- if __name__ == "__main__":
- asyncio.run(main())
|