orca-proxy 0.3.0-rc.6

Reverse proxy with HTTP routing and Wasm trigger dispatch
Documentation
//! Reverse proxy with HTTP routing for containers and Wasm trigger dispatch.
//!
//! Routes HTTP traffic by `Host` header to container backends (round-robin),
//! and by path pattern to Wasm component invocations via a callback.
//! Supports automatic TLS via ACME/Let's Encrypt (Caddy-style zero-config).

pub mod acme;
mod acme_proxy;
mod backend_timeouts;
mod body;
mod error_page;
mod forward;
mod handler;
pub mod rate_limit;
mod routing;
mod security_headers;
pub mod sni;
pub mod tls;
mod websocket;

pub use orca_core::config::{FallbackConfig, SecurityHeadersConfig};
/// Install the proxy's security-header policy (call once at startup). See
/// [`security_headers`].
pub use security_headers::init as init_security_headers;

use std::collections::HashMap;
use std::sync::Arc;
use std::sync::atomic::AtomicUsize;

use hyper::Request;
use hyper::body::Incoming;
use hyper::server::conn::http1;
use hyper::service::service_fn;
use hyper_util::rt::{TokioExecutor, TokioIo};
use hyper_util::server::conn::auto;
use tokio::net::TcpListener;
use tokio::sync::RwLock;
use tracing::{debug, info, warn};

use acme::AcmeManager;
pub use acme_proxy::{run_proxy_with_acme, run_proxy_with_acme_and_fallback};
use handler::{handle_acme_challenge, handle_request};
use rate_limit::RateLimiter;

/// Largest HTTP/2 header list the TLS listener accepts. Matches what the
/// HTTP/1.1 path tolerates in practice, so large SSO cookie jars work on both.
const H2_MAX_HEADER_LIST: u32 = 64 * 1024;

/// A backend target for container routing.
#[derive(Debug, Clone)]
pub struct RouteTarget {
    /// Address in the form `ip:port`.
    pub address: String,
    /// Owning service name.
    pub service_name: String,
    /// Optional path pattern (e.g., `"/api/*"`). When `None`, this target is a
    /// catch-all for the domain. When `Some`, only requests whose path matches
    /// the pattern are routed here. Longest-prefix match wins.
    pub path_pattern: Option<String>,
    /// Traffic weight (1-100, default 100). Used for weighted routing
    /// during canary deployments. Higher weight = more traffic.
    pub weight: u32,
    /// Prefix to strip from the request path before forwarding upstream,
    /// e.g. `"/admin"`. With `path_pattern = "/admin/*"` and
    /// `strip_prefix = Some("/admin")`, a request for `/admin/users` is
    /// forwarded as `/users` — same semantics as Caddy's `handle_path`.
    pub strip_prefix: Option<String>,
}

/// A Wasm HTTP trigger: maps a path pattern to a Wasm runtime instance.
#[derive(Debug, Clone)]
pub struct WasmTrigger {
    /// Path pattern (e.g., "/api/edge/*").
    pub pattern: String,
    /// Wasm runtime instance ID.
    pub runtime_id: String,
    /// Service name for logging.
    pub service_name: String,
}

/// Callback invoked when a request matches a Wasm trigger.
/// Receives (runtime_id, method, path, body) and returns the response body string.
pub type WasmInvoker =
    Arc<dyn Fn(String, String, String, String) -> WasmInvokeFuture + Send + Sync>;

/// Future type returned by the Wasm invoker.
pub type WasmInvokeFuture =
    std::pin::Pin<Box<dyn std::future::Future<Output = Result<String, String>> + Send>>;

/// Shared Wasm trigger table type.
pub type SharedWasmTriggers = Arc<RwLock<Vec<WasmTrigger>>>;

