# 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.background import background 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 _notify_gateway_status(is_open: bool) -> None: from devplacepy.database import get_table from devplacepy.utils import create_notification message = ( "AI gateway circuit breaker opened - upstream calls are being rejected" if is_open else "AI gateway circuit breaker closed - upstream calls have resumed" ) for admin in get_table("users").find(role="Admin"): create_notification(admin["uid"], "system", message, "gateway", "/admin/services/openai") 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 was_open = self._breaker.is_open self._breaker.record_failure() if not was_open and self._breaker.is_open: background.submit(_notify_gateway_status, True) 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 was_open = self._breaker.is_open self._breaker.record_failure() if not was_open and self._breaker.is_open: background.submit(_notify_gateway_status, True) else: was_open = self._breaker.is_open self._breaker.record_success() if was_open and not self._breaker.is_open: background.submit(_notify_gateway_status, False) 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, }