Skip to main content

ironflow_auth_proxy/
lib.rs

1//! # ironflow-auth-proxy
2//!
3//! HTTP service that keeps the Claude credential out of ironflow agent pods.
4//!
5//! The worker (`K8sEphemeralProvider::auth_proxy`) asks the admin API for an
6//! opaque token bound to one run and one step, and hands only that token to
7//! the pod. Claude Code sends it as `Authorization: Bearer` to this proxy,
8//! which swaps it for the real credential and relays the request to
9//! `api.anthropic.com`, streaming the answer back.
10//!
11//! The same mechanism keeps other secrets (a GitHub or GitLab token, an API
12//! key) out of the pod: a grant can hold a proxied secret with a host
13//! allowlist instead of the Claude credential. The pod calls
14//! `/r/<host>/<path>` with its opaque token (as `Authorization: Bearer`,
15//! `x-api-key`, `Private-Token` or the password of `Authorization: Basic`),
16//! and the proxy relays to `https://<host>/<path>` with the real secret
17//! injected the way the grant says (bearer, `Private-Token`, `x-api-key`, a
18//! named header or Basic). A host outside the allowlist (exact names or a
19//! leading `*.`, no port, no IP) gets a 403, so does a Claude token on `/r/`
20//! and a secret token on the Anthropic API. Redirects are returned to the
21//! pod, never followed.
22//!
23//! One listener serves:
24//!
25//! * `GET /healthz` - liveness;
26//! * `GET /metrics` - Prometheus metrics, when a recorder handle is given
27//!   ([`AuthProxyState::with_metrics`]): [`REQUESTS_TOTAL`] counts the `/r/`
28//!   requests by `secret` name and `result`;
29//! * `/admin/v1/...` - token issuance and revocation, behind
30//!   `Authorization: Bearer <IRONFLOW_AUTH_PROXY_ADMIN_KEY>`;
31//! * `/r/<host>/<path>` - the relay of proxied secrets: an unknown, expired
32//!   or revoked token gets a 401, a host outside the allowlist or a path with
33//!   `..`, `//` or percent-encoding a 403, CONNECT, TRACE or OPTIONS a 405.
34//!   Each request logs one `secret relay` event with its `result`
35//!   (`relayed`, `forbidden_host`, `forbidden_path`, `forbidden_method`,
36//!   `unknown_token`, `expired`, `revoked`, `upstream_error`,
37//!   `unavailable`);
38//! * anything else - the Anthropic relay: an unknown, expired or revoked
39//!   token gets a 401, a path outside `/v1/` or a request for another host a
40//!   403, a method other than GET/POST a 405.
41//!
42//! Grants live in a registry with two backends:
43//!
44//! * in memory (default): a single replica, tokens lost on restart;
45//! * PostgreSQL, when [`DATABASE_URL_ENV`] is set ([`registry_from_config`]):
46//!   shared by several replicas and surviving restarts. Only the token
47//!   SHA-256 is stored, never the token, and the credential is AES-256-GCM
48//!   encrypted at rest with the `IRONFLOW_SECRET_KEYS` key ring.
49//!
50//! Logs never contain a token, a credential, a secret, a header or a query
51//! string, only the short token id.
52//!
53//! # Examples
54//!
55//! ```no_run
56//! use std::env::var;
57//! use std::time::Duration;
58//!
59//! use ironflow_auth_proxy::{
60//!     AuthProxyConfig, AuthProxyState, DATABASE_URL_ENV, registry_from_config, serve, spawn_purge,
61//! };
62//! use ironflow_store::crypto::KeyRing;
63//! use tokio::net::TcpListener;
64//!
65//! # async fn example() -> Result<(), Box<dyn std::error::Error>> {
66//! let database_url = var(DATABASE_URL_ENV).ok();
67//! let registry = registry_from_config(database_url.as_deref(), KeyRing::from_env()?).await?;
68//! let config = AuthProxyConfig::new("0123456789abcdef0123456789abcdef");
69//! let state = AuthProxyState::with_registry(config, registry)?;
70//! let purge = spawn_purge(state.registry().clone(), Duration::from_secs(60));
71//! let listener = TcpListener::bind("0.0.0.0:8080").await?;
72//! serve(listener, state).await?;
73//! purge.abort();
74//! # Ok(())
75//! # }
76//! ```
77
78use std::fmt;
79use std::future::pending;
80use std::io;
81use std::sync::Arc;
82use std::time::{Duration, SystemTime, SystemTimeError, UNIX_EPOCH};
83
84use axum::body::{Body, Bytes, to_bytes};
85use axum::extract::{DefaultBodyLimit, Path, Request, State};
86use axum::http::header::{AUTHORIZATION, CONTENT_TYPE};
87use axum::http::{HeaderMap, Method, StatusCode};
88use axum::response::{IntoResponse, Response};
89use axum::routing::{delete, get, post};
90use axum::serve as serve_router;
91use axum::{Json, Router};
92use metrics::counter;
93use metrics_exporter_prometheus::PrometheusHandle;
94use reqwest::redirect::Policy;
95use reqwest::{Client, Error as ReqwestError};
96use serde_json::{from_slice, json};
97use thiserror::Error;
98use tokio::net::TcpListener;
99use tokio::signal::ctrl_c;
100#[cfg(unix)]
101use tokio::signal::unix::{SignalKind, signal};
102use tokio::task::JoinHandle;
103use tokio::{select, spawn, time};
104use tracing::{info, warn};
105use url::Url;
106
107use ironflow_core::auth_proxy::{
108    AuthProxyError, AuthProxyRegistry, DEFAULT_UPSTREAM, Grant, GrantCredential, RELAY_PREFIX,
109    TokenRejection, TokenRequest, admin_key_matches, downstream_headers, error_body,
110    extract_opaque_token, is_allowed_method, is_allowed_path, is_relay_method, is_relay_path,
111    is_valid_request_host, secret_upstream_headers, upstream_headers,
112};
113use ironflow_store::crypto::{CryptoError, KeyRing};
114use ironflow_store::error::StoreError;
115use ironflow_store::postgres::PostgresStore;
116
117/// Environment variable holding the PostgreSQL URL of the shared token
118/// registry. Unset (or empty), tokens live in memory.
119pub const DATABASE_URL_ENV: &str = "IRONFLOW_AUTH_PROXY_DATABASE_URL";
120
121/// Shortest admin key accepted.
122pub const MIN_ADMIN_KEY_LEN: usize = 32;
123
124/// Default largest request body relayed: 32 MiB.
125pub const DEFAULT_MAX_BODY_BYTES: usize = 32 * 1024 * 1024;
126
127/// Prometheus counter of the requests to the `/r/` relay, labelled by
128/// `secret` (its name, empty for an unknown token) and `result`.
129pub const REQUESTS_TOTAL: &str = "ironflow_auth_proxy_requests_total";
130
131/// Content type of the Prometheus text exposition format.
132const METRICS_CONTENT_TYPE: &str = "text/plain; version=0.0.4";
133
134/// Base of the URL of a `/r/` request, whose host is then replaced by the
135/// requested one.
136const SECRET_UPSTREAM_BASE: &str = "https://localhost";
137
138/// Length of the token id prefix written to logs.
139const SHORT_ID_LEN: usize = 12;
140
141/// Connect timeout towards the upstream. There is no total timeout: an
142/// agent turn streams for as long as the model answers.
143const UPSTREAM_CONNECT_TIMEOUT: Duration = Duration::from_secs(10);
144
145/// Configuration of the proxy. Its [`Debug`] output never shows the admin key.
146///
147/// # Examples
148///
149/// ```
150/// use ironflow_auth_proxy::AuthProxyConfig;
151///
152/// let config = AuthProxyConfig::new("0123456789abcdef0123456789abcdef");
153/// assert_eq!(config.upstream.as_str(), "https://api.anthropic.com/");
154/// assert!(config.secret_upstream.is_none());
155/// assert!(!format!("{config:?}").contains("0123456789abcdef"));
156/// ```
157#[derive(Clone)]
158pub struct AuthProxyConfig {
159    /// Where requests are relayed. Always `https://api.anthropic.com` in the binary.
160    pub upstream: Url,
161    /// Where `/r/<host>/` requests are relayed instead of `https://<host>`.
162    /// Always `None` in the binary.
163    pub secret_upstream: Option<Url>,
164    /// Key protecting the admin API.
165    pub admin_key: String,
166    /// Largest request body relayed, in bytes.
167    pub max_body_bytes: usize,
168}
169
170impl fmt::Debug for AuthProxyConfig {
171    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
172        f.debug_struct("AuthProxyConfig")
173            .field("upstream", &self.upstream.as_str())
174            .field(
175                "secret_upstream",
176                &self.secret_upstream.as_ref().map(Url::as_str),
177            )
178            .field("admin_key", &"<redacted>")
179            .field("max_body_bytes", &self.max_body_bytes)
180            .finish()
181    }
182}
183
184impl AuthProxyConfig {
185    /// Configuration relaying to [`DEFAULT_UPSTREAM`] with the given admin key.
186    ///
187    /// # Panics
188    ///
189    /// Panics if `admin_key` is shorter than [`MIN_ADMIN_KEY_LEN`].
190    ///
191    /// # Examples
192    ///
193    /// See [`AuthProxyConfig`].
194    pub fn new(admin_key: &str) -> Self {
195        assert!(
196            admin_key.len() >= MIN_ADMIN_KEY_LEN,
197            "the auth proxy admin key must be at least {MIN_ADMIN_KEY_LEN} characters"
198        );
199        Self {
200            upstream: Url::parse(DEFAULT_UPSTREAM).expect("DEFAULT_UPSTREAM is a valid URL"),
201            secret_upstream: None,
202            admin_key: admin_key.to_string(),
203            max_body_bytes: DEFAULT_MAX_BODY_BYTES,
204        }
205    }
206
207    /// Relay to another upstream. **For tests only**: the binary never calls
208    /// it, so a deployed proxy only reaches `api.anthropic.com`.
209    ///
210    /// # Examples
211    ///
212    /// ```
213    /// use ironflow_auth_proxy::AuthProxyConfig;
214    /// use url::Url;
215    ///
216    /// # fn example() -> Result<(), url::ParseError> {
217    /// let config = AuthProxyConfig::new("0123456789abcdef0123456789abcdef")
218    ///     .with_upstream(Url::parse("http://127.0.0.1:9000")?);
219    /// assert_eq!(config.upstream.port(), Some(9000));
220    /// # Ok(())
221    /// # }
222    /// ```
223    pub fn with_upstream(mut self, upstream: Url) -> Self {
224        self.upstream = upstream;
225        self
226    }
227
228    /// Relay every `/r/<host>/` request to `upstream` instead of
229    /// `https://<host>`. The host allowlist is still checked on `<host>`.
230    /// **For tests only**: the binary never calls it, so a deployed proxy
231    /// only reaches the allowlisted hosts, over https.
232    ///
233    /// # Examples
234    ///
235    /// ```
236    /// use ironflow_auth_proxy::AuthProxyConfig;
237    /// use url::Url;
238    ///
239    /// # fn example() -> Result<(), url::ParseError> {
240    /// let config = AuthProxyConfig::new("0123456789abcdef0123456789abcdef")
241    ///     .with_secret_upstream(Url::parse("http://127.0.0.1:9001")?);
242    /// assert_eq!(config.secret_upstream.and_then(|url| url.port()), Some(9001));
243    /// # Ok(())
244    /// # }
245    /// ```
246    pub fn with_secret_upstream(mut self, upstream: Url) -> Self {
247        self.secret_upstream = Some(upstream);
248        self
249    }
250}
251
252/// Shared state of the proxy: the token registry, the configuration and the
253/// upstream HTTP client. Cheap to clone.
254///
255/// # Examples
256///
257/// ```
258/// use ironflow_auth_proxy::{AuthProxyConfig, AuthProxyState};
259///
260/// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
261/// let state = AuthProxyState::new(AuthProxyConfig::new("0123456789abcdef0123456789abcdef"))?;
262/// assert!(state.registry().is_empty().await?);
263/// # Ok(())
264/// # }
265/// ```
266#[derive(Clone)]
267pub struct AuthProxyState {
268    registry: AuthProxyRegistry,
269    config: Arc<AuthProxyConfig>,
270    http: Client,
271    metrics: Option<PrometheusHandle>,
272}
273
274impl AuthProxyState {
275    /// Build the state with an empty in-memory registry.
276    ///
277    /// # Errors
278    ///
279    /// Returns the [`reqwest::Error`] raised when the HTTP client cannot be
280    /// built (TLS backend initialisation).
281    ///
282    /// # Examples
283    ///
284    /// See [`AuthProxyState`].
285    pub fn new(config: AuthProxyConfig) -> Result<Self, ReqwestError> {
286        Self::with_registry(config, AuthProxyRegistry::default())
287    }
288
289    /// Build the state over `registry`, for instance the shared PostgreSQL
290    /// registry returned by [`registry_from_config`]. Proxies built over
291    /// registries sharing a backend serve the same tokens.
292    ///
293    /// The upstream client never follows redirects and only bounds the
294    /// connection time.
295    ///
296    /// # Errors
297    ///
298    /// Returns the [`reqwest::Error`] raised when the HTTP client cannot be
299    /// built (TLS backend initialisation).
300    ///
301    /// # Examples
302    ///
303    /// ```
304    /// use ironflow_auth_proxy::{AuthProxyConfig, AuthProxyState};
305    /// use ironflow_core::auth_proxy::AuthProxyRegistry;
306    ///
307    /// # fn example() -> Result<(), reqwest::Error> {
308    /// let registry = AuthProxyRegistry::default();
309    /// let config = AuthProxyConfig::new("0123456789abcdef0123456789abcdef");
310    /// let a = AuthProxyState::with_registry(config.clone(), registry.clone())?;
311    /// let b = AuthProxyState::with_registry(config, registry)?;
312    /// # let _ = (a, b);
313    /// # Ok(())
314    /// # }
315    /// ```
316    pub fn with_registry(
317        config: AuthProxyConfig,
318        registry: AuthProxyRegistry,
319    ) -> Result<Self, ReqwestError> {
320        let http = Client::builder()
321            .connect_timeout(UPSTREAM_CONNECT_TIMEOUT)
322            .redirect(Policy::none())
323            .build()?;
324        Ok(Self {
325            registry,
326            config: Arc::new(config),
327            http,
328            metrics: None,
329        })
330    }
331
332    /// The token registry.
333    ///
334    /// # Examples
335    ///
336    /// See [`AuthProxyState`].
337    pub fn registry(&self) -> &AuthProxyRegistry {
338        &self.registry
339    }
340
341    /// Serve `GET /metrics` from `handle`, the handle of the installed
342    /// Prometheus recorder. Without it, `/metrics` answers 404.
343    ///
344    /// # Examples
345    ///
346    /// ```no_run
347    /// use ironflow_auth_proxy::{AuthProxyConfig, AuthProxyState};
348    /// use metrics_exporter_prometheus::PrometheusBuilder;
349    ///
350    /// # fn example() -> Result<(), Box<dyn std::error::Error>> {
351    /// let handle = PrometheusBuilder::new().install_recorder()?;
352    /// let state = AuthProxyState::new(AuthProxyConfig::new("0123456789abcdef0123456789abcdef"))?
353    ///     .with_metrics(handle);
354    /// # let _ = state;
355    /// # Ok(())
356    /// # }
357    /// ```
358    pub fn with_metrics(mut self, handle: PrometheusHandle) -> Self {
359        self.metrics = Some(handle);
360        self
361    }
362}
363
364/// Why the token registry could not be built. No variant carries the
365/// database URL, a key or a credential.
366///
367/// # Examples
368///
369/// ```
370/// use ironflow_auth_proxy::RegistryConfigError;
371///
372/// assert!(RegistryConfigError::MissingKeyRing.to_string().contains("IRONFLOW_SECRET_KEYS"));
373/// ```
374#[derive(Debug, Error)]
375pub enum RegistryConfigError {
376    /// A database is configured but no key to encrypt the credentials with.
377    #[error(
378        "{DATABASE_URL_ENV} is set but no encryption key is configured: set IRONFLOW_SECRET_KEYS (or IRONFLOW_SECRET_KEY)"
379    )]
380    MissingKeyRing,
381    /// The encryption key configuration is invalid.
382    #[error("invalid encryption key: {0}")]
383    Crypto(#[from] CryptoError),
384    /// The database cannot be opened or migrated.
385    #[error("cannot open the token registry database: {0}")]
386    Store(#[from] StoreError),
387}
388
389/// Build the token registry: in memory when `database_url` is `None` or
390/// blank, otherwise shared in PostgreSQL, the credentials encrypted with
391/// `key_ring`. Opening the database runs the store migrations.
392///
393/// # Errors
394///
395/// Returns [`RegistryConfigError::MissingKeyRing`] when a database is given
396/// without a key ring (checked before any connection), and
397/// [`RegistryConfigError::Store`] when the database cannot be reached or
398/// migrated.
399///
400/// # Examples
401///
402/// ```no_run
403/// use ironflow_auth_proxy::registry_from_config;
404/// use ironflow_store::crypto::KeyRing;
405///
406/// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
407/// let memory = registry_from_config(None, None).await?;
408/// assert!(memory.is_empty().await?);
409///
410/// let ring = KeyRing::from_spec(&format!("1:{}", "aa".repeat(32)), None)?;
411/// let shared = registry_from_config(Some("postgres://localhost/ironflow"), Some(ring)).await?;
412/// # let _ = shared;
413/// # Ok(())
414/// # }
415/// ```
416pub async fn registry_from_config(
417    database_url: Option<&str>,
418    key_ring: Option<KeyRing>,
419) -> Result<AuthProxyRegistry, RegistryConfigError> {
420    let Some(url) = database_url.map(str::trim).filter(|url| !url.is_empty()) else {
421        return Ok(AuthProxyRegistry::default());
422    };
423    let ring = key_ring.ok_or(RegistryConfigError::MissingKeyRing)?;
424    let mut store = PostgresStore::new(url).await?;
425    store.set_key_ring(ring);
426    Ok(AuthProxyRegistry::with_backend(Arc::new(store)))
427}
428
429/// The proxy router: health, metrics, admin API and relays.
430///
431/// # Examples
432///
433/// ```no_run
434/// use axum::serve;
435/// use ironflow_auth_proxy::{AuthProxyConfig, AuthProxyState, router};
436/// use tokio::net::TcpListener;
437///
438/// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
439/// let state = AuthProxyState::new(AuthProxyConfig::new("0123456789abcdef0123456789abcdef"))?;
440/// let listener = TcpListener::bind("127.0.0.1:0").await?;
441/// serve(listener, router(state)).await?;
442/// # Ok(())
443/// # }
444/// ```
445pub fn router(state: AuthProxyState) -> Router {
446    let limit = state.config.max_body_bytes;
447    Router::new()
448        .route("/healthz", get(healthz))
449        .route("/metrics", get(serve_metrics))
450        .route("/admin/v1/tokens", post(issue_token))
451        .route("/admin/v1/tokens/{id}", delete(revoke_token))
452        .route("/admin/v1/runs/{run_id}/tokens", delete(revoke_run))
453        .fallback(relay)
454        .layer(DefaultBodyLimit::max(limit))
455        .with_state(state)
456}
457
458/// Serve the proxy on `listener` until ctrl-c or SIGTERM, letting in-flight
459/// requests finish.
460///
461/// # Errors
462///
463/// Returns the I/O error that stopped the server.
464///
465/// # Examples
466///
467/// See the [crate documentation](crate).
468pub async fn serve(listener: TcpListener, state: AuthProxyState) -> io::Result<()> {
469    serve_router(listener, router(state))
470        .with_graceful_shutdown(shutdown_signal())
471        .await
472}
473
474/// Spawn a task dropping the expired grants of `registry` every `interval`.
475///
476/// Abort the returned handle to stop it.
477///
478/// # Panics
479///
480/// Panics if `interval` is zero, or when called outside a Tokio runtime.
481///
482/// # Examples
483///
484/// ```no_run
485/// use std::time::Duration;
486///
487/// use ironflow_auth_proxy::spawn_purge;
488/// use ironflow_core::auth_proxy::AuthProxyRegistry;
489///
490/// # async fn example() {
491/// let purge = spawn_purge(AuthProxyRegistry::default(), Duration::from_secs(60));
492/// purge.abort();
493/// # }
494/// ```
495pub fn spawn_purge(registry: AuthProxyRegistry, interval: Duration) -> JoinHandle<()> {
496    assert!(
497        !interval.is_zero(),
498        "purge interval must be greater than zero"
499    );
500    spawn(async move {
501        let mut ticker = time::interval(interval);
502        loop {
503            ticker.tick().await;
504            match now_unix() {
505                // Every replica purges: the deletes are idempotent.
506                Ok(now) => match registry.purge_expired(now).await {
507                    Ok(purged) if purged > 0 => info!(purged, "expired auth proxy tokens purged"),
508                    Ok(_) => {}
509                    Err(e) => warn!(error = %e, "expired auth proxy tokens purge failed"),
510                },
511                Err(e) => warn!(error = %e, "clock before the unix epoch; purge skipped"),
512            }
513        }
514    })
515}
516
517async fn shutdown_signal() {
518    let interrupt = async {
519        if let Err(e) = ctrl_c().await {
520            warn!(error = %e, "cannot listen for ctrl-c");
521            pending::<()>().await;
522        }
523    };
524    #[cfg(unix)]
525    let terminate = async {
526        match signal(SignalKind::terminate()) {
527            Ok(mut sigterm) => {
528                sigterm.recv().await;
529            }
530            Err(e) => {
531                warn!(error = %e, "cannot listen for SIGTERM");
532                pending::<()>().await;
533            }
534        }
535    };
536    #[cfg(not(unix))]
537    let terminate = pending::<()>();
538    select! {
539        () = interrupt => {},
540        () = terminate => {},
541    }
542    info!("shutting down");
543}
544
545fn now_unix() -> Result<u64, SystemTimeError> {
546    SystemTime::now()
547        .duration_since(UNIX_EPOCH)
548        .map(|d| d.as_secs())
549}
550
551fn short_id(id: &str) -> &str {
552    id.get(..SHORT_ID_LEN).unwrap_or(id)
553}
554
555fn error_response(status: StatusCode, kind: &str, message: &str) -> Response {
556    (status, Json(error_body(kind, message))).into_response()
557}
558
559fn clock_error(e: &SystemTimeError) -> Response {
560    warn!(error = %e, "system clock is before the unix epoch");
561    error_response(
562        StatusCode::INTERNAL_SERVER_ERROR,
563        "api_error",
564        "auth proxy clock error",
565    )
566}
567
568/// Whether the request carries the admin key as a bearer token.
569fn admin_authorized(state: &AuthProxyState, headers: &HeaderMap) -> bool {
570    headers
571        .get(AUTHORIZATION)
572        .and_then(|value| value.to_str().ok())
573        .and_then(|value| value.strip_prefix("Bearer "))
574        .is_some_and(|key| admin_key_matches(&state.config.admin_key, key.trim()))
575}
576
577/// The answer to a pod presenting no valid token. Logs the reason, never the token.
578fn invalid_token(reason: &str, path: &str) -> Response {
579    warn!(reason, path = %path, "request with an invalid token rejected");
580    unauthorized()
581}
582
583/// The 401 answer to an unknown, expired or revoked token.
584fn unauthorized() -> Response {
585    error_response(
586        StatusCode::UNAUTHORIZED,
587        "authentication_error",
588        "invalid or expired ironflow auth proxy token",
589    )
590}
591
592/// The answer when the registry backend fails. A 503 is retryable, and does
593/// not make a valid token look revoked as a 401 would.
594fn registry_unavailable(message: &str) -> Response {
595    error_response(StatusCode::SERVICE_UNAVAILABLE, "api_error", message)
596}
597
598fn admin_unauthorized() -> Response {
599    warn!("admin request without a valid admin key");
600    error_response(
601        StatusCode::UNAUTHORIZED,
602        "authentication_error",
603        "invalid or missing admin key",
604    )
605}
606
607async fn healthz() -> &'static str {
608    "ok"
609}
610
611async fn serve_metrics(State(state): State<AuthProxyState>) -> Response {
612    match &state.metrics {
613        Some(handle) => ([(CONTENT_TYPE, METRICS_CONTENT_TYPE)], handle.render()).into_response(),
614        None => error_response(
615            StatusCode::NOT_FOUND,
616            "not_found_error",
617            "metrics are disabled",
618        ),
619    }
620}
621
622async fn issue_token(
623    State(state): State<AuthProxyState>,
624    headers: HeaderMap,
625    body: Bytes,
626) -> Response {
627    if !admin_authorized(&state, &headers) {
628        return admin_unauthorized();
629    }
630    // The body carries a credential: neither the parse error nor the body
631    // is ever echoed or logged.
632    let Ok(request) = from_slice::<TokenRequest>(&body) else {
633        warn!("token request body rejected");
634        return error_response(
635            StatusCode::BAD_REQUEST,
636            "invalid_request_error",
637            "invalid token request body",
638        );
639    };
640    let now = match now_unix() {
641        Ok(now) => now,
642        Err(e) => return clock_error(&e),
643    };
644    let run_id = request.run_id.clone();
645    let step = request.step.clone();
646    // The name only: the value never reaches a log.
647    let secret = match &request.credential {
648        GrantCredential::Secret(secret) => Some(secret.name().to_string()),
649        GrantCredential::Claude(_) => None,
650    };
651    match state.registry.issue(request, now).await {
652        Ok(issued) => {
653            info!(
654                token = %issued.short_id(),
655                run_id = %run_id,
656                step = %step,
657                secret = secret.as_deref(),
658                "token issued"
659            );
660            (StatusCode::CREATED, Json(issued)).into_response()
661        }
662        Err(AuthProxyError::InvalidRequest(message)) => {
663            warn!(run_id = %run_id, step = %step, reason = %message, "token request refused");
664            error_response(StatusCode::BAD_REQUEST, "invalid_request_error", &message)
665        }
666        Err(AuthProxyError::Backend(e)) => {
667            warn!(run_id = %run_id, step = %step, error = %e, "token registry unavailable");
668            registry_unavailable("token registry unavailable")
669        }
670        Err(e) => {
671            warn!(error = %e, "token issuance failed");
672            error_response(
673                StatusCode::INTERNAL_SERVER_ERROR,
674                "api_error",
675                "token issuance failed",
676            )
677        }
678    }
679}
680
681async fn revoke_token(
682    State(state): State<AuthProxyState>,
683    Path(id): Path<String>,
684    headers: HeaderMap,
685) -> Response {
686    if !admin_authorized(&state, &headers) {
687        return admin_unauthorized();
688    }
689    match state.registry.revoke(&id).await {
690        Ok(true) => {
691            info!(token = %short_id(&id), "token revoked");
692            StatusCode::NO_CONTENT.into_response()
693        }
694        Ok(false) => error_response(StatusCode::NOT_FOUND, "not_found_error", "unknown token"),
695        Err(e) => {
696            warn!(token = %short_id(&id), error = %e, "token registry unavailable");
697            registry_unavailable("token registry unavailable")
698        }
699    }
700}
701
702async fn revoke_run(
703    State(state): State<AuthProxyState>,
704    Path(run_id): Path<String>,
705    headers: HeaderMap,
706) -> Response {
707    if !admin_authorized(&state, &headers) {
708        return admin_unauthorized();
709    }
710    match state.registry.revoke_run(&run_id).await {
711        Ok(revoked) => {
712            info!(run_id = %run_id, revoked, "run tokens revoked");
713            (StatusCode::OK, Json(json!({ "revoked": revoked }))).into_response()
714        }
715        Err(e) => {
716            warn!(run_id = %run_id, error = %e, "token registry unavailable");
717            registry_unavailable("token registry unavailable")
718        }
719    }
720}
721
722async fn relay(State(state): State<AuthProxyState>, request: Request) -> Response {
723    let (parts, body) = request.into_parts();
724    let method = parts.method;
725    let path = parts.uri.path().to_string();
726
727    // Absolute-form (`GET http://host/..`) and CONNECT ask the proxy to reach
728    // another host.
729    if parts.uri.authority().is_some() || method == Method::CONNECT {
730        warn!(method = %method, "request for another host refused");
731        return error_response(
732            StatusCode::FORBIDDEN,
733            "permission_error",
734            "only api.anthropic.com is reachable through this proxy",
735        );
736    }
737    if let Some(rest) = path
738        .strip_prefix(RELAY_PREFIX)
739        .and_then(|rest| rest.strip_prefix('/'))
740    {
741        let query = parts.uri.query();
742        return relay_secret(&state, method, &parts.headers, query, body, rest).await;
743    }
744
745    let now = match now_unix() {
746        Ok(now) => now,
747        Err(e) => return clock_error(&e),
748    };
749    let Some(token) = extract_opaque_token(&parts.headers) else {
750        return invalid_token("missing", &path);
751    };
752    let grant = match state.registry.resolve(&token, now).await {
753        Ok(grant) => grant,
754        Err(TokenRejection::Unknown) => return invalid_token("unknown", &path),
755        Err(TokenRejection::Expired) => return invalid_token("expired", &path),
756        Err(TokenRejection::Revoked) => return invalid_token("revoked", &path),
757        Err(TokenRejection::Unavailable(e)) => {
758            warn!(error = %e, path = %path, "token registry unavailable");
759            return registry_unavailable("auth proxy token registry unavailable");
760        }
761    };
762    let token = short_id(&grant.id);
763    let GrantCredential::Claude(credential) = &grant.credential else {
764        warn!(token = %token, path = %path, "secret token on the Anthropic API refused");
765        return error_response(
766            StatusCode::FORBIDDEN,
767            "permission_error",
768            "this token only reaches its allowlisted hosts under /r/",
769        );
770    };
771
772    if !is_allowed_path(&path) {
773        warn!(token = %token, path = %path, "path outside the API refused");
774        return error_response(
775            StatusCode::FORBIDDEN,
776            "permission_error",
777            "only the Anthropic API under /v1/ is reachable through this proxy",
778        );
779    }
780    if !is_allowed_method(&method) {
781        warn!(token = %token, method = %method, path = %path, "method refused");
782        return error_response(
783            StatusCode::METHOD_NOT_ALLOWED,
784            "invalid_request_error",
785            "only GET and POST are relayed",
786        );
787    }
788
789    let bytes = match to_bytes(body, state.config.max_body_bytes).await {
790        Ok(bytes) => bytes,
791        Err(e) => {
792            warn!(token = %token, error = %e, "request body rejected");
793            return error_response(
794                StatusCode::PAYLOAD_TOO_LARGE,
795                "request_too_large",
796                "request body too large or unreadable",
797            );
798        }
799    };
800
801    let mut url = state.config.upstream.clone();
802    url.set_path(&path);
803    url.set_query(parts.uri.query());
804    let upstream = state
805        .http
806        .request(method.clone(), url)
807        .headers(upstream_headers(&parts.headers, credential))
808        .body(bytes)
809        .send()
810        .await;
811    let upstream = match upstream {
812        Ok(upstream) => upstream,
813        Err(e) => {
814            warn!(token = %token, error = %e.without_url(), "upstream request failed");
815            return error_response(
816                StatusCode::BAD_GATEWAY,
817                "api_error",
818                "the Anthropic API is unreachable",
819            );
820        }
821    };
822
823    let status = upstream.status();
824    info!(
825        token = %token,
826        run_id = %grant.run_id,
827        step = %grant.step,
828        method = %method,
829        path = %path,
830        status = status.as_u16(),
831        "relayed"
832    );
833    let headers = downstream_headers(upstream.headers());
834    let mut response = Response::new(Body::from_stream(upstream.bytes_stream()));
835    *response.status_mut() = status;
836    *response.headers_mut() = headers;
837    response
838}
839
840/// Outcome of a `/r/` request: the `result` label of [`REQUESTS_TOTAL`]
841/// and the `result` field of its log event.
842#[derive(Clone, Copy, Debug, PartialEq, Eq)]
843enum RelayResult {
844    Relayed,
845    ForbiddenHost,
846    ForbiddenPath,
847    ForbiddenMethod,
848    UnknownToken,
849    Expired,
850    Revoked,
851    UpstreamError,
852    Unavailable,
853}
854
855impl RelayResult {
856    fn as_str(self) -> &'static str {
857        match self {
858            Self::Relayed => "relayed",
859            Self::ForbiddenHost => "forbidden_host",
860            Self::ForbiddenPath => "forbidden_path",
861            Self::ForbiddenMethod => "forbidden_method",
862            Self::UnknownToken => "unknown_token",
863            Self::Expired => "expired",
864            Self::Revoked => "revoked",
865            Self::UpstreamError => "upstream_error",
866            Self::Unavailable => "unavailable",
867        }
868    }
869}
870
871/// Count a `/r/` request and log it. `secret` is empty when the token is
872/// unknown. Never given the token, the secret value, a header or the query.
873fn record(
874    result: RelayResult,
875    host: &str,
876    secret: &str,
877    grant: Option<&Grant>,
878    upstream_status: Option<u16>,
879) {
880    counter!(REQUESTS_TOTAL, "secret" => secret.to_string(), "result" => result.as_str())
881        .increment(1);
882    let token = grant.map(|grant| short_id(&grant.id));
883    let run_id = grant.map(|grant| grant.run_id.as_str());
884    let step = grant.map(|grant| grant.step.as_str());
885    if result == RelayResult::Relayed {
886        info!(
887            host,
888            secret,
889            run_id,
890            step,
891            token,
892            result = result.as_str(),
893            upstream_status,
894            "secret relay"
895        );
896    } else {
897        warn!(
898            host,
899            secret,
900            run_id,
901            step,
902            token,
903            result = result.as_str(),
904            upstream_status,
905            "secret relay"
906        );
907    }
908}
909
910/// Relay `/r/<host>/<path>` (`rest` is `<host>/<path>`) to `https://<host>`
911/// with the proxied secret of the token's grant.
912async fn relay_secret(
913    state: &AuthProxyState,
914    method: Method,
915    headers: &HeaderMap,
916    query: Option<&str>,
917    body: Body,
918    rest: &str,
919) -> Response {
920    let (host, path) = match rest.split_once('/') {
921        Some((host, path)) => (host.to_ascii_lowercase(), format!("/{path}")),
922        None => (rest.to_ascii_lowercase(), "/".to_string()),
923    };
924    // Never echo a host that failed validation into a log.
925    let logged_host = if is_valid_request_host(&host) {
926        host.as_str()
927    } else {
928        "<invalid>"
929    };
930
931    let now = match now_unix() {
932        Ok(now) => now,
933        Err(e) => return clock_error(&e),
934    };
935    let Some(token) = extract_opaque_token(headers) else {
936        record(RelayResult::UnknownToken, logged_host, "", None, None);
937        return unauthorized();
938    };
939    let grant = match state.registry.resolve(&token, now).await {
940        Ok(grant) => grant,
941        Err(rejection) => {
942            let result = match rejection {
943                TokenRejection::Unknown => RelayResult::UnknownToken,
944                TokenRejection::Expired => RelayResult::Expired,
945                TokenRejection::Revoked => RelayResult::Revoked,
946                TokenRejection::Unavailable(e) => {
947                    warn!(error = %e, "token registry unavailable");
948                    record(RelayResult::Unavailable, logged_host, "", None, None);
949                    return registry_unavailable("auth proxy token registry unavailable");
950                }
951            };
952            record(result, logged_host, "", None, None);
953            return unauthorized();
954        }
955    };
956
957    let GrantCredential::Secret(secret) = &grant.credential else {
958        record(
959            RelayResult::ForbiddenHost,
960            logged_host,
961            "",
962            Some(&grant),
963            None,
964        );
965        return error_response(
966            StatusCode::FORBIDDEN,
967            "permission_error",
968            "a Claude token only reaches the Anthropic API",
969        );
970    };
971    let name = secret.name();
972    let refuse = |result| record(result, logged_host, name, Some(&grant), None);
973    let forbidden_host = || {
974        refuse(RelayResult::ForbiddenHost);
975        error_response(
976            StatusCode::FORBIDDEN,
977            "permission_error",
978            "host not allowed for this secret",
979        )
980    };
981    if !secret.allows_host(&host) {
982        return forbidden_host();
983    }
984    if !is_relay_path(&path) {
985        refuse(RelayResult::ForbiddenPath);
986        return error_response(
987            StatusCode::FORBIDDEN,
988            "permission_error",
989            "path not allowed: no '..', '//' or percent-encoding",
990        );
991    }
992    if !is_relay_method(&method) {
993        refuse(RelayResult::ForbiddenMethod);
994        return error_response(
995            StatusCode::METHOD_NOT_ALLOWED,
996            "invalid_request_error",
997            "only GET, HEAD, POST, PUT, PATCH and DELETE are relayed",
998        );
999    }
1000
1001    let bytes = match to_bytes(body, state.config.max_body_bytes).await {
1002        Ok(bytes) => bytes,
1003        Err(e) => {
1004            warn!(
1005                token = %short_id(&grant.id),
1006                secret = name,
1007                error = %e,
1008                "request body rejected"
1009            );
1010            return error_response(
1011                StatusCode::PAYLOAD_TOO_LARGE,
1012                "request_too_large",
1013                "request body too large or unreadable",
1014            );
1015        }
1016    };
1017
1018    let mut url = match &state.config.secret_upstream {
1019        Some(upstream) => upstream.clone(),
1020        None => {
1021            // `host` passed the allowlist, so it is a DNS name: this only
1022            // fails if the URL parser disagrees, and then refuses it.
1023            let url = Url::parse(SECRET_UPSTREAM_BASE).ok().and_then(|mut url| {
1024                url.set_host(Some(&host)).ok()?;
1025                Some(url)
1026            });
1027            let Some(url) = url else {
1028                return forbidden_host();
1029            };
1030            url
1031        }
1032    };
1033    url.set_path(&path);
1034    url.set_query(query);
1035    let upstream = state
1036        .http
1037        .request(method, url)
1038        .headers(secret_upstream_headers(headers, secret))
1039        .body(bytes)
1040        .send()
1041        .await;
1042    let upstream = match upstream {
1043        Ok(upstream) => upstream,
1044        Err(e) => {
1045            warn!(
1046                token = %short_id(&grant.id),
1047                secret = name,
1048                error = %e.without_url(),
1049                "upstream request failed"
1050            );
1051            refuse(RelayResult::UpstreamError);
1052            return error_response(
1053                StatusCode::BAD_GATEWAY,
1054                "upstream_error",
1055                "upstream unreachable",
1056            );
1057        }
1058    };
1059
1060    let status = upstream.status();
1061    record(
1062        RelayResult::Relayed,
1063        logged_host,
1064        name,
1065        Some(&grant),
1066        Some(status.as_u16()),
1067    );
1068    let headers = downstream_headers(upstream.headers());
1069    let mut response = Response::new(Body::from_stream(upstream.bytes_stream()));
1070    *response.status_mut() = status;
1071    *response.headers_mut() = headers;
1072    response
1073}
1074
1075#[cfg(test)]
1076mod tests {
1077    use super::*;
1078
1079    fn key_ring() -> KeyRing {
1080        KeyRing::from_spec(&format!("1:{}", "aa".repeat(32)), None).unwrap()
1081    }
1082
1083    #[tokio::test]
1084    async fn registry_from_config_defaults_to_memory() {
1085        let registry = registry_from_config(None, None).await.unwrap();
1086        assert!(registry.is_empty().await.unwrap());
1087
1088        let registry = registry_from_config(Some(""), None).await.unwrap();
1089        assert!(registry.is_empty().await.unwrap());
1090
1091        let registry = registry_from_config(Some("  \n"), Some(key_ring()))
1092            .await
1093            .unwrap();
1094        assert!(registry.is_empty().await.unwrap());
1095    }
1096
1097    #[tokio::test]
1098    async fn registry_from_config_requires_key_ring_with_database() {
1099        let result = registry_from_config(Some("postgres://localhost/x"), None).await;
1100        assert!(
1101            matches!(result, Err(RegistryConfigError::MissingKeyRing)),
1102            "{result:?}"
1103        );
1104    }
1105
1106    #[tokio::test]
1107    async fn registry_from_config_rejects_invalid_database_url() {
1108        let result = registry_from_config(Some("not a url"), Some(key_ring())).await;
1109        assert!(
1110            matches!(result, Err(RegistryConfigError::Store(_))),
1111            "{result:?}"
1112        );
1113    }
1114    #[test]
1115    fn relay_results_have_stable_labels() {
1116        let labels: Vec<&str> = [
1117            RelayResult::Relayed,
1118            RelayResult::ForbiddenHost,
1119            RelayResult::ForbiddenPath,
1120            RelayResult::ForbiddenMethod,
1121            RelayResult::UnknownToken,
1122            RelayResult::Expired,
1123            RelayResult::Revoked,
1124            RelayResult::UpstreamError,
1125            RelayResult::Unavailable,
1126        ]
1127        .into_iter()
1128        .map(RelayResult::as_str)
1129        .collect();
1130        assert_eq!(
1131            labels,
1132            vec![
1133                "relayed",
1134                "forbidden_host",
1135                "forbidden_path",
1136                "forbidden_method",
1137                "unknown_token",
1138                "expired",
1139                "revoked",
1140                "upstream_error",
1141                "unavailable",
1142            ]
1143        );
1144    }
1145
1146    #[test]
1147    fn config_debug_shows_secret_upstream_not_admin_key() {
1148        let config = AuthProxyConfig::new("0123456789abcdef0123456789abcdef")
1149            .with_secret_upstream(Url::parse("http://127.0.0.1:9001").unwrap());
1150        let debug = format!("{config:?}");
1151        assert!(debug.contains("http://127.0.0.1:9001/"), "{debug}");
1152        assert!(!debug.contains("0123456789abcdef"), "{debug}");
1153    }
1154}