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//! One listener serves:
12//!
13//! * `GET /healthz` - liveness;
14//! * `/admin/v1/...` - token issuance and revocation, behind
15//!   `Authorization: Bearer <IRONFLOW_AUTH_PROXY_ADMIN_KEY>`;
16//! * anything else - the relay: an unknown, expired or revoked token gets a
17//!   401, a path outside `/v1/` or a request for another host a 403, a method
18//!   other than GET/POST a 405.
19//!
20//! Grants live in a registry with two backends:
21//!
22//! * in memory (default): a single replica, tokens lost on restart;
23//! * PostgreSQL, when [`DATABASE_URL_ENV`] is set ([`registry_from_config`]):
24//!   shared by several replicas and surviving restarts. Only the token
25//!   SHA-256 is stored, never the token, and the credential is AES-256-GCM
26//!   encrypted at rest with the `IRONFLOW_SECRET_KEYS` key ring.
27//!
28//! Logs never contain a token or a credential, only the short token id.
29//!
30//! # Examples
31//!
32//! ```no_run
33//! use std::env::var;
34//! use std::time::Duration;
35//!
36//! use ironflow_auth_proxy::{
37//!     AuthProxyConfig, AuthProxyState, DATABASE_URL_ENV, registry_from_config, serve, spawn_purge,
38//! };
39//! use ironflow_store::crypto::KeyRing;
40//! use tokio::net::TcpListener;
41//!
42//! # async fn example() -> Result<(), Box<dyn std::error::Error>> {
43//! let database_url = var(DATABASE_URL_ENV).ok();
44//! let registry = registry_from_config(database_url.as_deref(), KeyRing::from_env()?).await?;
45//! let config = AuthProxyConfig::new("0123456789abcdef0123456789abcdef");
46//! let state = AuthProxyState::with_registry(config, registry)?;
47//! let purge = spawn_purge(state.registry().clone(), Duration::from_secs(60));
48//! let listener = TcpListener::bind("0.0.0.0:8080").await?;
49//! serve(listener, state).await?;
50//! purge.abort();
51//! # Ok(())
52//! # }
53//! ```
54
55use std::fmt;
56use std::future::pending;
57use std::io;
58use std::sync::Arc;
59use std::time::{Duration, SystemTime, SystemTimeError, UNIX_EPOCH};
60
61use axum::body::{Body, Bytes, to_bytes};
62use axum::extract::{DefaultBodyLimit, Path, Request, State};
63use axum::http::header::AUTHORIZATION;
64use axum::http::{HeaderMap, Method, StatusCode};
65use axum::response::{IntoResponse, Response};
66use axum::routing::{delete, get, post};
67use axum::serve as serve_router;
68use axum::{Json, Router};
69use reqwest::redirect::Policy;
70use reqwest::{Client, Error as ReqwestError};
71use serde_json::{from_slice, json};
72use thiserror::Error;
73use tokio::net::TcpListener;
74use tokio::signal::ctrl_c;
75#[cfg(unix)]
76use tokio::signal::unix::{SignalKind, signal};
77use tokio::task::JoinHandle;
78use tokio::{select, spawn, time};
79use tracing::{info, warn};
80use url::Url;
81
82use ironflow_core::auth_proxy::{
83    AuthProxyError, AuthProxyRegistry, DEFAULT_UPSTREAM, TokenRejection, TokenRequest,
84    admin_key_matches, downstream_headers, error_body, extract_opaque_token, is_allowed_method,
85    is_allowed_path, upstream_headers,
86};
87use ironflow_store::crypto::{CryptoError, KeyRing};
88use ironflow_store::error::StoreError;
89use ironflow_store::postgres::PostgresStore;
90
91/// Environment variable holding the PostgreSQL URL of the shared token
92/// registry. Unset (or empty), tokens live in memory.
93pub const DATABASE_URL_ENV: &str = "IRONFLOW_AUTH_PROXY_DATABASE_URL";
94
95/// Shortest admin key accepted.
96pub const MIN_ADMIN_KEY_LEN: usize = 32;
97
98/// Default largest request body relayed: 32 MiB.
99pub const DEFAULT_MAX_BODY_BYTES: usize = 32 * 1024 * 1024;
100
101/// Length of the token id prefix written to logs.
102const SHORT_ID_LEN: usize = 12;
103
104/// Connect timeout towards the upstream. There is no total timeout: an
105/// agent turn streams for as long as the model answers.
106const UPSTREAM_CONNECT_TIMEOUT: Duration = Duration::from_secs(10);
107
108/// Configuration of the proxy. Its [`Debug`] output never shows the admin key.
109///
110/// # Examples
111///
112/// ```
113/// use ironflow_auth_proxy::AuthProxyConfig;
114///
115/// let config = AuthProxyConfig::new("0123456789abcdef0123456789abcdef");
116/// assert_eq!(config.upstream.as_str(), "https://api.anthropic.com/");
117/// assert!(!format!("{config:?}").contains("0123456789abcdef"));
118/// ```
119#[derive(Clone)]
120pub struct AuthProxyConfig {
121    /// Where requests are relayed. Always `https://api.anthropic.com` in the binary.
122    pub upstream: Url,
123    /// Key protecting the admin API.
124    pub admin_key: String,
125    /// Largest request body relayed, in bytes.
126    pub max_body_bytes: usize,
127}
128
129impl fmt::Debug for AuthProxyConfig {
130    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
131        f.debug_struct("AuthProxyConfig")
132            .field("upstream", &self.upstream.as_str())
133            .field("admin_key", &"<redacted>")
134            .field("max_body_bytes", &self.max_body_bytes)
135            .finish()
136    }
137}
138
139impl AuthProxyConfig {
140    /// Configuration relaying to [`DEFAULT_UPSTREAM`] with the given admin key.
141    ///
142    /// # Panics
143    ///
144    /// Panics if `admin_key` is shorter than [`MIN_ADMIN_KEY_LEN`].
145    ///
146    /// # Examples
147    ///
148    /// See [`AuthProxyConfig`].
149    pub fn new(admin_key: &str) -> Self {
150        assert!(
151            admin_key.len() >= MIN_ADMIN_KEY_LEN,
152            "the auth proxy admin key must be at least {MIN_ADMIN_KEY_LEN} characters"
153        );
154        Self {
155            upstream: Url::parse(DEFAULT_UPSTREAM).expect("DEFAULT_UPSTREAM is a valid URL"),
156            admin_key: admin_key.to_string(),
157            max_body_bytes: DEFAULT_MAX_BODY_BYTES,
158        }
159    }
160
161    /// Relay to another upstream. **For tests only**: the binary never calls
162    /// it, so a deployed proxy only reaches `api.anthropic.com`.
163    ///
164    /// # Examples
165    ///
166    /// ```
167    /// use ironflow_auth_proxy::AuthProxyConfig;
168    /// use url::Url;
169    ///
170    /// # fn example() -> Result<(), url::ParseError> {
171    /// let config = AuthProxyConfig::new("0123456789abcdef0123456789abcdef")
172    ///     .with_upstream(Url::parse("http://127.0.0.1:9000")?);
173    /// assert_eq!(config.upstream.port(), Some(9000));
174    /// # Ok(())
175    /// # }
176    /// ```
177    pub fn with_upstream(mut self, upstream: Url) -> Self {
178        self.upstream = upstream;
179        self
180    }
181}
182
183/// Shared state of the proxy: the token registry, the configuration and the
184/// upstream HTTP client. Cheap to clone.
185///
186/// # Examples
187///
188/// ```
189/// use ironflow_auth_proxy::{AuthProxyConfig, AuthProxyState};
190///
191/// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
192/// let state = AuthProxyState::new(AuthProxyConfig::new("0123456789abcdef0123456789abcdef"))?;
193/// assert!(state.registry().is_empty().await?);
194/// # Ok(())
195/// # }
196/// ```
197#[derive(Clone)]
198pub struct AuthProxyState {
199    registry: AuthProxyRegistry,
200    config: Arc<AuthProxyConfig>,
201    http: Client,
202}
203
204impl AuthProxyState {
205    /// Build the state with an empty in-memory registry.
206    ///
207    /// # Errors
208    ///
209    /// Returns the [`reqwest::Error`] raised when the HTTP client cannot be
210    /// built (TLS backend initialisation).
211    ///
212    /// # Examples
213    ///
214    /// See [`AuthProxyState`].
215    pub fn new(config: AuthProxyConfig) -> Result<Self, ReqwestError> {
216        Self::with_registry(config, AuthProxyRegistry::default())
217    }
218
219    /// Build the state over `registry`, for instance the shared PostgreSQL
220    /// registry returned by [`registry_from_config`]. Proxies built over
221    /// registries sharing a backend serve the same tokens.
222    ///
223    /// The upstream client never follows redirects and only bounds the
224    /// connection time.
225    ///
226    /// # Errors
227    ///
228    /// Returns the [`reqwest::Error`] raised when the HTTP client cannot be
229    /// built (TLS backend initialisation).
230    ///
231    /// # Examples
232    ///
233    /// ```
234    /// use ironflow_auth_proxy::{AuthProxyConfig, AuthProxyState};
235    /// use ironflow_core::auth_proxy::AuthProxyRegistry;
236    ///
237    /// # fn example() -> Result<(), reqwest::Error> {
238    /// let registry = AuthProxyRegistry::default();
239    /// let config = AuthProxyConfig::new("0123456789abcdef0123456789abcdef");
240    /// let a = AuthProxyState::with_registry(config.clone(), registry.clone())?;
241    /// let b = AuthProxyState::with_registry(config, registry)?;
242    /// # let _ = (a, b);
243    /// # Ok(())
244    /// # }
245    /// ```
246    pub fn with_registry(
247        config: AuthProxyConfig,
248        registry: AuthProxyRegistry,
249    ) -> Result<Self, ReqwestError> {
250        let http = Client::builder()
251            .connect_timeout(UPSTREAM_CONNECT_TIMEOUT)
252            .redirect(Policy::none())
253            .build()?;
254        Ok(Self {
255            registry,
256            config: Arc::new(config),
257            http,
258        })
259    }
260
261    /// The token registry.
262    ///
263    /// # Examples
264    ///
265    /// See [`AuthProxyState`].
266    pub fn registry(&self) -> &AuthProxyRegistry {
267        &self.registry
268    }
269}
270
271/// Why the token registry could not be built. No variant carries the
272/// database URL, a key or a credential.
273///
274/// # Examples
275///
276/// ```
277/// use ironflow_auth_proxy::RegistryConfigError;
278///
279/// assert!(RegistryConfigError::MissingKeyRing.to_string().contains("IRONFLOW_SECRET_KEYS"));
280/// ```
281#[derive(Debug, Error)]
282pub enum RegistryConfigError {
283    /// A database is configured but no key to encrypt the credentials with.
284    #[error(
285        "{DATABASE_URL_ENV} is set but no encryption key is configured: set IRONFLOW_SECRET_KEYS (or IRONFLOW_SECRET_KEY)"
286    )]
287    MissingKeyRing,
288    /// The encryption key configuration is invalid.
289    #[error("invalid encryption key: {0}")]
290    Crypto(#[from] CryptoError),
291    /// The database cannot be opened or migrated.
292    #[error("cannot open the token registry database: {0}")]
293    Store(#[from] StoreError),
294}
295
296/// Build the token registry: in memory when `database_url` is `None` or
297/// blank, otherwise shared in PostgreSQL, the credentials encrypted with
298/// `key_ring`. Opening the database runs the store migrations.
299///
300/// # Errors
301///
302/// Returns [`RegistryConfigError::MissingKeyRing`] when a database is given
303/// without a key ring (checked before any connection), and
304/// [`RegistryConfigError::Store`] when the database cannot be reached or
305/// migrated.
306///
307/// # Examples
308///
309/// ```no_run
310/// use ironflow_auth_proxy::registry_from_config;
311/// use ironflow_store::crypto::KeyRing;
312///
313/// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
314/// let memory = registry_from_config(None, None).await?;
315/// assert!(memory.is_empty().await?);
316///
317/// let ring = KeyRing::from_spec(&format!("1:{}", "aa".repeat(32)), None)?;
318/// let shared = registry_from_config(Some("postgres://localhost/ironflow"), Some(ring)).await?;
319/// # let _ = shared;
320/// # Ok(())
321/// # }
322/// ```
323pub async fn registry_from_config(
324    database_url: Option<&str>,
325    key_ring: Option<KeyRing>,
326) -> Result<AuthProxyRegistry, RegistryConfigError> {
327    let Some(url) = database_url.map(str::trim).filter(|url| !url.is_empty()) else {
328        return Ok(AuthProxyRegistry::default());
329    };
330    let ring = key_ring.ok_or(RegistryConfigError::MissingKeyRing)?;
331    let mut store = PostgresStore::new(url).await?;
332    store.set_key_ring(ring);
333    Ok(AuthProxyRegistry::with_backend(Arc::new(store)))
334}
335
336/// The proxy router: health, admin API and relay.
337///
338/// # Examples
339///
340/// ```no_run
341/// use axum::serve;
342/// use ironflow_auth_proxy::{AuthProxyConfig, AuthProxyState, router};
343/// use tokio::net::TcpListener;
344///
345/// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
346/// let state = AuthProxyState::new(AuthProxyConfig::new("0123456789abcdef0123456789abcdef"))?;
347/// let listener = TcpListener::bind("127.0.0.1:0").await?;
348/// serve(listener, router(state)).await?;
349/// # Ok(())
350/// # }
351/// ```
352pub fn router(state: AuthProxyState) -> Router {
353    let limit = state.config.max_body_bytes;
354    Router::new()
355        .route("/healthz", get(healthz))
356        .route("/admin/v1/tokens", post(issue_token))
357        .route("/admin/v1/tokens/{id}", delete(revoke_token))
358        .route("/admin/v1/runs/{run_id}/tokens", delete(revoke_run))
359        .fallback(relay)
360        .layer(DefaultBodyLimit::max(limit))
361        .with_state(state)
362}
363
364/// Serve the proxy on `listener` until ctrl-c or SIGTERM, letting in-flight
365/// requests finish.
366///
367/// # Errors
368///
369/// Returns the I/O error that stopped the server.
370///
371/// # Examples
372///
373/// See the [crate documentation](crate).
374pub async fn serve(listener: TcpListener, state: AuthProxyState) -> io::Result<()> {
375    serve_router(listener, router(state))
376        .with_graceful_shutdown(shutdown_signal())
377        .await
378}
379
380/// Spawn a task dropping the expired grants of `registry` every `interval`.
381///
382/// Abort the returned handle to stop it.
383///
384/// # Panics
385///
386/// Panics if `interval` is zero, or when called outside a Tokio runtime.
387///
388/// # Examples
389///
390/// ```no_run
391/// use std::time::Duration;
392///
393/// use ironflow_auth_proxy::spawn_purge;
394/// use ironflow_core::auth_proxy::AuthProxyRegistry;
395///
396/// # async fn example() {
397/// let purge = spawn_purge(AuthProxyRegistry::default(), Duration::from_secs(60));
398/// purge.abort();
399/// # }
400/// ```
401pub fn spawn_purge(registry: AuthProxyRegistry, interval: Duration) -> JoinHandle<()> {
402    assert!(
403        !interval.is_zero(),
404        "purge interval must be greater than zero"
405    );
406    spawn(async move {
407        let mut ticker = time::interval(interval);
408        loop {
409            ticker.tick().await;
410            match now_unix() {
411                // Every replica purges: the deletes are idempotent.
412                Ok(now) => match registry.purge_expired(now).await {
413                    Ok(purged) if purged > 0 => info!(purged, "expired auth proxy tokens purged"),
414                    Ok(_) => {}
415                    Err(e) => warn!(error = %e, "expired auth proxy tokens purge failed"),
416                },
417                Err(e) => warn!(error = %e, "clock before the unix epoch; purge skipped"),
418            }
419        }
420    })
421}
422
423async fn shutdown_signal() {
424    let interrupt = async {
425        if let Err(e) = ctrl_c().await {
426            warn!(error = %e, "cannot listen for ctrl-c");
427            pending::<()>().await;
428        }
429    };
430    #[cfg(unix)]
431    let terminate = async {
432        match signal(SignalKind::terminate()) {
433            Ok(mut sigterm) => {
434                sigterm.recv().await;
435            }
436            Err(e) => {
437                warn!(error = %e, "cannot listen for SIGTERM");
438                pending::<()>().await;
439            }
440        }
441    };
442    #[cfg(not(unix))]
443    let terminate = pending::<()>();
444    select! {
445        () = interrupt => {},
446        () = terminate => {},
447    }
448    info!("shutting down");
449}
450
451fn now_unix() -> Result<u64, SystemTimeError> {
452    SystemTime::now()
453        .duration_since(UNIX_EPOCH)
454        .map(|d| d.as_secs())
455}
456
457fn short_id(id: &str) -> &str {
458    id.get(..SHORT_ID_LEN).unwrap_or(id)
459}
460
461fn error_response(status: StatusCode, kind: &str, message: &str) -> Response {
462    (status, Json(error_body(kind, message))).into_response()
463}
464
465fn clock_error(e: &SystemTimeError) -> Response {
466    warn!(error = %e, "system clock is before the unix epoch");
467    error_response(
468        StatusCode::INTERNAL_SERVER_ERROR,
469        "api_error",
470        "auth proxy clock error",
471    )
472}
473
474/// Whether the request carries the admin key as a bearer token.
475fn admin_authorized(state: &AuthProxyState, headers: &HeaderMap) -> bool {
476    headers
477        .get(AUTHORIZATION)
478        .and_then(|value| value.to_str().ok())
479        .and_then(|value| value.strip_prefix("Bearer "))
480        .is_some_and(|key| admin_key_matches(&state.config.admin_key, key.trim()))
481}
482
483/// The answer to a pod presenting no valid token. Logs the reason, never the token.
484fn invalid_token(reason: &str, path: &str) -> Response {
485    warn!(reason, path = %path, "request with an invalid token rejected");
486    error_response(
487        StatusCode::UNAUTHORIZED,
488        "authentication_error",
489        "invalid or expired ironflow auth proxy token",
490    )
491}
492
493/// The answer when the registry backend fails. A 503 is retryable, and does
494/// not make a valid token look revoked as a 401 would.
495fn registry_unavailable(message: &str) -> Response {
496    error_response(StatusCode::SERVICE_UNAVAILABLE, "api_error", message)
497}
498
499fn admin_unauthorized() -> Response {
500    warn!("admin request without a valid admin key");
501    error_response(
502        StatusCode::UNAUTHORIZED,
503        "authentication_error",
504        "invalid or missing admin key",
505    )
506}
507
508async fn healthz() -> &'static str {
509    "ok"
510}
511
512async fn issue_token(
513    State(state): State<AuthProxyState>,
514    headers: HeaderMap,
515    body: Bytes,
516) -> Response {
517    if !admin_authorized(&state, &headers) {
518        return admin_unauthorized();
519    }
520    // The body carries a credential: neither the parse error nor the body
521    // is ever echoed or logged.
522    let Ok(request) = from_slice::<TokenRequest>(&body) else {
523        warn!("token request body rejected");
524        return error_response(
525            StatusCode::BAD_REQUEST,
526            "invalid_request_error",
527            "invalid token request body",
528        );
529    };
530    let now = match now_unix() {
531        Ok(now) => now,
532        Err(e) => return clock_error(&e),
533    };
534    let run_id = request.run_id.clone();
535    let step = request.step.clone();
536    match state.registry.issue(request, now).await {
537        Ok(issued) => {
538            info!(
539                token = %issued.short_id(),
540                run_id = %run_id,
541                step = %step,
542                "token issued"
543            );
544            (StatusCode::CREATED, Json(issued)).into_response()
545        }
546        Err(AuthProxyError::InvalidRequest(message)) => {
547            warn!(run_id = %run_id, step = %step, reason = %message, "token request refused");
548            error_response(StatusCode::BAD_REQUEST, "invalid_request_error", &message)
549        }
550        Err(AuthProxyError::Backend(e)) => {
551            warn!(run_id = %run_id, step = %step, error = %e, "token registry unavailable");
552            registry_unavailable("token registry unavailable")
553        }
554        Err(e) => {
555            warn!(error = %e, "token issuance failed");
556            error_response(
557                StatusCode::INTERNAL_SERVER_ERROR,
558                "api_error",
559                "token issuance failed",
560            )
561        }
562    }
563}
564
565async fn revoke_token(
566    State(state): State<AuthProxyState>,
567    Path(id): Path<String>,
568    headers: HeaderMap,
569) -> Response {
570    if !admin_authorized(&state, &headers) {
571        return admin_unauthorized();
572    }
573    match state.registry.revoke(&id).await {
574        Ok(true) => {
575            info!(token = %short_id(&id), "token revoked");
576            StatusCode::NO_CONTENT.into_response()
577        }
578        Ok(false) => error_response(StatusCode::NOT_FOUND, "not_found_error", "unknown token"),
579        Err(e) => {
580            warn!(token = %short_id(&id), error = %e, "token registry unavailable");
581            registry_unavailable("token registry unavailable")
582        }
583    }
584}
585
586async fn revoke_run(
587    State(state): State<AuthProxyState>,
588    Path(run_id): Path<String>,
589    headers: HeaderMap,
590) -> Response {
591    if !admin_authorized(&state, &headers) {
592        return admin_unauthorized();
593    }
594    match state.registry.revoke_run(&run_id).await {
595        Ok(revoked) => {
596            info!(run_id = %run_id, revoked, "run tokens revoked");
597            (StatusCode::OK, Json(json!({ "revoked": revoked }))).into_response()
598        }
599        Err(e) => {
600            warn!(run_id = %run_id, error = %e, "token registry unavailable");
601            registry_unavailable("token registry unavailable")
602        }
603    }
604}
605
606async fn relay(State(state): State<AuthProxyState>, request: Request) -> Response {
607    let (parts, body) = request.into_parts();
608    let method = parts.method;
609    let path = parts.uri.path().to_string();
610
611    // Absolute-form (`GET http://host/..`) and CONNECT ask the proxy to reach
612    // another host.
613    if parts.uri.authority().is_some() || method == Method::CONNECT {
614        warn!(method = %method, "request for another host refused");
615        return error_response(
616            StatusCode::FORBIDDEN,
617            "permission_error",
618            "only api.anthropic.com is reachable through this proxy",
619        );
620    }
621
622    let now = match now_unix() {
623        Ok(now) => now,
624        Err(e) => return clock_error(&e),
625    };
626    let Some(token) = extract_opaque_token(&parts.headers) else {
627        return invalid_token("missing", &path);
628    };
629    let grant = match state.registry.resolve(&token, now).await {
630        Ok(grant) => grant,
631        Err(TokenRejection::Unknown) => return invalid_token("unknown", &path),
632        Err(TokenRejection::Expired) => return invalid_token("expired", &path),
633        Err(TokenRejection::Unavailable(e)) => {
634            warn!(error = %e, path = %path, "token registry unavailable");
635            return registry_unavailable("auth proxy token registry unavailable");
636        }
637    };
638    let token = short_id(&grant.id);
639
640    if !is_allowed_path(&path) {
641        warn!(token = %token, path = %path, "path outside the API refused");
642        return error_response(
643            StatusCode::FORBIDDEN,
644            "permission_error",
645            "only the Anthropic API under /v1/ is reachable through this proxy",
646        );
647    }
648    if !is_allowed_method(&method) {
649        warn!(token = %token, method = %method, path = %path, "method refused");
650        return error_response(
651            StatusCode::METHOD_NOT_ALLOWED,
652            "invalid_request_error",
653            "only GET and POST are relayed",
654        );
655    }
656
657    let bytes = match to_bytes(body, state.config.max_body_bytes).await {
658        Ok(bytes) => bytes,
659        Err(e) => {
660            warn!(token = %token, error = %e, "request body rejected");
661            return error_response(
662                StatusCode::PAYLOAD_TOO_LARGE,
663                "request_too_large",
664                "request body too large or unreadable",
665            );
666        }
667    };
668
669    let mut url = state.config.upstream.clone();
670    url.set_path(&path);
671    url.set_query(parts.uri.query());
672    let upstream = state
673        .http
674        .request(method.clone(), url)
675        .headers(upstream_headers(&parts.headers, &grant.credential))
676        .body(bytes)
677        .send()
678        .await;
679    let upstream = match upstream {
680        Ok(upstream) => upstream,
681        Err(e) => {
682            warn!(token = %token, error = %e.without_url(), "upstream request failed");
683            return error_response(
684                StatusCode::BAD_GATEWAY,
685                "api_error",
686                "the Anthropic API is unreachable",
687            );
688        }
689    };
690
691    let status = upstream.status();
692    info!(
693        token = %token,
694        run_id = %grant.run_id,
695        step = %grant.step,
696        method = %method,
697        path = %path,
698        status = status.as_u16(),
699        "relayed"
700    );
701    let headers = downstream_headers(upstream.headers());
702    let mut response = Response::new(Body::from_stream(upstream.bytes_stream()));
703    *response.status_mut() = status;
704    *response.headers_mut() = headers;
705    response
706}
707
708#[cfg(test)]
709mod tests {
710    use super::*;
711
712    fn key_ring() -> KeyRing {
713        KeyRing::from_spec(&format!("1:{}", "aa".repeat(32)), None).unwrap()
714    }
715
716    #[tokio::test]
717    async fn registry_from_config_defaults_to_memory() {
718        let registry = registry_from_config(None, None).await.unwrap();
719        assert!(registry.is_empty().await.unwrap());
720
721        let registry = registry_from_config(Some(""), None).await.unwrap();
722        assert!(registry.is_empty().await.unwrap());
723
724        let registry = registry_from_config(Some("  \n"), Some(key_ring()))
725            .await
726            .unwrap();
727        assert!(registry.is_empty().await.unwrap());
728    }
729
730    #[tokio::test]
731    async fn registry_from_config_requires_key_ring_with_database() {
732        let result = registry_from_config(Some("postgres://localhost/x"), None).await;
733        assert!(
734            matches!(result, Err(RegistryConfigError::MissingKeyRing)),
735            "{result:?}"
736        );
737    }
738
739    #[tokio::test]
740    async fn registry_from_config_rejects_invalid_database_url() {
741        let result = registry_from_config(Some("not a url"), Some(key_ring())).await;
742        assert!(
743            matches!(result, Err(RegistryConfigError::Store(_))),
744            "{result:?}"
745        );
746    }
747}