""" =================================================================== 모듈: common/opencode_service.py 설명: OpenCode 서버(opencode serve)와의 HTTP 통신 - JSON API (세션/메시지/프로바이더/에이전트) → 동기 httpx (DRF 뷰용) - SSE 이벤트 스트림 → 비동기 (SSE 뷰·파일 워처용) =================================================================== """ import asyncio import logging from typing import Any, AsyncIterator import httpx from django.conf import settings log = logging.getLogger(__name__) # 메시지 전송은 LLM 응답 완료까지 걸리므로 타임아웃을 길게 잡는다 (10분) LONG_TIMEOUT = httpx.Timeout(600.0, connect=10.0) DEFAULT_TIMEOUT = httpx.Timeout(30.0, connect=10.0) class OpencodeService: """OpenCode 서버 프록시.""" def __init__(self) -> None: self.base_url = settings.OPENCODE_BASE_URL # 모든 요청에 ?directory= 를 붙여 CodeAssist workspace 프로젝트로 고정 (1.18 멀티 프로젝트 서버) self.directory = getattr(settings, "OPENCODE_DIRECTORY", "") or None def _params(self) -> dict | None: return {"directory": self.directory} if self.directory else None def _request( self, method: str, path: str, json: Any | None = None, timeout: httpx.Timeout = DEFAULT_TIMEOUT, ) -> Any: with httpx.Client(base_url=self.base_url, timeout=timeout) as client: response = client.request(method, path, json=json, params=self._params()) response.raise_for_status() return response.json() if response.content else None # ── 세션 def list_sessions(self) -> Any: return self._request("GET", "/session") def create_session(self) -> Any: return self._request("POST", "/session", json={}) def delete_session(self, session_id: str) -> Any: return self._request("DELETE", f"/session/{session_id}") def rename_session(self, session_id: str, title: str) -> Any: return self._request("PATCH", f"/session/{session_id}", json={"title": title}) # ── 메시지 def list_messages(self, session_id: str) -> Any: return self._request("GET", f"/session/{session_id}/message") def abort_session(self, session_id: str) -> Any: log.info("응답 중단: session=%s", session_id) return self._request("POST", f"/session/{session_id}/abort") def send_message(self, session_id: str, payload: dict) -> Any: log.info( "메시지 전송: session=%s provider=%s model=%s", session_id, payload.get("providerID"), payload.get("modelID"), ) # OpenCode 1.18+ 는 최상위 providerID/modelID 를 무시한다(세션 기본 모델 사용). # model:{providerID,modelID} 오브젝트 형식으로 변환해서 보내야 선택이 먹는다. provider_id = payload.pop("providerID", None) model_id = payload.pop("modelID", None) if provider_id and model_id: payload["model"] = {"providerID": provider_id, "modelID": model_id} return self._request( "POST", f"/session/{session_id}/message", json=payload, timeout=LONG_TIMEOUT ) # ── 설정 def get_providers(self) -> Any: return self._request("GET", "/config/providers") def list_agents(self) -> Any: return self._request("GET", "/agent") # ── 비동기 버전 (스트림 어댑터용 — 이벤트 루프 안에서 블록 없이) async def _arequest(self, method: str, path: str, json: Any | None = None) -> Any: async with httpx.AsyncClient(base_url=self.base_url, timeout=DEFAULT_TIMEOUT) as client: response = await client.request(method, path, json=json, params=self._params()) response.raise_for_status() return response.json() if response.content else None async def prompt_async(self, session_id: str, payload: dict) -> None: """POST /session/{id}/prompt_async — 204 즉시 반환, 답변은 /event 로 흘러옴.""" await self._arequest("POST", f"/session/{session_id}/prompt_async", json=payload) async def get_session_a(self, session_id: str) -> Any: return await self._arequest("GET", f"/session/{session_id}") async def list_messages_a(self, session_id: str) -> Any: return await self._arequest("GET", f"/session/{session_id}/message") async def abort_session_a(self, session_id: str) -> Any: return await self._arequest("POST", f"/session/{session_id}/abort") # ── SSE 이벤트 스트림 (비동기) async def iter_event_jsons(self) -> AsyncIterator[str | None]: """SSE 프레임을 파싱해 이벤트 JSON 문자열 단위로 순회한다. 자가 치유: 업스트림이 60초간 조용하면 연결을 재수립한다 — 서버 재시작 등으로 반쯤 죽은 소켓이 남아도 스스로 복구된다. 무활동/재연결 시점에는 None(하트비트)을 내보낸다.""" stream_timeout = httpx.Timeout(None, connect=10.0) while True: try: async with httpx.AsyncClient( base_url=self.base_url, timeout=stream_timeout ) as client: async with client.stream("GET", "/event", params=self._params()) as response: response.raise_for_status() data_lines: list[str] = [] lines = response.aiter_lines() while True: try: line = await asyncio.wait_for(lines.__anext__(), timeout=60.0) except asyncio.TimeoutError: yield None # 60초 무활동 — 하트비트 후 재연결 break except StopAsyncIteration: break if line == "": if data_lines: yield "\n".join(data_lines) data_lines = [] elif line.startswith("data:"): data_lines.append(line[5:].lstrip()) # 그 외 필드(id:, event:, 코멘트)는 무시 except Exception: log.warning("OpenCode SSE 업스트림 끊김 — 2초 후 재연결") yield None await asyncio.sleep(2) opencode_service = OpencodeService() def upstream_error(e: Exception): """OpenCode 서버 에러 → DRF APIException (FE 는 detail 필드를 읽는다).""" from rest_framework.exceptions import APIException if isinstance(e, httpx.ConnectError): exc = APIException("OpenCode 서버에 연결할 수 없습니다. opencode serve 실행 여부를 확인하세요.") exc.status_code = 503 return exc if isinstance(e, httpx.HTTPStatusError): exc = APIException(f"OpenCode 서버 오류: {e.response.text[:500]}") exc.status_code = e.response.status_code return exc log.exception("OpenCode 중계 오류") exc = APIException(f"OpenCode 중계 오류: {e}") exc.status_code = 502 return exc