import json import unittest from datetime import datetime, timezone from unittest.mock import AsyncMock, Mock, patch from app.contracts.engine_gateway import EngineMessage, GenerateRequest from engine_gateway import provider_registry class _FakeStreamReader: def __init__(self, lines: list[bytes] | None = None, body: bytes = b""): self.lines = list(lines or []) self.body = body async def readline(self) -> bytes: return self.lines.pop(0) if self.lines else b"" async def read(self) -> bytes: return self.body class _FakeStreamWriter: def __init__(self): self.writes: list[bytes] = [] self.closed = False def write(self, data: bytes) -> None: self.writes.append(data) async def drain(self) -> None: return None def is_closing(self) -> bool: return self.closed def close(self) -> None: self.closed = True async def wait_closed(self) -> None: return None class _FakeAgyProcess: def __init__(self, events: list[dict]): self.stdin = _FakeStreamWriter() self.stdout = _FakeStreamReader( [(json.dumps(event, ensure_ascii=False) + "\n").encode("utf-8") for event in events] ) self.stderr = _FakeStreamReader() self.returncode = None async def wait(self) -> int: if self.returncode is None: self.returncode = 0 return self.returncode def kill(self) -> None: self.returncode = -9 class ProviderRegistryTest(unittest.IsolatedAsyncioTestCase): def test_cli_subprocess_env_backfills_windows_essentials(self): # 런타임 승격 체인의 psutil 환경 이식에서 SystemRoot 가 유실되면 agy(Go)가 # 인증서 풀/홈 해석에 실패해 빈 모델 목록을 내놓는다(2026-08-18 실측). scrubbed = { "PATH": r"C:\Windows\System32", "USERPROFILE": r"C:\Users\encep", } with patch.dict(provider_registry.os.environ, scrubbed, clear=True): env = provider_registry._cli_subprocess_env() self.assertEqual(env["SystemRoot"], r"C:\Windows") self.assertEqual(env["SystemDrive"], "C:") self.assertEqual(env["ComSpec"], r"C:\Windows\system32\cmd.exe") self.assertEqual(env["PATH"], r"C:\Windows\System32") def test_cli_subprocess_env_preserves_existing_essentials(self): # os.environ은 Windows에서 키를 대문자로 정규화한다. 이미 값이 있으면 # 기본 케이스 키를 덧붙이지 않고 기존값을 그대로 둔다. scrubbed = { "SYSTEMROOT": r"D:\Win", "PATH": "x", } with patch.dict(provider_registry.os.environ, scrubbed, clear=True): env = provider_registry._cli_subprocess_env() self.assertEqual(env["SYSTEMROOT"], r"D:\Win") self.assertNotIn("SystemRoot", env) def setUp(self): provider_registry.clear_capability_cache() def tearDown(self): provider_registry.clear_capability_cache() async def test_codex_catalog_uses_live_models_with_terra_medium_default(self): payload = { "data": [ { "id": "gpt-5.6-sol", "model": "gpt-5.6-sol", "displayName": "GPT-5.6-Sol", "description": "Frontier", "hidden": False, "isDefault": True, "defaultReasoningEffort": "low", "supportedReasoningEfforts": [ {"reasoningEffort": "low"}, {"reasoningEffort": "medium"}, ], }, { "id": "gpt-5.6-terra", "model": "gpt-5.6-terra", "displayName": "GPT-5.6-Terra", "description": "Balanced", "hidden": False, "isDefault": False, "defaultReasoningEffort": "medium", "supportedReasoningEfforts": [ {"reasoningEffort": "low"}, {"reasoningEffort": "medium"}, {"reasoningEffort": "high"}, ], }, ] } long_payload = "x" * 24_001 with ( patch.object(provider_registry, "_binary", return_value="codex.exe"), patch.object( provider_registry, "_codex_model_list", AsyncMock(return_value=payload), ), ): result = await provider_registry.discover_capabilities("codex_cli") self.assertTrue(result.available) self.assertEqual(result.source, "live_cli") self.assertEqual(result.default_model, "gpt-5.6-terra") self.assertEqual(result.default_reasoning_effort, "medium") terra = next(model for model in result.models if model.id == "gpt-5.6-terra") self.assertTrue(terra.is_default) self.assertEqual(terra.reasoning_efforts, ["low", "medium", "high"]) async def test_agy_catalog_uses_cli_list_with_flash_high_default(self): stdout = "\n".join( [ "gemini-3.6-flash-high", "gemini-3.6-flash-medium", "claude-sonnet-4-6", ] ) with ( patch.object(provider_registry, "_binary", return_value="agy.exe"), patch.object( provider_registry, "_run_process", AsyncMock(return_value=(stdout, "")), ), ): result = await provider_registry.discover_capabilities("agy_cli") self.assertTrue(result.available) self.assertEqual(result.default_model, "gemini-3.6-flash-high") self.assertEqual(result.default_reasoning_effort, "high") selected = next(model for model in result.models if model.is_default) self.assertEqual(selected.reasoning_efforts, ["high"]) self.assertEqual(selected.label, "Gemini 3.6 Flash (High)") async def test_agy_catalog_parses_tab_separated_id_label_lines(self): # 2026-08 agy CLI는 `agy models`를 "model_id\t표시 라벨" 형태로 출력한다. # 탭이 포함된 줄을 통째로 버리면 카탈로그가 비어 available=false가 된다. stdout = "\n".join( [ "gemini-3.7-flash-high\tGemini 3.7 Flash (High)", "gemini-3.7-flash-medium\tGemini 3.7 Flash (Medium)", "claude-sonnet-4-6\tClaude Sonnet 4.6 (Thinking)", "Fetching available models...", ] ) with ( patch.object(provider_registry, "_binary", return_value="agy.exe"), patch.object( provider_registry, "_run_process", AsyncMock(return_value=(stdout, "")), ), ): result = await provider_registry.discover_capabilities("agy_cli") self.assertTrue(result.available) self.assertEqual( [model.id for model in result.models], ["gemini-3.7-flash-high", "gemini-3.7-flash-medium", "claude-sonnet-4-6"], ) # 기본 모델 gemini-3.6-flash-high 가 목록에 없으면 첫 모델로 폴백한다. self.assertEqual(result.default_model, "gemini-3.7-flash-high") self.assertEqual(result.default_reasoning_effort, "high") selected = next(model for model in result.models if model.id == result.default_model) self.assertEqual(selected.reasoning_efforts, ["high"]) self.assertEqual(selected.label, "Gemini 3.7 Flash (High)") async def test_claude_cli_catalog_is_explicit_static_alias_fallback(self): with patch.object(provider_registry, "_binary", return_value="claude.exe"): result = await provider_registry.discover_capabilities("claude_cli") self.assertTrue(result.available) self.assertEqual(result.source, "static_cli") self.assertEqual(result.default_model, "gateway-default") self.assertEqual([model.id for model in result.models], ["gateway-default", "opus", "sonnet", "fable"]) async def test_anthropic_catalog_fails_closed_without_api_key(self): with patch.dict(provider_registry.os.environ, {}, clear=True): result = await provider_registry.discover_capabilities( "claude_api", force=True ) self.assertFalse(result.available) self.assertEqual(result.source, "unavailable") self.assertEqual(result.models, []) self.assertIn("ANTHROPIC_API_KEY", result.detail) async def test_openai_catalog_intersects_live_models_with_explicit_allowlist(self): response = Mock() response.raise_for_status.return_value = None response.json.return_value = { "data": [ {"id": "gpt-5.6-sol"}, {"id": "gpt-5.6-terra"}, {"id": "gpt-4.1"}, {"id": "gpt-4o-mini-tts"}, ] } client = AsyncMock() client.__aenter__.return_value = client client.__aexit__.return_value = False client.get.return_value = response with ( patch.dict( provider_registry.os.environ, { "OPENAI_API_KEY": "test-openai-key", "OPENAI_ENGINE_MODEL": "gpt-5.6-terra", "OPENAI_ENGINE_MODELS": "gpt-5.6-terra,gpt-4.1", }, clear=True, ), patch.object( provider_registry.httpx, "AsyncClient", return_value=client, ), ): result = await provider_registry.discover_capabilities( "openai", force=True ) self.assertTrue(result.available) self.assertEqual(result.source, "live_api") self.assertEqual(result.default_model, "gpt-5.6-terra") self.assertEqual( [model.id for model in result.models], ["gpt-5.6-terra", "gpt-4.1"], ) self.assertEqual(result.models[0].reasoning_efforts, ["low", "medium", "high"]) self.assertEqual(result.models[1].reasoning_efforts, []) request = client.get.await_args self.assertEqual(request.args[0], f"{provider_registry.OPENAI_API_BASE}/models") self.assertEqual( request.kwargs["headers"]["Authorization"], "Bearer test-openai-key", ) async def test_openai_generation_uses_responses_api_without_temperature_for_reasoning(self): capabilities = provider_registry.EngineCapabilitiesResponse( provider="openai", available=True, source="live_api", models=[ provider_registry.EngineModelOption( id="gpt-5.6-terra", label="gpt-5.6-terra", reasoning_efforts=["low", "medium", "high"], default_reasoning_effort="medium", is_default=True, ) ], default_model="gpt-5.6-terra", default_reasoning_effort="medium", fetched_at=1, ) response = Mock() response.raise_for_status.return_value = None response.json.return_value = { "model": "gpt-5.6-terra", "output": [ { "type": "message", "content": [{"type": "output_text", "text": "OK"}], } ], "usage": { "input_tokens": 12, "output_tokens": 2, "input_tokens_details": {"cached_tokens": 3}, }, } client = AsyncMock() client.__aenter__.return_value = client client.__aexit__.return_value = False client.post.return_value = response request = GenerateRequest( provider="openai", model="gpt-5.6-terra", reasoning_effort="medium", messages=[EngineMessage(role="user", content="hello")], ) with ( patch.dict( provider_registry.os.environ, {"OPENAI_API_KEY": "test-openai-key"}, clear=True, ), patch.object( provider_registry, "discover_capabilities", AsyncMock(return_value=capabilities), ), patch.object( provider_registry.httpx, "AsyncClient", return_value=client, ), ): result = await provider_registry.generate_with_provider( request, system_prompt="system", user_payload="current turn", ) self.assertEqual(result.text, "OK") self.assertEqual(result.provider, "openai") self.assertEqual(result.tokens_in, 12) self.assertEqual(result.tokens_out, 2) call = client.post.await_args self.assertEqual(call.args[0], f"{provider_registry.OPENAI_API_BASE}/responses") payload = call.kwargs["json"] self.assertEqual(payload["model"], "gpt-5.6-terra") self.assertEqual(payload["instructions"], "system") self.assertEqual(payload["input"], [{"role": "user", "content": "current turn"}]) self.assertIs(payload["store"], False) self.assertEqual(payload["reasoning"], {"effort": "medium"}) async def test_codex_generation_uses_model_and_reasoning_from_selection(self): capabilities = provider_registry.EngineCapabilitiesResponse( provider="codex_cli", available=True, source="live_cli", models=[ provider_registry.EngineModelOption( id="gpt-5.6-terra", label="GPT-5.6-Terra", reasoning_efforts=["low", "medium", "high"], default_reasoning_effort="medium", is_default=True, ) ], default_model="gpt-5.6-terra", default_reasoning_effort="medium", fetched_at=1, ) stdout = "\n".join( [ json.dumps( { "type": "item.completed", "item": {"type": "agent_message", "text": "OK"}, } ), json.dumps( { "type": "turn.completed", "usage": { "input_tokens": 12, "cached_input_tokens": 2, "output_tokens": 3, }, } ), ] ) runner = AsyncMock(return_value=(stdout, "")) request = GenerateRequest( provider="codex_cli", model="gpt-5.6-terra", reasoning_effort="medium", messages=[EngineMessage(role="user", content="hello")], ) with ( patch.object(provider_registry, "_binary", return_value="codex.exe"), patch.object( provider_registry, "discover_capabilities", AsyncMock(return_value=capabilities), ), patch.object(provider_registry, "_run_process", runner), ): result = await provider_registry.generate_with_provider( request, system_prompt="system", user_payload="hello", ) self.assertEqual(result.text, "OK") self.assertEqual(result.tokens_in, 12) self.assertEqual(result.tokens_out, 3) self.assertEqual(result.cost_usd, 0.0000705) args = runner.await_args.args[0] self.assertIn("gpt-5.6-terra", args) self.assertIn('model_reasoning_effort="medium"', args) self.assertEqual(args[-1], "-") self.assertIn("[시스템 지침]", runner.await_args.kwargs["input_text"]) async def test_agy_generation_uses_stream_json_usage_and_reference_cost(self): capabilities = provider_registry.EngineCapabilitiesResponse( provider="agy_cli", available=True, source="live_cli", models=[ provider_registry.EngineModelOption( id="gemini-3.7-flash-high", label="Gemini 3.7 Flash (High)", reasoning_efforts=["high"], default_reasoning_effort="high", ) ], default_model="gemini-3.7-flash-high", default_reasoning_effort="high", fetched_at=1, ) request = GenerateRequest( provider="agy_cli", model="gemini-3.7-flash-high", reasoning_effort="high", messages=[EngineMessage(role="user", content="hello")], ) process = _FakeAgyProcess( [ { "event": "step_update", "step_update": { "step_type": "agent_response", "state": "DONE", "text_delta": "OK", }, }, { "event": "result", "result": { "status": "SUCCESS", "response": "OK", "usage": { "input_tokens": 12, "cache_read_tokens": 2, "output_tokens": 2, }, }, }, ] ) captured: list[tuple] = [] async def fake_create_subprocess_exec(*args, **kwargs): captured.append(args) return process long_payload = "x" * 24_001 with ( patch.object(provider_registry, "_binary", return_value="agy"), patch.object( provider_registry, "_utcnow", return_value=datetime(2026, 8, 28, tzinfo=timezone.utc), ), patch.object( provider_registry, "discover_capabilities", AsyncMock(return_value=capabilities), ), patch.object( provider_registry.asyncio, "create_subprocess_exec", fake_create_subprocess_exec, ), ): result = await provider_registry.generate_with_provider( request, system_prompt="system", user_payload=long_payload, ) self.assertEqual(result.text, "OK") self.assertEqual(result.tokens_in, 12) self.assertEqual(result.tokens_out, 2) self.assertEqual(result.cost_usd, 0.00001515) args = captured[0] self.assertIn("--input-format", args) self.assertEqual(args[args.index("--input-format") + 1], "stream-json") self.assertEqual(args[args.index("--output-format") + 1], "stream-json") self.assertNotIn("--print", args) self.assertTrue(all(long_payload not in str(arg) for arg in args)) self.assertTrue(process.stdin.closed) sent = json.loads(b"".join(process.stdin.writes).decode("utf-8")) self.assertEqual(sent["event"], "user") self.assertIn("[시스템 지침]", sent["message"]["content"]) self.assertIn(long_payload, sent["message"]["content"]) async def test_agy_stream_forwards_live_deltas_without_repeating_final_response(self): capabilities = provider_registry.EngineCapabilitiesResponse( provider="agy_cli", available=True, source="live_cli", models=[ provider_registry.EngineModelOption( id="gemini-3.6-flash-high", label="Gemini 3.6 Flash (High)", reasoning_efforts=["high"], default_reasoning_effort="high", is_default=True, ) ], default_model="gemini-3.6-flash-high", default_reasoning_effort="high", fetched_at=1, ) request = GenerateRequest( provider="agy_cli", model="gemini-3.6-flash-high", reasoning_effort="high", messages=[EngineMessage(role="user", content="hello")], ) process = _FakeAgyProcess( [ { "event": "step_update", "step_update": { "step_type": "agent_response", "state": "ACTIVE", "text_delta": "안", }, }, { "event": "step_update", "step_update": { "step_type": "agent_response", "state": "DONE", "text_delta": "녕", }, }, { "event": "result", "result": { "status": "SUCCESS", "response": "안녕", "usage": { "input_tokens": 12, "cache_read_tokens": 2, "output_tokens": 2, }, }, }, ] ) captured: list[tuple] = [] async def fake_create_subprocess_exec(*args, **kwargs): captured.append(args) return process with ( patch.object(provider_registry, "_binary", return_value="agy.exe"), patch.object( provider_registry, "_utcnow", return_value=datetime(2026, 8, 28, tzinfo=timezone.utc), ), patch.object( provider_registry, "discover_capabilities", AsyncMock(return_value=capabilities), ), patch.object( provider_registry.asyncio, "create_subprocess_exec", fake_create_subprocess_exec, ), ): events = [ event async for event in provider_registry.stream_with_provider( request, system_prompt="system", user_payload="hello", ) ] self.assertEqual([event.type for event in events], ["delta", "delta", "done"]) self.assertEqual("".join(event.text for event in events), "안녕") self.assertEqual(events[-1].result.text, "안녕") self.assertEqual(events[-1].result.tokens_in, 12) self.assertEqual(events[-1].result.tokens_out, 2) self.assertEqual(events[-1].result.cost_usd, 0.00001515) args = captured[0] self.assertIn("--input-format", args) self.assertEqual(args[args.index("--input-format") + 1], "stream-json") self.assertIn("--output-format", args) self.assertEqual(args[args.index("--output-format") + 1], "stream-json") self.assertNotIn("--print", args) async def test_agy_stdin_rejection_reaps_child_process(self): request = GenerateRequest( provider="agy_cli", model="gemini-3.6-flash-high", reasoning_effort="high", messages=[EngineMessage(role="user", content="hello")], ) process = _FakeAgyProcess([]) async def broken_drain() -> None: raise BrokenPipeError() process.stdin.drain = broken_drain # type: ignore[method-assign] async def fake_create_subprocess_exec(*args, **kwargs): return process with ( patch.object(provider_registry, "_binary", return_value="agy.exe"), patch.object( provider_registry, "_resolve_selection", AsyncMock(return_value=("gemini-3.6-flash-high", "high")), ), patch.object( provider_registry.asyncio, "create_subprocess_exec", fake_create_subprocess_exec, ), ): with self.assertRaisesRegex(provider_registry.ProviderError, "stdin 평가 입력"): async for _event in provider_registry._stream_agy( request, system_prompt="system", user_payload="hello", ): pass self.assertTrue(process.stdin.closed) self.assertEqual(process.returncode, -9) async def test_generation_rejects_model_effort_not_returned_by_provider(self): capabilities = provider_registry.EngineCapabilitiesResponse( provider="agy_cli", available=True, source="live_cli", models=[ provider_registry.EngineModelOption( id="gemini-3.6-flash-high", label="Gemini 3.6 Flash (High)", reasoning_efforts=["high"], default_reasoning_effort="high", ) ], default_model="gemini-3.6-flash-high", default_reasoning_effort="high", fetched_at=1, ) request = GenerateRequest( provider="agy_cli", model="gemini-3.6-flash-high", reasoning_effort="low", messages=[EngineMessage(role="user", content="hello")], ) with patch.object( provider_registry, "discover_capabilities", AsyncMock(return_value=capabilities), ): with self.assertRaisesRegex( provider_registry.ProviderError, "사용할 수 없는 추론 강도" ): await provider_registry._resolve_selection(request, "agy_cli") if __name__ == "__main__": unittest.main()