910 lines
33 KiB
Python
910 lines
33 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("/")
|
|
|
|
|
|
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()
|
|
|
|
|
|
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 _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()
|
|
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={
|
|
"x-api-key": api_key,
|
|
"anthropic-version": "2023-06-01",
|
|
},
|
|
)
|
|
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(provider: EngineProvider) -> EngineCapabilitiesResponse:
|
|
if provider == "claude_cli":
|
|
return await _discover_claude_cli()
|
|
if provider == "claude_api":
|
|
return await _discover_claude_api()
|
|
if provider == "codex_cli":
|
|
return await _discover_codex_cli()
|
|
if provider == "agy_cli":
|
|
return await _discover_agy_cli()
|
|
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)
|
|
if os.name == "nt" and len(prompt) > 24_000:
|
|
raise ProviderError(
|
|
"Agy CLI 프롬프트가 Windows 명령줄 안전 한도(24,000자)를 초과했습니다."
|
|
)
|
|
args = [binary, "--model", model, "--sandbox"]
|
|
if effort:
|
|
args += ["--effort", effort]
|
|
args += [
|
|
"--print-timeout",
|
|
f"{int(CLI_TIMEOUT_SECONDS)}s",
|
|
"--output-format",
|
|
"stream-json",
|
|
"--print",
|
|
prompt,
|
|
]
|
|
proc = await asyncio.create_subprocess_exec(
|
|
*args,
|
|
cwd=str(_cli_runtime_cwd()),
|
|
stdout=asyncio.subprocess.PIPE,
|
|
stderr=asyncio.subprocess.PIPE,
|
|
env=_cli_subprocess_env(),
|
|
)
|
|
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:
|
|
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()
|
|
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={
|
|
"x-api-key": api_key,
|
|
"anthropic-version": "2023-06-01",
|
|
},
|
|
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),
|
|
)
|
|
|
|
|
|
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)
|
|
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)
|