Skip to main content

openkind_api/
http.rs

1//! HTTP layer — axum 0.8.
2//!
3//! Routes:
4//! - `POST /v1/systemone`  → the Jev evaluation endpoint
5//! - `GET  /health`        → liveness
6//! - `GET  /v1/models`     → list available model aliases (Jev shape)
7//! - `GET  /metrics`       → Prometheus scrape
8//! - `GET  /playground`    → embedded web UI, only with `--playground on`
9//! - `POST /v1/arrow`      → unofficial bulk Arrow IPC endpoint, only with
10//!   `--arrow on`; outside the TypeSafe wire contract (see [`crate::arrow`])
11//!
12//! Every response is stamped with `x-typesafe-request-id` (the SDK reads
13//! this to log per-request correlation), and `/v1/*` is gated by an
14//! optional bearer token when `OPENKIND_API_KEY` is set.
15
16use std::sync::Arc;
17
18use axum::{
19    extract::State,
20    response::IntoResponse,
21    routing::{get, post},
22    Json, Router,
23};
24use openkind_core::{ResponseContract, SystemRequest};
25use openkind_engine::{dispatch, EngineRegistry};
26use tower_http::trace::{DefaultMakeSpan, TraceLayer};
27use tracing::Level;
28
29use crate::error::ApiError;
30use crate::middleware::AuthConfig;
31use crate::models::ModelsResponse;
32use crate::AppState;
33
34/// Maximum allowed request payload size in bytes (16 MB) to protect against DoS memory exhaustion.
35pub const MAX_PAYLOAD_SIZE_BYTES: usize = 16 * 1024 * 1024;
36
37/// Build the HTTP router with explicit payload size limit and the default
38/// per-IP rate limit (see [`crate::middleware::RateLimitConfig::default`]).
39pub fn router_with_state_and_limit(
40    state: AppState,
41    auth: AuthConfig,
42    max_payload_bytes: usize,
43) -> Router {
44    router_full(
45        state,
46        auth,
47        max_payload_bytes,
48        crate::middleware::RateLimiter::new(crate::middleware::RateLimitConfig::default()),
49        false,
50        false,
51    )
52}
53
54/// Build the HTTP router with explicit payload size limit and rate limiting.
55pub fn router_with_state_auth_rate_limit(
56    state: AppState,
57    auth: AuthConfig,
58    max_payload_bytes: usize,
59    rate_limiter: crate::middleware::RateLimiter,
60) -> Router {
61    router_full(state, auth, max_payload_bytes, rate_limiter, false, false)
62}
63
64/// Build the daemon HTTP router: explicit payload size limit, rate limiting,
65/// and the embedded playground UI when `playground` is set
66/// (`openkindd --playground on`).
67pub fn router_daemon(
68    state: AppState,
69    auth: AuthConfig,
70    max_payload_bytes: usize,
71    rate_limiter: crate::middleware::RateLimiter,
72    playground: bool,
73) -> Router {
74    router_full(
75        state,
76        auth,
77        max_payload_bytes,
78        rate_limiter,
79        playground,
80        false,
81    )
82}
83
84/// Build the daemon router with the unofficial Arrow endpoint explicitly enabled.
85/// `arrow = false` has the same behavior as [`router_daemon`].
86pub fn router_daemon_with_arrow(
87    state: AppState,
88    auth: AuthConfig,
89    max_payload_bytes: usize,
90    rate_limiter: crate::middleware::RateLimiter,
91    playground: bool,
92    arrow: bool,
93) -> Router {
94    router_full(
95        state,
96        auth,
97        max_payload_bytes,
98        rate_limiter,
99        playground,
100        arrow,
101    )
102}
103
104fn router_full(
105    state: AppState,
106    auth: AuthConfig,
107    max_payload_bytes: usize,
108    rate_limiter: crate::middleware::RateLimiter,
109    playground: bool,
110    arrow: bool,
111) -> Router {
112    let routes = Router::new()
113        // POST /v1/systemone — canonical Jev decision evaluation endpoint.
114        .route("/v1/systemone", post(systemone))
115        // POST /v1/system_one — SDK alias for decision evaluation endpoint.
116        .route("/v1/system_one", post(systemone))
117        // GET /v1/models — list registered models and their capabilities.
118        .route("/v1/models", get(list_models))
119        // GET /health — unauthenticated service liveness probe.
120        .route("/health", get(health))
121        // GET /metrics — Prometheus text-format scrape target.
122        .route("/metrics", get(prometheus_metrics));
123    // GET /playground — embedded web UI. Opt-in because it is a developer
124    // convenience outside the wire contract. The HTML shell is public;
125    // model lifecycle routes stay behind auth.
126    let routes = if playground {
127        routes
128            .route("/playground", get(crate::playground::playground_page))
129            .route(
130                "/playground/api/models",
131                get(crate::playground::list_models).post(crate::playground::change_model),
132            )
133    } else {
134        routes
135    };
136    // POST /v1/arrow: unofficial bulk Arrow IPC endpoint. Opt-in like the
137    // playground because it is outside the TypeSafe wire contract; it still
138    // sits behind the /v1 auth gate and rate limiter.
139    let routes = if arrow {
140        routes.route("/v1/arrow", post(crate::arrow::arrow_batch))
141    } else {
142        routes
143    };
144    // A disabled limiter has no observable effect. Leave its middleware
145    // off the router so it cannot allocate or dispatch on every request.
146    let routes = if rate_limiter.is_enabled() {
147        routes.layer(axum::middleware::from_fn_with_state(
148            rate_limiter,
149            crate::middleware::rate_limit_layer,
150        ))
151    } else {
152        routes
153    };
154    routes
155        // Order matters: layers added LATER are OUTERMOST. We want
156        // request_id outermost so it stamps the response on every code
157        // path, including 401s from auth_layer and 429s from the rate
158        // limiter (both short-circuit before any handler middleware
159        // fires). Authentication sits outside rate limiting so rejected
160        // credentials cannot exhaust the budget shared by requests from
161        // the same TCP peer (for example, a reverse proxy).
162        .layer(axum::middleware::from_fn_with_state(
163            auth,
164            crate::middleware::auth_layer,
165        ))
166        .layer(axum::middleware::from_fn(
167            crate::middleware::request_id_layer,
168        ))
169        .layer(axum::extract::DefaultBodyLimit::max(max_payload_bytes))
170        // Keep request headers out of telemetry. In particular, auth
171        // credentials and caller-supplied request IDs must never become span
172        // fields; native-engine metrics use only fixed outcome labels.
173        .layer(
174            TraceLayer::new_for_http().make_span_with(
175                DefaultMakeSpan::new()
176                    .level(Level::INFO)
177                    .include_headers(false),
178            ),
179        )
180        .with_state(Arc::new(state))
181}
182
183/// Build the HTTP router with a pre-built state and default 16MB payload limit. Used by the daemon.
184pub fn router_with_state(state: AppState, auth: AuthConfig) -> Router {
185    router_with_state_and_limit(state, auth, MAX_PAYLOAD_SIZE_BYTES)
186}
187
188/// Build the HTTP router with default authentication configuration (no API key required).
189pub fn router(registry: EngineRegistry) -> Router {
190    router_with_state(AppState::new(registry), AuthConfig::default())
191}
192
193/// Build the HTTP router with explicit authentication configuration.
194pub fn router_with_auth(registry: EngineRegistry, auth: AuthConfig) -> Router {
195    router_with_state(AppState::new(registry), auth)
196}
197
198// Re-exported at the crate root for tests.
199pub use router_with_state as build_router_with_state;
200
201/// Canonical evaluation handler for POST `/v1/systemone` and `/v1/system_one`.
202async fn systemone(
203    State(state): State<Arc<AppState>>,
204    headers: axum::http::HeaderMap,
205    req: Result<Json<SystemRequest>, axum::extract::rejection::JsonRejection>,
206) -> Result<axum::response::Response, ApiError> {
207    let Json(req) = match req {
208        Ok(j) => j,
209        Err(rejection) => match rejection {
210            axum::extract::rejection::JsonRejection::BytesRejection(e) => {
211                return Err(ApiError::PayloadTooLarge(e.to_string()));
212            }
213            axum::extract::rejection::JsonRejection::JsonSyntaxError(e) => {
214                return Err(ApiError::BadJson(e.to_string()));
215            }
216            axum::extract::rejection::JsonRejection::JsonDataError(e) => {
217                return Err(ApiError::InvalidBody(e.to_string()));
218            }
219            other => {
220                return Err(ApiError::InvalidBody(other.to_string()));
221            }
222        },
223    };
224    // Proxy-cache mode: the hook decides whether this alias is proxied and
225    // answers from the distilling cache or forwards upstream with the
226    // caller's own credentials. Everything else dispatches locally.
227    if let Some(proxy) = &state.proxy {
228        if proxy.wants(&req) {
229            let contract = ResponseContract::from_request(&req)
230                .map_err(|error| ApiError::InvalidBody(error.to_string()))?;
231            let caller_key = if proxy.forwards_caller_credentials() {
232                bearer_of(&headers)
233            } else {
234                None
235            };
236            let outcome = proxy.evaluate(req, caller_key).await?;
237            contract.validate(&outcome.response).map_err(|error| {
238                ApiError::BadGateway(format!("upstream returned an invalid response: {error}"))
239            })?;
240            let mut response = (axum::http::StatusCode::OK, Json(outcome.response)).into_response();
241            let headers = response.headers_mut();
242            if let Ok(value) = axum::http::HeaderValue::from_str(outcome.source.as_str()) {
243                headers.insert(
244                    axum::http::HeaderName::from_static("x-openkind-cache"),
245                    value,
246                );
247            }
248            if let Some(detail) = &outcome.detail {
249                if let Ok(text) = serde_json::to_string(detail) {
250                    if let Ok(value) = axum::http::HeaderValue::from_str(&text) {
251                        headers.insert(
252                            axum::http::HeaderName::from_static("x-openkind-cache-detail"),
253                            value,
254                        );
255                    }
256                }
257            }
258            return Ok(response);
259        }
260    }
261    let resp = dispatch(req, &state.registry).await?;
262    Ok((axum::http::StatusCode::OK, Json(resp)).into_response())
263}
264
265/// Extract the caller's bearer credential (the upstream Jev key) without
266/// logging or storing it beyond the proxy's salted hash.
267fn bearer_of(headers: &axum::http::HeaderMap) -> Option<String> {
268    let value = headers.get(axum::http::header::AUTHORIZATION)?;
269    let value = value.to_str().ok()?;
270    let token = value
271        .strip_prefix("Bearer ")
272        .or_else(|| value.strip_prefix("bearer "))?;
273    let token = token.trim();
274    if token.is_empty() {
275        None
276    } else {
277        Some(token.to_owned())
278    }
279}
280
281/// Model listing handler for GET `/v1/models`.
282async fn list_models(State(state): State<Arc<AppState>>) -> impl IntoResponse {
283    // In proxy mode the upstream listing is authoritative for proxied
284    // aliases; fall back to the local registry when the upstream cannot be
285    // reached.
286    if let Some(proxy) = &state.proxy {
287        if let Some(models) = proxy.models().await {
288            return Json(models).into_response();
289        }
290    }
291    let models = state.registry.list_models();
292    Json(ModelsResponse::new(models)).into_response()
293}
294
295/// Service liveness probe handler for GET `/health`.
296///
297/// The body is constant, so it is pre-encoded once instead of rebuilding a
298/// `serde_json::Value` and re-serializing it on every probe. The bytes and
299/// content type match what `Json(json!({"status":"ok"}))` produced.
300async fn health() -> impl IntoResponse {
301    static HEALTH_BODY: &str = "{\"status\":\"ok\"}";
302    (
303        [(
304            axum::http::header::CONTENT_TYPE,
305            axum::http::HeaderValue::from_static("application/json"),
306        )],
307        HEALTH_BODY,
308    )
309}
310
311static HANDLE: std::sync::OnceLock<metrics_exporter_prometheus::PrometheusHandle> =
312    std::sync::OnceLock::new();
313
314/// Prometheus text-format metrics exporter. Returns a valid empty body
315/// even when the recorder hasn't been installed — scrapers should
316/// always see 200. Real metrics fire once the server binary calls
317/// `install_metrics_recorder()` on startup.
318async fn prometheus_metrics() -> impl IntoResponse {
319    if let Some(h) = HANDLE.get() {
320        (
321            axum::http::StatusCode::OK,
322            [("content-type", "text/plain; version=0.0.4")],
323            h.render(),
324        )
325    } else {
326        (
327            axum::http::StatusCode::OK,
328            [("content-type", "text/plain; version=0.0.4")],
329            "# metrics recorder not installed\n".to_string(),
330        )
331    }
332}
333
334/// Install the Prometheus recorder. Called by the server binary on startup.
335/// Exposed here so the binary doesn't have to depend on
336/// metrics-exporter-prometheus directly.
337pub fn install_metrics_recorder() -> anyhow::Result<()> {
338    if HANDLE.get().is_some() {
339        return Ok(());
340    }
341    use metrics_exporter_prometheus::PrometheusBuilder;
342    let handle = PrometheusBuilder::new()
343        .install_recorder()
344        .map_err(|e| anyhow::anyhow!("install metrics recorder: {e}"))?;
345    let _ = HANDLE.set(handle);
346    Ok(())
347}
348
349#[cfg(test)]
350#[path = "http_tests.rs"]
351mod tests;