|
|
@@ -0,0 +1,103 @@
|
|
|
+"""
|
|
|
+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")
|
|
|
+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 conversation_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 new unique room and spawn worker."""
|
|
|
+ try:
|
|
|
+ body = await request.json()
|
|
|
+ app_identity = body.get("identity", f"user-{int(time.time()*1000)%100000}")
|
|
|
+ except Exception:
|
|
|
+ app_identity = f"user-{int(time.time()*1000)%100000}"
|
|
|
+
|
|
|
+ room_name = f"room-{int(time.time_ns())}"
|
|
|
+
|
|
|
+ # Spawn worker (runs in background, connects to room)
|
|
|
+ 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": 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", 9000)
|
|
|
+ await site.start()
|
|
|
+ logger.info("Dispatcher listening on :9000")
|
|
|
+ await asyncio.Future()
|
|
|
+
|
|
|
+
|
|
|
+if __name__ == "__main__":
|
|
|
+ asyncio.run(main())
|