/// Run the reverse proxy on the given port.
pub async fn run_proxy(
    route_table: Arc<RwLock<HashMap<String, Vec<RouteTarget>>>>,
    wasm_triggers: SharedWasmTriggers,
    wasm_invoker: Option<WasmInvoker>,
    port: u16,
    tls_acceptor: Option<tokio_rustls::TlsAcceptor>,
    acme_manager: Option<AcmeManager>,
) -> anyhow::Result<()> {
    let addr = format!("0.0.0.0:{port}");
    let listener = TcpListener::bind(&addr).await?;
    let proto = if tls_acceptor.is_some() {
        "HTTPS"
    } else {
        "HTTP"
    };
    info!("Reverse proxy listening on {addr} ({proto})");

    serve_loop(
        listener,
        route_table,
        wasm_triggers,
        wasm_invoker,
        tls_acceptor,
        acme_manager,
    )
    .await
}

/// Run the proxy with optional fallback support.
#[allow(clippy::too_many_arguments)]
pub async fn run_proxy_with_fallback(
    route_table: Arc<RwLock<HashMap<String, Vec<RouteTarget>>>>,
    wasm_triggers: SharedWasmTriggers,
    wasm_invoker: Option<WasmInvoker>,
    port: u16,
    tls_acceptor: Option<tokio_rustls::TlsAcceptor>,
    acme_manager: Option<AcmeManager>,
    fallback: Option<FallbackConfig>,
) -> anyhow::Result<()> {
    let addr = format!("0.0.0.0:{port}");
    let listener = TcpListener::bind(&addr).await?;
    let proto = if tls_acceptor.is_some() {
        "HTTPS"
    } else {
        "HTTP"
    };
    info!("Reverse proxy listening on {addr} ({proto})");

    serve_loop_with_fallback(
        listener,
        route_table,
        wasm_triggers,
        wasm_invoker,
        tls_acceptor,
        acme_manager,
        fallback,
    )
    .await
}

/// Shared dynamic cert resolver for hot-provisioning.
pub type SharedCertResolver = Arc<acme::DynCertResolver>;

/// Core accept loop shared by HTTP and HTTPS listeners.
async fn serve_loop(
    listener: TcpListener,
    route_table: Arc<RwLock<HashMap<String, Vec<RouteTarget>>>>,
    wasm_triggers: SharedWasmTriggers,
    wasm_invoker: Option<WasmInvoker>,
    tls_acceptor: Option<tokio_rustls::TlsAcceptor>,
    acme_manager: Option<AcmeManager>,
) -> anyhow::Result<()> {
    serve_loop_with_fallback(
        listener,
        route_table,
        wasm_triggers,
        wasm_invoker,
        tls_acceptor,
        acme_manager,
        None,
    )
    .await
}

