vignette/apps/api/app/engine_client.py
2026-06-29 08:12:14 +09:00

177 lines
6.7 KiB
Python

"""엔진 게이트웨이 HTTP 클라이언트.
⚠️ 게이트웨이 자체(apps/api/engine_gateway/)는 람다가 직접 만든다. 여기는 *호출부*만.
게이트웨이가 provider 라우팅(claude_api/claude_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}/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,
EngineMessage,
EngineGatewaySseLineDecoder,
EngineGatewaySsePacket,
GenerateRequest,
GenerateResponse,
StreamRequest,
normalize_engine_gateway_model,
)
class EngineError(RuntimeError):
"""게이트웨이 호출 실패. 라우트가 503/502 로 변환."""
class EngineClient:
"""ENGINE_URL 게이트웨이 비동기 클라이언트. 앱 수명주기 동안 1 인스턴스 재사용."""
def __init__(self, base_url: Optional[str] = None) -> None:
self.base_url = (base_url or settings.engine_url).rstrip("/")
self.engine_mode = settings.engine_mode
self.default_model: Optional[str] = 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,
timeout=httpx.Timeout(
settings.engine_timeout,
connect=settings.engine_connect_timeout,
),
)
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: str,
default_model: Optional[str] = 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
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)
if self.default_model and "model" not in payload:
payload["model"] = self.default_model
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:
r = await self.client.get("/ready")
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 generate(self, req: GenerateRequest) -> GenerateResponse:
"""단발 생성."""
try:
r = await self.client.post("/v1/generate", json=self._payload(req))
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
# 앱 전역 싱글톤 (main lifespan 에서 startup/shutdown)
engine_client = EngineClient()