Skip to content
Merged
Show file tree
Hide file tree
Changes from 2 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 24 additions & 0 deletions docs/api-reference.md
Original file line number Diff line number Diff line change
Expand Up @@ -249,6 +249,30 @@ Gateway-side failures use fixed public messages. Diagnose them with protected
logs and safe metadata such as request ID, provider, model, and status. Do not
log provider keys, prompts, responses, or raw upstream bodies.

## Error codes

A refusal a caller is expected to act on carries a stable code, both as an
`Otari-Error-Code` header and as `code` in the body beside the human-readable
`detail`: `{"detail": "...", "code": "budget_exceeded"}`. Map refusals by the
code: it keeps its meaning across releases, while the `detail` text may be
reworded.

| `Otari-Error-Code` | Status | Meaning | Also sent |
|---|---|---|---|
| `budget_exceeded` | 403 | A budget refused the request | `Otari-Budget-Scope`: `user` for the billed user's own budget, otherwise the ceiling's scope (`api_token`, `workspace`, `organization`, ...) |
| `user_blocked` | 403 | The billed user is blocked | |
| `user_not_found` | 404 | The billed user does not exist | |
| `rate_limited` | 429 | A gateway rate limit is full | `Otari-Rate-Limit-Rule` for a `rate_limits` rule; `Retry-After` when waiting helps |
| `upstream_rate_limited` | 429 | The provider rate limited the gateway | `Retry-After` when the provider sent one |
| `invalid_model` | 400 | The model selector names no configured provider | |
| `model_not_allowed` | 403 | The key may not use the model | |
| `context_length_exceeded` | 400 | The prompt is too long for the model | |
| `pricing_required` | 402 | `require_pricing` is on and the model has no price | |

A failure after a stream has started arrives as an error event, which carries
the code as `error.code` on Chat Completions and Responses:
`{"error": {"message": "...", "type": "server_error", "code": "upstream_rate_limited"}}`.

## Caller-orchestrated MCP

Two stored-server endpoints let an application own its own MCP tool loop, as an
Expand Down
61 changes: 48 additions & 13 deletions src/gateway/api/routes/_pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@
from urllib.parse import ParseResult, urlparse

