Files
openapi-first/openapi_first/templates/vet_app/sse.py

45 lines
1.3 KiB
Python

import asyncio
import random
import json
_sounds_by_species = {"dog": ["woof"], "cat": ["meow"], "bird": ["coo"]}
_subscribers: dict[int, list[asyncio.Queue]] = {}
_worker_tasks: dict[int, asyncio.Task] = {}
async def _sound_worker(pet_id: int, species: str):
sounds = _sounds_by_species.get(species, ["woof"])
while True:
sound = random.choice(sounds)
data = json.dumps({"sound": sound})
queues = _subscribers.get(pet_id, [])
for q in queues:
await q.put(data)
await asyncio.sleep(random.uniform(1, 5))
def _ensure_worker(pet_id: int, species: str):
if pet_id not in _worker_tasks or _worker_tasks[pet_id].done():
_subscribers.setdefault(pet_id, [])
_worker_tasks[pet_id] = asyncio.create_task(
_sound_worker(pet_id, species)
)
async def subscribe(pet_id: int, species: str) -> asyncio.Queue:
q: asyncio.Queue = asyncio.Queue()
_subscribers.setdefault(pet_id, []).append(q)
_ensure_worker(pet_id, species)
return q
def unsubscribe(pet_id: int, q: asyncio.Queue):
queues = _subscribers.get(pet_id, [])
if q in queues:
queues.remove(q)
if not queues:
task = _worker_tasks.pop(pet_id, None)
if task and not task.done():
task.cancel()
_subscribers.pop(pet_id, None)