"""엔진 게이트웨이 HTTP 클라이언트. ⚠️ 게이트웨이 자체(apps/api/engine_gateway/)는 람다가 직접 만든다. 여기는 *호출부*만. 게이트웨이가 provider 라우팅(claude_cli/claude_api/codex_cli/agy_cli/openai/solar)· 모델 탐색·캐싱·상주 claude -p 풀을 흡수한다(마스터플랜 §0, R1). 이 백엔드는 ENGINE_URL로 HTTP 호출만 한다. 계약 (람다와 합의할 게이트웨이 API): POST {ENGINE_URL}/v1/generate — 단발 생성 (평가 deep-loop 등) POST {ENGINE_URL}/v1/stream — SSE 토큰 스트림 (내담자 AI 응답) GET {ENGINE_URL}/v1/capabilities — provider별 사용 가능 모델·추론 강도 GET {ENGINE_URL}/health 요청 바디는 3-AI 역할별 system 레이어(L0~L6, 설계서 §1.2)를 게이트웨이에 넘기되, text 는 *PII 마스킹 후(text_masked)* 만 보낸다 (R7/F-03, 마스킹은 호출 전 가드레일이 완료). """ from __future__ import annotations import asyncio from typing import Any, AsyncIterator, Optional import httpx from .config import settings from .contracts.engine_gateway import ( AIRole as AIRole, EngineCapabilitiesResponse, EngineMessage as EngineMessage, EngineGatewaySseLineDecoder, EngineGatewaySsePacket, EngineProvider, GenerateRequest, GenerateResponse, ReasoningEffort, StreamRequest, normalize_engine_gateway_model, ) class EngineError(RuntimeError): """게이트웨이 호출 실패. 라우트가 503/502 로 변환.""" class EngineClient: """ENGINE_URL 게이트웨이 비동기 클라이언트. 앱 수명주기 동안 1 인스턴스 재사용.""" def __init__( self, base_url: Optional[str] = None, *, shared_secret: Optional[str] = None, ) -> None: self.base_url = (base_url or settings.engine_url).rstrip("/") configured_secret = settings.engine_gateway_shared_secret.get_secret_value() self._shared_secret = ( configured_secret if shared_secret is None else shared_secret ).strip() self.engine_mode = settings.engine_mode self.live_client_provider: Optional[EngineProvider] = settings.live_client_provider self.default_model: Optional[str] = None self.default_reasoning_effort: Optional[ReasoningEffort] = None self._client: Optional[httpx.AsyncClient] = None self._lock = asyncio.Lock() async def startup(self) -> None: async with self._lock: if self._client is None: self._client = self._new_client() def _new_client(self) -> httpx.AsyncClient: return httpx.AsyncClient( base_url=self.base_url, headers=self._auth_headers(), timeout=httpx.Timeout( settings.engine_timeout, connect=settings.engine_connect_timeout, ), ) def _auth_headers(self) -> dict[str, str]: if not self._shared_secret: return {} return {"X-Vignette-Engine-Token": self._shared_secret} async def shutdown(self) -> None: async with self._lock: if self._client is not None: await self._client.aclose() self._client = None async def configure( self, *, base_url: str, engine_mode: EngineProvider, default_model: Optional[str] = None, default_reasoning_effort: Optional[ReasoningEffort] = None, ) -> None: next_url = base_url.rstrip("/") next_model = normalize_engine_gateway_model(default_model) async with self._lock: url_changed = next_url != self.base_url self.base_url = next_url self.engine_mode = engine_mode self.default_model = next_model self.default_reasoning_effort = default_reasoning_effort if self._client is not None and url_changed: old_client = self._client self._client = self._new_client() await old_client.aclose() def _payload(self, req: GenerateRequest) -> dict[str, Any]: payload = req.model_dump(exclude_none=True) provider = self.engine_mode if req.ai_role == "client" and req.session_id and self.live_client_provider: provider = self.live_client_provider if "provider" not in payload: payload["provider"] = provider # 관리자 기본 모델/추론 강도는 그 provider에 속한 값이다. 실시간 lane이 # 다른 provider면 잘못된 모델 slug를 넘기지 않고 해당 provider 기본값을 쓴다. same_provider = payload["provider"] == self.engine_mode if same_provider and self.default_model and "model" not in payload: payload["model"] = self.default_model if ( same_provider and self.default_reasoning_effort and "reasoning_effort" not in payload ): payload["reasoning_effort"] = self.default_reasoning_effort return payload @property def client(self) -> httpx.AsyncClient: if self._client is None: raise EngineError("EngineClient not started — call startup() in lifespan") return self._client async def health(self) -> bool: return bool((await self.health_detail()).get("ok")) async def health_detail(self) -> dict[str, Any]: try: params: dict[str, str] = {"provider": self.engine_mode} if self.default_model: params["model"] = self.default_model if self.default_reasoning_effort: params["reasoning_effort"] = self.default_reasoning_effort r = await self.client.get("/ready", params=params) if r.status_code == 404: live = await self.client.get("/health") return { "ok": live.status_code == 200, "detail": "gateway liveness only; readiness endpoint unavailable", "status_code": live.status_code, } payload: dict[str, Any] = {} try: payload = r.json() except ValueError: payload = {} return { "ok": r.status_code == 200 and bool(payload.get("ok", False)), "detail": str(payload.get("detail") or r.text or "engine readiness failed"), "status_code": r.status_code, "cached": bool(payload.get("cached", False)), } except httpx.HTTPError as exc: return { "ok": False, "detail": f"engine readiness transport error: {exc}", "status_code": None, } async def capabilities( self, *, provider: EngineProvider, base_url: str | None = None, force: bool = False, ) -> EngineCapabilitiesResponse: target_url = (base_url or self.base_url).rstrip("/") params = {"provider": provider, "force": str(force).lower()} try: if target_url == self.base_url: response = await self.client.get("/v1/capabilities", params=params) else: async with httpx.AsyncClient( base_url=target_url, headers=self._auth_headers(), timeout=httpx.Timeout(30, connect=settings.engine_connect_timeout), ) as client: response = await client.get("/v1/capabilities", params=params) response.raise_for_status() return EngineCapabilitiesResponse.model_validate(response.json()) except httpx.HTTPStatusError as exc: raise EngineError( f"engine capabilities {exc.response.status_code}: {exc.response.text}" ) from exc except (httpx.HTTPError, ValueError) as exc: raise EngineError(f"engine capabilities unavailable: {exc}") from exc async def generate( self, req: GenerateRequest, *, timeout: float | None = None, ) -> GenerateResponse: """단발 생성.""" try: kwargs: dict[str, Any] = {} if timeout is not None: kwargs["timeout"] = timeout r = await self.client.post( "/v1/generate", json=self._payload(req), **kwargs, ) r.raise_for_status() except httpx.HTTPStatusError as e: raise EngineError(f"engine generate {e.response.status_code}: {e.response.text}") from e except httpx.HTTPError as e: raise EngineError(f"engine generate transport error: {e}") from e return GenerateResponse.model_validate(r.json()) async def stream(self, req: StreamRequest) -> AsyncIterator[str]: """SSE 토큰 스트림 프록시. 게이트웨이 SSE(`text/event-stream`) 의 원시 non-empty line을 yield한다. 새 호출부는 raw line 대신 `stream_packets()`를 사용한다. """ try: async with self.client.stream( "POST", "/v1/stream", json=self._payload(req) ) as r: r.raise_for_status() async for line in r.aiter_lines(): if line: yield line except httpx.HTTPStatusError as e: raise EngineError(f"engine stream {e.response.status_code}") from e except httpx.HTTPError as e: raise EngineError(f"engine stream transport error: {e}") from e async def stream_packets(self, req: StreamRequest) -> AsyncIterator[EngineGatewaySsePacket]: """Decode gateway SSE into the shared token/done/error contract. Keep wire-format parsing at the gateway client boundary so app services do not depend on raw SSE line structure. A future Node.js gateway should only need to preserve `app.contracts.engine_gateway`. """ decoder = EngineGatewaySseLineDecoder() async for raw in self.stream(req): packet = decoder.feed_line(raw) if packet is not None: yield packet async def close_session(self, session_id: str) -> bool: """회기 종료 시 게이트웨이의 상주 페르소나 프로세스를 회수한다.""" if not session_id: return False try: response = await self.client.delete(f"/session/{session_id}") response.raise_for_status() return bool(response.json().get("closed")) except (httpx.HTTPError, ValueError): # DB 회기 종료 성공을 게이트웨이 정리 실패 때문에 되돌리지는 않는다. return False # 앱 전역 싱글톤 (main lifespan 에서 startup/shutdown) engine_client = EngineClient()