import asyncio import json import os import tempfile import unittest from pathlib import Path from unittest.mock import AsyncMock, patch from aiohttp import web from aiohttp.test_utils import TestClient, TestServer from decision_server.connections import Connection, Connections, endpoint from decision_server.protocol import ( PRECONDITIONS, SKILLS, VERSION, DecisionError, loads, validate, validate_plan, ) from decision_server.providers.codex import CodexAccount from decision_server.server import PREFIX, STATE, create_app def plan_value(): return { "version": VERSION, "objectId": "block", "goalId": "placement", "summary": "搬运方块", "steps": [ {"skill": s, "precondition": PRECONDITIONS[s], "onFailure": "stop"} for s in SKILLS ], } def observation(request_id="r1"): return { "version": VERSION, "stamp": { "runId": "run1", "sceneRevision": 1, "sequence": 0, "planRevision": 0, "requestId": request_id, }, "source": "mujoco-ground-truth", "units": "SI", "frame": "world-z-up", "time": 0, "phase": "open", "base": {"position": [0, 0, 0.09], "yaw": 0}, "joints": [0] * 5, "opening": 1, "tcp": [0.2, 0, 0.2], "object": {"id": "block", "position": [0.257, 0.015, 0.128], "speed": 0}, "goal": {"id": "placement", "position": [0.257, 0.615, 0.128]}, "evidence": { "fingerForces": [0, 0], "supported": True, "onGoalSupport": False, "secure": False, "transported": 0, }, "safety": [], } def request_value(ident="r1"): return {"observation": observation(ident), "instruction": "把方块搬到目标", "remaining": SKILLS} class ProtocolTests(unittest.TestCase): def test_shared_schema_and_semantics(self): self.assertEqual(validate("Observation", observation()), observation()) self.assertEqual(validate_plan(plan_value(), SKILLS), plan_value()) for mutate in ( lambda p: p.update(objectId="other"), lambda p: p["steps"].reverse(), lambda p: p["steps"][0].update(precondition="released"), lambda p: p["steps"].append(p["steps"][0]), lambda p: p.update(command="shell"), ): value = plan_value() mutate(value) with self.assertRaises(DecisionError): validate_plan(value, SKILLS) for number in (float("nan"), float("inf"), True, 1e10): value = observation() value["time"] = number with self.assertRaises(DecisionError): validate("Observation", value) for text in ('{"a":1,"a":2}', '{"a":NaN}', "no json"): with self.assertRaises(DecisionError): loads(text) def test_endpoint_and_credentials(self): for url in ( "http://evil.test/v1", "https://host/?key=secret", "https://user:key@host", "file:///etc/passwd", "http://[bad", "https://x:99999", "https://x/\\evil", ): with self.assertRaises(DecisionError, msg=url): endpoint(url) for url in ("http://127.0.0.1:9000/v1", "http://localhost/v1", "https://api.openai.com/v1"): self.assertEqual(endpoint(url), url) with tempfile.TemporaryDirectory() as directory: store = Connections(Path(directory)) data = { "role": "llm", "protocol": "responses", "baseUrl": "https://api.openai.com/v1", "model": "test-model", "apiKey": "test-secret", } store.set(data) self.assertNotIn("test-secret", store.path.read_text()) self.assertEqual(store.path.stat().st_mode & 0o777, 0o600) self.assertFalse(Connections(Path(directory)).values["llm"].key) store.set({**data, "baseUrl": "http://localhost:9000", "apiKey": ""}) with self.assertRaises(DecisionError): store.get("llm") with self.assertRaises(DecisionError): store.set({**data, "protocol": "codex"}) class ServerTests(unittest.IsolatedAsyncioTestCase): async def asyncSetUp(self): self.temp = tempfile.TemporaryDirectory() self.responses = [] self.received = [] self.started = asyncio.Event() self.release = asyncio.Event() self.block = False async def upstream(request): self.received.append({"path": request.path, "body": await request.json()}) self.started.set() if self.block: await self.release.wait() if self.responses: return self.responses.pop(0) return web.json_response( { "status": "completed", "output": [ { "type": "message", "content": [{"type": "output_text", "text": json.dumps(plan_value())}], } ], "usage": {"input_tokens": 10, "output_tokens": 20, "secret": "test-secret"}, } ) upstream_app = web.Application() upstream_app.router.add_post("/{path:.*}", upstream) self.upstream = TestServer(upstream_app) await self.upstream.start_server() app = create_app(Path(self.temp.name), "token", ["http://localhost:5173"]) self.client = TestClient(TestServer(app)) await self.client.start_server() self.service = app[STATE] self.service.hosts = {f"127.0.0.1:{self.client.port}"} self.headers = {"Authorization": "Bearer token", "Origin": "http://localhost:5173"} self.conn = { "role": "llm", "protocol": "responses", "baseUrl": str(self.upstream.make_url("/v1")), "model": "fixture", "apiKey": "test-secret", } self.service.connections.set(self.conn) async def asyncTearDown(self): self.release.set() await self.client.close() await self.upstream.close() self.temp.cleanup() async def post(self, path, data): return await self.client.post(PREFIX + path, json=data, headers=self.headers) async def test_host_origin_token_and_body(self): cases = [ ({}, 401), ({**self.headers, "Host": "evil.test"}, 403), ({**self.headers, "Origin": "https://evil.test"}, 403), (self.headers, 200), ] for headers, status in cases: response = await self.client.get(PREFIX + "/status", headers=headers) self.assertEqual(response.status, status) response = await self.client.options( PREFIX + "/plan", headers={"Origin": "http://localhost:5173"} ) self.assertEqual(response.status, 204) self.assertEqual(response.headers["Access-Control-Allow-Origin"], "http://localhost:5173") response = await self.client.post( PREFIX + "/plan", data="x" * 70000, headers={**self.headers, "Content-Type": "application/json"}, ) self.assertEqual(response.status, 413) response = await self.post("/codex/turn", {}) self.assertEqual(response.status, 405) async def test_responses_stamp_usage_and_duplicate(self): response = await self.post("/plan", request_value()) self.assertEqual(response.status, 200, await response.text()) value = await response.json() self.assertEqual(value["stamp"], observation()["stamp"]) self.assertEqual(value["value"], plan_value()) self.assertEqual(value["usage"], {"input_tokens": 10, "output_tokens": 20}) self.assertEqual(self.received[0]["body"]["tools"], []) self.assertEqual(self.received[0]["body"]["tool_choice"], "none") self.assertEqual((await self.post("/plan", request_value())).status, 409) state = await (await self.client.get(PREFIX + "/status", headers=self.headers)).text() self.assertNotIn("test-secret", state) self.assertNotIn("instruction", state) async def test_explicit_chat_protocol(self): self.service.connections.set({**self.conn, "protocol": "chat-completions"}) self.responses.append( web.json_response( { "choices": [ {"finish_reason": "stop", "message": {"content": json.dumps(plan_value())}} ] } ) ) response = await self.post("/plan", request_value()) self.assertEqual(response.status, 200, await response.text()) self.assertEqual(self.received[0]["path"], "/v1/chat/completions") self.assertIn("response_format", self.received[0]["body"]) async def test_typesafe_choices_and_probabilities(self): self.service.connections.set( { **self.conn, "role": "jev", "protocol": "typesafe", "baseUrl": str(self.upstream.make_url("/v1/systemone")), } ) chosen = { "choice": "open", "grasp": "uncertain", "diagnosis": "none", "recovery": "continue", } self.responses.append( web.json_response({"answers": {k: {"choice": v} for k, v in chosen.items()}}) ) response = await self.post( "/decide", {"observation": observation(), "candidates": ["open", "stop"]} ) self.assertEqual(response.status, 200, await response.text()) self.assertEqual((await response.json())["value"], {"version": VERSION, **chosen}) self.assertIn("questions", self.received[0]["body"]) self.assertNotIn("messages", self.received[0]["body"]) bad = {k: {"choice": v} for k, v in chosen.items()} bad["choice"] = {"choice": "carry"} self.responses.append(web.json_response({"answers": bad})) response = await self.post( "/decide", {"observation": observation("r2"), "candidates": ["open", "stop"]} ) self.assertEqual((await response.json())["error"], "jev_invalid_choice") async def test_openrouter_explicit_decisions_and_real_usage_only(self): self.service.connections.set( { **self.conn, "role": "jev", "protocol": "openrouter-decisions", "baseUrl": str(self.upstream.make_url("/api/alpha/decisions")), } ) selected = {"choice": "open", "grasp": "empty", "diagnosis": "none", "recovery": "continue"} self.responses.append( web.json_response( { "answers": {k: {"choice": v} for k, v in selected.items()}, "usage": {"cost": 0.00004, "input_tokens": 100, "untrusted": "secret"}, } ) ) response = await self.post( "/decide", {"observation": observation(), "candidates": ["open", "stop"]} ) result = await response.json() self.assertEqual(response.status, 200, result) self.assertEqual(result["usage"], {"cost": 0.00004, "input_tokens": 100}) self.assertEqual(self.received[0]["body"]["provider"], {"allow_fallbacks": False}) self.assertEqual(self.received[0]["path"], "/api/alpha/decisions") async def test_credentials_not_in_prompts_and_tools_not_executed(self): payload = request_value() payload["instruction"] = "do not disclose test-secret" await self.post("/plan", payload) self.assertNotIn("test-secret", json.dumps(self.received[0]["body"])) self.responses.append( web.json_response( { "status": "completed", "output": [ {"type": "function_call", "name": "shell", "arguments": "untrusted"} ], } ) ) result = await (await self.post("/plan", request_value("r2"))).json() self.assertEqual(result["error"], "llm_tool_or_unknown_output") async def test_http_and_bad_json_no_retry_no_secret_echo(self): for index, status in enumerate((401, 429, 302)): self.responses.append( web.Response(status=status, text="test-secret", headers={"Location": "/stolen"}) ) response = await self.post("/plan", request_value(str(index))) self.assertEqual((await response.json())["error"], f"upstream_http_{status}") self.assertEqual(len(self.received), index + 1) self.assertEqual((await self.post("/plan", request_value("budget"))).status, 429) self.assertEqual(len(self.received), 3) async def test_bad_contract_no_fallback(self): self.responses.append(web.Response(text="not JSON test-secret")) response = await self.post("/plan", request_value()) self.assertEqual((await response.json())["error"], "invalid_json") value = plan_value() value["steps"].reverse() self.responses.append( web.json_response( { "status": "completed", "output": [ { "type": "message", "content": [{"type": "output_text", "text": json.dumps(value)}], } ], } ) ) response = await self.post("/plan", request_value("r2")) self.assertEqual((await response.json())["error"], "invalid_plan_order") self.assertEqual(len(self.received), 2) async def test_cancel_reconfigure_and_count_failed_requests(self): self.block = True pending = asyncio.create_task(self.post("/plan", request_value())) await asyncio.wait_for(self.started.wait(), 3) self.assertEqual((await self.post("/plan", request_value("parallel"))).status, 409) response = await self.post("/cancel", {"runId": "run1", "requestId": "r1"}) self.assertTrue((await response.json())["cancelled"]) result = await pending self.assertEqual((await result.json())["error"], "request_cancelled") self.assertEqual(self.service.runs["run1"]["llm"], 1) self.assertFalse(self.service.active) self.started.clear() pending = asyncio.create_task(self.post("/plan", request_value("r2"))) await asyncio.wait_for(self.started.wait(), 3) response = await self.client.put( PREFIX + "/connections", json=self.conn, headers=self.headers ) self.assertEqual(response.status, 200) self.assertEqual((await (await pending).json())["error"], "request_cancelled") async def test_cancel_before_post_prevents_late_launch(self): response = await self.post("/cancel", {"runId": "run1", "requestId": "r1"}) self.assertEqual(response.status, 200) response = await self.post("/plan", request_value()) self.assertEqual((await response.json())["error"], "request_cancelled") self.assertEqual(self.received, []) async def test_timeout_budgets_and_redaction(self): with patch( "decision_server.providers.openai.plan", new=AsyncMock(side_effect=TimeoutError) ): response = await self.post("/plan", request_value()) self.assertEqual((await response.json())["error"], "request_timeout") value = plan_value() value["summary"] = "test-secret" with patch( "decision_server.providers.openai.plan", new=AsyncMock(return_value=(value, {})) ): response = await self.post("/plan", request_value("r2")) self.assertEqual((await response.json())["value"]["summary"], "[redacted]") for i in range(60): stamp = {**observation()["stamp"], "runId": "jev-run", "requestId": str(i)} self.service.admit("jev", stamp) with self.assertRaises(DecisionError): self.service.admit("jev", {**stamp, "requestId": "61"}) async def test_codex_planning_fails_closed(self): self.service.connections.values["llm"] = Connection("codex", "", "account-model") with patch.object( self.service.codex, "status", AsyncMock(return_value={"loggedIn": False}) ): response = await self.post("/plan", request_value()) self.assertEqual((await response.json())["error"], "codex_chatgpt_login_required") self.assertEqual(self.received, []) with self.assertRaises(DecisionError): await self.service.codex.rpc("command/exec", {}) class NativeCodexTests(unittest.IsolatedAsyncioTestCase): @unittest.skipUnless( os.environ.get("DECISION_CODEX_SMOKE") == "1", "opt-in: isolated native CLI, no login/turn" ) async def test_isolated_status_and_models(self): with tempfile.TemporaryDirectory() as directory: client = CodexAccount(Path(directory)) try: self.assertFalse((await client.status())["loggedIn"]) self.assertFalse((await client.models())["planningAvailable"]) self.assertFalse((Path(directory) / "home/auth.json").exists()) finally: await client.close()