/// Serve loop variant with fallback support for SNI passthrough and HTTP forwarding.
#[allow(clippy::too_many_arguments)]
pub(crate) async fn serve_loop_with_fallback(
    listener: TcpListener,
    route_table: Arc<RwLock<HashMap<String, Vec<RouteTarget>>>>,
    wasm_triggers: SharedWasmTriggers,
    wasm_invoker: Option<WasmInvoker>,
    tls_acceptor: Option<tokio_rustls::TlsAcceptor>,
    acme_manager: Option<AcmeManager>,
    fallback: Option<FallbackConfig>,
) -> anyhow::Result<()> {
    let counter = Arc::new(AtomicUsize::new(0));
    // Upload- and response-aware timeouts (#187), not reqwest's read_timeout.
    let client = Arc::new(backend_timeouts::client());
    // A TLS endpoint exists if this listener terminates TLS itself, or if
    // an ACME manager is present — the plain-HTTP listener of the ACME
    // dual-listener setup carries one for HTTP-01 challenges, which is
    // exactly the "HTTPS runs alongside on 443" signal. With neither, the
    // HTTP→HTTPS redirect has nowhere to send clients and must not fire
    // (#123: ACME unconfigured meant every routed host redirected into a
    // closed port).
    let https_enabled = tls_acceptor.is_some() || acme_manager.is_some();
    let acme = acme_manager.map(Arc::new);
    let is_tls = tls_acceptor.is_some();
    let rate_limiter = RateLimiter::new();

    let fallback = Arc::new(fallback);
    loop {
        let (stream, peer) = match listener.accept().await {
            Ok(conn) => conn,
            Err(e) => {
                warn!("Proxy accept error: {e}");
                continue;
            }
        };

        let routes = route_table.clone();
        let triggers = wasm_triggers.clone();
        let invoker = wasm_invoker.clone();
        let counter = counter.clone();
        let client = client.clone();
        let acme = acme.clone();
        let tls = tls_acceptor.clone();
        let rl = rate_limiter.clone();
        let fb = fallback.clone();
        let routes_for_sni = routes.clone();

        let fb_for_service = fb.clone();
        tokio::spawn(async move {
            let service = service_fn(move |req: Request<Incoming>| {
                let routes = routes.clone();
                let triggers = triggers.clone();
                let invoker = invoker.clone();
                let counter = counter.clone();
                let client = client.clone();
                let acme = acme.clone();
                let rl = rl.clone();
                let fb = fb_for_service.clone();
                async move {
                    if let Some(resp) = handle_acme_challenge(&req, acme.as_deref()).await {
                        return Ok(resp);
                    }
                    let mut resp = handle_request(
                        req,
                        &routes,
                        &triggers,
                        invoker.as_ref(),
                        &counter,
                        &client,
                        is_tls,
                        https_enabled,
                        &rl,
                        peer,
                        fb.as_ref().as_ref(),
                    )
                    .await?;
                    // Inject baseline security headers (add-if-absent; HSTS only
                    // over TLS). No-op unless a policy was installed at startup.
                    security_headers::apply(&mut resp, is_tls);
                    Ok::<_, hyper::Error>(resp)
                }
            });
            if let Some(acceptor) = tls {
                let mut stream = stream;
                // Peek SNI to decide between local TLS termination and pass-through
                let sni = sni::peek_sni(&mut stream).await;
                let should_passthrough = if let Some(ref host) = sni {
                    let routes_lock = routes_for_sni.read().await;
                    let known = routes_lock.contains_key(host);
                    drop(routes_lock);
                    !known && fb.as_ref().as_ref().and_then(|f| f.tls.as_ref()).is_some()
                } else {
                    false
                };

                if should_passthrough {
                    let target = fb
                        .as_ref()
                        .as_ref()
                        .and_then(|f| f.tls.clone())
                        .expect("checked above");
                    debug!(?sni, %target, "SNI passthrough");
                    match tokio::net::TcpStream::connect(&target).await {
                        Ok(mut backend) => {
                            if let Err(e) =
                                tokio::io::copy_bidirectional(&mut stream, &mut backend).await
                            {
                                debug!("Passthrough copy error from {peer}: {e}");
                            }
                        }
                        Err(e) => warn!("Failed to connect to TLS fallback {target}: {e}"),
                    }
                    return;
                }

                match acceptor.accept(stream).await {
                    Ok(tls_stream) => {
                        let io = TokioIo::new(tls_stream);
                        // h1 or h2, picked by the ALPN result. WebSocket
                        // upgrades stay on h1: extended CONNECT (RFC 8441)
                        // is not enabled, so browsers open a separate
                        // HTTP/1.1 connection for them, as before.
                        // hyper's h2 default caps the header list at 16 KiB and
                        // answers 431 above it; SSO cookie jars (Keycloak plus
                        // the app's own) exceed that where h1 never refused.
                        let mut builder = auto::Builder::new(TokioExecutor::new());
                        builder.http2().max_header_list_size(H2_MAX_HEADER_LIST);
                        if let Err(e) = builder.serve_connection_with_upgrades(io, service).await {
                            debug!("TLS proxy error from {peer}: {e}");
                        }
                    }
                    Err(e) => debug!("TLS handshake failed from {peer}: {e}"),
                }
            } else {
                // The plain listener stays HTTP/1-only: browsers never speak
                // cleartext h2 (h2c), and this port only serves redirects,
                // ACME challenges and plain-HTTP routes, so accepting h2c
                // would add parser surface for no client.
                let io = TokioIo::new(stream);
                if let Err(e) = http1::Builder::new()
                    .serve_connection(io, service)
                    .with_upgrades()
                    .await
                {
                    debug!("Proxy connection error from {peer}: {e}");
                }
            }
        });
    }
}