#![cfg(fuzzing)]
#![allow(dead_code, unused_imports)]
mod admin;
mod admission;
mod aliases;
mod availability;
mod backends;
mod budget;
mod config;
mod convergence;
mod credentials;
mod desired_state;
mod error;
mod key_material;
mod mint;
mod ops;
mod policy;
mod pricing;
mod principals;
mod rate_limit;
mod redis_support;
mod reload;
mod revocation;
mod routes;
mod shutdown;
mod state;
mod status;
mod streaming;
#[allow(unused_imports)]
mod telemetry;
mod usage;
use std::collections::HashMap;
use std::sync::{Arc, Mutex, OnceLock};
use std::time::{Duration, UNIX_EPOCH};
use crate::backends::catalog::{
self, Admission, CatalogContent, CatalogDiff, CatalogSnapshot, CatalogSource, JsonPointer,
LastKnownGoodCatalog, ModelField, Refusable, Refusal, RefusalReason, SourceValidators,
};
use crate::backends::models_dev::{
self, ModelsDevAdapter, ModelsDevError, SEED_PAYLOAD, seed_snapshot,
};
use crate::budget::NoBudget;
use crate::config::{Config, Model};
use crate::mint::{MintAlgorithm, MintRequest};
use crate::principals::{
Presented, PrincipalStore, PrincipalStoreError, TokenVerificationError, TokenVerifier,
};
use crate::rate_limit::NoLimit;
use crate::revocation::NoDenylist;
use crate::state::{AppState, ConfigSnapshot, ReplicaObservability, SnapshotError};
use crate::usage::{UsageDelivery, UsageFanout};
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Rejection {
Load(String),
Invalid(String),
BadRequest(String),
Unauthenticated(&'static str),
Unauthorized(&'static str),
Catalog {
code: &'static str,
message: String,
pointer: Option<String>,
},
Unavailable,
}
impl Refusable for Rejection {
fn refusal(&self) -> Refusal {
let Self::Catalog { code, pointer, .. } = self else {
return Refusal::new(RefusalReason::Unknown);
};
let reason = RefusalReason::ALL
.iter()
.copied()
.find(|reason| reason.as_str() == *code)
.unwrap_or(RefusalReason::Unknown);
pointer.as_ref().map_or_else(
|| Refusal::new(reason),
|pointer| Refusal::at(reason, JsonPointer::new(pointer.clone())),
)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ConfigShape {
pub stateful: bool,
pub namespaces: usize,
pub providers: usize,
pub models: usize,
pub credentials: usize,
pub gateway_keys: usize,
pub verifiers: usize,
}
pub fn config_from_toml_str(input: &str) -> Result<ConfigShape, Rejection> {
match Config::from_toml_str(input) {
Ok(config) => Ok(ConfigShape {
stateful: config.mode == config::Mode::Stateful,
namespaces: config.namespace.len(),
providers: config.provider.len(),
models: config.model.len(),
credentials: config.credential.len(),
gateway_keys: config.gateway_key.len(),
verifiers: config.gateway_verifier.len(),
}),
Err(config::ConfigError::Load(message)) => Err(Rejection::Load(message)),
Err(config::ConfigError::Invalid(message)) => Err(Rejection::Invalid(message)),
}
}
pub fn credentials_query_namespaces(raw_query: Option<&str>) -> Result<Option<String>, Rejection> {
routes::fuzz_parse_credential_query(raw_query)
.map_err(|error| Rejection::BadRequest(error.to_string()))
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct VerifiedToken {
pub namespace: String,
pub subject: String,
pub capabilities: usize,
pub scoped_aliases: bool,
pub max_request_microdollars: Option<u64>,
}
pub fn verify_token(credential: &str) -> Result<Option<VerifiedToken>, Rejection> {
let presented = Presented { credential };
let resolved = futures::executor::block_on(verifier().resolve(&presented));
match resolved {
Ok(None) => Ok(None),
Ok(Some(key)) => Ok(Some(VerifiedToken {
namespace: key.namespace,
subject: key.subject,
capabilities: key.scope.map_or(0, |scope| scope.len()),
scoped_aliases: key.alias_scope.is_some(),
max_request_microdollars: key.max_request_microdollars,
})),
Err(PrincipalStoreError::Unauthorized(error)) => {
Err(Rejection::Unauthenticated(code(&error)))
}
Err(PrincipalStoreError::Forbidden(error)) => Err(Rejection::Unauthorized(code(&error))),
Err(PrincipalStoreError::Unavailable) => Err(Rejection::Unavailable),
}
}
pub fn mint_hs256_token(
namespace: &str,
subject: &str,
audience: &str,
ttl_seconds: u64,
issued_at: Option<u64>,
scope: Option<Vec<String>>,
aliases: Option<Vec<String>>,
) -> Option<String> {
mint::fuzz_mint_token_with_raw_scope(
MintRequest {
kid: HS256_KID,
algorithm: MintAlgorithm::Hs256,
key_material: HS256_MATERIAL,
namespace,
subject,
audience,
ttl: Duration::from_secs(ttl_seconds),
aliases,
max_request_microdollars: None,
scope: None,
},
issued_at,
scope,
)
.ok()
.map(|minted| minted.token)
}
pub fn resign_seed_onto_this_run(token: &str) -> Option<String> {
use base64::{Engine as _, engine::general_purpose::URL_SAFE_NO_PAD};
let mut segments = token.strip_prefix("axt1.")?.split('.');
let header: jsonwebtoken::Header =
serde_json::from_slice(&URL_SAFE_NO_PAD.decode(segments.next()?).ok()?).ok()?;
let mut claims: serde_json::Map<String, serde_json::Value> =
serde_json::from_slice(&URL_SAFE_NO_PAD.decode(segments.next()?).ok()?).ok()?;
let iat = claims.get("iat")?.as_u64()?;
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.ok()?
.as_secs();
let offset = i128::from(now) - i128::from(iat);
let shift = |value: &mut serde_json::Value| {
if let Some(seconds) = value.as_u64() {
let shifted = (i128::from(seconds) + offset).clamp(0, i128::from(u64::MAX));
*value = serde_json::Value::from(u64::try_from(shifted).unwrap_or(0));
}
};
for claim in ["iat", "exp", "nbf"] {
if let Some(value) = claims.get_mut(claim) {
shift(value);
}
}
let kid = header.kid.clone().unwrap_or_else(|| HS256_KID.to_owned());
mint::fuzz_sign_claims(
&header,
&serde_json::Value::Object(claims),
MintAlgorithm::Hs256,
HS256_MATERIAL,
&kid,
)
.ok()
}
pub const AUDIENCE: &str = "fuzz.axond.invalid";
pub const NAMESPACES: [&str; 2] = ["fuzz", "denied"];
pub const CAPABILITY_COUNT: usize = principals::Capability::ALL.len();
pub const MAX_TTL_SECONDS: u64 = 900;
pub fn epoch_min_iat() -> u64 {
static MIN_IAT: OnceLock<u64> = OnceLock::new();
*MIN_IAT.get_or_init(|| {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |since| since.as_secs())
.saturating_sub(EPOCH_LOOKBACK_SECONDS)
})
}
const EPOCH_LOOKBACK_SECONDS: u64 = 300;
pub const HS256_KID: &str = "fuzz-hs256";
pub const EDDSA_KID: &str = "fuzz-eddsa";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CatalogImport {
pub content_id: String,
pub source_url: String,
pub schema_version: &'static str,
pub providers: usize,
pub models: usize,
pub offerings: usize,
pub priced_offerings: usize,
pub overrides: usize,
pub model_ids: Vec<String>,
pub offering_keys: Vec<String>,
pub override_pointers: Vec<String>,
pub overrides_are_contradictions: bool,
pub overrides_point_into_offerings: bool,
pub raw_digest: String,
pub raw_bytes: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct CatalogDiffShape {
pub changes: usize,
pub providers_added: usize,
pub providers_removed: usize,
pub providers_changed: usize,
pub models_added: usize,
pub models_removed: usize,
pub offerings_added: usize,
pub offerings_removed: usize,
pub neutral_changed: usize,
pub lifecycle_changed: usize,
pub capabilities_changed: usize,
pub metadata_changed: usize,
pub prices_changed: usize,
}
impl CatalogDiffShape {
fn of(diff: &CatalogDiff) -> Self {
let counts = diff.counts();
Self {
changes: diff.changes().len(),
providers_added: counts.providers_added,
providers_removed: counts.providers_removed,
providers_changed: counts.providers_changed,
models_added: counts.models_added,
models_removed: counts.models_removed,
offerings_added: counts.offerings_added,
offerings_removed: counts.offerings_removed,
neutral_changed: counts.neutral_changed,
lifecycle_changed: counts.lifecycle_changed,
capabilities_changed: counts.capabilities_changed,
metadata_changed: counts.metadata_changed,
prices_changed: counts.prices_changed,
}
}
pub const fn is_price_only(self) -> bool {
self.prices_changed > 0 && self.prices_changed == self.changes
}
pub const fn is_metadata_only(self) -> bool {
self.prices_changed == 0
&& self.changes > 0
&& self.changes
== self.metadata_changed
+ self.capabilities_changed
+ self.lifecycle_changed
+ self.neutral_changed
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CatalogAdmission {
pub outcome: &'static str,
pub refusal: Option<Rejection>,
pub fetched: Vec<String>,
pub active_content_id: String,
pub active_models: usize,
pub active_is_seed: bool,
pub diff: Option<CatalogDiffShape>,
pub import: Option<CatalogImport>,
}
pub const CATALOG_SEED_PAYLOAD: &str = SEED_PAYLOAD;
pub fn catalog_source_url() -> String {
ModelsDevAdapter::default().source_url().to_owned()
}
pub fn catalog_seed_content_id() -> String {
seed_snapshot().content.content_id().to_string()
}
pub fn catalog_parse(
payload: &[u8],
fetched_at_secs: u64,
etag: Option<&str>,
) -> Result<CatalogImport, Rejection> {
parse_catalog(payload, fetched_at_secs, etag).map(|snapshot| describe_catalog(&snapshot))
}
pub fn catalog_import_over_seed(payload: &[u8], etag: Option<&str>) -> CatalogAdmission {
let seed = seed_snapshot();
let seed_content_id = seed.content.content_id();
let mut active = LastKnownGoodCatalog::default();
active.admit(seed);
let requested = Arc::new(Mutex::new(Vec::new()));
let fetch = RecordingFetch {
payload: payload.to_vec(),
etag: etag.map(ToOwned::to_owned),
requested: Arc::clone(&requested),
};
let source = models_dev::ModelsDevSource::new(ModelsDevAdapter::default(), fetch);
let refreshed = futures::executor::block_on(source.refresh(None));
let fetched = requested
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
let parsed = parse_catalog(payload, CATALOG_FETCHED_AT_SECS, etag);
assert_eq!(
parsed.is_ok(),
matches!(refreshed, Ok(catalog::CatalogRefresh::Updated { .. })),
"the source and the parser disagree about whether a payload is usable"
);
let import = parsed.as_ref().ok().map(describe_catalog);
let (outcome, refusal, diff) = match active.admit_result(parsed) {
Ok(Admission::Initial { .. }) => ("initial", None, None),
Ok(Admission::Unchanged { .. }) => ("unchanged", None, None),
Ok(Admission::Updated { diff, .. }) => ("updated", None, Some(CatalogDiffShape::of(&diff))),
Err((rejection, _)) => ("refused", Some(rejection), None),
};
let active = active.active().expect("the seed stays active");
CatalogAdmission {
outcome,
refusal,
fetched,
active_content_id: active.content.content_id().to_string(),
active_models: active.content.models().len(),
active_is_seed: active.content.content_id() == seed_content_id,
diff,
import,
}
}
pub fn catalog_diff(
previous: &[u8],
current: &[u8],
) -> Result<Option<CatalogDiffShape>, Rejection> {
let mut active = LastKnownGoodCatalog::default();
active.admit(parse_catalog(previous, CATALOG_FETCHED_AT_SECS, None)?);
match active.admit(parse_catalog(current, CATALOG_FETCHED_AT_SECS + 60, None)?) {
Admission::Unchanged { .. } => Ok(None),
Admission::Updated { diff, .. } => Ok(Some(CatalogDiffShape::of(&diff))),
Admission::Initial { .. } => unreachable!("the previous payload was admitted first"),
}
}
pub fn runtime_routes() -> Vec<String> {
routes_of(&runtime_state().config())
}
pub fn publication_moves_runtime_routes() -> bool {
let state = build_state(seam_config().clone()).expect("the seam's own config resolves");
let before = routes_of(&state.config());
let mut config = seam_config().clone();
let Some(extra) = config.model.first().cloned() else {
return false;
};
config.model.push(Model {
name: format!("{}-published", extra.name),
..extra
});
let snapshot = ConfigSnapshot::build(config, &seam_env(), 1)
.expect("the seam's config resolves with one more alias");
state.publish(snapshot);
routes_of(&state.config()) != before
}
pub fn override_oracle_notices_a_missing_override() -> bool {
let seed = seed_snapshot();
let mut models = seed.content.models().to_vec();
let mut removed = false;
for model in &mut models {
if model.neutral.is_none() {
continue;
}
for offering in &mut model.offerings {
if !offering.overrides.is_empty() {
offering.overrides.remove(0);
removed = true;
break;
}
}
if removed {
break;
}
}
if !removed {
return false;
}
let Ok(content) = CatalogContent::new(seed.content.providers().to_vec(), models) else {
return false;
};
!describe_catalog(&CatalogSnapshot {
source: seed.source,
content,
})
.overrides_are_contradictions
}
fn routes_of(snapshot: &ConfigSnapshot) -> Vec<String> {
snapshot
.config
.model
.iter()
.flat_map(|model| {
model.targets.iter().map(move |target| {
format!("{} => {}/{}", model.name, target.provider, target.model)
})
})
.collect()
}
pub const CATALOG_FETCHED_AT_SECS: u64 = 1_767_225_600;
fn parse_catalog(
payload: &[u8],
fetched_at_secs: u64,
etag: Option<&str>,
) -> Result<CatalogSnapshot, Rejection> {
let validators = etag.map_or_else(SourceValidators::default, SourceValidators::etag);
let fetched_at = UNIX_EPOCH + Duration::from_secs(fetched_at_secs);
ModelsDevAdapter::default()
.parse(payload, validators, fetched_at)
.map_err(|error| catalog_rejection(&error))
}
fn catalog_rejection(error: &ModelsDevError) -> Rejection {
use ModelsDevError as E;
let (code, pointer) = match error {
E::UnsupportedEndpoint { .. } => ("unsupported_endpoint", None),
E::NotJson { .. } => ("not_json", None),
E::Schema { pointer, .. } => ("schema", pointer.as_ref()),
E::IdMismatch { pointer, .. } => ("id_mismatch", Some(pointer)),
E::Identifier { pointer, .. } => ("identifier", Some(pointer)),
E::UnknownStatus { pointer, .. } => ("unknown_status", Some(pointer)),
E::UnknownModality { pointer, .. } => ("unknown_modality", Some(pointer)),
E::Price { pointer, .. } => ("price", Some(pointer)),
E::UnknownTierType { pointer, .. } => ("unknown_tier_type", Some(pointer)),
E::DuplicateTier { pointer } => ("duplicate_tier", Some(pointer)),
E::NeutralPrice { pointer } => ("neutral_price", Some(pointer)),
E::UncanonicalizableText { pointer, .. } => ("uncanonicalizable_text", Some(pointer)),
E::AmbiguousModelKey { pointer, .. } => ("ambiguous_model_key", Some(pointer)),
E::Content { .. } => ("content", None),
};
Rejection::Catalog {
code,
message: error.to_string(),
pointer: pointer.map(|pointer| pointer.as_str().to_owned()),
}
}
fn describe_catalog(snapshot: &CatalogSnapshot) -> CatalogImport {
let content = &snapshot.content;
let mut model_ids = Vec::with_capacity(content.models().len());
let mut offering_keys = Vec::new();
let mut override_pointers = Vec::new();
let mut priced_offerings = 0;
let mut overrides = 0;
let mut overrides_are_contradictions = true;
let mut overrides_point_into_offerings = true;
for model in content.models() {
model_ids.push(model.id.as_str().to_owned());
for offering in &model.offerings {
offering_keys.push(format!(
"{}|{}|{}",
model.id, offering.provider, offering.published_model_id
));
priced_offerings += usize::from(offering.price.is_some());
overrides += offering.overrides.len();
let stated: Vec<ModelField> =
offering.overrides.iter().map(|(field, _)| *field).collect();
match model.neutral.as_ref() {
Some(neutral) => {
let contradicted: Vec<ModelField> = offering.facts.differences(neutral);
if stated != contradicted {
overrides_are_contradictions = false;
}
}
None => {
if !stated.is_empty() {
overrides_are_contradictions = false;
}
}
}
for (field, pointer) in &offering.overrides {
if !pointer.as_str().starts_with(offering.pointer.as_str()) {
overrides_point_into_offerings = false;
}
override_pointers.push(format!(
"{}|{}|{}|{}",
model.id,
offering.provider,
field.as_str(),
pointer
));
}
}
}
CatalogImport {
content_id: content.content_id().to_string(),
source_url: snapshot.source.source_url.clone(),
schema_version: snapshot.source.schema_version.as_str(),
providers: content.providers().len(),
models: content.models().len(),
offerings: content.offering_count(),
priced_offerings,
overrides,
model_ids,
offering_keys,
override_pointers,
overrides_are_contradictions,
overrides_point_into_offerings,
raw_digest: snapshot.source.raw.digest.to_string(),
raw_bytes: snapshot.source.raw.size_bytes,
}
}
fn runtime_state() -> &'static AppState {
static STATE: OnceLock<AppState> = OnceLock::new();
STATE
.get_or_init(|| build_state(seam_config().clone()).expect("the seam's own config resolves"))
}
fn build_state(config: Config) -> Result<AppState, SnapshotError> {
AppState::with_resources(
config,
&seam_env(),
Arc::new(UsageDelivery::telemetry(UsageFanout::new(Vec::new()))),
Box::new(NoBudget),
Box::new(NoLimit),
Box::new(NoDenylist),
ReplicaObservability::stateless(),
)
}
struct RecordingFetch {
payload: Vec<u8>,
etag: Option<String>,
requested: Arc<Mutex<Vec<String>>>,
}
#[async_trait::async_trait]
impl models_dev::CatalogFetch for RecordingFetch {
async fn get(
&self,
url: &str,
_validators: Option<&SourceValidators>,
) -> Result<models_dev::FetchResponse, models_dev::FetchError> {
self.requested
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(url.to_owned());
Ok(models_dev::FetchResponse::Payload {
bytes: self.payload.clone(),
validators: self
.etag
.as_deref()
.map_or_else(SourceValidators::default, SourceValidators::etag),
})
}
}
const HS256_MATERIAL: &str = "axond-fuzz-hs256-material-not-a-secret";
const EDDSA_PUBLIC_BASE64: &str = "ZnV6ei1heG9uZC1lZDI1NTE5LXB1YmxpYy1rZXktMzI=";
const CONFIG: &str = r#"
[[namespace]]
id = "fuzz"
default = true
[[namespace]]
id = "denied"
[[gateway_key]]
env = "AXOND_FUZZ_STATIC_KEY"
namespace = "fuzz"
[gateway_token]
audience = "fuzz.axond.invalid"
[[gateway_verifier]]
kid = "fuzz-hs256"
alg = "HS256"
env = "AXOND_FUZZ_HS256"
namespaces = ["fuzz"]
max_ttl = "15m"
[[gateway_verifier]]
kid = "fuzz-eddsa"
alg = "EdDSA"
env = "AXOND_FUZZ_EDDSA"
namespaces = ["fuzz", "denied"]
max_ttl = "15m"
[[gateway_token_epoch]]
namespace = "fuzz"
min_iat = {MIN_IAT}
# One provider and one alias, so the request path this seam compiles has a
# routing table with something in it: `catalog_import` asserts that importing a
# catalogue leaves that table byte-for-byte alone, and a table that was empty to
# begin with would assert nothing. Neither is reachable — the base URL resolves
# nowhere and no target is ever dispatched.
[[provider]]
id = "fuzz-provider"
kind = "openai-compatible"
base_url = "https://provider.fuzz.axond.invalid/v1"
[[credential]]
namespace = "fuzz"
provider = "fuzz-provider"
env = "AXOND_FUZZ_PROVIDER_KEY"
[[model]]
name = "fuzz-alias"
targets = [
{ provider = "fuzz-provider", model = "fuzz-upstream-model", price = { input_microdollars_per_million = 1, output_microdollars_per_million = 1 } },
]
"#;
fn seam_config() -> &'static Config {
static CONFIG_ONCE: OnceLock<Config> = OnceLock::new();
CONFIG_ONCE.get_or_init(|| {
let text = CONFIG.replace("{MIN_IAT}", &epoch_min_iat().to_string());
Config::from_toml_str(&text).expect("the seam's own config is valid")
})
}
fn seam_env() -> HashMap<String, String> {
HashMap::from([
(
"AXOND_FUZZ_STATIC_KEY".to_owned(),
"axond-fuzz-static-key-not-a-secret".to_owned(),
),
("AXOND_FUZZ_HS256".to_owned(), HS256_MATERIAL.to_owned()),
(
"AXOND_FUZZ_EDDSA".to_owned(),
EDDSA_PUBLIC_BASE64.to_owned(),
),
(
"AXOND_FUZZ_PROVIDER_KEY".to_owned(),
"axond-fuzz-provider-key-not-a-secret".to_owned(),
),
])
}
fn verifier() -> &'static TokenVerifier {
static VERIFIER: OnceLock<TokenVerifier> = OnceLock::new();
VERIFIER.get_or_init(|| {
let config = seam_config();
let env = seam_env();
TokenVerifier::build(config, &env)
.expect("the seam's own verifiers build")
.expect("the seam configures verifiers")
})
}
fn code(error: &TokenVerificationError) -> &'static str {
error.code()
}