Files
CODE_ASSISTANT/5_django_backend/apps/chat/events.py
T
leehc991028andClaude Fable 5.1 2869985139 feat(backend): Django + OpenCode 백엔드 추가 — 5_django_backend/
앱이 기대하는 base-backend 계약(envelope·JWT·/chat/stream SSE)을 그대로 구현.
Django 는 로그인·세션 미러(SQLite/PG)·OpenCode 이벤트 번역만 맡고, 답변은 전용
OpenCode 인스턴스(opencode/ workspace, codeassist 에이전트)가 만듦.
- accounts: 이메일 로그인, access 60분 / refresh 14일, entra/config 는 501
- chat: 세션 목록·검색·메시지·취소 + POST /chat/stream 어댑터(part.delta→token,
  idle→usage/done, 클라 끊겨도 턴 감시 태스크가 DB 마무리, 첫 이벤트 60초 타임아웃)
- OpenCode 1.18 멀티 프로젝트라 모든 요청에 ?directory= 부착
- docs-lib 에 OpenCode SDK 1.18.6 타입 원본, pytest 32

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-09-16 21:28:31 +09:00

77 lines
2.6 KiB
Python

"""
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()