Files
devplacepy/devplacepy/services/openai_gateway/gateway.py
T
retoorandClaude Sonnet 5 67c85e4184 Fix gateway model fallback silently skipped for unrouted client model names
The per-route fallback_model mechanism (used to fail over to a different
provider when the primary upstream errors, e.g. insufficient balance)
looked up the fallback keyed on the client's raw, literal "model" string.
That string only matches a configured route when the caller sends the
exact alias ("molodetz"/"molodetz~embed"/"molodetz-img-small") or another
exact route name - any other value (the common case for external agents,
which rarely echo DevPlace's own alias) resolves no route at all, so
resolve_fallback() returned None and a real 402/5xx from the primary
upstream propagated straight to the caller even with a fallback configured
on the default route.

Fixed in all three handlers (chat/embed/image): the fallback lookup key
now reflects whether a route actually matched the raw request (chat_overlay
result), independent of gateway_force_model - which is always forced true
by a successful overlay and therefore cannot be used to tell "matched a
specific route" apart from "used the default". An unmatched request now
normalizes to the default alias for fallback purposes, so its configured
fallback_model is consulted instead of silently skipped.

Added a regression test that reproduces the exact failure (an unrouted
client model name, primary upstream returns 402, fallback configured on
the default route) and confirms it now recovers instead of surfacing the
402 to the caller.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-09-03 18:56:05 +00:00

1114 lines
41 KiB
Python

# retoor <retoor@molodetz.nl>
import asyncio
import json
import logging
import time
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,
resolve_fallback,
)
from devplacepy.services.openai_gateway.system_message import apply_system_directives
from devplacepy.services.openai_gateway.thinking import apply_thinking
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 _call_failed(resp, exc, timing: dict) -> bool:
if timing.get("circuit_open") or exc is not None:
return True
return resp is not None and resp.status_code >= 400
def _fallback_headers(cfg: dict, key_field: str) -> dict:
headers = {"Content-Type": "application/json"}
if cfg.get(key_field):
headers["Authorization"] = f"Bearer {cfg[key_field]}"
return headers
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 _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,
stream=False,
):
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
connect_holder = {"ms": 0.0}
attempts = 1
resp = None
exc = None
try:
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, stream=stream)
send_start = time.monotonic()
resp, exc, attempts, queue_wait_ms = await retry_send(
do_call,
sem,
cfg["gateway_max_retries"],
cfg["gateway_retry_backoff_ms"],
log,
)
timing["upstream_latency_ms"] = round(
(time.monotonic() - send_start) * 1000, 3
)
timing["queue_wait_ms"] = round(queue_wait_ms, 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)
base_cfg = 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}")
# Fallback lookup is keyed on whichever route (if any) actually matched
# the client's raw request, never on the *effective* model: overlay
# unconditionally forces gateway_force_model, so "was a specific route
# matched" cannot be read off that flag - it is always true once any
# route (including "molodetz" itself) has matched. An unmatched request
# normalizes to "molodetz" so its own configured fallback still protects
# the default path, instead of silently never being consulted.
fallback_key = requested if overlay is not None else "molodetz"
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"] = stream
if stream:
# Always ask the upstream for usage on its final chunk, regardless of
# whether the client itself requested it, so the ledger always has real
# cost/token numbers for a streamed call; `include_usage` (the client's
# own ask) only controls whether that chunk is relayed to the client.
stream_options = dict(body.get("stream_options") or {})
stream_options["include_usage"] = True
payload["stream_options"] = stream_options
else:
payload.pop("stream_options", None)
apply_thinking(
payload,
cfg.get("gateway_upstream_url", ""),
default_enabled=bool(cfg.get("gateway_thinking", False)),
)
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"
)
send_start = time.monotonic()
resp, exc, timing = await self._send(
client,
sem,
"POST",
cfg["gateway_upstream_url"],
headers,
cfg,
log,
json_body=payload,
stream=stream,
)
if resp is not None and stream and resp.status_code != 200:
# The streamed response body hasn't been read yet; the shared error
# handling below needs resp.text/.content, which requires an explicit
# read for a stream=True response (a no-op if already buffered).
await resp.aread()
requested_model = requested
if _call_failed(resp, exc, timing):
fallback_route = resolve_fallback(fallback_key, "chat")
fallback_overlay = (
chat_overlay(fallback_route.source_model, base_cfg)
if fallback_route is not None
else None
)
if fallback_overlay:
log(
f"model {requested!r} failed, falling back to "
f"{fallback_route.source_model!r} -> {fallback_overlay['gateway_model']!r}"
)
cfg = {**base_cfg, **fallback_overlay}
model = cfg["gateway_model"]
payload = dict(payload)
payload["model"] = model
headers = _fallback_headers(cfg, "gateway_api_key")
send_start = time.monotonic()
resp, exc, timing = await self._send(
client,
sem,
"POST",
cfg["gateway_upstream_url"],
headers,
cfg,
log,
json_body=payload,
stream=stream,
)
if resp is not None and stream and resp.status_code != 200:
await resp.aread()
pricing = pricing_from_cfg(cfg)
context_map = parse_context_map(cfg.get("gateway_model_context_map"))
base = {
"owner_kind": owner[0],
"owner_id": owner[1],
"backend": "chat",
"endpoint": "chat/completions",
"requested_model": requested_model,
"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,
)
if stream:
# Real upstream streaming: forward SSE chunks to the client as they
# arrive (so TTFT/inter-token latency are genuine), and finalize the
# ledger row from inside the generator once the stream ends - cost and
# token counts are not known yet at this point, so unlike the
# non-streaming response below, this response carries no
# X-Gateway-Cost-USD/token headers (HTTP headers must precede the body).
log(f"chat POST -> 200 stream ({timing['upstream_latency_ms']:.0f}ms to first byte)")
return StreamingResponse(
self._stream_chat_response(
resp,
base,
pricing,
context_map,
include_usage,
handle_start,
send_start,
log,
),
media_type="text/event-stream",
headers={
"X-Gateway-Model": model,
"X-Gateway-Backend": "chat",
"X-App-Reference": app_reference or "default",
},
)
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)")
return Response(
content=resp.content,
media_type="application/json",
headers=resp_headers,
)
async def _stream_chat_response(
self,
resp: httpx.Response,
base: dict,
pricing,
context_map: dict,
include_usage: bool,
handle_start: float,
send_start: float,
log,
):
first_at: Optional[float] = None
last_at: Optional[float] = None
content_chunks = 0
usage_captured: Optional[dict] = None
saw_done = False
success = True
error_category: Optional[str] = None
try:
async for raw_line in resp.aiter_lines():
line = raw_line.strip()
if not line or not line.startswith("data:"):
continue
now = time.monotonic()
if first_at is None:
first_at = now
last_at = now
data_str = line[len("data:") :].strip()
if data_str == "[DONE]":
saw_done = True
yield "data: [DONE]\n\n"
break
try:
chunk = json.loads(data_str)
except ValueError:
continue
choices = chunk.get("choices") or []
delta = (choices[0].get("delta") if choices else None) or {}
if delta.get("content") or delta.get("reasoning_content"):
content_chunks += 1
usage = chunk.get("usage")
if usage:
usage_captured = usage
if not include_usage:
# The client did not ask for the usage chunk; swallow it,
# we still captured it above for the ledger.
continue
yield f"data: {json.dumps(chunk)}\n\n"
except GeneratorExit:
success = False
error_category = "client_disconnected"
raise
except Exception as exc: # noqa: BLE001 - an upstream stream can drop mid-flight; never crash the response
success = False
error_category = classify_error(0, exc)
log(f"chat stream interrupted: {exc}")
finally:
try:
await resp.aclose()
except Exception: # noqa: BLE001 - releasing the connection must never mask the real outcome
pass
ttft_ms = (
round((first_at - send_start) * 1000, 3)
if first_at is not None
else None
)
inter_token_ms = (
round((last_at - first_at) * 1000 / (content_chunks - 1), 3)
if first_at is not None and last_at is not None and content_chunks > 1
else None
)
base["total_latency_ms"] = round(
(time.monotonic() - handle_start) * 1000, 3
)
base["gateway_overhead_ms"] = round(
max(
base["total_latency_ms"]
- base.get("upstream_latency_ms", 0.0)
- base.get("queue_wait_ms", 0.0),
0.0,
),
3,
)
base["status_code"] = 200
base["success"] = success
base["error_category"] = error_category
base["usage"] = usage_captured
base["ttft_ms"] = ttft_ms
base["inter_token_ms"] = inter_token_ms
self._ledger.record(base, pricing, context_map)
if not saw_done:
yield "data: [DONE]\n\n"
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)
base_cfg = 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)
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}")
# See handle_chat's matching comment: fallback lookup must key off
# whether a route actually matched, not off gateway_force_model, which
# a matched overlay always forces true.
fallback_key = requested if overlay is not None else "molodetz~embed"
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,
)
requested_model = requested
if _call_failed(resp, exc, timing):
fallback_route = resolve_fallback(fallback_key, "embed")
fallback_overlay = (
embed_overlay(fallback_route.source_model, base_cfg)
if fallback_route is not None
else None
)
if fallback_overlay:
log(
f"embed model {requested!r} failed, falling back to "
f"{fallback_route.source_model!r} -> {fallback_overlay['gateway_embed_model']!r}"
)
cfg = {**base_cfg, **fallback_overlay}
model = cfg["gateway_embed_model"]
payload = dict(payload)
payload["model"] = model
headers = _fallback_headers(cfg, "gateway_embed_key")
resp, exc, timing = await self._send(
client,
sem,
"POST",
cfg["gateway_embed_url"],
headers,
cfg,
log,
json_body=payload,
)
pricing = pricing_from_cfg(cfg)
context_map = parse_context_map(cfg.get("gateway_model_context_map"))
base = {
"owner_kind": owner[0],
"owner_id": owner[1],
"backend": "embed",
"endpoint": "embeddings",
"requested_model": requested_model,
"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)
base_cfg = 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)
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}")
# See handle_chat's matching comment: fallback lookup must key off
# whether a route actually matched, not off gateway_force_model, which
# a matched overlay always forces true.
fallback_key = requested if overlay is not None else "molodetz-img-small"
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,
)
requested_model = requested
if _call_failed(resp, exc, timing):
fallback_route = resolve_fallback(fallback_key, "image")
fallback_overlay = (
image_overlay(fallback_route.source_model, base_cfg)
if fallback_route is not None
else None
)
if fallback_overlay:
log(
f"image model {requested!r} failed, falling back to "
f"{fallback_route.source_model!r} -> {fallback_overlay['gateway_image_model']!r}"
)
cfg = {**base_cfg, **fallback_overlay}
model = cfg["gateway_image_model"]
payload = dict(payload)
payload["model"] = model
headers = _fallback_headers(cfg, "gateway_image_key")
resp, exc, timing = await self._send(
client,
sem,
"POST",
cfg["gateway_image_url"],
headers,
cfg,
log,
json_body=payload,
)
pricing = pricing_from_cfg(cfg)
context_map = parse_context_map(cfg.get("gateway_model_context_map"))
base = {
"owner_kind": owner[0],
"owner_id": owner[1],
"backend": "image",
"endpoint": "images/generations",
"requested_model": requested_model,
"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,
}