#![allow(clippy::unused_async)]
use axum::body::Body;
use axum::extract::rejection::JsonRejection;
use axum::extract::{Query, Request, State};
use axum::http::{HeaderMap, HeaderValue, StatusCode};
use axum::response::{IntoResponse, Response};
use futures_util::StreamExt;
use log_lazy::LogLazy;
use std::collections::{BTreeMap, HashSet};
use crate::accounts::RoutingContext;
pub(crate) use crate::api_error::error_response;
use crate::api_error::malformed_json_response;
pub use crate::app_state::AppState;
use crate::config::UpstreamProvider;
pub use crate::model_routing::models as openai_models;
pub use crate::monitoring_api::{accounts_endpoint, metrics_endpoint, usage_endpoint};
use crate::openai;
pub(crate) use crate::request_routing::{request_routing_context, retry_after_duration};
use crate::responses;
use crate::subscription::SubscriptionProvider;
pub const API_PREFIX: &str = "/api/latest/anthropic/";
pub const REQUIRED_FORWARD_HEADERS: &[&str] = &[
"anthropic-beta",
"anthropic-version",
"x-claude-code-session-id",
];
pub const DEFAULT_ANTHROPIC_VERSION: &str = "2023-06-01";
pub const OAUTH_BETA_FLAG: &str = "oauth-2025-04-20";
pub const MAX_PROXY_REQUEST_BYTES: usize = crate::config::DEFAULT_MAX_PROXY_REQUEST_BYTES;
#[must_use]
pub fn merge_oauth_beta(existing: Option<&str>) -> String {
match existing {
Some(v) if v.split(',').map(str::trim).any(|f| f == OAUTH_BETA_FLAG) => v.to_string(),
Some(v) if !v.trim().is_empty() => format!("{v},{OAUTH_BETA_FLAG}"),
_ => OAUTH_BETA_FLAG.to_string(),
}
}
const HOP_BY_HOP_HEADERS: &[&str] = &[
"host",
"connection",
"keep-alive",
"proxy-authenticate",
"proxy-authorization",
"te",
"trailer",
"transfer-encoding",
"upgrade",
];
const RESPONSE_CREDENTIAL_HEADERS: &[&str] = &[
"authorization",
"proxy-authorization",
"set-cookie",
"set-cookie2",
"x-api-key",
];
pub(crate) fn relay_response_headers(headers: &HeaderMap) -> HeaderMap {
let connection_headers: HashSet<String> = headers
.get_all("connection")
.iter()
.filter_map(|value| value.to_str().ok())
.flat_map(|value| value.split(','))
.map(|name| name.trim().to_ascii_lowercase())
.filter(|name| !name.is_empty())
.collect();
let mut relayed = HeaderMap::new();
for (name, value) in headers {
let name_lower = name.as_str();
if HOP_BY_HOP_HEADERS.contains(&name_lower)
|| RESPONSE_CREDENTIAL_HEADERS.contains(&name_lower)
|| is_operator_subscription_header(name_lower)
|| name_lower == "content-length"
|| connection_headers.contains(name_lower)
{
continue;
}
relayed.append(name.clone(), value.clone());
}
relayed
}
fn is_operator_subscription_header(name: &str) -> bool {
matches!(
name,
"x-codex-plan-type"
| "x-codex-active-limit"
| "x-codex-credits-balance"
| "x-codex-entitlement"
| "x-codex-entitlements"
| "x-codex-state-token"
) || name.starts_with("x-codex-credits-")
|| name.starts_with("x-codex-entitlement-")
}
#[allow(clippy::unused_async)]
pub async fn health() -> impl IntoResponse {
(StatusCode::OK, "ok")
}
pub(crate) fn is_admin_authorised(state: &AppState, headers: &HeaderMap) -> bool {
let provided = extract_bearer_token(headers);
if provided.is_some_and(|token| state.admin.verify(token)) {
return true;
}
crate::admin_auth::admin_access_granted(
&state.token_manager,
provided,
state.admin_key.as_deref(),
state.allow_anonymous_admin,
)
}
pub(crate) fn extract_admin_bearer(headers: &HeaderMap) -> Option<&str> {
extract_bearer_token(headers)
}
fn extract_bearer_token(headers: &HeaderMap) -> Option<&str> {
headers
.get("authorization")
.and_then(|v| v.to_str().ok())
.and_then(|value| {
let (scheme, token) = value.split_once(' ')?;
(scheme.eq_ignore_ascii_case("bearer") || scheme.eq_ignore_ascii_case("token"))
.then_some(token)
})
.filter(|token| !token.is_empty())
}
pub(crate) fn extract_client_token(headers: &HeaderMap) -> Option<&str> {
extract_bearer_token(headers).or_else(|| {
headers
.get("x-api-key")
.and_then(|value| value.to_str().ok())
.filter(|value| !value.is_empty())
})
}
pub(crate) fn authenticate_client(
state: &AppState,
headers: &HeaderMap,
) -> Result<crate::token::TokenClaims, Box<Response>> {
let Some(token) = extract_client_token(headers) else {
state.logger.debug(|| "Missing Authorization header");
return Err(Box::new(error_response(
StatusCode::UNAUTHORIZED,
"authentication_error",
"Missing Authorization Bearer token or x-api-key",
)));
};
state.token_manager.validate_token(token).map_err(|error| {
let status = if matches!(error, crate::token::TokenError::Revoked) {
StatusCode::FORBIDDEN
} else {
StatusCode::UNAUTHORIZED
};
state
.logger
.debug(|| format!("Token validation failed: {error}"));
Box::new(error_response(
status,
"authentication_error",
error.client_message(),
))
})
}
pub async fn proxy_handler(State(state): State<AppState>, req: Request) -> Response {
let path = req.uri().path().to_string();
let method = req.method().clone();
let incoming_headers = req.headers().clone();
if state.upstream_provider == UpstreamProvider::Auto {
if let Err(response) = authenticate_client(&state, &incoming_headers) {
return *response;
}
let (routed, request) =
match crate::model_routing::route_anthropic_request(&state, req).await {
Ok(routed) => routed,
Err(response) => return response,
};
return Box::pin(proxy_handler(State(routed), request)).await;
}
state.logger.verbose(|| format!("Incoming {method} {path}"));
let upstream_path = resolve_upstream_path(&path);
state
.logger
.debug(|| format!("Resolved upstream path: {upstream_path}"));
let upstream_url = format!(
"{}{}",
state.upstream_base_url.trim_end_matches('/'),
upstream_path
);
let upstream_url = if let Some(query) = req.uri().query() {
format!("{upstream_url}?{query}")
} else {
upstream_url
};
if let Some(session_id) = incoming_headers.get("x-claude-code-session-id") {
state
.logger
.verbose(|| format!("Session: {}", session_id.to_str().unwrap_or("<invalid>")));
}
if crate::anthropic_bridge::is_bridged(state.upstream_provider) {
let body_bytes =
match axum::body::to_bytes(req.into_body(), state.max_proxy_request_bytes).await {
Ok(bytes) => bytes,
Err(e) => {
return error_response(
StatusCode::PAYLOAD_TOO_LARGE,
"invalid_request_error",
&format!(
"request body exceeds the {} byte proxy limit: {e}",
state.max_proxy_request_bytes
),
);
}
};
let body = match serde_json::from_slice(&body_bytes) {
Ok(body) => body,
Err(error) => return malformed_json_response(&error.to_string()),
};
return crate::anthropic_bridge::handle_anthropic_surface(
&state,
&incoming_headers,
&path,
body,
)
.await;
}
let claims = match authenticate_client(&state, &incoming_headers) {
Ok(claims) => claims,
Err(response) => return *response,
};
if let Err(e) = state.token_manager.enforce_request_budget(&claims.sub) {
state
.logger
.debug(|| format!("Token budget check failed: {e}"));
return crate::token_http::budget_error_response(&e);
}
let mut body_bytes =
match axum::body::to_bytes(req.into_body(), state.max_proxy_request_bytes).await {
Ok(bytes) => bytes,
Err(e) => {
return error_response(
StatusCode::PAYLOAD_TOO_LARGE,
"invalid_request_error",
&format!(
"request body exceeds the {} byte proxy limit: {e}",
state.max_proxy_request_bytes
),
);
}
};
let routing_body = match serde_json::from_slice(&body_bytes) {
Ok(body) => body,
Err(error) => return malformed_json_response(&error.to_string()),
};
crate::audit::record_authorised_request(
&state,
&claims,
crate::metrics::Surface::Anthropic,
&path,
Some(&routing_body),
);
let pinned_account = match state.token_manager.account_for(&claims.sub) {
Ok(account) => account,
Err(error) => {
return error_response(
StatusCode::INTERNAL_SERVER_ERROR,
"api_error",
&format!("failed to resolve token account binding: {error}"),
);
}
};
let routing_context = request_routing_context(&incoming_headers, &routing_body, pinned_account);
let (oauth_token, selected_account) =
match resolve_upstream_credentials(&state, &routing_context).await {
Ok(pair) => pair,
Err(e) => {
tracing::error!("Failed to resolve upstream credentials: {e}");
return error_response(
StatusCode::BAD_GATEWAY,
"api_error",
"Upstream authentication unavailable",
);
}
};
let mut upstream_body = routing_body.clone();
openai::reconcile_subscription_parameters(SubscriptionProvider::Claude, &mut upstream_body);
if upstream_body != routing_body {
body_bytes = serde_json::to_vec(&upstream_body)
.map(bytes::Bytes::from)
.unwrap_or(body_bytes);
}
let body_bytes = if crate::claude_identity::is_oauth_credential(&oauth_token) {
crate::claude_identity::ensure_claude_code_system_bytes(&upstream_body, body_bytes)
} else {
body_bytes
};
let upstream_headers = build_upstream_headers(&incoming_headers, &oauth_token, &state.logger);
state.logger.verbose(|| {
format!(
"Forwarding {method} {upstream_url} ({} bytes)",
body_bytes.len()
)
});
let upstream_req = state
.client
.request(method, &upstream_url)
.headers(upstream_headers)
.body(body_bytes);
let correlation_id = crate::request_log::correlation_id(&incoming_headers);
let upstream_resp = match state
.request_log
.send_upstream(&correlation_id, &state.client, upstream_req)
.await
{
Ok(resp) => resp,
Err(e) => {
tracing::error!("Upstream request failed: {e}");
return error_response(
StatusCode::BAD_GATEWAY,
"api_error",
&format!("Upstream request failed: {e}"),
);
}
};
let status = StatusCode::from_u16(upstream_resp.status().as_u16())
.unwrap_or(StatusCode::INTERNAL_SERVER_ERROR);
let retry_after = retry_after_duration(upstream_resp.headers());
crate::request_routing::record_claude_evidence(&state, status.as_u16());
state
.logger
.verbose(|| format!("Upstream responded: {status}"));
state.metrics.record_request(
crate::metrics::Surface::Anthropic,
status.as_u16(),
selected_account.as_deref(),
);
if status.as_u16() == 429 {
if let (Some(router), Some(name)) =
(state.account_router.as_ref(), selected_account.as_deref())
{
router.report_failure_with_retry_after(name, "upstream returned 429", retry_after);
}
}
let response_headers = relay_response_headers(upstream_resp.headers());
let response_log = std::sync::Arc::clone(&state.request_log);
let mut usage = status
.is_success()
.then(|| crate::usage::UsageTracker::new(state.token_manager.clone(), claims.sub));
let stream = upstream_resp.bytes_stream().map(move |chunk| {
if let Ok(bytes) = &chunk {
response_log.record_upstream_body(&correlation_id, bytes);
if let Some(tracker) = &mut usage {
tracker.feed(bytes);
}
}
chunk.map_err(std::io::Error::other)
});
let body = Body::from_stream(stream);
let mut response = Response::new(body);
*response.status_mut() = status;
*response.headers_mut() = response_headers;
response
}
#[must_use]
pub fn resolve_upstream_path(path: &str) -> String {
if let Some(rest) = path.strip_prefix("/api/anthropic") {
return rest.to_string();
}
if let Some(rest) = path.strip_prefix("/api/latest/anthropic") {
return rest.to_string();
}
path.to_string()
}
pub(crate) fn build_upstream_headers(
incoming: &HeaderMap,
oauth_token: &str,
logger: &LogLazy,
) -> HeaderMap {
let mut headers = HeaderMap::new();
for (name, value) in incoming {
let name_lower = name.as_str().to_lowercase();
if matches!(
name_lower.as_str(),
"authorization" | "x-api-key" | "content-length"
) || HOP_BY_HOP_HEADERS.contains(&name_lower.as_str())
{
continue;
}
headers.insert(name.clone(), value.clone());
}
if let Ok(auth_val) = HeaderValue::from_str(&format!("Bearer {oauth_token}")) {
headers.insert("authorization", auth_val);
}
if !headers.contains_key("anthropic-version") {
headers.insert(
"anthropic-version",
HeaderValue::from_static(DEFAULT_ANTHROPIC_VERSION),
);
}
let existing_beta = headers
.get("anthropic-beta")
.and_then(|v| v.to_str().ok())
.map(String::from);
if let Ok(beta_val) = HeaderValue::from_str(&merge_oauth_beta(existing_beta.as_deref())) {
headers.insert("anthropic-beta", beta_val);
}
for &header_name in REQUIRED_FORWARD_HEADERS {
if let Some(val) = headers.get(header_name) {
logger.trace(|| {
format!(
"Forwarding {header_name}: {}",
val.to_str().unwrap_or("<non-utf8>")
)
});
}
}
headers
}
async fn resolve_upstream_credentials(
state: &AppState,
context: &RoutingContext,
) -> Result<(String, Option<String>), Box<dyn std::error::Error + Send + Sync>> {
if let Some(router) = state.account_router.as_ref() {
let sel = router.select_subscription(context)?;
let now_ms = chrono::Utc::now().timestamp_millis();
let token = state
.subscription_cache
.get_fresh_for(
&state.client,
router.provider(),
&sel.name,
sel.token,
now_ms,
)
.await;
return Ok((token.access_token, Some(sel.name)));
}
let token = state
.oauth_provider
.get_fresh_token(&state.client, &state.subscription_cache)
.await?;
Ok((token, None))
}
#[path = "proxy_openai.rs"]
mod openai_handlers;
pub use openai_handlers::{openai_chat_completions, openai_responses};
#[derive(Clone, Copy)]
enum OpenAIShape {
Chat,
Response,
}
async fn forward_openai(
state: &AppState,
headers: &HeaderMap,
body: serde_json::Value,
routing_body: &serde_json::Value,
surface: crate::metrics::Surface,
stream_options: (bool, OpenAIShape, bool),
) -> Response {
let (stream_requested, shape, include_usage) = stream_options;
let served_model = body["model"].as_str().unwrap_or_default().to_string();
let path = match shape {
OpenAIShape::Chat => "/v1/chat/completions",
OpenAIShape::Response => "/v1/responses",
};
if let Some(resp) = maybe_mpp_challenge(state, headers, path) {
return resp;
}
let claims = match authenticate_client(state, headers) {
Ok(claims) => claims,
Err(response) => return *response,
};
if let Err(e) = state.token_manager.enforce_request_budget(&claims.sub) {
return crate::token_http::budget_error_response(&e);
}
crate::audit::record_authorised_request(state, &claims, surface, path, Some(routing_body));
let pinned_account = match state.token_manager.account_for(&claims.sub) {
Ok(account) => account,
Err(error) => {
return error_response(
StatusCode::INTERNAL_SERVER_ERROR,
"api_error",
&format!("failed to resolve token account binding: {error}"),
);
}
};
let routing_context = request_routing_context(headers, routing_body, pinned_account);
let (oauth_token, selected_account) =
match resolve_upstream_credentials(state, &routing_context).await {
Ok(p) => p,
Err(e) => {
tracing::error!("openai: upstream credentials unavailable: {e}");
return error_response(
StatusCode::BAD_GATEWAY,
"api_error",
"Upstream authentication unavailable",
);
}
};
let upstream_url = format!(
"{}/v1/messages",
state.upstream_base_url.trim_end_matches('/')
);
let mut body = body;
if crate::claude_identity::is_oauth_credential(&oauth_token) {
crate::claude_identity::ensure_claude_code_system(&mut body);
}
let serialized = match serde_json::to_vec(&body) {
Ok(v) => v,
Err(e) => {
return error_response(
StatusCode::INTERNAL_SERVER_ERROR,
"api_error",
&format!("failed to serialize translated body: {e}"),
);
}
};
let bytes_sent = serialized.len() as u64;
let mut req_builder = state
.client
.post(&upstream_url)
.header("authorization", format!("Bearer {oauth_token}"))
.header("content-type", "application/json")
.header("anthropic-version", DEFAULT_ANTHROPIC_VERSION)
.body(serialized);
let merged_beta = merge_oauth_beta(headers.get("anthropic-beta").and_then(|v| v.to_str().ok()));
req_builder = req_builder.header("anthropic-beta", merged_beta);
let correlation_id = crate::request_log::correlation_id(headers);
let upstream_resp = match state
.request_log
.send_upstream(&correlation_id, &state.client, req_builder)
.await
{
Ok(r) => r,
Err(e) => {
state
.metrics
.record_request(surface, 502, selected_account.as_deref());
return error_response(
StatusCode::BAD_GATEWAY,
"api_error",
&format!("upstream request failed: {e}"),
);
}
};
let upstream_status = upstream_resp.status();
let retry_after = retry_after_duration(upstream_resp.headers());
let response_headers = relay_response_headers(upstream_resp.headers());
if stream_requested && upstream_status.is_success() {
state
.metrics
.record_request(surface, 200, selected_account.as_deref());
let stream_shape = match shape {
OpenAIShape::Chat => openai::OpenAIStreamShape::ChatCompletion,
OpenAIShape::Response => openai::OpenAIStreamShape::Response,
};
let mut translator = openai::OpenAIStreamTranslator::new(stream_shape, &served_model)
.with_include_usage(include_usage);
let response_log = std::sync::Arc::clone(&state.request_log);
let mut usage =
crate::usage::UsageTracker::new(state.token_manager.clone(), claims.sub.clone());
let stream = upstream_resp.bytes_stream().map(move |chunk| match chunk {
Ok(bytes) => {
response_log.record_upstream_body(&correlation_id, &bytes);
usage.feed(&bytes);
Ok::<bytes::Bytes, std::io::Error>(bytes::Bytes::from(
translator.push(&bytes).join(""),
))
}
Err(e) => Err(std::io::Error::other(e)),
});
let mut response = Response::new(Body::from_stream(stream));
*response.status_mut() = StatusCode::OK;
*response.headers_mut() = response_headers;
response.headers_mut().insert(
"content-type",
HeaderValue::from_static("text/event-stream; charset=utf-8"),
);
return response;
}
let upstream_body = match upstream_resp.bytes().await {
Ok(b) => b,
Err(e) => {
state
.metrics
.record_request(surface, 502, selected_account.as_deref());
return error_response(
StatusCode::BAD_GATEWAY,
"api_error",
&format!("upstream body read failed: {e}"),
);
}
};
state
.request_log
.record_upstream_body(&correlation_id, &upstream_body);
let bytes_received = upstream_body.len() as u64;
state.metrics.record_bytes(bytes_sent, bytes_received);
if !upstream_status.is_success() {
if upstream_status.as_u16() == 429 {
if let (Some(router), Some(name)) =
(state.account_router.as_ref(), selected_account.as_deref())
{
router.report_failure_with_retry_after(name, "upstream returned 429", retry_after);
}
}
state.metrics.record_request(
surface,
upstream_status.as_u16(),
selected_account.as_deref(),
);
let parsed: serde_json::Value =
serde_json::from_slice(&upstream_body).unwrap_or_else(|_| serde_json::json!({}));
let mut resp = (
StatusCode::from_u16(upstream_status.as_u16()).unwrap_or(StatusCode::BAD_GATEWAY),
axum::Json(parsed),
)
.into_response();
*resp.headers_mut() = response_headers;
resp.headers_mut()
.insert("content-type", HeaderValue::from_static("application/json"));
return resp;
}
let anthropic: serde_json::Value = match serde_json::from_slice(&upstream_body) {
Ok(v) => v,
Err(e) => {
state
.metrics
.record_request(surface, 502, selected_account.as_deref());
return error_response(
StatusCode::BAD_GATEWAY,
"api_error",
&format!("upstream returned non-JSON: {e}"),
);
}
};
state
.token_manager
.record_token_usage(
&claims.sub,
crate::usage::token_count(&anthropic).unwrap_or(0),
)
.unwrap_or_else(|error| tracing::warn!("failed to persist token usage: {error}"));
let translated = match shape {
OpenAIShape::Chat => openai::anthropic_to_chat_completion(&anthropic, &served_model),
OpenAIShape::Response => responses::anthropic_to_response(&anthropic, &served_model),
};
state
.metrics
.record_request(surface, 200, selected_account.as_deref());
let mut response = (StatusCode::OK, axum::Json(translated)).into_response();
*response.headers_mut() = response_headers;
response
.headers_mut()
.insert("content-type", HeaderValue::from_static("application/json"));
response
}
pub(crate) fn maybe_mpp_challenge(
state: &AppState,
headers: &HeaderMap,
path: &str,
) -> Option<Response> {
if !state.mpp.is_configured() {
return None;
}
if crate::mpp::has_payment_credential(headers) {
return Some(crate::mpp::unsupported_payment_verification());
}
Some(crate::mpp::payment_required(&state.mpp, path))
}