""" OpenCode `/event` 구독 버스 — 프로세스당 하나. OpenCode 는 서버 전체 이벤트를 SSE 하나로 흘림. 여기서 한 번만 구독하고 세션별 asyncio.Queue 로 나눠줌. 구독자 없으면 태스크도 안 돌림. 큐 item: dict(이벤트) 또는 None(하트비트 — 60초 무활동/재연결). """ import asyncio import json import logging from common.opencode_service import opencode_service log = logging.getLogger(__name__) def session_id_of(event: dict) -> str | None: props = event.get("properties") or {} info = props.get("info") or {} part = props.get("part") or {} if str(event.get("type", "")).startswith("session."): return info.get("id") or props.get("sessionID") return part.get("sessionID") or props.get("sessionID") or info.get("sessionID") class EventBus: def __init__(self, source=None): # source: async iterator factory — 테스트에서 가짜로 바꿈 self._source = source or opencode_service.iter_event_jsons self._subs: dict[str, set[asyncio.Queue]] = {} self._task: asyncio.Task | None = None def subscribe(self, session_id: str) -> asyncio.Queue: q: asyncio.Queue = asyncio.Queue() self._subs.setdefault(session_id, set()).add(q) self._ensure_running() return q def unsubscribe(self, session_id: str, q: asyncio.Queue) -> None: subs = self._subs.get(session_id) if subs: subs.discard(q) if not subs: del self._subs[session_id] def _ensure_running(self) -> None: if self._task is None or self._task.done(): self._task = asyncio.get_running_loop().create_task(self._pump(), name="opencode-event-pump") async def _pump(self) -> None: try: async for payload in self._source(): if payload is None: for subs in list(self._subs.values()): for q in list(subs): q.put_nowait(None) continue try: event = json.loads(payload) except ValueError: continue sid = session_id_of(event) if not sid: continue for q in list(self._subs.get(sid, ())): q.put_nowait(event) except asyncio.CancelledError: raise except Exception: # noqa: BLE001 log.exception("이벤트 펌프 죽음 — 다음 구독 때 재시작") finally: self._task = None event_bus = EventBus()