# retoor import asyncio import json import logging import time import uuid from typing import Optional import httpx from fastapi.responses import JSONResponse, Response, StreamingResponse from devplacepy import stealth from devplacepy.services.openai_gateway import config from devplacepy.services.openai_gateway.reliability import CircuitBreaker, retry_send from devplacepy.services.openai_gateway.routing import ( chat_overlay, embed_overlay, image_overlay, ) from devplacepy.services.openai_gateway.system_message import apply_system_directives from devplacepy.services.openai_gateway.usage import ( GatewayUsageLedger, classify_error, extract_image_usage, extract_params, parse_context_map, pricing_from_cfg, usage_response_headers, ) from devplacepy.services.openai_gateway.vision import VisionAugmenter, VisionCache logger = logging.getLogger(__name__) def _fake_stream(data: dict, model: str, include_usage: bool = False): chunk_id = data.get("id", f"chatcmpl-{uuid.uuid4().hex[:12]}") created = data.get("created", int(time.time())) out_model = data.get("model", model) try: msg = data["choices"][0]["message"] except (KeyError, IndexError): msg = {"content": ""} tool_calls = msg.get("tool_calls") content = msg.get("content") or "" reasoning_content = msg.get("reasoning_content") or "" def _chunk(delta: dict, finish: Optional[str] = None) -> str: return ( "data: " + json.dumps( { "id": chunk_id, "object": "chat.completion.chunk", "created": created, "model": out_model, "choices": [{"index": 0, "delta": delta, "finish_reason": finish}], } ) + "\n\n" ) async def gen(): yield _chunk({"role": "assistant"}) if reasoning_content: for i in range(0, len(reasoning_content), 50): yield _chunk({"reasoning_content": reasoning_content[i : i + 50]}) if tool_calls: yield _chunk({"tool_calls": tool_calls}) elif content: for i in range(0, len(content), 50): yield _chunk({"content": content[i : i + 50]}) yield _chunk({}, finish="tool_calls" if tool_calls else "stop") if include_usage and data.get("usage"): yield ( "data: " + json.dumps( { "id": chunk_id, "object": "chat.completion.chunk", "created": created, "model": out_model, "choices": [], "usage": data["usage"], } ) + "\n\n" ) yield "data: [DONE]\n\n" return gen() def _connect_tracer(holder: dict): started: dict = {} async def trace(name: str, info: dict) -> None: if name.endswith("connect_tcp.started") or name.endswith("start_tls.started"): started[name] = time.monotonic() elif name.endswith("connect_tcp.complete") or name.endswith( "start_tls.complete" ): begin = started.get(name.replace(".complete", ".started")) if begin is not None: holder["ms"] += (time.monotonic() - begin) * 1000 return trace class GatewayRuntime: def __init__(self): self._client: Optional[httpx.AsyncClient] = None self._sem: Optional[asyncio.Semaphore] = None self._instances = 0 self._timeout = 0 self._vision_cache: Optional[VisionCache] = None self._vision_cache_size = -1 self._ledger = GatewayUsageLedger() self._breaker = CircuitBreaker( config.CIRCUIT_THRESHOLD_DEFAULT, config.CIRCUIT_COOLDOWN_SECONDS_DEFAULT ) self.requests = 0 self.errors = 0 self.in_flight = 0 self.peak_in_flight = 0 self.vision_calls = 0 self.embed_calls = 0 self.image_calls = 0 self.last_status = 0 self.last_latency_ms = 0 def _ensure(self, cfg: dict): instances = max(1, cfg["gateway_instances"]) timeout = max(1, cfg["gateway_timeout"]) if ( self._client is None or instances != self._instances or timeout != self._timeout ): old = self._client limits = httpx.Limits( max_connections=instances, max_keepalive_connections=instances ) self._client = stealth.stealth_async_client( timeout=float(timeout), limits=limits ) self._sem = asyncio.Semaphore(instances) self._instances = instances self._timeout = timeout if old is not None: asyncio.create_task(old.aclose()) size = cfg["gateway_vision_cache_size"] if self._vision_cache is None or size != self._vision_cache_size: self._vision_cache = VisionCache(size) self._vision_cache_size = size self._breaker.configure( cfg["gateway_circuit_threshold"], cfg["gateway_circuit_cooldown_seconds"] ) return self._client, self._sem async def aclose(self) -> None: if self._client is not None: await self._client.aclose() self._client = None self._instances = 0 async def _send( self, client, sem, method, url, headers, cfg, log, json_body=None, content=None ): timing = { "queue_wait_ms": 0.0, "upstream_latency_ms": 0.0, "connect_ms": 0.0, "retries_attempted": 0, "retry_succeeded": False, "circuit_open": False, } if not self._breaker.allow(): timing["circuit_open"] = True log("circuit breaker open, rejecting upstream call") return None, None, timing self.requests += 1 self.in_flight += 1 if self.in_flight > self.peak_in_flight: self.peak_in_flight = self.in_flight wait_start = time.monotonic() connect_holder = {"ms": 0.0} attempts = 1 resp = None exc = None try: async with sem: timing["queue_wait_ms"] = round( (time.monotonic() - wait_start) * 1000, 3 ) async def do_call(): request = client.build_request( method, url, headers=headers, json=json_body, content=content ) request.extensions["trace"] = _connect_tracer(connect_holder) return await client.send(request) send_start = time.monotonic() resp, exc, attempts = await retry_send( do_call, cfg["gateway_max_retries"], cfg["gateway_retry_backoff_ms"], log, ) timing["upstream_latency_ms"] = round( (time.monotonic() - send_start) * 1000, 3 ) finally: self.in_flight -= 1 timing["connect_ms"] = round(connect_holder["ms"], 3) timing["retries_attempted"] = max(attempts - 1, 0) self.last_latency_ms = int(timing["upstream_latency_ms"]) if exc is not None: self.errors += 1 self._breaker.record_failure() log(f"{method} {url} connection failed after {attempts} attempt(s): {exc}") return None, exc, timing self.last_status = resp.status_code if resp.status_code >= 500: self.errors += 1 self._breaker.record_failure() else: self._breaker.record_success() timing["retry_succeeded"] = attempts > 1 return resp, None, timing async def handle_chat( self, body: dict, cfg: dict, owner: tuple, user_agent: str, app_reference: str, log=None ): log = log or (lambda message: None) overlay = chat_overlay(body.get("model"), cfg) if overlay: cfg = {**cfg, **overlay} log(f"routed model {body.get('model')!r} -> {cfg['gateway_model']!r}") client, sem = self._ensure(cfg) pricing = pricing_from_cfg(cfg) context_map = parse_context_map(cfg.get("gateway_model_context_map")) params = extract_params(body) handle_start = time.monotonic() messages = body.get("messages", []) or [] vision_cost = 0.0 if cfg["gateway_vision_enabled"]: augmenter = VisionAugmenter( cfg["gateway_vision_url"], cfg["gateway_vision_model"], cfg["gateway_vision_key"], self._vision_cache, ledger=self._ledger, owner=owner, pricing=pricing, context_map=context_map, app_reference=app_reference, ) messages = await augmenter.augment_messages(client, messages) self.vision_calls += augmenter.calls vision_cost = augmenter.cost_usd messages = apply_system_directives(messages, cfg.get("gateway_system_preamble", "")) requested = body.get("model") if cfg["gateway_force_model"] or not requested or requested == "molodetz": model = cfg["gateway_model"] elif overlay is not None: model = requested else: model = cfg["gateway_model"] log(f"requested model {requested!r} has no route, falling back to {model!r}") stream = bool(body.get("stream")) include_usage = bool((body.get("stream_options") or {}).get("include_usage")) payload = dict(body) payload["model"] = model payload["messages"] = messages payload["stream"] = False payload.pop("stream_options", None) headers = {"Content-Type": "application/json"} if cfg["gateway_api_key"]: headers["Authorization"] = f"Bearer {cfg['gateway_api_key']}" else: log( "No upstream API key configured (gateway_api_key / DEEPSEEK_API_KEY / OPENROUTER_API_KEY); upstream will likely reject the request" ) resp, exc, timing = await self._send( client, sem, "POST", cfg["gateway_upstream_url"], headers, cfg, log, json_body=payload, ) base = { "owner_kind": owner[0], "owner_id": owner[1], "backend": "chat", "endpoint": "chat/completions", "model": model, "user_agent": user_agent, "app_reference": app_reference, **params, **timing, } def finalize(status_code, success, category, usage=None): base["total_latency_ms"] = round( (time.monotonic() - handle_start) * 1000, 3 ) base["gateway_overhead_ms"] = round( max( base["total_latency_ms"] - timing["upstream_latency_ms"] - timing["queue_wait_ms"], 0.0, ), 3, ) base["status_code"] = status_code base["success"] = success base["error_category"] = category base["usage"] = usage row = self._ledger.record(base, pricing, context_map) headers = usage_response_headers(row) if vision_cost and headers: chat_cost = float((row or {}).get("cost_usd") or 0.0) headers["X-Gateway-Cost-USD"] = f"{chat_cost + vision_cost:.8f}" headers["X-Gateway-Vision-Cost-USD"] = f"{vision_cost:.8f}" return headers if timing["circuit_open"]: resp_headers = finalize(503, False, "circuit_open") return JSONResponse( status_code=503, content={ "error": { "message": "Upstream temporarily unavailable", "type": "circuit_open", } }, headers=resp_headers, ) if exc is not None: resp_headers = finalize(502, False, classify_error(0, exc)) return JSONResponse( status_code=502, content={ "error": { "message": f"Upstream connection failed: {exc}", "type": "upstream_error", } }, headers=resp_headers, ) if resp.status_code != 200: resp_headers = finalize( resp.status_code, False, classify_error(resp.status_code, None, resp.text), ) log(f"chat upstream POST -> {resp.status_code}: {resp.text[:300]}") return JSONResponse( status_code=resp.status_code, content={"error": {"message": resp.text, "type": "upstream_error"}}, headers=resp_headers, ) try: data = resp.json() except ValueError: self.errors += 1 resp_headers = finalize(502, False, "gateway") log("chat upstream returned 200 but body was not valid JSON") return JSONResponse( status_code=502, content={ "error": { "message": "invalid upstream response", "type": "upstream_error", } }, headers=resp_headers, ) resp_headers = finalize(200, True, None, data.get("usage")) log(f"chat POST -> 200 ({timing['upstream_latency_ms']:.0f}ms)") if stream: return StreamingResponse( _fake_stream(data, model, include_usage), media_type="text/event-stream", headers=resp_headers, ) return Response( content=resp.content, media_type="application/json", headers=resp_headers, ) async def handle_embeddings( self, body: dict, cfg: dict, owner: tuple, user_agent: str, app_reference: str, log=None ): log = log or (lambda message: None) vision_cost = 0.0 overlay = embed_overlay(body.get("model"), cfg) if overlay: cfg = {**cfg, **overlay} log(f"routed embed model {body.get('model')!r} -> {cfg['gateway_embed_model']!r}") if not cfg["gateway_embed_enabled"]: from devplacepy.services.audit import record as audit from devplacepy.services.openai_gateway.usage import audit_actor_for actor_kind, actor_uid, actor_role = audit_actor_for(owner[0], owner[1]) audit.record_system( "ai.gateway.call", actor_kind=actor_kind, actor_uid=actor_uid, actor_role=actor_role, origin="api", result="denied", summary="embeddings disabled", metadata={ "backend": "embed", "endpoint": "embeddings", "owner_kind": owner[0], "owner_id": owner[1], }, ) return JSONResponse( status_code=503, content={ "error": { "message": "Embeddings are disabled", "type": "embeddings_disabled", } }, ) client, sem = self._ensure(cfg) pricing = pricing_from_cfg(cfg) context_map = parse_context_map(cfg.get("gateway_model_context_map")) params = extract_params(body) handle_start = time.monotonic() requested = body.get("model") if cfg["gateway_force_model"] or not requested or requested == "molodetz~embed": model = cfg["gateway_embed_model"] elif overlay is not None: model = requested else: model = cfg["gateway_embed_model"] log(f"requested embed model {requested!r} has no route, falling back to {model!r}") payload = dict(body) payload["model"] = model headers = {"Content-Type": "application/json"} if cfg["gateway_embed_key"]: headers["Authorization"] = f"Bearer {cfg['gateway_embed_key']}" else: log( "No upstream embeddings API key configured (gateway_embed_key / gateway_vision_key / OPENROUTER_API_KEY); upstream will likely reject the request" ) resp, exc, timing = await self._send( client, sem, "POST", cfg["gateway_embed_url"], headers, cfg, log, json_body=payload, ) base = { "owner_kind": owner[0], "owner_id": owner[1], "backend": "embed", "endpoint": "embeddings", "model": model, "user_agent": user_agent, "app_reference": app_reference, **params, **timing, } def finalize(status_code, success, category, usage=None): base["total_latency_ms"] = round( (time.monotonic() - handle_start) * 1000, 3 ) base["gateway_overhead_ms"] = round( max( base["total_latency_ms"] - timing["upstream_latency_ms"] - timing["queue_wait_ms"], 0.0, ), 3, ) base["status_code"] = status_code base["success"] = success base["error_category"] = category base["usage"] = usage row = self._ledger.record(base, pricing, context_map) headers = usage_response_headers(row) if vision_cost and headers: chat_cost = float((row or {}).get("cost_usd") or 0.0) headers["X-Gateway-Cost-USD"] = f"{chat_cost + vision_cost:.8f}" headers["X-Gateway-Vision-Cost-USD"] = f"{vision_cost:.8f}" return headers if timing["circuit_open"]: resp_headers = finalize(503, False, "circuit_open") return JSONResponse( status_code=503, content={ "error": { "message": "Upstream temporarily unavailable", "type": "circuit_open", } }, headers=resp_headers, ) if exc is not None: resp_headers = finalize(502, False, classify_error(0, exc)) return JSONResponse( status_code=502, content={ "error": { "message": f"Upstream connection failed: {exc}", "type": "upstream_error", } }, headers=resp_headers, ) if resp.status_code != 200: resp_headers = finalize( resp.status_code, False, classify_error(resp.status_code, None, resp.text), ) log(f"embed upstream POST -> {resp.status_code}: {resp.text[:300]}") return JSONResponse( status_code=resp.status_code, content={"error": {"message": resp.text, "type": "upstream_error"}}, headers=resp_headers, ) try: data = resp.json() except ValueError: self.errors += 1 resp_headers = finalize(502, False, "gateway") log("embed upstream returned 200 but body was not valid JSON") return JSONResponse( status_code=502, content={ "error": { "message": "invalid upstream response", "type": "upstream_error", } }, headers=resp_headers, ) self.embed_calls += 1 resp_headers = finalize(200, True, None, data.get("usage")) log(f"embed POST -> 200 ({timing['upstream_latency_ms']:.0f}ms)") return JSONResponse(content=data, headers=resp_headers) async def handle_images( self, body: dict, cfg: dict, owner: tuple, user_agent: str, app_reference: str, log=None ): log = log or (lambda message: None) overlay = image_overlay(body.get("model"), cfg) if overlay: cfg = {**cfg, **overlay} log( f"routed image model {body.get('model')!r} -> {cfg['gateway_image_model']!r}" ) if not cfg["gateway_image_enabled"]: from devplacepy.services.audit import record as audit from devplacepy.services.openai_gateway.usage import audit_actor_for actor_kind, actor_uid, actor_role = audit_actor_for(owner[0], owner[1]) audit.record_system( "ai.gateway.call", actor_kind=actor_kind, actor_uid=actor_uid, actor_role=actor_role, origin="api", result="denied", summary="image generation disabled", metadata={ "backend": "image", "endpoint": "images/generations", "owner_kind": owner[0], "owner_id": owner[1], }, ) return JSONResponse( status_code=503, content={ "error": { "message": "Image generation is disabled", "type": "images_disabled", } }, ) client, sem = self._ensure(cfg) pricing = pricing_from_cfg(cfg) context_map = parse_context_map(cfg.get("gateway_model_context_map")) params = extract_params(body) handle_start = time.monotonic() requested = body.get("model") if ( cfg["gateway_force_model"] or not requested or requested in ("molodetz-img-small", "molodetz-img") ): model = cfg["gateway_image_model"] elif overlay is not None: model = requested else: model = cfg["gateway_image_model"] log(f"requested image model {requested!r} has no route, falling back to {model!r}") payload = dict(body) payload["model"] = model headers = {"Content-Type": "application/json"} if cfg["gateway_image_key"]: headers["Authorization"] = f"Bearer {cfg['gateway_image_key']}" else: log( "No upstream image API key configured (gateway_image_key / gateway_api_key / OPENROUTER_API_KEY); upstream will likely reject the request" ) resp, exc, timing = await self._send( client, sem, "POST", cfg["gateway_image_url"], headers, cfg, log, json_body=payload, ) base = { "owner_kind": owner[0], "owner_id": owner[1], "backend": "image", "endpoint": "images/generations", "model": model, "user_agent": user_agent, "app_reference": app_reference, **params, **timing, } def finalize(status_code, success, category, usage=None): base["total_latency_ms"] = round( (time.monotonic() - handle_start) * 1000, 3 ) base["gateway_overhead_ms"] = round( max( base["total_latency_ms"] - timing["upstream_latency_ms"] - timing["queue_wait_ms"], 0.0, ), 3, ) base["status_code"] = status_code base["success"] = success base["error_category"] = category base["usage"] = usage row = self._ledger.record(base, pricing, context_map) return usage_response_headers(row) if timing["circuit_open"]: resp_headers = finalize(503, False, "circuit_open") return JSONResponse( status_code=503, content={ "error": { "message": "Upstream temporarily unavailable", "type": "circuit_open", } }, headers=resp_headers, ) if exc is not None: resp_headers = finalize(502, False, classify_error(0, exc)) return JSONResponse( status_code=502, content={ "error": { "message": f"Upstream connection failed: {exc}", "type": "upstream_error", } }, headers=resp_headers, ) if resp.status_code != 200: resp_headers = finalize( resp.status_code, False, classify_error(resp.status_code, None, resp.text), ) log(f"image upstream POST -> {resp.status_code}: {resp.text[:300]}") return JSONResponse( status_code=resp.status_code, content={"error": {"message": resp.text, "type": "upstream_error"}}, headers=resp_headers, ) try: data = resp.json() except ValueError: self.errors += 1 resp_headers = finalize(502, False, "gateway") log("image upstream returned 200 but body was not valid JSON") return JSONResponse( status_code=502, content={ "error": { "message": "invalid upstream response", "type": "upstream_error", } }, headers=resp_headers, ) self.image_calls += 1 resp_headers = finalize(200, True, None, extract_image_usage(data)) log(f"image POST -> 200 ({timing['upstream_latency_ms']:.0f}ms)") return JSONResponse(content=data, headers=resp_headers) async def handle_passthrough( self, method: str, subpath: str, content_type: str, body: bytes, cfg: dict, owner: tuple, user_agent: str, app_reference: str, log=None, ): log = log or (lambda message: None) vision_cost = 0.0 client, sem = self._ensure(cfg) pricing = pricing_from_cfg(cfg) context_map = parse_context_map(cfg.get("gateway_model_context_map")) handle_start = time.monotonic() base_url = cfg["gateway_upstream_url"] if base_url.endswith("/chat/completions"): base_url = base_url[: -len("/chat/completions")] url = f"{base_url.rstrip('/')}/{subpath}" headers = {} if cfg["gateway_api_key"]: headers["Authorization"] = f"Bearer {cfg['gateway_api_key']}" if content_type: headers["Content-Type"] = content_type resp, exc, timing = await self._send( client, sem, method, url, headers, cfg, log, content=body ) base = { "owner_kind": owner[0], "owner_id": owner[1], "backend": "chat", "endpoint": subpath, "model": cfg["gateway_model"], "user_agent": user_agent, "app_reference": app_reference, **timing, } def finalize(status_code, success, category, usage=None): base["total_latency_ms"] = round( (time.monotonic() - handle_start) * 1000, 3 ) base["gateway_overhead_ms"] = round( max( base["total_latency_ms"] - timing["upstream_latency_ms"] - timing["queue_wait_ms"], 0.0, ), 3, ) base["status_code"] = status_code base["success"] = success base["error_category"] = category base["usage"] = usage row = self._ledger.record(base, pricing, context_map) headers = usage_response_headers(row) if vision_cost and headers: chat_cost = float((row or {}).get("cost_usd") or 0.0) headers["X-Gateway-Cost-USD"] = f"{chat_cost + vision_cost:.8f}" headers["X-Gateway-Vision-Cost-USD"] = f"{vision_cost:.8f}" return headers if timing["circuit_open"]: resp_headers = finalize(503, False, "circuit_open") return JSONResponse( status_code=503, content={ "error": { "message": "Upstream temporarily unavailable", "type": "circuit_open", } }, headers=resp_headers, ) if exc is not None: resp_headers = finalize(502, False, classify_error(0, exc)) return JSONResponse( status_code=502, content={ "error": { "message": f"Upstream connection failed: {exc}", "type": "upstream_error", } }, headers=resp_headers, ) usage = None if resp.status_code < 400 and "application/json" in ( resp.headers.get("content-type") or "" ): try: usage = resp.json().get("usage") except ValueError: usage = None resp_headers = finalize( resp.status_code, resp.status_code < 400, None if resp.status_code < 400 else classify_error(resp.status_code, None, resp.text), usage, ) log( f"passthrough {method} {url} -> {resp.status_code} ({timing['upstream_latency_ms']:.0f}ms)" ) return Response( content=resp.content, status_code=resp.status_code, media_type=resp.headers.get("content-type"), headers=resp_headers, ) def metrics(self) -> dict: return { "requests": self.requests, "errors": self.errors, "in_flight": self.in_flight, "peak_in_flight": self.peak_in_flight, "vision_calls": self.vision_calls, "embed_calls": self.embed_calls, "image_calls": self.image_calls, "last_status": self.last_status, "last_latency_ms": self.last_latency_ms, "pool": self._instances, "circuit_open": self._breaker.is_open, }