use axum::body::{Body, Bytes};
use axum::extract::{Request, State};
use axum::http::StatusCode;
use axum::http::header::{self, HeaderValue};
use axum::response::Response;
use http_body_util::BodyExt as _;
use r402_extensions::{
EIP2612_GAS_SPONSORING_KEY, ERC20_APPROVAL_GAS_SPONSORING_KEY, Eip2612GasSponsoringExtension,
Erc20ApprovalGasSponsoringExtension,
};
use r402_http::server::{SettlementOverrides, reason_to_status, settlement_to_header};
use r402_http::{
PAYMENT_REQUIRED, PAYMENT_RESPONSE, PAYMENT_SIGNATURE, ensure_expose_headers, merge_private,
set_no_store,
};
use r402_protocol::error::{ErrorReason, FacilitatorError};
use r402_protocol::extension::{AdvertiseContext, Extension};
use r402_protocol::payment::{
Base64Bytes, Extensions, PaymentRequired, PaymentRequirements, ResourceInfo, SettleResponse,
SupportedResponse,
};
use r402_server::{
PaymentFlowName, PaymentRequiredBuildContext, ResourceServer, SettlePhase, WirePaymentPayload,
validate_accepts_against_supported,
};
use serde_json::json;
use super::inflight::{HashReservation, SettleOnce};
use super::price::{self, LoadedRates, Scheme};
use super::tee::{self, SettleJob, TeeBody};
use super::{tags, usage};
use crate::http::bill::Bill;
use crate::http::{RequestContext, openai_error};
use crate::proxy;
use crate::state::AppState;
pub(crate) async fn gate(State(state): State<AppState>, request: Request) -> Response {
let Some(ctx) = request.extensions().get::<RequestContext>().cloned() else {
return openai_error::server_error("missing request context");
};
let Some(server) = state.resource_server().cloned() else {
return openai_error::server_error("payment is enabled without a resource server");
};
if state.payment().accepts.is_empty() {
return json_status(
StatusCode::INTERNAL_SERVER_ERROR,
&json!({ "error": "payment accepts are empty" }),
);
}
dispatch(&state, server, &ctx, request).await
}
async fn dispatch(
state: &AppState,
server: ResourceServer,
ctx: &RequestContext,
request: Request,
) -> Response {
let scheme = match ctx.bill {
Bill::Exact => Scheme::Exact,
Bill::Upto => Scheme::Upto,
Bill::Unpaid | Bill::Reject => {
return openai_error::server_error("payment gate on unbilled route");
}
};
let Some(rates) = state.rates_for(&ctx.model, scheme) else {
return json_status(
StatusCode::INTERNAL_SERVER_ERROR,
&json!({ "error": "missing loaded rates" }),
);
};
let embeddings = crate::http::bill::openai_path(request.uri().path()) == "/v1/embeddings";
let max_out = price::max_out(
rates,
ctx.max_tokens,
ctx.max_completion_tokens,
ctx.max_output_tokens,
);
let amount = price::ceiling_or_price(rates, max_out, embeddings);
let tags = match tags::price_tags_for(state.payment(), amount, Some(scheme.as_str())) {
Ok(tags) if !tags.is_empty() => tags,
Ok(_) => {
return json_status(
StatusCode::INTERNAL_SERVER_ERROR,
&json!({ "error": "payment accepts are empty" }),
);
}
Err(error) => {
return json_status(
StatusCode::INTERNAL_SERVER_ERROR,
&json!({ "error": error.to_string() }),
);
}
};
let reqs: Vec<PaymentRequirements> = tags.into_iter().map(|tag| tag.requirements).collect();
let resource = resource_info(state, request.uri().path());
let payment_required = match build_payment_required(&server, reqs, resource).await {
Ok(body) => body,
Err(response) => return response,
};
settle_after_verify(state, server, ctx, rates, amount, payment_required, request).await
}
async fn settle_after_verify(
state: &AppState,
server: ResourceServer,
ctx: &RequestContext,
rates: &LoadedRates,
ceiling: u128,
payment_required: PaymentRequired,
request: Request,
) -> Response {
let Some(header) = request.headers().get(PAYMENT_SIGNATURE) else {
return challenge(
&payment_required,
"Payment-Signature header is required",
StatusCode::PAYMENT_REQUIRED,
);
};
let signature = header.as_bytes().to_vec();
let Some(payload) = decode_signature(&signature) else {
return challenge(
&payment_required,
"Invalid or malformed payment header",
StatusCode::PAYMENT_REQUIRED,
);
};
let Some(requirements) = server
.find_matching_requirements(&payment_required.accepts, &payload)
.cloned()
else {
return challenge(
&payment_required,
"Unable to find matching payment requirements",
StatusCode::PAYMENT_REQUIRED,
);
};
match server.get_payment_flow(&requirements) {
Ok(PaymentFlowName::Authorization) => {}
Ok(flow) => {
return json_status(
StatusCode::INTERNAL_SERVER_ERROR,
&json!({ "error": "unsupported payment flow", "flow": flow.as_str() }),
);
}
Err(error) => {
return json_status(
StatusCode::INTERNAL_SERVER_ERROR,
&json!({ "error": error.to_string() }),
);
}
}
let Some(reservation) = state.inflight().reserve(signature) else {
return challenge(
&payment_required,
"payment already in flight",
StatusCode::PAYMENT_REQUIRED,
);
};
if let Err(response) =
verify_or_reject(&server, &payload, &requirements, &payment_required).await
{
return response;
}
let stream = ctx.stream;
let verified = VerifiedSettle {
state,
server,
model_id: ctx.model.as_str(),
rates,
ceiling,
payload,
requirements,
advertised: payment_required.extensions.clone(),
resource_url: payment_required.resource.url.to_string(),
reservation,
};
if stream {
if rates.scheme == Scheme::Upto {
stream_then_settle(verified, request).await
} else {
wait_2xx_then_spawn(verified, request).await
}
} else {
sequential_wait_settle(verified, request).await
}
}
struct VerifiedSettle<'a> {
state: &'a AppState,
server: ResourceServer,
model_id: &'a str,
rates: &'a LoadedRates,
ceiling: u128,
payload: WirePaymentPayload,
requirements: PaymentRequirements,
advertised: Extensions,
resource_url: String,
reservation: HashReservation,
}
async fn sequential_wait_settle(ctx: VerifiedSettle<'_>, request: Request) -> Response {
let Some((upstream, client, catalog)) = ctx.state.resolve_upstream(Some(ctx.model_id)) else {
return openai_error::bad_gateway("upstream is not configured");
};
let upstream = match proxy::send_upstream(client, upstream, catalog, request).await {
Ok(response) => response,
Err(response) => return response,
};
let buffered = match proxy::buffer_upstream(upstream).await {
Ok(buffered) => buffered,
Err(response) => return response,
};
if !buffered.status.is_success() {
return proxy::buffered_response(buffered.status, buffered.headers, buffered.body);
}
let once = SettleOnce::new();
if !once.take() {
return openai_error::server_error("settle already taken");
}
attach_settlement(ctx, buffered).await
}
async fn wait_2xx_then_spawn(ctx: VerifiedSettle<'_>, request: Request) -> Response {
let Some((upstream, client, catalog)) = ctx.state.resolve_upstream(Some(ctx.model_id)) else {
return openai_error::bad_gateway("upstream is not configured");
};
let upstream = match proxy::send_upstream(client, upstream, catalog, request).await {
Ok(response) => response,
Err(response) => return response,
};
if !proxy::upstream_status(&upstream).is_success() {
let buffered = match proxy::buffer_upstream(upstream).await {
Ok(buffered) => buffered,
Err(response) => return response,
};
return proxy::buffered_response(buffered.status, buffered.headers, buffered.body);
}
let once = SettleOnce::new();
if !once.take() {
return openai_error::server_error("settle already taken");
}
spawn_exact_settle(ctx);
proxy::stream_upstream(upstream)
}
async fn stream_then_settle(ctx: VerifiedSettle<'_>, request: Request) -> Response {
let Some((upstream, client, catalog)) = ctx.state.resolve_upstream(Some(ctx.model_id)) else {
return openai_error::bad_gateway("upstream is not configured");
};
let request = match force_stream_usage(request).await {
Ok(request) => request,
Err(response) => return response,
};
let missing_usage = ctx.state.payment().missing_usage;
let abort_usage = ctx.state.payment().abort_usage;
let rates = *ctx.rates;
let ceiling = ctx.ceiling;
let upstream = match proxy::send_upstream(client, upstream, catalog, request).await {
Ok(response) => response,
Err(response) => return response,
};
if !proxy::upstream_status(&upstream).is_success() {
let buffered = match proxy::buffer_upstream(upstream).await {
Ok(buffered) => buffered,
Err(response) => return response,
};
return proxy::buffered_response(buffered.status, buffered.headers, buffered.body);
}
let status = proxy::upstream_status(&upstream);
let headers = proxy::response_headers(&upstream);
let VerifiedSettle {
server,
payload,
requirements,
advertised,
resource_url,
reservation,
..
} = ctx;
let (inflight, signature) = reservation.disarm();
let tee = TeeBody::new(
upstream.bytes_stream(),
SettleJob {
server,
payload,
requirements,
advertised,
resource_url,
inflight,
signature,
rates,
ceiling,
missing_usage,
abort_usage,
},
);
let mut response = Response::new(Body::from_stream(tee));
*response.status_mut() = status;
*response.headers_mut() = headers;
response
}
fn spawn_exact_settle(ctx: VerifiedSettle<'_>) {
let VerifiedSettle {
server,
payload,
requirements,
advertised,
resource_url,
reservation,
..
} = ctx;
let (inflight, signature) = reservation.disarm();
tee::spawn_settle_task(tee::SpawnSettle {
server,
payload,
requirements,
advertised,
resource_url,
inflight,
signature,
overrides: None,
});
}
#[allow(
clippy::result_large_err,
reason = "gate returns HTTP responses as Err"
)]
async fn force_stream_usage(request: Request) -> Result<Request, Response> {
let (parts, body) = request.into_parts();
let bytes = match body.collect().await {
Ok(collected) => collected.to_bytes(),
Err(_) => return Err(openai_error::invalid_request("failed to read request body")),
};
let bytes = usage::force_include_usage(&bytes).map_or(bytes, Bytes::from);
Ok(Request::from_parts(parts, Body::from(bytes)))
}
async fn attach_settlement(ctx: VerifiedSettle<'_>, buffered: proxy::BufferedUpstream) -> Response {
let usage = usage::from_json(&buffered.body);
let actual = price::actual(
ctx.rates,
ctx.ceiling,
usage.as_ref(),
ctx.state.payment().missing_usage,
);
let overrides =
(ctx.rates.scheme == Scheme::Upto).then(|| SettlementOverrides::amount(actual.to_string()));
let settlement = match ctx
.server
.settle_payment(
&ctx.payload,
&ctx.requirements,
overrides.as_ref(),
SettlePhase::AfterHandler,
Some(ctx.resource_url.as_str()),
Some(&ctx.advertised),
)
.await
{
Ok(settlement) => settlement,
Err(error) => return settle_facilitator_error(error),
};
drop(ctx.reservation);
if !settlement.is_success() {
return settlement_failure(&settlement);
}
let mut response = proxy::buffered_response(buffered.status, buffered.headers, buffered.body);
match settlement_to_header(&settlement) {
Ok(value) => {
let _ = response.headers_mut().insert(PAYMENT_RESPONSE, value);
ensure_expose_headers(response.headers_mut());
merge_private(response.headers_mut());
}
Err(error) => {
return json_status(
StatusCode::INTERNAL_SERVER_ERROR,
&json!({ "error": error.to_string() }),
);
}
}
response
}
#[allow(
clippy::result_large_err,
reason = "gate returns HTTP responses as Err"
)]
async fn build_payment_required(
server: &ResourceServer,
reqs: Vec<PaymentRequirements>,
resource: ResourceInfo,
) -> Result<PaymentRequired, Response> {
let supported = match server.facilitator().supported().await {
Ok(supported) => supported,
Err(error) => return Err(transport_or_internal(error)),
};
if let Err(error) = validate_accepts_against_supported(server, &reqs, &supported) {
return Err(json_status(
StatusCode::INTERNAL_SERVER_ERROR,
&json!({ "error": "facilitator support", "reason": error.to_string() }),
));
}
let mut body = match server
.create_payment_required_response(
reqs,
PaymentRequiredBuildContext {
resource,
error: None,
extensions: Extensions::new(),
supported: supported.clone(),
payment_payload: None,
},
)
.await
{
Ok(body) => body,
Err(error) => {
return Err(json_status(
StatusCode::INTERNAL_SERVER_ERROR,
&json!({ "error": error.to_string() }),
));
}
};
advertise_permit2_gas_sponsoring(&mut body, &supported, server);
Ok(body)
}
fn advertise_permit2_gas_sponsoring(
body: &mut PaymentRequired,
supported: &SupportedResponse,
server: &ResourceServer,
) {
if !body.accepts.iter().any(|req| is_evm_permit2(req, server)) {
return;
}
if body.extensions.get(EIP2612_GAS_SPONSORING_KEY).is_none() {
let entry = Eip2612GasSponsoringExtension::new().advertise(
&AdvertiseContext::for_payment_required(&body.resource, &body.accepts, None),
);
if let Some(entry) = entry {
body.extensions.insert(EIP2612_GAS_SPONSORING_KEY, entry);
}
}
let erc20_listed = supported
.extensions
.iter()
.any(|ext| ext.as_str() == ERC20_APPROVAL_GAS_SPONSORING_KEY);
if erc20_listed
&& body
.extensions
.get(ERC20_APPROVAL_GAS_SPONSORING_KEY)
.is_none()
{
let entry = Erc20ApprovalGasSponsoringExtension::new().advertise(
&AdvertiseContext::for_payment_required(&body.resource, &body.accepts, None),
);
if let Some(entry) = entry {
body.extensions
.insert(ERC20_APPROVAL_GAS_SPONSORING_KEY, entry);
}
}
}
fn is_evm_permit2(req: &PaymentRequirements, server: &ResourceServer) -> bool {
if extra_atm(req) == Some("permit2") {
return true;
}
server
.registered_scheme(req.scheme.as_str(), &req.network)
.is_some_and(|scheme| scheme.default_asset_transfer_method() == "permit2")
}
fn extra_atm(req: &PaymentRequirements) -> Option<&str> {
req.extra
.as_ref()
.and_then(|extra| extra.get("assetTransferMethod"))
.and_then(serde_json::Value::as_str)
}
#[allow(
clippy::result_large_err,
reason = "gate returns HTTP responses as Err"
)]
async fn verify_or_reject(
server: &ResourceServer,
payload: &WirePaymentPayload,
requirements: &PaymentRequirements,
payment_required: &PaymentRequired,
) -> Result<(), Response> {
let outcome = server
.verify_payment(payload, requirements, Some(&payment_required.extensions))
.await
.map_err(|error| verify_facilitator_error(payment_required, error))?;
match outcome.response {
r402_protocol::payment::VerifyResponse::Invalid {
reason, message, ..
} => {
let reason = reason.unwrap_or(ErrorReason::UnexpectedVerifyError);
let status = reason_to_status(&reason);
let detail = message.map_or_else(|| reason.to_string(), |text| text.to_string());
Err(challenge(
payment_required,
&format!("Verification failed: {detail}"),
status,
))
}
_ => Ok(()),
}
}
fn decode_signature(header: &[u8]) -> Option<WirePaymentPayload> {
let decoded = Base64Bytes::from(header).decode().ok()?;
serde_json::from_slice(&decoded).ok()
}
fn resource_info(state: &AppState, path: &str) -> ResourceInfo {
let url = state.base_url().map_or_else(
|| path.to_owned(),
|base| proxy::join_origin(base, path, None).to_string(),
);
ResourceInfo::new(url)
.with_description("o402 LLM gateway")
.with_mime_type("application/json")
}
fn challenge(payment_required: &PaymentRequired, error: &str, status: StatusCode) -> Response {
let body = payment_required.clone().with_error(error);
let Ok(body_bytes) = serde_json::to_vec(&body) else {
return json_status(
StatusCode::INTERNAL_SERVER_ERROR,
&json!({ "error": "payment-required serialization failed" }),
);
};
let Ok(header_value) = HeaderValue::from_bytes(Base64Bytes::encode(&body_bytes).as_ref())
else {
return json_status(
StatusCode::INTERNAL_SERVER_ERROR,
&json!({ "error": "payment-required header encoding failed" }),
);
};
let mut response = Response::new(Body::from(body_bytes));
*response.status_mut() = status;
let _ = response
.headers_mut()
.insert(PAYMENT_REQUIRED, header_value);
let _ = response.headers_mut().insert(
header::CONTENT_TYPE,
HeaderValue::from_static("application/json"),
);
ensure_expose_headers(response.headers_mut());
set_no_store(response.headers_mut());
response
}
fn settlement_failure(failure: &SettleResponse) -> Response {
let body_bytes = serde_json::to_vec(failure).unwrap_or_else(|_| b"{}".to_vec());
let header_value = failure
.encode_base64_any()
.and_then(|b64| HeaderValue::from_bytes(b64.as_ref()).ok());
let mut response = Response::new(Body::from(body_bytes));
*response.status_mut() = StatusCode::PAYMENT_REQUIRED;
let _ = response.headers_mut().insert(
header::CONTENT_TYPE,
HeaderValue::from_static("application/json"),
);
if let Some(header_value) = header_value {
let _ = response
.headers_mut()
.insert(PAYMENT_RESPONSE, header_value);
}
ensure_expose_headers(response.headers_mut());
set_no_store(response.headers_mut());
response
}
fn verify_facilitator_error(
payment_required: &PaymentRequired,
error: FacilitatorError,
) -> Response {
match error {
FacilitatorError::Transport { kind } => transport_response(kind),
other => {
let reason = other
.as_payment_problem()
.map_or(ErrorReason::UnexpectedVerifyError, |problem| {
problem.reason()
});
challenge(
payment_required,
&format!("Verification failed: {other}"),
reason_to_status(&reason),
)
}
}
}
fn settle_facilitator_error(error: FacilitatorError) -> Response {
match error {
FacilitatorError::Transport { kind } => transport_response(kind),
other => {
let mut response = json_status(
StatusCode::PAYMENT_REQUIRED,
&json!({
"error": "settlement aborted",
"details": other.to_string(),
}),
);
set_no_store(response.headers_mut());
response
}
}
}
fn transport_or_internal(error: FacilitatorError) -> Response {
match error {
FacilitatorError::Transport { kind } => transport_response(kind),
other => json_status(
StatusCode::INTERNAL_SERVER_ERROR,
&json!({ "error": other.to_string() }),
),
}
}
fn transport_response(kind: r402_protocol::error::FacilitatorTransportKind) -> Response {
json_status(
StatusCode::BAD_GATEWAY,
&json!({
"error": "facilitator transport",
"kind": kind.to_string(),
}),
)
}
fn json_status(status: StatusCode, body: &serde_json::Value) -> Response {
let mut response = Response::new(Body::from(body.to_string()));
*response.status_mut() = status;
let _ = response.headers_mut().insert(
header::CONTENT_TYPE,
HeaderValue::from_static("application/json"),
);
ensure_expose_headers(response.headers_mut());
response
}