# retoor <retoor@molodetz.nl>
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,
}