from any_llm import LLMProvider
from any_llm.exceptions import AnyLLMError, InvalidRequestError, UnsupportedParameterError
from any_llm.exceptions import AnyLLMError, ContextLengthExceededError, InvalidRequestError, UnsupportedParameterError
from any_llm.types.completion import (
ChatCompletion,
ChatCompletionChunk,
Expand Down Expand Up @@ -103,6 +103,15 @@
from gateway.core.config import ATTEMPT_ID_HEADER, REQUEST_ID_HEADER, GatewayConfig
from gateway.core.database import DATABASE_ERRORS, release_session
from gateway.core.env import otari_env
from gateway.core.error_codes import (
CONTEXT_LENGTH_EXCEEDED,
INVALID_MODEL,
MODEL_NOT_ALLOWED,
PRICING_REQUIRED,
UPSTREAM_RATE_LIMITED,
error_code_of,
error_headers,
)
from gateway.core.metered_pricing import calculate_metered_cost, quantize_cost
from gateway.core.unit_of_work import UnitOfWork
from gateway.core.usage import (
Expand Down Expand Up @@ -637,20 +646,41 @@ def classify_provider_error(exc: BaseException) -> ProviderErrorMapping | None:
def provider_error_headers(exc: BaseException, status_code: int) -> dict[str, str] | None:
"""Response headers for a classified provider failure, or ``None``.

Forwards the upstream ``Retry-After`` on a 429, which is the one header a
rate-limited caller can act on and the one piece of a provider's rate-limit
response that its message body cannot always carry. Restricted to the 429:
on the statuses that surface as a fixed-detail 502 the header would describe
the gateway's own upstream account, which is not the caller's to read.

Returns ``None`` rather than an empty dict when there is nothing to send, so
``HTTPException(headers=...)`` stays unset instead of being handed a dict
that adds nothing.
On a 429, ``Otari-Error-Code: upstream_rate_limited`` (so a caller can tell
the provider's limit from the gateway's own) and the upstream ``Retry-After``,
the one piece of a provider's rate-limit response that its message body cannot
always carry. Restricted to the 429: on the statuses that surface as a
fixed-detail 502 the header would describe the gateway's own upstream
account, which is not the caller's to read.
"""
if status_code == status.HTTP_400_BAD_REQUEST and _is_context_length_error(exc):
return error_headers(CONTEXT_LENGTH_EXCEEDED)
if status_code != status.HTTP_429_TOO_MANY_REQUESTS:
return None
headers = error_headers(UPSTREAM_RATE_LIMITED)
retry_after = upstream_retry_after(exc)
return {"Retry-After": retry_after} if retry_after is not None else None
if retry_after is not None:
headers["Retry-After"] = retry_after
return headers


def _is_context_length_error(exc: BaseException) -> bool:
"""Whether any-llm classified the failure as a prompt too long for the model."""
return any(isinstance(current, ContextLengthExceededError) for current in upstream_exception_chain(exc))


def refusal_code(exc: BaseException) -> str | None:
"""The ``Otari-Error-Code`` an exception ending a request stands for, or None.

Read by the stream error events, which go out after the headers that would
otherwise carry it.
"""
if isinstance(exc, HTTPException):
return error_code_of(exc.headers)
mapping = classify_provider_error(exc)
if mapping is None:
return None
return error_code_of(provider_error_headers(exc, mapping.status_code))


def failure_status_code(exc: BaseException) -> int:
Expand Down Expand Up @@ -1042,6 +1072,7 @@ def _raise_for_unresolvable_model(model_selector: str, exc: Exception) -> NoRetu
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=unresolvable_model_detail(model_selector),
headers=error_headers(INVALID_MODEL),
) from exc


Expand Down Expand Up @@ -1482,6 +1513,7 @@ async def top_up_reservation_for_attempt(ctx: RequestContext, attempt: Attempt)
raise HTTPException(
status_code=status.HTTP_402_PAYMENT_REQUIRED,
detail=no_pricing_error_detail(f"{attempt.instance}:{attempt.model}"),
headers=error_headers(PRICING_REQUIRED),
)
repriced = estimate_cost(
pricing,
Expand Down Expand Up @@ -1983,7 +2015,7 @@ async def resolve_request_context(
started_at=started_at,
request_id=request_id,
)
raise adapter.error(403, not_allowed_detail, ErrorKind.PERMISSION)
raise adapter.error(403, not_allowed_detail, ErrorKind.PERMISSION, headers=error_headers(MODEL_NOT_ALLOWED))

# Organization-scoped model restriction (otari#643): the org key
# resolved for this workspace+provider may narrow which models it
Expand Down Expand Up @@ -2014,7 +2046,9 @@ async def resolve_request_context(
started_at=started_at,
request_id=request_id,
)
raise adapter.error(403, not_allowed_detail, ErrorKind.PERMISSION)
raise adapter.error(
403, not_allowed_detail, ErrorKind.PERMISSION, headers=error_headers(MODEL_NOT_ALLOWED)
)

if idempotency is not None and session_principal is None:
try:
Expand Down Expand Up @@ -2156,6 +2190,7 @@ async def resolve_request_context(
402,
no_pricing_detail,
ErrorKind.INVALID_REQUEST,
headers=error_headers(PRICING_REQUIRED),
)

# Resolve uploaded attachments only once the request is authorized
Expand Down
5 changes: 3 additions & 2 deletions src/gateway/api/routes/chat.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@
provider_error_headers,
raise_all_streaming_attempts_failed,
rate_limit_headers,
refusal_code,
resolve_dispatch_provider,
resolve_request_context,
run_platform_non_stream,
Expand Down Expand Up @@ -70,7 +71,7 @@
mcp_tool_loop_stream,
)
from gateway.services.tools import CODE_EXECUTION_HEADER, WEB_SEARCH_HEADER, Dialect, ToolUseBudget
from gateway.streaming import OPENAI_STREAM_FORMAT, StreamFormat
from gateway.streaming import OPENAI_STREAM_FORMAT, StreamFormat, openai_error_event
from gateway.types.attempt import Attempt
from gateway.types.normalization_target import NormalizationTarget
from gateway.types.session_principal import SessionPrincipal
Expand Down Expand Up @@ -208,7 +209,7 @@ def provider_error(self, exc: BaseException) -> HTTPException:
)

def stream_error_payload(self, exc: BaseException) -> str:
return self.stream_format.error_payload
return openai_error_event(self.stream_format, refusal_code(exc))

