1use 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
91pub const DATABASE_URL_ENV: &str = "IRONFLOW_AUTH_PROXY_DATABASE_URL";
94
95pub const MIN_ADMIN_KEY_LEN: usize = 32;
97
98pub const DEFAULT_MAX_BODY_BYTES: usize = 32 * 1024 * 1024;
100
101const SHORT_ID_LEN: usize = 12;
103
104const UPSTREAM_CONNECT_TIMEOUT: Duration = Duration::from_secs(10);
107
108#[derive(Clone)]
120pub struct AuthProxyConfig {
121 pub upstream: Url,
123 pub admin_key: String,
125 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 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 pub fn with_upstream(mut self, upstream: Url) -> Self {
178 self.upstream = upstream;
179 self
180 }
181}
182
183#[derive(Clone)]
198pub struct AuthProxyState {
199 registry: AuthProxyRegistry,
200 config: Arc<AuthProxyConfig>,
201 http: Client,
202}
203
204impl AuthProxyState {
205 pub fn new(config: AuthProxyConfig) -> Result<Self, ReqwestError> {
216 Self::with_registry(config, AuthProxyRegistry::default())
217 }
218
219 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 pub fn registry(&self) -> &AuthProxyRegistry {
267 &self.registry
268 }
269}
270
271#[derive(Debug, Error)]
282pub enum RegistryConfigError {
283 #[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 #[error("invalid encryption key: {0}")]
290 Crypto(#[from] CryptoError),
291 #[error("cannot open the token registry database: {0}")]
293 Store(#[from] StoreError),
294}
295
296pub 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
336pub 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
364pub 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
380pub 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 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
474fn 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
483fn 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
493fn 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 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 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}