dispatcher.py 3.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118
  1. """
  2. LiveKit Worker Dispatcher — per-room ASR worker management.
  3. App calls POST /connect → gets tokens for a unique room → connects.
  4. Worker is spawned for that room automatically.
  5. """
  6. from __future__ import annotations
  7. import asyncio
  8. import logging
  9. import os
  10. import sys
  11. import time
  12. import jwt
  13. from aiohttp import web
  14. sys.path.insert(0, os.path.join(os.path.dirname(os.path.abspath(__file__)), "whisper_asr"))
  15. logger = logging.getLogger("dispatcher")
  16. LIVEKIT_URL = os.environ.get("LIVEKIT_URL", "ws://localhost:7880")
  17. PUBLIC_LIVEKIT_URL = os.environ.get("PUBLIC_LIVEKIT_URL", LIVEKIT_URL)
  18. API_KEY = "devkey"
  19. API_SECRET = "secretsecretsecretsecretsecret12"
  20. def _room_token(room_name: str, identity: str, *, can_create: bool = False) -> str:
  21. n = int(time.time())
  22. return jwt.encode({
  23. "iss": API_KEY, "sub": identity,
  24. "nbf": n - 60, "exp": n + 24 * 3600,
  25. "video": {
  26. "roomJoin": True, "roomCreate": can_create,
  27. "room": room_name,
  28. "canPublish": True, "canSubscribe": True, "canPublishData": True,
  29. },
  30. }, API_SECRET, algorithm="HS256")
  31. _workers: dict[str, asyncio.Task] = {}
  32. async def _spawn_worker(room_name: str):
  33. """Run a worker for one room. Exits when room is empty or cancelled."""
  34. from worker import Worker
  35. try:
  36. logger.info("Worker starting room=%s", room_name)
  37. w = Worker(room_name, identity="asr-bot")
  38. await w.run()
  39. except asyncio.CancelledError:
  40. pass
  41. except Exception:
  42. logger.exception("Worker %s failed", room_name)
  43. finally:
  44. logger.info("Worker stopped room=%s", room_name)
  45. async def handle_connect(request: web.Request) -> web.Response:
  46. """POST /connect — return tokens for a room and spawn worker.
  47. Body: {"identity": "user-xxx", "room": "room-xxx" (optional)}
  48. If room is provided and a worker already exists, reuse it.
  49. """
  50. try:
  51. body = await request.json()
  52. app_identity = body.get("identity", f"user-{int(time.time()*1000)%100000}")
  53. resume_room = body.get("room")
  54. except Exception:
  55. app_identity = f"user-{int(time.time()*1000)%100000}"
  56. resume_room = None
  57. if resume_room:
  58. room_name = resume_room
  59. if room_name in _workers and not _workers[room_name].done():
  60. logger.info("Session resume: %s (%s)", room_name, app_identity)
  61. else:
  62. # Recreate worker for the same room so App can reuse the room name
  63. # even after the previous worker exited.
  64. task = asyncio.create_task(_spawn_worker(room_name))
  65. _workers[room_name] = task
  66. logger.info("Room %s recreated for %s (total workers: %d)", room_name, app_identity, len(_workers))
  67. else:
  68. room_name = f"room-{int(time.time_ns())}"
  69. task = asyncio.create_task(_spawn_worker(room_name))
  70. _workers[room_name] = task
  71. logger.info("Room %s created for %s (total workers: %d)", room_name, app_identity, len(_workers))
  72. return web.json_response({
  73. "room": room_name,
  74. "url": PUBLIC_LIVEKIT_URL,
  75. "token": _room_token(room_name, app_identity, can_create=True),
  76. })
  77. async def handle_health(request: web.Request) -> web.Response:
  78. # Clean up finished workers
  79. for name, task in list(_workers.items()):
  80. if task.done():
  81. _workers.pop(name, None)
  82. return web.json_response({"status": "ok", "active_workers": len(_workers)})
  83. async def main():
  84. logging.basicConfig(level=logging.INFO)
  85. app = web.Application()
  86. app.router.add_post("/connect", handle_connect)
  87. app.router.add_get("/health", handle_health)
  88. runner = web.AppRunner(app)
  89. await runner.setup()
  90. site = web.TCPSite(runner, "0.0.0.0", 9100)
  91. await site.start()
  92. logger.info("Dispatcher listening on :9100")
  93. await asyncio.Future()
  94. if __name__ == "__main__":
  95. asyncio.run(main())