def format_chunk(self, chunk: ChatCompletionChunk) -> str:
return f"data: {chunk.model_dump_json()}\n\n"
Expand Down
5 changes: 3 additions & 2 deletions src/gateway/api/routes/responses.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@
prepare_gateway_tools,
provider_error_headers,
raise_all_streaming_attempts_failed,
refusal_code,
release_reservation,
resolve_dispatch_provider,
resolve_request_context,
Expand Down Expand Up @@ -68,7 +69,7 @@
)
from gateway.services.tool_format import inject_purpose_hints_responses, openai_to_responses_tools
from gateway.services.tools import CODE_EXECUTION_HEADER, WEB_SEARCH_HEADER, Dialect, ToolUseBudget
from gateway.streaming import RESPONSES_STREAM_FORMAT, StreamFormat
from gateway.streaming import RESPONSES_STREAM_FORMAT, StreamFormat, openai_error_event
from gateway.types.attempt import Attempt
from gateway.types.normalization_target import NormalizationTarget

Expand Down Expand Up @@ -342,7 +343,7 @@ def provider_error(self, exc: BaseException) -> HTTPException:
)

def stream_error_payload(self, exc: BaseException) -> str:
return self.stream_format.error_payload
return openai_error_event(self.stream_format, refusal_code(exc))

def format_chunk(self, chunk: ResponseStreamEvent) -> str:
return f"event: {chunk.type}\ndata: {chunk.model_dump_json(exclude_none=True)}\n\n"
Expand Down
39 changes: 39 additions & 0 deletions src/gateway/core/error_codes.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
"""Stable, machine-readable codes for the refusals a caller is expected to act on.

Sent as the ``Otari-Error-Code`` response header and as ``code`` in the error
body, beside the human-readable ``detail``, so a client maps a refusal by its
code rather than by matching text that may be reworded. A streamed error event
carries it as ``error.code``. A code, once sent, keeps its meaning.
"""

from collections.abc import Mapping

ERROR_CODE_HEADER = "Otari-Error-Code"
# Which budget refused: ``user`` for the billed user's own budget, otherwise the
# scope of the ceiling (``api_token``, ``workspace``, ``organization``, ...).
BUDGET_SCOPE_HEADER = "Otari-Budget-Scope"
# The ``rate_limits`` rule a 429 names.
RATE_LIMIT_RULE_HEADER = "Otari-Rate-Limit-Rule"

BUDGET_EXCEEDED = "budget_exceeded"
USER_BLOCKED = "user_blocked"
USER_NOT_FOUND = "user_not_found"
RATE_LIMITED = "rate_limited"
UPSTREAM_RATE_LIMITED = "upstream_rate_limited"
INVALID_MODEL = "invalid_model"
MODEL_NOT_ALLOWED = "model_not_allowed"
CONTEXT_LENGTH_EXCEEDED = "context_length_exceeded"
PRICING_REQUIRED = "pricing_required"


def error_headers(code: str, **extra: str | None) -> dict[str, str]:
"""``Otari-Error-Code`` plus any of the extra headers that have a value."""
headers = {ERROR_CODE_HEADER: code}
names = {"budget_scope": BUDGET_SCOPE_HEADER, "rule": RATE_LIMIT_RULE_HEADER}
headers.update({names[name]: value for name, value in extra.items() if value is not None})
return headers


def error_code_of(headers: Mapping[str, str] | None) -> str | None:
"""The code a refusal's headers carry, or None."""
return (headers or {}).get(ERROR_CODE_HEADER)
22 changes: 22 additions & 0 deletions src/gateway/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,11 +7,13 @@
from urllib.parse import urlsplit

from fastapi import FastAPI, Request, Response, status
from fastapi.exception_handlers import http_exception_handler
from fastapi.exceptions import RequestValidationError
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import FileResponse, HTMLResponse, JSONResponse, RedirectResponse
from fastapi.routing import APIRoute
from fastapi.staticfiles import StaticFiles
from starlette.exceptions import HTTPException as StarletteHTTPException
from starlette.middleware.base import BaseHTTPMiddleware, RequestResponseEndpoint
from typing_extensions import override

Expand All @@ -22,6 +24,7 @@
from gateway.context_propagation import TraceContextPropagationMiddleware
from gateway.core.config import API_KEY_HEADER, API_ROOT, GATEWAY_TOKEN_HEADER, X_API_KEY_HEADER, GatewayConfig
from gateway.core.database import create_session, dispose_db, init_db
from gateway.core.error_codes import error_code_of
from gateway.core.feature import Worker
from gateway.dashboard import DASHBOARD_PACKAGE_PATH, get_dashboard_build_id, get_dashboard_dir
from gateway.exceptions import TenancyError
Expand Down Expand Up @@ -733,6 +736,24 @@ async def lifespan(app: FastAPI) -> AsyncGenerator[None, None]:
return lifespan


