use std::collections::BTreeSet;
use std::sync::{Arc, Mutex};
use std::time::Instant;
use sha2::{Digest, Sha256};
use crate::auth::AuthContext;
use crate::errors::RpcError;
pub const INTROSPECT_ENDPOINT: &str = "/__introspect_token__";
pub const INTROSPECT_ENABLED_HEADER: &str = "vgi-token-introspection";
pub const MAX_INTROSPECT_BODY_BYTES: usize = 8192;
const MAX_TOKEN_CHARS: usize = 4096;
pub const DEFAULT_INTROSPECT_TTL_SECONDS: u64 = 300;
pub const DEFAULT_INTROSPECT_RATE_LIMIT: u32 = 20;
pub fn token_digest(token: &str) -> String {
let mut h = Sha256::new();
h.update(token.as_bytes());
format!("{:x}", h.finalize())
}
pub fn is_jws_shaped(token: &str) -> bool {
let mut parts = token.split('.');
let (Some(a), Some(b), Some(c), None) =
(parts.next(), parts.next(), parts.next(), parts.next())
else {
return false;
};
let b64url = |s: &str| {
s.bytes()
.all(|c| c.is_ascii_alphanumeric() || c == b'-' || c == b'_')
};
!a.is_empty() && !b.is_empty() && b64url(a) && b64url(b) && b64url(c)
}
#[derive(Clone, Debug)]
pub struct TokenIdentity {
pub principal: String,
pub token_name: String,
pub ttl_seconds: Option<u64>,
}
impl TokenIdentity {
pub fn new(principal: impl Into<String>) -> Self {
Self {
principal: principal.into(),
token_name: String::new(),
ttl_seconds: None,
}
}
pub fn with_token_name(mut self, name: impl Into<String>) -> Self {
self.token_name = name.into();
self
}
pub fn with_ttl_seconds(mut self, ttl: u64) -> Self {
self.ttl_seconds = Some(ttl);
self
}
}
pub type TokenResolver =
Arc<dyn Fn(&str) -> std::result::Result<Option<TokenIdentity>, RpcError> + Send + Sync>;
#[derive(Debug)]
pub enum IntrospectOutcome {
Resolved {
principal: String,
token_name: String,
ttl_seconds: u64,
},
NotAnIntrospector,
Unresolved,
RateLimited,
Unavailable { retry_after_seconds: u32 },
}
struct RateLimiter {
per_window: u32,
window: std::time::Duration,
state: Mutex<(Instant, std::collections::HashMap<String, u32>)>,
}
impl RateLimiter {
fn new(per_window: u32) -> Self {
Self {
per_window,
window: std::time::Duration::from_secs(1),
state: Mutex::new((Instant::now(), std::collections::HashMap::new())),
}
}
fn allow(&self, key: &str) -> bool {
let now = Instant::now();
let mut guard = self.state.lock().unwrap_or_else(|e| e.into_inner());
let (start, counts) = &mut *guard;
if now.duration_since(*start) >= self.window {
counts.clear();
*start = now;
}
let count = counts.entry(key.to_string()).or_insert(0);
if *count >= self.per_window {
return false;
}
*count += 1;
true
}
}
pub struct TokenIntrospector {
resolver: TokenResolver,
principals: BTreeSet<String>,
default_ttl_seconds: u64,
limiter: RateLimiter,
}
impl std::fmt::Debug for TokenIntrospector {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("TokenIntrospector")
.field("principals", &self.principals.len())
.field("default_ttl_seconds", &self.default_ttl_seconds)
.finish_non_exhaustive()
}
}
impl TokenIntrospector {
pub fn new<I, S>(
resolver: TokenResolver,
principals: I,
default_ttl_seconds: u64,
rate_limit_per_second: u32,
) -> Self
where
I: IntoIterator<Item = S>,
S: Into<String>,
{
let principals: BTreeSet<String> = principals
.into_iter()
.map(Into::into)
.filter(|p| !p.is_empty())
.collect();
assert!(
!principals.is_empty(),
"introspect_principals must name at least one principal. Introspection \
is a distinct capability from authentication: allowing any \
authenticated caller lets any user resolve any other user's \
credential to its owner."
);
assert!(
default_ttl_seconds > 0,
"introspect_default_ttl must be positive: a zero or absent TTL silently \
disables the caller's cache and turns every request into a round trip."
);
Self {
resolver,
principals,
default_ttl_seconds,
limiter: RateLimiter::new(rate_limit_per_second),
}
}
pub fn introspect(&self, auth: &AuthContext, body: &[u8]) -> IntrospectOutcome {
if !auth.authenticated || !self.principals.contains(&auth.principal) {
tracing::warn!(
target: "vgi_rpc.http.introspect",
principal = %auth.principal,
authenticated = auth.authenticated,
"introspection refused: caller is not an introspector"
);
return IntrospectOutcome::NotAnIntrospector;
}
if !self.limiter.allow(&auth.principal) {
tracing::warn!(
target: "vgi_rpc.http.introspect",
principal = %auth.principal,
"introspection rate limit exceeded"
);
return IntrospectOutcome::RateLimited;
}
let Some(token) = parse_token(body) else {
return IntrospectOutcome::Unresolved;
};
let digest = token_digest(&token);
if is_jws_shaped(&token) {
tracing::warn!(
target: "vgi_rpc.http.introspect",
principal = %auth.principal,
token_digest = %digest,
"introspection refused: JWS-shaped subject"
);
return IntrospectOutcome::Unresolved;
}
match (self.resolver)(&token) {
Ok(Some(identity)) => {
tracing::info!(
target: "vgi_rpc.http.introspect",
principal = %auth.principal,
token_digest = %digest,
resolved_principal = %identity.principal,
"introspection: resolved"
);
IntrospectOutcome::Resolved {
principal: identity.principal,
token_name: identity.token_name,
ttl_seconds: identity.ttl_seconds.unwrap_or(self.default_ttl_seconds),
}
}
Ok(None) => {
tracing::info!(
target: "vgi_rpc.http.introspect",
principal = %auth.principal,
token_digest = %digest,
"introspection: credential did not resolve"
);
IntrospectOutcome::Unresolved
}
Err(err) => {
tracing::error!(
target: "vgi_rpc.http.introspect",
principal = %auth.principal,
token_digest = %digest,
error = %err.message,
"introspection unavailable"
);
IntrospectOutcome::Unavailable {
retry_after_seconds: err
.retry_after_seconds
.unwrap_or(crate::errors::DEFAULT_AUTH_RETRY_AFTER_SECONDS),
}
}
}
}
}
fn parse_token(body: &[u8]) -> Option<String> {
if body.len() > MAX_INTROSPECT_BODY_BYTES {
return None;
}
let value: serde_json::Value = serde_json::from_slice(body).ok()?;
let token = value.get("token")?.as_str()?;
if token.is_empty() || token.len() > MAX_TOKEN_CHARS {
return None;
}
Some(token.to_string())
}
#[cfg(test)]
mod tests {
use super::*;
const SUBJECT: &str = "opaque-subject-token";
fn introspector() -> TokenIntrospector {
TokenIntrospector::new(
Arc::new(|token: &str| {
Ok((token == SUBJECT || is_jws_shaped(token))
.then(|| TokenIdentity::new("subject@example").with_token_name("laptop")))
}),
["proxy"],
DEFAULT_INTROSPECT_TTL_SECONDS,
DEFAULT_INTROSPECT_RATE_LIMIT,
)
}
fn caller(principal: &str) -> AuthContext {
AuthContext::for_principal("conformance", principal)
}
fn body(token: &str) -> Vec<u8> {
serde_json::json!({ "token": token })
.to_string()
.into_bytes()
}
#[test]
fn resolves_a_known_credential() {
let outcome = introspector().introspect(&caller("proxy"), &body(SUBJECT));
let IntrospectOutcome::Resolved {
principal,
token_name,
ttl_seconds,
} = outcome
else {
panic!("expected a resolution, got {outcome:?}");
};
assert_eq!(principal, "subject@example");
assert_eq!(token_name, "laptop");
assert_eq!(ttl_seconds, DEFAULT_INTROSPECT_TTL_SECONDS);
}
#[test]
fn authentication_alone_does_not_grant_introspection() {
let outcome = introspector().introspect(&caller("someone-else"), &body(SUBJECT));
assert!(matches!(outcome, IntrospectOutcome::NotAnIntrospector));
let outcome = introspector().introspect(&AuthContext::anonymous(), &body(SUBJECT));
assert!(matches!(outcome, IntrospectOutcome::NotAnIntrospector));
}
#[test]
#[should_panic(expected = "at least one principal")]
fn empty_allowlist_is_not_a_permissive_default() {
let _ = TokenIntrospector::new(
Arc::new(|_: &str| Ok(None)),
Vec::<String>::new(),
DEFAULT_INTROSPECT_TTL_SECONDS,
DEFAULT_INTROSPECT_RATE_LIMIT,
);
}
#[test]
fn jws_shaped_subject_never_reaches_the_resolver() {
let jws = "eyJhbGciOiJIUzI1NiJ9.eyJzdWIiOiJhbGljZSJ9.c2lnbmF0dXJl";
assert!(is_jws_shaped(jws));
let resolver_ran = Arc::new(std::sync::atomic::AtomicBool::new(false));
let flag = resolver_ran.clone();
let it = TokenIntrospector::new(
Arc::new(move |_: &str| {
flag.store(true, std::sync::atomic::Ordering::SeqCst);
Ok(Some(TokenIdentity::new("subject@example")))
}),
["proxy"],
DEFAULT_INTROSPECT_TTL_SECONDS,
DEFAULT_INTROSPECT_RATE_LIMIT,
);
assert!(matches!(
it.introspect(&caller("proxy"), &body(jws)),
IntrospectOutcome::Unresolved
));
assert!(
!resolver_ran.load(std::sync::atomic::Ordering::SeqCst),
"the resolver was handed a JWS"
);
}
#[test]
fn jws_shape_test_does_not_catch_opaque_credentials() {
assert!(!is_jws_shaped("conformance-opaque-subject-token"));
assert!(!is_jws_shaped("a.b"));
assert!(!is_jws_shaped("a.b.c.d"));
assert!(!is_jws_shaped("a.b.c!"));
assert!(!is_jws_shaped(".b.c"));
assert!(
is_jws_shaped("a.b."),
"an unsecured JWS still has the shape"
);
}
#[test]
fn unknown_expired_and_malformed_are_one_answer() {
let it = introspector();
for probe in [
body("no-such-credential"),
body("expired-credential"),
body("!!malformed!!"),
b"not json at all".to_vec(),
b"{}".to_vec(),
serde_json::json!({ "token": 7 }).to_string().into_bytes(),
serde_json::json!({ "token": "x".repeat(MAX_TOKEN_CHARS + 1) })
.to_string()
.into_bytes(),
] {
assert!(matches!(
it.introspect(&caller("proxy"), &probe),
IntrospectOutcome::Unresolved
));
}
}
#[test]
fn an_oversized_body_is_refused_without_being_parsed() {
let huge =
serde_json::json!({ "token": "x", "pad": "p".repeat(MAX_INTROSPECT_BODY_BYTES) })
.to_string()
.into_bytes();
assert!(matches!(
introspector().introspect(&caller("proxy"), &huge),
IntrospectOutcome::Unresolved
));
}
#[test]
fn a_resolver_outage_is_transient_not_a_rejection() {
let it = TokenIntrospector::new(
Arc::new(|_: &str| {
Err(RpcError::auth_unavailable("token store down").with_retry_after(7))
}),
["proxy"],
DEFAULT_INTROSPECT_TTL_SECONDS,
DEFAULT_INTROSPECT_RATE_LIMIT,
);
assert!(matches!(
it.introspect(&caller("proxy"), &body(SUBJECT)),
IntrospectOutcome::Unavailable {
retry_after_seconds: 7
}
));
}
#[test]
fn rate_limit_bounds_the_oracle() {
let it = TokenIntrospector::new(
Arc::new(|_: &str| Ok(None)),
["proxy"],
DEFAULT_INTROSPECT_TTL_SECONDS,
2,
);
let probe = body("guess");
assert!(matches!(
it.introspect(&caller("proxy"), &probe),
IntrospectOutcome::Unresolved
));
assert!(matches!(
it.introspect(&caller("proxy"), &probe),
IntrospectOutcome::Unresolved
));
assert!(matches!(
it.introspect(&caller("proxy"), &probe),
IntrospectOutcome::RateLimited
));
}
#[test]
fn token_digest_is_stable_and_is_not_the_credential() {
let d = token_digest(SUBJECT);
assert_eq!(d, token_digest(SUBJECT));
assert_eq!(d.len(), 64);
assert!(!d.contains(SUBJECT));
}
}