"""Loopback-only bounded model gateway. No physics, arbitrary URL proxy or arbitrary RPC.""" import asyncio import hmac import secrets import time from collections import deque from pathlib import Path import aiohttp from aiohttp import web from .connections import Connections from .protocol import ( LANGUAGE_VERSION, VERSION, DecisionError, candidates, fields, loads, remaining_skills, schema_for, validate, ) from .providers import jev, openai from .providers.codex import CodexAccount from .providers.http import post, usage PREFIX = "/api/decision/v1" STATE = web.AppKey("decision_state", object) WEB_SERVICE = web.RequestKey("website_service", object) class Service: def __init__(self, directory, token, origins, port, *, connections=None): self.connections = connections if connections is not None else Connections(directory) self.token = token self.origins = set(origins) self.hosts = {f"127.0.0.1:{port}", f"localhost:{port}", f"[::1]:{port}"} self.codex = CodexAccount(directory / "codex") self.session = None self.active = {} self.runs = {} self.calls = deque() self.cancelled = {} self.records = deque(maxlen=100) self.epoch = 0 def invalidate(self): self.epoch += 1 for task in self.active.values(): task.cancel() def admit(self, role, stamp): now = time.monotonic() self.cancelled = {k: t for k, t in self.cancelled.items() if now - t < 3600} if (stamp["runId"], stamp["requestId"]) in self.cancelled: raise DecisionError("request_cancelled", 409) while self.calls and now - self.calls[0] > 3600: self.calls.popleft() if len(self.calls) >= 120: raise DecisionError("session_hourly_budget_exceeded", 429) # Bounded tombstones prevent reused request IDs or cancelled calls resetting budgets. self.runs = {k: v for k, v in self.runs.items() if now - v["start"] < 3600} run_id = stamp["runId"] run = self.runs.get(run_id) if run is None: if len(self.runs) >= 32: raise DecisionError("too_many_runs", 429) run = { "start": now, "llm": 0, "jev": 0, "ids": set(), "scene": stamp["sceneRevision"], "sequence": -1, "revision": -1, } self.runs[run_id] = run if now - run["start"] > 1200: raise DecisionError("run_wall_deadline_exceeded", 408) if ( stamp["requestId"] in run["ids"] or stamp["sceneRevision"] != run["scene"] or stamp["sequence"] < run["sequence"] or stamp["planRevision"] < run["revision"] ): raise DecisionError("stale_or_duplicate_request", 409) if run[role] >= (60 if role == "jev" else 3): raise DecisionError("run_request_budget_exceeded", 429) run[role] += 1 run["ids"].add(stamp["requestId"]) run["sequence"], run["revision"] = stamp["sequence"], stamp["planRevision"] self.calls.append(now) async def request(self, role, data): if self.active: raise DecisionError("request_already_running", 409) common = ["observation"] fields( data, common + (["instruction", "remaining"] if role == "llm" else ["candidates"]), [] if role == "llm" else ["failure"], ) data = self.connections.redact(data) raw = data["observation"] version = raw.get("version") if isinstance(raw, dict) else VERSION obs = validate("Observation", raw, version) if role == "llm": instruction = data["instruction"] if not isinstance(instruction, str) or not 1 <= len(instruction) <= 2000: raise DecisionError("invalid_instruction") remaining_skills(data["remaining"], version) else: candidates(data["candidates"], version) codes = schema_for(version)["$defs"]["SkillResult"]["properties"]["code"]["enum"] if data.get("failure", "none") not in codes: raise DecisionError("invalid_failure_code") conn = self.connections.get(role) async def invoke(): if conn.protocol == "codex": if version == LANGUAGE_VERSION: raise DecisionError("language_requires_configured_api", 409) return await self.codex.plan(data, conn.model) if role == "llm": return await openai.plan(self.session, conn, data) return await jev.decide(self.session, conn, data) return await self.execute(role, obs["stamp"], conn, invoke) async def execute(self, role, stamp, conn, invoke): if self.active: raise DecisionError("request_already_running", 409) self.admit(role, stamp) key = (stamp["runId"], stamp["requestId"]) epoch = self.epoch start = time.monotonic() code = "completed" task = asyncio.create_task(invoke()) self.active[key] = task try: value, metrics = await asyncio.wait_for(task, 60) if self.epoch != epoch: raise DecisionError("connection_changed", 409) # Summaries are the only free-form upstream strings; redact current credential values. value = self.connections.redact(value) return { "stamp": stamp, "value": value, "provider": conn.protocol, "model": conn.model, "usage": metrics, "elapsedMs": (time.monotonic() - start) * 1000, } except asyncio.CancelledError: code = "cancelled" raise DecisionError("request_cancelled", 409) from None except TimeoutError: code = "timeout" raise DecisionError("request_timeout", 504) from None except DecisionError as exc: code = exc.code raise except Exception: code = "internal_error" raise finally: if not task.done(): task.cancel() await asyncio.gather(task, return_exceptions=True) self.active.pop(key, None) self.records.append( { "role": role, "status": code, "elapsedMs": round((time.monotonic() - start) * 1000), } ) def service_for(request): return request.get(WEB_SERVICE) or request.app[STATE] @web.middleware async def boundary(request, handler): service = service_for(request) # Exact Host and Origin checks precede authentication and even OPTIONS; no DNS wildcard. if request.headers.get("Host", "") not in service.hosts: return web.json_response({"error": "host_forbidden"}, status=403) origin = request.headers.get("Origin") if origin and origin not in service.origins: return web.json_response({"error": "origin_forbidden"}, status=403) if request.method == "OPTIONS": response = web.Response(status=204) elif not hmac.compare_digest( request.headers.get("Authorization", "").encode(), ("Bearer " + service.token).encode() ): response = web.json_response({"error": "token_required"}, status=401) else: try: response = await handler(request) except DecisionError as exc: response = web.json_response({"error": exc.code}, status=exc.status) except web.HTTPException as exc: response = web.json_response({"error": "http_request_rejected"}, status=exc.status) except Exception: # Never echo raw provider responses, config, tracebacks, account URLs or headers. response = web.json_response({"error": "internal_error"}, status=500) response.headers.update({"Cache-Control": "no-store", "X-Content-Type-Options": "nosniff"}) if origin: response.headers.update( { "Access-Control-Allow-Origin": origin, "Vary": "Origin", "Access-Control-Allow-Headers": "Authorization,Content-Type", "Access-Control-Allow-Methods": "GET,POST,PUT,OPTIONS", } ) return response async def body(request): if request.content_type != "application/json": raise DecisionError("json_content_type_required", 415) try: return loads(await request.text()) except UnicodeError: raise DecisionError("invalid_encoding") from None async def status(request): s = service_for(request) return web.json_response( { "version": "lekiwi-agent-v1", "keyStorage": "memory-only", "connections": {k: v.public() for k, v in s.connections.values.items()}, "active": len(s.active), "records": list(s.records), "codexCheckedModels": sorted(s.codex.checked), } ) async def configure(request): s = service_for(request) data = await body(request) result = s.connections.set(data) s.invalidate() return web.json_response(result) async def plan(request): return web.json_response(await service_for(request).request("llm", await body(request))) async def decide(request): return web.json_response(await service_for(request).request("jev", await body(request))) async def cancel(request): data = fields(await body(request), ["runId", "requestId"]) if any(not isinstance(v, str) or len(v) > 128 for v in data.values()): raise DecisionError("invalid_request_id") service = service_for(request) key = (data["runId"], data["requestId"]) service.cancelled[key] = time.monotonic() if len(service.cancelled) > 256: service.cancelled.pop(next(iter(service.cancelled))) task = service.active.get(key) if task: task.cancel() return web.json_response({"cancelled": task is not None}) async def test_connection(request): s = service_for(request) data = fields(await body(request), ["role"]) if data["role"] not in ("llm", "jev"): raise DecisionError("invalid_role") if s.active: raise DecisionError("request_already_running", 409) conn = s.connections.get(data["role"]) # Tests are explicit billable calls, count against the same session budget, no retries. stamp = { "runId": "connection-tests", "requestId": secrets.token_hex(12), "sceneRevision": 0, "sequence": 0, "planRevision": 0, } async def invoke(): if conn.protocol == "codex": return await s.codex.check_model(conn.model), {} if data["role"] == "llm": value, metrics = await openai.structured( s.session, conn, 'Return {"ok":true}.', { "type": "object", "additionalProperties": False, "required": ["ok"], "properties": {"ok": {"type": "boolean"}}, }, ) if value != {"ok": True}: raise DecisionError("connection_test_invalid_response", 502) else: result = await post( s.session, conn, "", { "model": conn.model, "state": "Connection test. Choose ok.", **( {"provider": {"allow_fallbacks": False}} if conn.protocol == "openrouter-decisions" else {} ), "questions": {"test": jev.question(["ok"], "Choose ok.")}, }, ) jev.answer(result, "test", ["ok"]) metrics = usage(result) return {"ok": True}, metrics return web.json_response(await s.execute(data["role"], stamp, conn, invoke)) async def codex_operation(request): codex = service_for(request).codex operations = { "status": codex.status, "login": codex.login, "cancel": codex.cancel_login, "logout": codex.logout, "models": codex.models, "limits": codex.limits, } name = request.match_info["operation"] if name not in operations: raise DecisionError("unknown_codex_operation", 404) if request.method == "POST": fields(await body(request), []) return web.json_response(await operations[name]()) def create_app(directory=None, token=None, origins=None, port=8768): directory = (directory or Path.home() / ".local/state/mujoco-decision").expanduser().resolve() service = Service( directory, token or secrets.token_urlsafe(32), origins or { "http://localhost:5173", "http://127.0.0.1:5173", "http://localhost:4173", "http://127.0.0.1:4173", }, port, ) app = web.Application(middlewares=[boundary], client_max_size=65536) app[STATE] = service app.router.add_get(PREFIX + "/status", status) app.router.add_put(PREFIX + "/connections", configure) app.router.add_post(PREFIX + "/test", test_connection) async def command(request): from .commands import command as run_command return web.json_response(await run_command(service_for(request), await body(request))) app.router.add_post(PREFIX + "/command", command) app.router.add_post(PREFIX + "/plan", plan) app.router.add_post(PREFIX + "/decide", decide) app.router.add_post(PREFIX + "/cancel", cancel) app.router.add_get(PREFIX + "/codex/{operation:status|models|limits}", codex_operation) app.router.add_post(PREFIX + "/codex/{operation:login|cancel|logout}", codex_operation) async def options(_): return web.Response(status=204) app.router.add_route("OPTIONS", PREFIX + "/{path:.*}", options) async def lifecycle(_): async with aiohttp.ClientSession(trust_env=False) as session: service.session = session yield service.invalidate() await asyncio.gather(*service.active.values(), return_exceptions=True) await service.codex.close() app.cleanup_ctx.append(lifecycle) return app