async def _http_exception_handler(request: Request, exc: Exception) -> Response:
"""FastAPI's own HTTPException response, with the refusal's ``Otari-Error-Code`` as ``code`` in the body too.

A body field survives where a header does not (a proxy that drops it, an SDK
that surfaces only the body), so a client can map the refusal from either.
"""
if not isinstance(exc, StarletteHTTPException):
raise exc
code = error_code_of(exc.headers)
if code is None:
return await http_exception_handler(request, exc)
return JSONResponse(
{"detail": exc.detail, "code": code},
status_code=exc.status_code,
headers=exc.headers,
)


async def _tenancy_error_handler(_: Request, exc: Exception) -> Response:
"""Render a tenancy domain error as the status it carries.

Expand Down Expand Up @@ -1060,6 +1081,7 @@ async def root_index() -> str:
install_rate_limits(app, config)

register_routers(app, config)
app.add_exception_handler(StarletteHTTPException, _http_exception_handler)
app.add_exception_handler(TenancyError, _tenancy_error_handler)
app.add_exception_handler(ControlPlaneError, _control_plane_error_handler)
app.add_exception_handler(RequestValidationError, _validation_error_handler)
Expand Down
19 changes: 11 additions & 8 deletions src/gateway/rate_limit.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@

from fastapi import HTTPException, Request, status

from gateway.core.error_codes import RATE_LIMITED, error_headers
from gateway.log_config import logger
from gateway.metrics import REGISTRY, Counter
from gateway.ports.rate_limit_store_port import RateLimitStorePort, RateLimitWindow
Expand Down Expand Up @@ -115,7 +116,7 @@ def _info_or_raise(window: RateLimitWindow, limit: int) -> RateLimitInfo:
raise HTTPException(
status_code=status.HTTP_429_TOO_MANY_REQUESTS,
detail="Rate limit exceeded",
headers={"Retry-After": str(math.ceil(window.reset_after))},
headers={"Retry-After": str(math.ceil(window.reset_after)), **error_headers(RATE_LIMITED)},
)
# Wall-clock time for the externally facing reset header.
return RateLimitInfo(limit=limit, remaining=limit - window.count, reset=time.time() + window.reset_after)
Expand Down Expand Up @@ -293,14 +294,13 @@ def _count(n: int, noun: str) -> str:
return f"{n:,} {noun}" if n == 1 else f"{n:,} {noun}s"


def _refused(detail: str, retry_after: float | None) -> HTTPException:
def _refused(detail: str, retry_after: float | None, rule: str) -> HTTPException:
"""A 429 with ``detail``, without ``Retry-After`` when no wait would let the request in."""
RATE_LIMIT_HITS.inc()
return HTTPException(
status_code=status.HTTP_429_TOO_MANY_REQUESTS,
detail=detail,
headers={"Retry-After": str(max(math.ceil(retry_after), 1))} if retry_after is not None else None,
)
headers = error_headers(RATE_LIMITED, rule=rule)
if retry_after is not None:
headers["Retry-After"] = str(max(math.ceil(retry_after), 1))
return HTTPException(status_code=status.HTTP_429_TOO_MANY_REQUESTS, detail=detail, headers=headers)


async def _count_rule(
Expand All @@ -327,6 +327,7 @@ async def _count_rule(
raise _refused(
f"Rate limit '{rule.name}'{label} exceeded: {_count(rule.rpm, 'request')} per minute",
window.reset_after,
rule.name,
)
hold.entries.append((f"{base}:rpm", window.handle))
if rule.tpm is not None:
Expand All @@ -337,12 +338,13 @@ async def _count_rule(
f"Request needs an estimated {_count(cost, 'token')}; "
f"rate limit '{rule.name}'{label} allows {rule.tpm:,} per minute"
)
raise _refused(msg, None)
raise _refused(msg, None, rule.name)
window = await store.hit(f"{base}:tpm", rule.tpm, _RULE_WINDOW_SEC, cost=cost)
if window.handle is None:
raise _refused(
f"Rate limit '{rule.name}'{label} exceeded: {_count(rule.tpm, 'token')} per minute",
window.reset_after,
rule.name,
)
hold.entries.append((f"{base}:tpm", window.handle))
hold.estimates.append((f"{base}:tpm", window.handle))
Expand All @@ -352,6 +354,7 @@ async def _count_rule(
raise _refused(
f"Rate limit '{rule.name}'{label} exceeded: {_count(rule.max_concurrent, 'request')} in flight",
_CONCURRENCY_RETRY_AFTER_SEC,
rule.name,
)
hold.leases.append((f"{base}:concurrent", lease))

Expand Down
Loading
Loading