vignette/apps/api/engine_gateway/provider_registry.py
Yun Chan 45b84faa0d
Some checks failed
API contract / OpenAPI type drift (push) Failing after 1m0s
codex·agy OAuth 계정 연결 지원 (ChatGPT/Antigravity 네이티브 어댑터)
omniroute와 동일하게 공개 클라이언트 자격증명으로 서버사이드 OAuth 교환을
제공한다. 관리자가 제공자 로그인 후 브라우저에 남는 code를 붙여넣으면
토큰·refresh token·메타데이터(account-id, Code Assist project)를 저장하고
게이트웨이에 push한다.

- codex_api 어댑터: chatgpt.com/backend-api/codex/responses (Responses SSE,
  401 시 refresh token으로 자가 갱신)
- antigravity_api 어댑터: cloudcode-pa v1internal:streamGenerateContent
  (Gemini 형식 SSE, loadCodeAssist로 프로젝트 발급, 401 자가 갱신)
- 자격증명 저장소에 refresh_token_encrypted·extra 컬럼 추가(부트스트랩
  SQL 포함), 게이트웨이 push가 구조화 자격증명을 전달
2026-09-11 19:01:35 +09:00

1718 lines
64 KiB
Python

"""Provider 탐색과 Claude CLI 이외 실행 어댑터.
Provider별 CLI/API 세부 구현은 게이트웨이가 소유한다. 애플리케이션과 맞닿는
wire 계약은 ``app.contracts.engine_gateway``에 유지한다.
"""
from __future__ import annotations
import asyncio
import json
import os
import shutil
import tempfile
import time
from dataclasses import dataclass
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, AsyncIterator, Iterable, Literal, cast
import httpx
from app.contracts.engine_gateway import (
ENGINE_PROVIDER_DEFAULTS,
ENGINE_REASONING_EFFORTS,
EngineCapabilitiesResponse,
EngineModelOption,
EngineProvider,
GenerateRequest,
ReasoningEffort,
normalize_engine_gateway_model,
)
from app.services.llm_pricing import estimate_reference_cost
CODEX_DEFAULT_MODEL, CODEX_DEFAULT_EFFORT = ENGINE_PROVIDER_DEFAULTS["codex_cli"]
AGY_DEFAULT_MODEL, AGY_DEFAULT_EFFORT = ENGINE_PROVIDER_DEFAULTS["agy_cli"]
CLAUDE_CLI_DEFAULT_MODEL, CLAUDE_DEFAULT_EFFORT = ENGINE_PROVIDER_DEFAULTS[
"claude_cli"
]
CAPABILITY_CACHE_TTL_SECONDS = float(
os.environ.get("ENGINE_CAPABILITY_CACHE_TTL_SECONDS", "60")
)
CLI_TIMEOUT_SECONDS = float(os.environ.get("ENGINE_CLI_TIMEOUT_SECONDS", "300"))
ANTHROPIC_API_BASE = os.environ.get(
"ANTHROPIC_API_BASE", "https://api.anthropic.com"
).rstrip("/")
OPENAI_API_BASE = os.environ.get(
"OPENAI_BASE_URL", "https://api.openai.com/v1"
).rstrip("/")
OPENROUTER_API_BASE = os.environ.get(
"OPENROUTER_API_BASE", "https://openrouter.ai/api/v1"
).rstrip("/")
_OPENAI_DEFAULT_ENGINE_MODELS = (
"gpt-5.6-terra",
"gpt-5.6-luna",
"gpt-4.1",
)
class ProviderError(RuntimeError):
"""자격 증명을 노출하지 않고 provider 탐색·생성 실패를 전달한다."""
@dataclass(frozen=True, slots=True)
class ProviderGenerateResult:
text: str
model: str
provider: EngineProvider
tokens_in: int = 0
tokens_out: int = 0
cost_usd: float = 0.0
inference_geo: str | None = None
structured: dict[str, Any] | None = None
@dataclass(frozen=True, slots=True)
class ProviderStreamEvent:
type: Literal["delta", "done"]
text: str = ""
result: ProviderGenerateResult | None = None
_CAPABILITY_CACHE: dict[EngineProvider, tuple[float, EngineCapabilitiesResponse]] = {}
_CAPABILITY_LOCK = asyncio.Lock()
# 관리자 연결 UI가 push한 자격증명의 종류(api_key | oauth_token). 토큰 값 자체는
# os.environ으로만 들어가고 여기엔 어떤 헤더로 보낼지 판정하는 힌트만 둔다.
_PROVIDER_AUTH_KINDS: dict[str, str] = {}
# OAuth 연결(claude/codex/agy)의 구조화 자격증명. 값은 내부 push 엔드포인트로만
# 들어오고 재시작 시 초기화된다(재push는 boot_id 불일치 감지로 자동). 토큰 만료 시
# refresh_token으로 자가 갱신한다.
_PROVIDER_CREDENTIALS: dict[str, dict[str, Any]] = {}
def set_provider_credential(
provider: str,
*,
token: str,
auth_kind: str = "api_key",
refresh_token: str = "",
extra: dict[str, Any] | None = None,
) -> None:
_PROVIDER_CREDENTIALS[provider] = {
"token": token,
"auth_kind": auth_kind,
"refresh_token": refresh_token,
"extra": dict(extra or {}),
}
set_provider_auth_kind(provider, auth_kind)
def get_provider_credential(provider: str) -> dict[str, Any] | None:
return _PROVIDER_CREDENTIALS.get(provider)
def set_provider_auth_kind(provider: str, auth_kind: str) -> None:
if auth_kind in {"api_key", "oauth_token"}:
_PROVIDER_AUTH_KINDS[provider] = auth_kind
else:
_PROVIDER_AUTH_KINDS.pop(provider, None)
def get_provider_auth_kind(provider: str) -> str:
return _PROVIDER_AUTH_KINDS.get(provider, "api_key")
def _anthropic_headers(api_key: str) -> dict[str, str]:
"""Anthropic 인증 헤더. OAuth 액세스 토큰은 x-api-key 대신 Bearer를 쓴다."""
if get_provider_auth_kind("claude") == "oauth_token":
return {
"Authorization": f"Bearer {api_key}",
"anthropic-version": "2023-06-01",
}
return {
"x-api-key": api_key,
"anthropic-version": "2023-06-01",
}
# ── OAuth 계정 어댑터(codex_api / antigravity_api) 공용 유틸 ───────────
async def _refresh_oauth_token(provider: str) -> str | None:
"""401 시 저장된 refresh_token으로 액세스 토큰을 갱신하고 새 값을 돌려준다."""
cred = _PROVIDER_CREDENTIALS.get(provider)
if not cred or not cred.get("refresh_token"):
return None
if provider == "codex":
url = "https://auth.openai.com/oauth/token"
data = {
"grant_type": "refresh_token",
"refresh_token": cred["refresh_token"],
"client_id": "app_EMoamEEZ73f0CkXaXp7hrann",
}
else:
url = "https://oauth2.googleapis.com/token"
data = {
"grant_type": "refresh_token",
"refresh_token": cred["refresh_token"],
"client_id": "1071006060591-tmhssin2h21lcre235vtolojh4g403ep.apps.googleusercontent.com",
"client_secret": "GOCSPX-K58FWR486LdLJ1mLB8sXC4zqDAf",
}
try:
async with httpx.AsyncClient(timeout=30) as client:
response = await client.post(url, data=data)
response.raise_for_status()
payload = response.json()
except (httpx.HTTPError, ValueError):
return None
token = str(payload.get("access_token") or "").strip()
if token:
cred["token"] = token
rotated = str(payload.get("refresh_token") or "").strip()
if rotated:
cred["refresh_token"] = rotated
return token or None
def clear_capability_cache() -> None:
_CAPABILITY_CACHE.clear()
def _now() -> float:
return time.time()
def _utcnow() -> datetime:
return datetime.now(timezone.utc)
def _efforts(values: Iterable[str]) -> list[ReasoningEffort]:
allowed = set(ENGINE_REASONING_EFFORTS)
return [cast(ReasoningEffort, value) for value in values if value in allowed]
def _configured_openai_models() -> list[str]:
"""운영 엔진에 노출할 OpenAI 텍스트 모델을 명시 allowlist로 제한한다."""
raw = os.environ.get("OPENAI_ENGINE_MODELS", "").strip()
candidates = (
[item.strip() for item in raw.split(",")]
if raw
else list(_OPENAI_DEFAULT_ENGINE_MODELS)
)
configured_default = os.environ.get("OPENAI_ENGINE_MODEL", "").strip()
if configured_default:
candidates.insert(0, configured_default)
result: list[str] = []
for model in candidates:
if not model or model in result:
continue
if any(character.isspace() for character in model) or "/" in model:
raise ProviderError("OPENAI_ENGINE_MODELS에 유효하지 않은 모델 식별자가 있습니다.")
result.append(model)
if not result:
raise ProviderError("OPENAI_ENGINE_MODELS가 비어 있습니다.")
return result
def _openai_reasoning_efforts(model: str) -> list[ReasoningEffort]:
# API 모델 목록은 추론 강도 메타데이터를 제공하지 않는다. 운영 기본값으로 쓰는
# GPT-5 계열에는 여러 세대가 공통 지원하는 보수적 교집합만 노출한다.
if model.startswith("gpt-5"):
return ["low", "medium", "high"]
return []
def _binary(env_name: str, fallback: str) -> str | None:
configured = os.environ.get(env_name, "").strip()
if configured:
path = Path(configured)
return str(path) if path.exists() else shutil.which(configured)
if os.name == "nt":
shim = shutil.which(fallback)
if fallback == "codex" and shim:
npm_vendor_root = (
Path(shim).parent
/ "node_modules"
/ "@openai"
/ "codex"
/ "node_modules"
/ "@openai"
)
native_candidates = sorted(
npm_vendor_root.glob("codex-win32-*/vendor/*/bin/codex.exe")
)
if native_candidates:
return str(native_candidates[0])
executable = shutil.which(f"{fallback}.exe")
if executable:
return executable
if fallback == "agy":
# agy 인스톨러의 Windows 표준 위치(%LOCALAPPDATA%\agy\bin). PATH에
# 올라 있지 않은 머신에서도 게이트웨이가 공급자를 잃지 않게 한다.
local_app_data = os.environ.get("LOCALAPPDATA", "")
if local_app_data:
well_known = Path(local_app_data) / "agy" / "bin" / "agy.exe"
if well_known.exists():
return str(well_known)
return shim
return shutil.which(fallback)
def _safe_process_error(stderr: bytes, fallback: str) -> str:
detail = stderr.decode("utf-8", errors="replace").strip()
if not detail:
return fallback
return detail[-1200:]
# Windows 필수 변수의 표준 기본값 — 승격 체인의 psutil 환경 이식에서 유실될 수 있다.
_WINDOWS_ESSENTIAL_ENV_DEFAULTS = {
"SystemRoot": r"C:\Windows",
"SystemDrive": "C:",
"ComSpec": r"C:\Windows\system32\cmd.exe",
}
def _cli_subprocess_env() -> dict[str, str]:
"""CLI 자식 프로세스에 물려줄 환경.
공개 런타임 승격은 이전 프로세스 환경을 psutil로 통째로 이식하는데, 이 캡처에서
SystemRoot 같은 Windows 필수 변수가 유실되면 Go 계열 CLI(agy)가 시스템 인증서
풀·홈 디렉터리 해석에 실패하고 빈 모델 목록을 조용히 내놓는다(2026-08-18 실측).
상속 환경에서 빠진 필수 키만 기본값으로 채운다.
"""
env = dict(os.environ)
if os.name == "nt":
# os.environ은 Windows에서 키를 대문자로 정규화하므로 대소문자 무시 조회한다.
upper_names = {key.upper(): key for key in env}
for name, default in _WINDOWS_ESSENTIAL_ENV_DEFAULTS.items():
existing = upper_names.get(name.upper())
if existing is None or not env[existing]:
env[name] = default
return env
async def _run_process(
args: list[str],
*,
input_text: str | None = None,
cwd: str | None = None,
timeout: float = CLI_TIMEOUT_SECONDS,
) -> tuple[str, str]:
proc = await asyncio.create_subprocess_exec(
*args,
stdin=asyncio.subprocess.PIPE if input_text is not None else asyncio.subprocess.DEVNULL,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
cwd=cwd,
env=_cli_subprocess_env(),
)
try:
stdout, stderr = await asyncio.wait_for(
proc.communicate(
input_text.encode("utf-8") if input_text is not None else None
),
timeout=timeout,
)
except TimeoutError as exc:
proc.kill()
await proc.wait()
raise ProviderError(f"provider 명령이 {timeout:.0f}초 안에 끝나지 않았습니다.") from exc
if proc.returncode != 0:
raise ProviderError(
_safe_process_error(stderr, f"provider 명령 실패: 종료 코드 {proc.returncode}")
)
return (
stdout.decode("utf-8", errors="replace"),
stderr.decode("utf-8", errors="replace"),
)
def _unavailable(provider: EngineProvider, detail: str) -> EngineCapabilitiesResponse:
return EngineCapabilitiesResponse(
provider=provider,
available=False,
source="unavailable",
detail=detail,
fetched_at=_now(),
)
def _display_model_name(model_id: str) -> str:
parts = model_id.split("-")
effort = parts[-1] if parts and parts[-1] in {"low", "medium", "high"} else None
if effort:
parts = parts[:-1]
words: list[str] = []
for part in parts:
if part.lower() in {"gpt", "oss"}:
words.append(part.upper())
elif any(char.isdigit() for char in part):
words.append(part)
else:
words.append(part.capitalize())
label = " ".join(words)
return f"{label} ({effort.capitalize()})" if effort else label
async def _discover_claude_cli() -> EngineCapabilitiesResponse:
if _binary("CLAUDE_BIN", "claude") is None:
return _unavailable("claude_cli", "Claude CLI를 찾을 수 없습니다.")
efforts = _efforts(("low", "medium", "high", "xhigh", "max"))
models = [
EngineModelOption(
id=CLAUDE_CLI_DEFAULT_MODEL,
label="Claude CLI 기본 모델",
description="로그인된 Claude CLI가 권장하는 기본 모델을 사용합니다.",
reasoning_efforts=efforts,
default_reasoning_effort=CLAUDE_DEFAULT_EFFORT,
is_default=True,
),
*[
EngineModelOption(
id=model,
label=f"Claude {model.capitalize()} 최신",
description="Claude CLI가 제공하는 안정 alias입니다.",
reasoning_efforts=efforts,
default_reasoning_effort=CLAUDE_DEFAULT_EFFORT,
)
for model in ("opus", "sonnet", "fable")
],
]
return EngineCapabilitiesResponse(
provider="claude_cli",
available=True,
source="static_cli",
models=models,
default_model=CLAUDE_CLI_DEFAULT_MODEL,
default_reasoning_effort=CLAUDE_DEFAULT_EFFORT,
detail="Claude CLI는 모델 목록 명령이 없어 공식 alias를 사용합니다.",
fetched_at=_now(),
)
async def _codex_model_list(binary: str) -> dict[str, Any]:
proc = await asyncio.create_subprocess_exec(
binary,
"app-server",
stdin=asyncio.subprocess.PIPE,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
env=_cli_subprocess_env(),
)
if proc.stdin is None or proc.stdout is None:
proc.kill()
await proc.wait()
raise ProviderError("Codex app-server stdio를 열 수 없습니다.")
messages = (
{
"method": "initialize",
"id": 0,
"params": {
"clientInfo": {
"name": "vignette_engine_gateway",
"title": "Vignette Engine Gateway",
"version": "1.0.0",
}
},
},
{"method": "initialized", "params": {}},
{
"method": "model/list",
"id": 6,
"params": {"limit": 100, "includeHidden": False},
},
)
for message in messages:
proc.stdin.write((json.dumps(message) + "\n").encode("utf-8"))
await proc.stdin.drain()
try:
while True:
raw = await asyncio.wait_for(proc.stdout.readline(), timeout=20)
if not raw:
raise ProviderError("Codex model/list 응답이 비어 있습니다.")
try:
message = json.loads(raw)
except json.JSONDecodeError:
continue
if message.get("id") == 6:
if message.get("error"):
raise ProviderError(str(message["error"].get("message") or message["error"]))
return cast(dict[str, Any], message.get("result") or {})
except TimeoutError as exc:
raise ProviderError("Codex model/list 응답 시간이 초과됐습니다.") from exc
finally:
if proc.stdin is not None and not proc.stdin.is_closing():
proc.stdin.close()
if proc.returncode is None:
try:
await asyncio.wait_for(proc.wait(), timeout=2)
except TimeoutError:
proc.kill()
await proc.wait()
async def _discover_codex_cli() -> EngineCapabilitiesResponse:
binary = _binary("CODEX_BIN", "codex")
if binary is None:
return _unavailable("codex_cli", "Codex CLI를 찾을 수 없습니다.")
try:
payload = await _codex_model_list(binary)
except (OSError, ProviderError) as exc:
return _unavailable("codex_cli", f"Codex 모델 조회 실패: {exc}")
raw_models = payload.get("data") if isinstance(payload, dict) else []
models: list[EngineModelOption] = []
for item in raw_models if isinstance(raw_models, list) else []:
if not isinstance(item, dict) or item.get("hidden"):
continue
model_id = str(item.get("model") or item.get("id") or "").strip()
if not model_id:
continue
supported = item.get("supportedReasoningEfforts") or []
efforts = _efforts(
str(entry.get("reasoningEffort") or "")
for entry in supported
if isinstance(entry, dict)
)
raw_default = str(item.get("defaultReasoningEffort") or "")
default_effort = (
cast(ReasoningEffort, raw_default)
if raw_default in efforts
else (efforts[0] if efforts else None)
)
models.append(
EngineModelOption(
id=model_id,
label=str(item.get("displayName") or model_id),
description=str(item.get("description") or ""),
reasoning_efforts=efforts,
default_reasoning_effort=default_effort,
is_default=model_id == CODEX_DEFAULT_MODEL,
)
)
if not models:
return _unavailable("codex_cli", "Codex가 선택 가능한 모델을 반환하지 않았습니다.")
default_model = (
CODEX_DEFAULT_MODEL
if any(model.id == CODEX_DEFAULT_MODEL for model in models)
else next((model.id for model in models if model.is_default), models[0].id)
)
selected = next(model for model in models if model.id == default_model)
default_effort = (
CODEX_DEFAULT_EFFORT
if CODEX_DEFAULT_EFFORT in selected.reasoning_efforts
else selected.default_reasoning_effort
)
return EngineCapabilitiesResponse(
provider="codex_cli",
available=True,
source="live_cli",
models=models,
default_model=default_model,
default_reasoning_effort=default_effort,
detail="Codex app-server model/list에서 실시간 조회했습니다.",
fetched_at=_now(),
)
async def _discover_agy_cli() -> EngineCapabilitiesResponse:
binary = _binary("AGY_BIN", "agy")
if binary is None:
return _unavailable("agy_cli", "Agy CLI를 찾을 수 없습니다.")
try:
stdout, _ = await _run_process([binary, "models"], timeout=30)
except (OSError, ProviderError) as exc:
return _unavailable("agy_cli", f"Agy 모델 조회 실패: {exc}")
models: list[EngineModelOption] = []
for line in stdout.splitlines():
# agy CLI는 2026-08부터 "model_id\t표시 라벨" 형태로 출력한다. 탭 왼쪽 토큰이
# 모델 id고 라벨은 CLI가 준 값을 우선한다. 탭 없이 공백이 섞인 줄은 상태/안내
# 문구이므로 건너뛴다(구형 베어 id 출력 계약은 그대로 유지).
label: str | None = None
if "\t" in line:
model_id, _, rest = line.partition("\t")
model_id = model_id.strip()
label = rest.strip() or None
if not model_id:
continue
else:
model_id = line.strip()
if not model_id or any(char.isspace() for char in model_id):
continue
suffix = model_id.rsplit("-", 1)[-1]
if suffix in {"low", "medium", "high"}:
efforts = _efforts((suffix,))
default_effort = cast(ReasoningEffort, suffix)
else:
efforts = _efforts(("low", "medium", "high"))
default_effort = AGY_DEFAULT_EFFORT if model_id == AGY_DEFAULT_MODEL else "medium"
models.append(
EngineModelOption(
id=model_id,
label=label or _display_model_name(model_id),
description="Agy CLI가 현재 계정에 노출한 모델입니다.",
reasoning_efforts=efforts,
default_reasoning_effort=default_effort,
is_default=model_id == AGY_DEFAULT_MODEL,
)
)
if not models:
return _unavailable("agy_cli", "Agy가 선택 가능한 모델을 반환하지 않았습니다.")
default_model = (
AGY_DEFAULT_MODEL
if any(model.id == AGY_DEFAULT_MODEL for model in models)
else models[0].id
)
selected = next(model for model in models if model.id == default_model)
return EngineCapabilitiesResponse(
provider="agy_cli",
available=True,
source="live_cli",
models=models,
default_model=default_model,
default_reasoning_effort=selected.default_reasoning_effort,
detail="agy models에서 실시간 조회했습니다.",
fetched_at=_now(),
)
async def _discover_claude_api() -> EngineCapabilitiesResponse:
api_key = os.environ.get("ANTHROPIC_API_KEY", "").strip() or os.environ.get(
"ANTHROPIC_AUTH_TOKEN", ""
).strip()
if not api_key:
return _unavailable("claude_api", "ANTHROPIC_API_KEY가 설정되지 않았습니다.")
try:
async with httpx.AsyncClient(timeout=20) as client:
response = await client.get(
f"{ANTHROPIC_API_BASE}/v1/models",
params={"limit": 100},
headers={**_anthropic_headers(api_key)},
)
response.raise_for_status()
payload = response.json()
except (httpx.HTTPError, ValueError) as exc:
return _unavailable("claude_api", f"Anthropic 모델 조회 실패: {exc}")
models: list[EngineModelOption] = []
for item in payload.get("data", []) if isinstance(payload, dict) else []:
if not isinstance(item, dict):
continue
model_id = str(item.get("id") or "").strip()
if not model_id:
continue
effort_capability = (item.get("capabilities") or {}).get("effort") or {}
efforts = _efforts(
effort
for effort in ENGINE_REASONING_EFFORTS
if isinstance(effort_capability.get(effort), dict)
and effort_capability[effort].get("supported")
)
default_effort: ReasoningEffort | None = (
CLAUDE_DEFAULT_EFFORT
if CLAUDE_DEFAULT_EFFORT in efforts
else (efforts[0] if efforts else None)
)
models.append(
EngineModelOption(
id=model_id,
label=str(item.get("display_name") or model_id),
description="Anthropic Models API가 현재 키에 노출한 모델입니다.",
reasoning_efforts=efforts,
default_reasoning_effort=default_effort,
)
)
if not models:
return _unavailable("claude_api", "Anthropic이 선택 가능한 모델을 반환하지 않았습니다.")
configured_default = os.environ.get("ANTHROPIC_MODEL", "").strip()
default_model = (
configured_default
if configured_default and any(model.id == configured_default for model in models)
else models[0].id
)
selected = next(model for model in models if model.id == default_model)
selected.is_default = True
return EngineCapabilitiesResponse(
provider="claude_api",
available=True,
source="live_api",
models=models,
default_model=default_model,
default_reasoning_effort=selected.default_reasoning_effort,
detail="Anthropic /v1/models에서 실시간 조회했습니다.",
fetched_at=_now(),
)
async def _discover_openai() -> EngineCapabilitiesResponse:
api_key = os.environ.get("OPENAI_API_KEY", "").strip()
if not api_key:
return _unavailable("openai", "OPENAI_API_KEY가 설정되지 않았습니다.")
try:
allowed_models = _configured_openai_models()
async with httpx.AsyncClient(timeout=20) as client:
response = await client.get(
f"{OPENAI_API_BASE}/models",
headers={"Authorization": f"Bearer {api_key}"},
)
response.raise_for_status()
payload = response.json()
except ProviderError as exc:
return _unavailable("openai", str(exc))
except (httpx.HTTPError, ValueError) as exc:
return _unavailable("openai", f"OpenAI 모델 조회 실패: {exc}")
live_ids = {
str(item.get("id") or "").strip()
for item in payload.get("data", [])
if isinstance(item, dict)
} if isinstance(payload, dict) else set()
available_ids = [model for model in allowed_models if model in live_ids]
if not available_ids:
return _unavailable(
"openai",
"OpenAI가 allowlist의 텍스트 모델을 반환하지 않았습니다.",
)
configured_default = os.environ.get("OPENAI_ENGINE_MODEL", "").strip()
default_model = (
configured_default
if configured_default in available_ids
else available_ids[0]
)
models: list[EngineModelOption] = []
for model_id in available_ids:
efforts = _openai_reasoning_efforts(model_id)
default_effort: ReasoningEffort | None = (
"medium" if "medium" in efforts else None
)
models.append(
EngineModelOption(
id=model_id,
label=model_id,
description="OpenAI Models API와 운영 allowlist가 함께 허용한 모델입니다.",
reasoning_efforts=efforts,
default_reasoning_effort=default_effort,
is_default=model_id == default_model,
)
)
selected = next(model for model in models if model.id == default_model)
return EngineCapabilitiesResponse(
provider="openai",
available=True,
source="live_api",
models=models,
default_model=default_model,
default_reasoning_effort=selected.default_reasoning_effort,
detail="OpenAI /v1/models와 운영 allowlist를 교차 확인했습니다.",
fetched_at=_now(),
)
def _configured_openrouter_models() -> list[str]:
"""OpenRouter allowlist. 미설정이면 계정이 노출한 전체 모델을 그대로 쓴다."""
raw = os.environ.get("OPENROUTER_ENGINE_MODELS", "").strip()
if not raw:
return []
result: list[str] = []
for model in (item.strip() for item in raw.split(",")):
if not model or model in result:
continue
if any(character.isspace() for character in model) or "/" not in model:
raise ProviderError(
"OPENROUTER_ENGINE_MODELS의 모델 식별자는 vendor/model 형태여야 합니다."
)
result.append(model)
return result
async def _discover_openrouter() -> EngineCapabilitiesResponse:
api_key = os.environ.get("OPENROUTER_API_KEY", "").strip()
if not api_key:
return _unavailable("openrouter", "OPENROUTER_API_KEY가 설정되지 않았습니다.")
try:
allowed_models = _configured_openrouter_models()
async with httpx.AsyncClient(timeout=20) as client:
response = await client.get(
f"{OPENROUTER_API_BASE}/models",
headers={"Authorization": f"Bearer {api_key}"},
)
response.raise_for_status()
payload = response.json()
except ProviderError as exc:
return _unavailable("openrouter", str(exc))
except (httpx.HTTPError, ValueError) as exc:
return _unavailable("openrouter", f"OpenRouter 모델 조회 실패: {exc}")
data = payload.get("data", []) if isinstance(payload, dict) else []
live: dict[str, dict[str, Any]] = {}
for item in data if isinstance(data, list) else []:
if not isinstance(item, dict):
continue
model_id = str(item.get("id") or "").strip()
if model_id:
live[model_id] = item
if allowed_models:
selected_ids = [model for model in allowed_models if model in live]
else:
selected_ids = list(live)
if not selected_ids:
return _unavailable(
"openrouter",
"OpenRouter가 사용 가능한 모델을 반환하지 않았습니다.",
)
configured_default = os.environ.get("OPENROUTER_ENGINE_MODEL", "").strip()
default_model = (
configured_default
if configured_default in selected_ids
else selected_ids[0]
)
models: list[EngineModelOption] = []
for model_id in selected_ids:
entry = live[model_id]
pricing = entry.get("pricing") or {}
is_free = str(pricing.get("prompt") or "") == "0" and str(
pricing.get("completion") or ""
) == "0"
models.append(
EngineModelOption(
id=model_id,
label=str(entry.get("name") or model_id),
description=(
"OpenRouter 무료 모델입니다."
if is_free
else "OpenRouter 계정에서 현재 사용할 수 있는 모델입니다."
),
reasoning_efforts=[],
default_reasoning_effort=None,
is_default=model_id == default_model,
)
)
return EngineCapabilitiesResponse(
provider="openrouter",
available=True,
source="live_api",
models=models,
default_model=default_model,
default_reasoning_effort=None,
detail="OpenRouter /models에서 실시간 조회했습니다.",
fetched_at=_now(),
)
async def _generate_openrouter(
req: GenerateRequest,
system_prompt: str,
user_payload: str,
) -> ProviderGenerateResult:
pricing_started_at = _utcnow()
api_key = os.environ.get("OPENROUTER_API_KEY", "").strip()
if not api_key:
raise ProviderError("OPENROUTER_API_KEY가 설정되지 않았습니다.")
model, effort = await _resolve_selection(req, "openrouter")
if req.ai_role == "client":
messages: list[dict[str, str]] = [
{"role": "user", "content": user_payload}
]
else:
messages = [
{"role": message.role, "content": message.content}
for message in req.messages
if message.role != "system"
]
if system_prompt:
messages = [{"role": "system", "content": system_prompt}, *messages]
payload: dict[str, Any] = {
"model": model,
"messages": messages,
"max_tokens": req.max_tokens,
"temperature": req.temperature,
}
if effort:
payload["reasoning"] = {"effort": effort}
try:
async with httpx.AsyncClient(timeout=CLI_TIMEOUT_SECONDS) as client:
response = await client.post(
f"{OPENROUTER_API_BASE}/chat/completions",
headers={
"Authorization": f"Bearer {api_key}",
"Content-Type": "application/json",
},
json=payload,
)
response.raise_for_status()
body = response.json()
except (httpx.HTTPError, ValueError) as exc:
raise ProviderError(f"OpenRouter Chat Completions 호출 실패: {exc}") from exc
choices = body.get("choices", []) if isinstance(body, dict) else []
text = ""
if isinstance(choices, list) and choices and isinstance(choices[0], dict):
message = choices[0].get("message") or {}
if isinstance(message, dict):
text = str(message.get("content") or "").strip()
if not text:
raise ProviderError("OpenRouter가 텍스트 응답을 반환하지 않았습니다.")
usage = body.get("usage") or {}
detail = usage.get("prompt_tokens_details") or {}
estimate = estimate_reference_cost(
provider="openrouter",
model=str(body.get("model") or model),
tokens_in=int(usage.get("prompt_tokens") or 0),
tokens_out=int(usage.get("completion_tokens") or 0),
priced_at=pricing_started_at,
cached_input_tokens=int(detail.get("cached_tokens") or 0),
)
return ProviderGenerateResult(
text=text,
model=str(body.get("model") or model),
provider="openrouter",
tokens_in=int(usage.get("prompt_tokens") or 0),
tokens_out=int(usage.get("completion_tokens") or 0),
cost_usd=estimate.cost_usd if estimate is not None else 0.0,
structured=_structured_or_none(text, req),
)
CODEX_API_BASE = "https://chatgpt.com/backend-api/codex"
ANTIGRAVITY_API_BASE = "https://cloudcode-pa.googleapis.com"
def _configured_codex_api_models() -> list[str]:
raw = os.environ.get("CODEX_API_ENGINE_MODELS", "").strip()
models = (
[item.strip() for item in raw.split(",") if item.strip()]
if raw
else ["gpt-5.1-codex", "gpt-5.1-codex-mini", "gpt-5.1", "gpt-5.1-mini"]
)
result: list[str] = []
for model in models:
if model not in result:
result.append(model)
return result
def _configured_antigravity_models() -> list[str]:
raw = os.environ.get("AGY_API_ENGINE_MODELS", "").strip()
models = (
[item.strip() for item in raw.split(",") if item.strip()]
if raw
else ["gemini-3.6-flash-high", "gemini-3.6-flash-low", "gemini-3.6-pro-high"]
)
result: list[str] = []
for model in models:
if model not in result:
result.append(model)
return result
async def _discover_codex_api() -> EngineCapabilitiesResponse:
if get_provider_credential("codex") is None:
return _unavailable("codex_api", "ChatGPT 계정 OAuth 연결이 없습니다.")
models: list[EngineModelOption] = []
allowed = _configured_codex_api_models()
default_model = allowed[0] if allowed else ""
for model in allowed:
efforts = _openai_reasoning_efforts(model)
models.append(
EngineModelOption(
id=model,
label=model,
description="관리자가 연결한 ChatGPT 계정(Codex 백엔드)에서 사용 가능한 모델입니다.",
reasoning_efforts=efforts,
default_reasoning_effort="medium" if "medium" in efforts else None,
is_default=model == default_model,
)
)
if not models:
return _unavailable("codex_api", "CODEX_API_ENGINE_MODELS 설정이 비어 있습니다.")
return EngineCapabilitiesResponse(
provider="codex_api",
available=True,
source="static_cli",
models=models,
default_model=default_model,
default_reasoning_effort=None,
detail="ChatGPT 계정 연결이 확인되어 정적 모델 목록을 제공합니다.",
fetched_at=_now(),
)
def _antigravity_headers(access_token: str) -> dict[str, str]:
return {
"Authorization": f"Bearer {access_token}",
"Content-Type": "application/json",
"User-Agent": "antigravity/cli/1.0.0 (aidev_client; os_type=windows; arch=amd64; auth_method=consumer)",
"Accept": "text/event-stream",
}
def _split_effort_suffix(model_id: str) -> tuple[str, ReasoningEffort | None]:
suffix = model_id.rsplit("-", 1)[-1]
if suffix in {"low", "medium", "high"}:
return model_id[: -(len(suffix) + 1)], cast(ReasoningEffort, suffix)
return model_id, None
async def _discover_antigravity_api() -> EngineCapabilitiesResponse:
if get_provider_credential("agy") is None:
return _unavailable("antigravity_api", "Agy(Google 계정) OAuth 연결이 없습니다.")
models: list[EngineModelOption] = []
allowed = _configured_antigravity_models()
default_model = allowed[0] if allowed else ""
for model in allowed:
base_model, effort = _split_effort_suffix(model)
models.append(
EngineModelOption(
id=model,
label=f"{base_model} ({effort.capitalize()})" if effort else base_model,
description="관리자가 연결한 Google 계정(Antigravity Code Assist) 모델입니다.",
reasoning_efforts=_efforts(("low", "medium", "high")),
default_reasoning_effort=effort,
is_default=model == default_model,
)
)
if not models:
return _unavailable("antigravity_api", "AGY_API_ENGINE_MODELS 설정이 비어 있습니다.")
return EngineCapabilitiesResponse(
provider="antigravity_api",
available=True,
source="static_cli",
models=models,
default_model=default_model,
default_reasoning_effort=None,
detail="Google 계정 연결이 확인되어 정적 모델 목록을 제공합니다.",
fetched_at=_now(),
)
async def _generate_codex_api(
req: GenerateRequest,
system_prompt: str,
user_payload: str,
) -> ProviderGenerateResult:
"""ChatGPT 계정(codex 백엔드)으로 Responses API 형식을 SSE로 소비한다."""
pricing_started_at = _utcnow()
cred = get_provider_credential("codex")
if cred is None:
raise ProviderError("ChatGPT 계정 OAuth 연결이 없습니다.")
model, effort = await _resolve_selection(req, "codex_api")
if req.ai_role == "client":
input_items: list[dict[str, str]] = [{"role": "user", "content": user_payload}]
else:
input_items = [
{"role": message.role, "content": message.content}
for message in req.messages
if message.role != "system"
]
payload: dict[str, Any] = {
"model": model,
"instructions": system_prompt or None,
"input": input_items,
"stream": True,
"store": False,
}
if effort:
payload["reasoning"] = {"effort": effort}
headers = {
"Authorization": f"Bearer {cred['token']}",
"Content-Type": "application/json",
"Accept": "text/event-stream",
"Openai-Beta": "responses=experimental",
"originator": "codex_cli_rs",
"User-Agent": "codex-cli/0.153.2 (Windows; amd64)",
}
account_id = (cred.get("extra") or {}).get("chatgpt_account_id")
if account_id:
headers["chatgpt-account-id"] = str(account_id)
if req.session_id:
headers["session_id"] = req.session_id
text, usage = await _post_sse_json(
f"{CODEX_API_BASE}/responses",
headers=headers,
json_body=payload,
provider="codex_api",
refresh_token_provider="codex",
refresh_headers=lambda token: {
**headers,
"Authorization": f"Bearer {token}",
},
extract=(lambda data: str(data.get("delta") or "") if data.get("type") == "response.output_text.delta" else None),
usage_extractor=(lambda data: (data.get("response") or {}).get("usage") if data.get("type") == "response.completed" else None),
)
if not text.strip():
raise ProviderError("Codex 백엔드가 텍스트 응답을 반환하지 않았습니다.")
tokens_in = int(usage.get("input_tokens") or 0)
tokens_out = int(usage.get("output_tokens") or 0)
estimate = estimate_reference_cost(
provider="openai",
model=model,
tokens_in=tokens_in,
tokens_out=tokens_out,
priced_at=pricing_started_at,
)
return ProviderGenerateResult(
text=text,
model=model,
provider="codex_api",
tokens_in=tokens_in,
tokens_out=tokens_out,
cost_usd=estimate.cost_usd if estimate is not None else 0.0,
structured=_structured_or_none(text, req),
)
async def _generate_antigravity_api(
req: GenerateRequest,
system_prompt: str,
user_payload: str,
) -> ProviderGenerateResult:
"""Google 계정(Antigravity Code Assist)으로 Gemini 형식을 SSE로 소비한다."""
pricing_started_at = _utcnow()
cred = get_provider_credential("agy")
if cred is None:
raise ProviderError("Agy(Google 계정) OAuth 연결이 없습니다.")
selected_model, effort = await _resolve_selection(req, "antigravity_api")
base_model, _suffix_effort = _split_effort_suffix(selected_model)
if req.ai_role == "client":
contents = [{"role": "user", "parts": [{"text": user_payload}]}]
else:
contents = [
{
"role": message.role,
"parts": [{"text": message.content}],
}
for message in req.messages
if message.role != "system"
]
project_id = str((cred.get("extra") or {}).get("project_id") or "")
payload: dict[str, Any] = {
"model": base_model,
"project": project_id,
"requestType": "agent",
"userAgent": "antigravity",
"contents": contents,
"enabledCreditTypes": ["GOOGLE_ONE_AI"],
}
if system_prompt:
payload["systemInstruction"] = {"parts": [{"text": system_prompt}]}
if effort:
payload["generationConfig"] = {"effort": effort}
def usage_extractor(data: dict[str, Any]) -> dict[str, Any] | None:
return data.get("usageMetadata") or None
text, usage = await _post_sse_json(
f"{ANTIGRAVITY_API_BASE}/v1internal:streamGenerateContent?alt=sse",
headers=_antigravity_headers(cred["token"]),
json_body=payload,
provider="antigravity_api",
refresh_token_provider="agy",
refresh_headers=lambda token: _antigravity_headers(token),
extract=lambda data: "".join(
str(part.get("text") or "")
for part in ((data.get("candidates") or [{}])[0].get("content") or {}).get("parts") or []
if isinstance(part, dict)
) or None,
usage_extractor=usage_extractor,
)
if not text.strip():
raise ProviderError("Antigravity 백엔드가 텍스트 응답을 반환하지 않았습니다.")
tokens_in = int(usage.get("promptTokenCount") or 0)
tokens_out = int(usage.get("candidatesTokenCount") or 0)
return ProviderGenerateResult(
text=text,
model=base_model,
provider="antigravity_api",
tokens_in=tokens_in,
tokens_out=tokens_out,
cost_usd=0.0,
structured=_structured_or_none(text, req),
)
async def _post_sse_json(
url: str,
*,
headers: dict[str, str],
json_body: dict[str, Any],
provider: str,
refresh_token_provider: str,
refresh_headers,
extract,
usage_extractor,
) -> tuple[str, dict[str, Any]]:
"""SSE(JSON 이벤트)를 소비해 누적 텍스트와 마지막 usage를 돌려준다. 401은 한 번 갱신해 재시도한다."""
text_parts: list[str] = []
usage: dict[str, Any] = {}
for attempt in range(2):
try:
async with httpx.AsyncClient(timeout=CLI_TIMEOUT_SECONDS) as client:
async with client.stream("POST", url, headers=headers, json=json_body) as response:
if response.status_code == 401 and attempt == 0:
refreshed = await _refresh_oauth_token(refresh_token_provider)
if refreshed:
headers = refresh_headers(refreshed)
continue
raise ProviderError(f"{provider} 인증이 만료되었습니다. 관리자 화면에서 다시 로그인해 주세요.")
response.raise_for_status()
async for line in response.aiter_lines():
if not line.startswith("data:"):
continue
raw = line[len("data:"):].strip()
if not raw or raw == "[DONE]":
continue
try:
data = json.loads(raw)
except json.JSONDecodeError:
continue
if not isinstance(data, dict):
continue
delta = extract(data)
if delta is not None:
text_parts.append(delta)
extracted_usage = usage_extractor(data)
if extracted_usage:
usage = extracted_usage
break
except httpx.HTTPStatusError as exc:
raise ProviderError(f"{provider} 백엔드 호출 실패: HTTP {exc.response.status_code}") from exc
except httpx.HTTPError as exc:
raise ProviderError(f"{provider} 백엔드 전송 실패: {exc}") from exc
return "".join(text_parts), usage
async def _discover(provider: EngineProvider) -> EngineCapabilitiesResponse:
if provider == "claude_cli":
return await _discover_claude_cli()
if provider == "claude_api":
return await _discover_claude_api()
if provider == "openai":
return await _discover_openai()
if provider == "openrouter":
return await _discover_openrouter()
if provider == "codex_cli":
return await _discover_codex_cli()
if provider == "codex_api":
return await _discover_codex_api()
if provider == "agy_cli":
return await _discover_agy_cli()
if provider == "antigravity_api":
return await _discover_antigravity_api()
return _unavailable(provider, f"{provider} 어댑터는 아직 모델 탐색을 지원하지 않습니다.")
async def discover_capabilities(
provider: EngineProvider, *, force: bool = False
) -> EngineCapabilitiesResponse:
cached = _CAPABILITY_CACHE.get(provider)
if (
not force
and cached is not None
and time.monotonic() - cached[0] < CAPABILITY_CACHE_TTL_SECONDS
):
return cached[1].model_copy(deep=True)
async with _CAPABILITY_LOCK:
cached = _CAPABILITY_CACHE.get(provider)
if (
not force
and cached is not None
and time.monotonic() - cached[0] < CAPABILITY_CACHE_TTL_SECONDS
):
return cached[1].model_copy(deep=True)
result = await _discover(provider)
_CAPABILITY_CACHE[provider] = (time.monotonic(), result)
return result.model_copy(deep=True)
def _cli_prompt(system_prompt: str, user_payload: str) -> str:
parts = []
if system_prompt.strip():
parts.append("[시스템 지침]\n" + system_prompt.strip())
parts.append("[응답할 입력]\n" + user_payload.strip())
return "\n\n".join(parts)
def _cli_runtime_cwd() -> Path:
path = Path(
os.environ.get(
"ENGINE_CLI_CWD",
str(Path(tempfile.gettempdir()) / "vignette-engine-runtime"),
)
)
path.mkdir(parents=True, exist_ok=True)
return path
async def _resolve_selection(
req: GenerateRequest, provider: EngineProvider
) -> tuple[str, ReasoningEffort | None]:
capabilities = await discover_capabilities(provider)
if not capabilities.available:
raise ProviderError(capabilities.detail or f"{provider}를 사용할 수 없습니다.")
requested_model = normalize_engine_gateway_model(req.model)
model = requested_model or capabilities.default_model
option = next((item for item in capabilities.models if item.id == model), None)
if option is None:
raise ProviderError(f"{provider}에서 사용할 수 없는 모델입니다: {model}")
effort = req.reasoning_effort or option.default_reasoning_effort
if effort is not None and effort not in option.reasoning_efforts:
raise ProviderError(f"{model}에서 사용할 수 없는 추론 강도입니다: {effort}")
return option.id, effort
def _structured_or_none(text: str, req: GenerateRequest) -> dict[str, Any] | None:
if not req.structured_schema:
return None
try:
parsed = json.loads(text)
except (json.JSONDecodeError, TypeError):
return None
return parsed if isinstance(parsed, dict) else None
async def _generate_codex(
req: GenerateRequest, system_prompt: str, user_payload: str
) -> ProviderGenerateResult:
pricing_started_at = _utcnow()
binary = _binary("CODEX_BIN", "codex")
if binary is None:
raise ProviderError("Codex CLI를 찾을 수 없습니다.")
model, effort = await _resolve_selection(req, "codex_cli")
cli_cwd = _cli_runtime_cwd()
args = [
binary,
"exec",
"--json",
"--ephemeral",
"--skip-git-repo-check",
"--ignore-user-config",
"--ignore-rules",
"--sandbox",
"read-only",
"-C",
str(cli_cwd),
"-m",
model,
]
if effort:
args += ["-c", f'model_reasoning_effort="{effort}"']
args.append("-")
stdout, _ = await _run_process(
args,
input_text=_cli_prompt(system_prompt, user_payload),
cwd=str(cli_cwd),
)
text = ""
tokens_in = 0
tokens_out = 0
cached_input_tokens = 0
for line in stdout.splitlines():
try:
event = json.loads(line)
except json.JSONDecodeError:
continue
if event.get("type") == "item.completed":
item = event.get("item") or {}
if item.get("type") == "agent_message":
text = str(item.get("text") or text)
elif event.get("type") == "turn.completed":
usage = event.get("usage") or {}
tokens_in = int(usage.get("input_tokens") or 0)
tokens_out = int(usage.get("output_tokens") or 0)
cached_input_tokens = int(usage.get("cached_input_tokens") or 0)
elif event.get("type") in {"turn.failed", "error"}:
raise ProviderError(str(event.get("message") or event))
if not text.strip():
raise ProviderError("Codex CLI가 최종 응답을 반환하지 않았습니다.")
estimate = estimate_reference_cost(
provider="codex_cli",
model=model,
tokens_in=tokens_in,
tokens_out=tokens_out,
priced_at=pricing_started_at,
cached_input_tokens=cached_input_tokens,
)
return ProviderGenerateResult(
text=text,
model=model,
provider="codex_cli",
tokens_in=tokens_in,
tokens_out=tokens_out,
cost_usd=estimate.cost_usd if estimate is not None else 0.0,
structured=_structured_or_none(text, req),
)
async def _generate_agy(
req: GenerateRequest, system_prompt: str, user_payload: str
) -> ProviderGenerateResult:
# text 출력은 토큰 사용량을 주지 않는다. stream-json의 terminal result를
# 동일하게 소비해 generate와 stream 모두 같은 token/cost 계약을 유지한다.
async for event in _stream_agy(req, system_prompt, user_payload):
if event.type == "done" and event.result is not None:
return event.result
raise ProviderError("Agy CLI가 최종 응답을 반환하지 않았습니다.")
async def _stream_agy(
req: GenerateRequest, system_prompt: str, user_payload: str
) -> AsyncIterator[ProviderStreamEvent]:
"""Agy stream-json의 agent_response delta를 게이트웨이 토큰으로 전달한다.
Agy print 모드는 대화 내용을 로컬 conversation 저장소에 남길 수 있으므로 여기서는
--continue/--conversation을 쓰지 않는다. 회기 메모리는 매 요청의 마스킹된 prompt가
소유하고, 프로세스는 응답 뒤 종료한다.
"""
pricing_started_at = _utcnow()
binary = _binary("AGY_BIN", "agy")
if binary is None:
raise ProviderError("Agy CLI를 찾을 수 없습니다.")
model, effort = await _resolve_selection(req, "agy_cli")
prompt = _cli_prompt(system_prompt, user_payload)
args = [binary, "--model", model, "--sandbox"]
if effort:
args += ["--effort", effort]
args += [
"--print-timeout",
f"{int(CLI_TIMEOUT_SECONDS)}s",
# 긴 deep-loop 축어록을 Windows argv에 싣지 않는다. Agy의 공식 stream-json
# 입력 계약은 prompt를 stdin의 단일 user 이벤트로 받으므로 명령줄 길이 한계를
# 피하면서 전체 마스킹 근거를 그대로 보존한다.
"--input-format",
"stream-json",
"--output-format",
"stream-json",
]
proc = await asyncio.create_subprocess_exec(
*args,
cwd=str(_cli_runtime_cwd()),
stdin=asyncio.subprocess.PIPE,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
env=_cli_subprocess_env(),
)
assert proc.stdin is not None
assert proc.stdout is not None
assert proc.stderr is not None
stderr_task = asyncio.create_task(proc.stderr.read())
emitted = ""
final_text = ""
tokens_in = 0
tokens_out = 0
cached_input_tokens = 0
result_status = ""
try:
# 공식 protocol: 한 줄에 한 user event. 마지막 turn 뒤 stdin을 닫아도 CLI는
# terminal result를 내보낸 뒤 종료한다. stdin 거절은 child와 stderr를 정리한 뒤
# provider 오류로 승격해 프로세스를 남기지 않는다.
try:
proc.stdin.write(
(
json.dumps(
{"event": "user", "message": {"content": prompt}},
ensure_ascii=False,
)
+ "\n"
).encode("utf-8")
)
await proc.stdin.drain()
except (BrokenPipeError, ConnectionResetError) as exc:
raise ProviderError("Agy CLI가 stdin 평가 입력을 수락하지 않았습니다.") from exc
finally:
if not proc.stdin.is_closing():
proc.stdin.close()
try:
await proc.stdin.wait_closed()
except (BrokenPipeError, ConnectionResetError):
# 이미 종료된 CLI가 close 직후 EOF를 끊어도 finally가 child를 회수한다.
pass
async with asyncio.timeout(CLI_TIMEOUT_SECONDS):
while True:
raw = await proc.stdout.readline()
if not raw:
break
try:
event = json.loads(raw.decode("utf-8", errors="replace"))
except json.JSONDecodeError:
continue
if event.get("event") == "step_update":
update = event.get("step_update") or {}
if update.get("step_type") == "agent_response":
delta = str(update.get("text_delta") or "")
if delta:
emitted += delta
yield ProviderStreamEvent(type="delta", text=delta)
elif event.get("event") == "result":
result = event.get("result") or {}
result_status = str(result.get("status") or "")
final_text = str(result.get("response") or "")
usage = result.get("usage") or {}
tokens_in = int(usage.get("input_tokens") or 0)
tokens_out = int(usage.get("output_tokens") or 0)
cached_input_tokens = int(
usage.get("cache_read_tokens")
or usage.get("cached_input_tokens")
or 0
)
returncode = await proc.wait()
except TimeoutError as exc:
raise ProviderError(
f"Agy CLI 응답 시간이 {int(CLI_TIMEOUT_SECONDS)}초를 넘었습니다."
) from exc
finally:
if proc.returncode is None:
proc.kill()
await proc.wait()
stderr = await stderr_task
if returncode != 0:
raise ProviderError(_safe_process_error(stderr, f"Agy CLI exit {returncode}"))
if result_status and result_status != "SUCCESS":
raise ProviderError(f"Agy CLI 생성 실패: {result_status}")
resolved_text = final_text or emitted
if not resolved_text.strip():
raise ProviderError("Agy CLI가 최종 응답을 반환하지 않았습니다.")
if final_text and final_text.startswith(emitted):
remainder = final_text[len(emitted) :]
if remainder:
emitted += remainder
yield ProviderStreamEvent(type="delta", text=remainder)
elif not emitted:
emitted = resolved_text
yield ProviderStreamEvent(type="delta", text=resolved_text)
estimate = estimate_reference_cost(
provider="agy_cli",
model=model,
tokens_in=tokens_in,
tokens_out=tokens_out,
priced_at=pricing_started_at,
cached_input_tokens=cached_input_tokens,
)
yield ProviderStreamEvent(
type="done",
result=ProviderGenerateResult(
text=resolved_text,
model=model,
provider="agy_cli",
tokens_in=tokens_in,
tokens_out=tokens_out,
cost_usd=estimate.cost_usd if estimate is not None else 0.0,
structured=_structured_or_none(resolved_text, req),
),
)
async def _generate_claude_api(
req: GenerateRequest, system_prompt: str
) -> ProviderGenerateResult:
pricing_started_at = _utcnow()
api_key = os.environ.get("ANTHROPIC_API_KEY", "").strip() or os.environ.get(
"ANTHROPIC_AUTH_TOKEN", ""
).strip()
if not api_key:
raise ProviderError("ANTHROPIC_API_KEY가 설정되지 않았습니다.")
model, effort = await _resolve_selection(req, "claude_api")
messages = [
{"role": message.role, "content": message.content}
for message in req.messages
if message.role != "system"
]
payload: dict[str, Any] = {
"model": model,
"max_tokens": req.max_tokens,
"temperature": req.temperature,
"messages": messages,
}
if system_prompt:
payload["system"] = system_prompt
if effort:
payload["output_config"] = {"effort": effort}
try:
async with httpx.AsyncClient(timeout=CLI_TIMEOUT_SECONDS) as client:
response = await client.post(
f"{ANTHROPIC_API_BASE}/v1/messages",
headers=_anthropic_headers(api_key),
json=payload,
)
response.raise_for_status()
body = response.json()
except (httpx.HTTPError, ValueError) as exc:
raise ProviderError(f"Anthropic Messages API 호출 실패: {exc}") from exc
text = "".join(
str(block.get("text") or "")
for block in body.get("content", [])
if isinstance(block, dict) and block.get("type") == "text"
)
if not text:
raise ProviderError("Anthropic Messages API가 텍스트 응답을 반환하지 않았습니다.")
usage = body.get("usage") or {}
inference_geo = body.get("inference_geo")
tokens_in = int(usage.get("input_tokens") or 0)
tokens_out = int(usage.get("output_tokens") or 0)
estimate = estimate_reference_cost(
provider="claude_api",
model=str(body.get("model") or model),
tokens_in=tokens_in,
tokens_out=tokens_out,
priced_at=pricing_started_at,
cached_input_tokens=int(usage.get("cache_read_input_tokens") or 0),
)
return ProviderGenerateResult(
text=text,
model=str(body.get("model") or model),
provider="claude_api",
tokens_in=tokens_in,
tokens_out=tokens_out,
cost_usd=estimate.cost_usd if estimate is not None else 0.0,
inference_geo=str(inference_geo) if inference_geo else None,
structured=_structured_or_none(text, req),
)
def _openai_response_text(body: dict[str, Any]) -> str:
direct = body.get("output_text")
if isinstance(direct, str) and direct.strip():
return direct
parts: list[str] = []
for item in body.get("output", []) if isinstance(body, dict) else []:
if not isinstance(item, dict) or item.get("type") != "message":
continue
for content in item.get("content", []):
if (
isinstance(content, dict)
and content.get("type") == "output_text"
and isinstance(content.get("text"), str)
):
parts.append(content["text"])
return "".join(parts)
async def _generate_openai(
req: GenerateRequest,
system_prompt: str,
user_payload: str,
) -> ProviderGenerateResult:
pricing_started_at = _utcnow()
api_key = os.environ.get("OPENAI_API_KEY", "").strip()
if not api_key:
raise ProviderError("OPENAI_API_KEY가 설정되지 않았습니다.")
model, effort = await _resolve_selection(req, "openai")
if req.ai_role == "client":
messages: list[dict[str, str]] = [
{"role": "user", "content": user_payload}
]
else:
messages = [
{"role": message.role, "content": message.content}
for message in req.messages
if message.role != "system"
]
payload: dict[str, Any] = {
"model": model,
"input": messages,
"max_output_tokens": req.max_tokens,
# 상담 시뮬레이션 입력을 OpenAI의 응답 상태 저장소에 남기지 않는다.
# 회기 기록의 SSOT는 Vignette의 NAS PostgreSQL뿐이다.
"store": False,
}
if system_prompt:
payload["instructions"] = system_prompt
if effort:
payload["reasoning"] = {"effort": effort}
else:
payload["temperature"] = req.temperature
try:
async with httpx.AsyncClient(timeout=CLI_TIMEOUT_SECONDS) as client:
response = await client.post(
f"{OPENAI_API_BASE}/responses",
headers={
"Authorization": f"Bearer {api_key}",
"Content-Type": "application/json",
},
json=payload,
)
response.raise_for_status()
body = response.json()
except (httpx.HTTPError, ValueError) as exc:
raise ProviderError(f"OpenAI Responses API 호출 실패: {exc}") from exc
if not isinstance(body, dict):
raise ProviderError("OpenAI Responses API가 객체 응답을 반환하지 않았습니다.")
text = _openai_response_text(body).strip()
if not text:
raise ProviderError("OpenAI Responses API가 텍스트 응답을 반환하지 않았습니다.")
usage = body.get("usage") or {}
tokens_in = int(usage.get("input_tokens") or 0)
tokens_out = int(usage.get("output_tokens") or 0)
input_details = usage.get("input_tokens_details") or {}
estimate = estimate_reference_cost(
provider="openai",
model=str(body.get("model") or model),
tokens_in=tokens_in,
tokens_out=tokens_out,
priced_at=pricing_started_at,
cached_input_tokens=int(input_details.get("cached_tokens") or 0),
)
inference_geo = body.get("inference_geo")
return ProviderGenerateResult(
text=text,
model=str(body.get("model") or model),
provider="openai",
tokens_in=tokens_in,
tokens_out=tokens_out,
cost_usd=estimate.cost_usd if estimate is not None else 0.0,
inference_geo=str(inference_geo) if inference_geo else None,
structured=_structured_or_none(text, req),
)
async def generate_with_provider(
req: GenerateRequest,
*,
system_prompt: str,
user_payload: str,
) -> ProviderGenerateResult:
provider = req.provider
if provider == "codex_cli":
return await _generate_codex(req, system_prompt, user_payload)
if provider == "agy_cli":
return await _generate_agy(req, system_prompt, user_payload)
if provider == "claude_api":
return await _generate_claude_api(req, system_prompt)
if provider == "openai":
return await _generate_openai(req, system_prompt, user_payload)
if provider == "codex_api":
return await _generate_codex_api(req, system_prompt, user_payload)
if provider == "antigravity_api":
return await _generate_antigravity_api(req, system_prompt, user_payload)
if provider == "openrouter":
return await _generate_openrouter(req, system_prompt, user_payload)
raise ProviderError(f"이 게이트웨이에서 실행할 수 없는 provider입니다: {provider}")
async def stream_with_provider(
req: GenerateRequest,
*,
system_prompt: str,
user_payload: str,
) -> AsyncIterator[ProviderStreamEvent]:
"""Provider가 제공하는 가장 이른 출력 단위를 공통 delta/done 계약으로 바꾼다."""
if req.provider == "agy_cli":
async for event in _stream_agy(req, system_prompt, user_payload):
yield event
return
result = await generate_with_provider(
req,
system_prompt=system_prompt,
user_payload=user_payload,
)
yield ProviderStreamEvent(type="delta", text=result.text)
yield ProviderStreamEvent(type="done", result=result)