diff --git a/README.md b/README.md index cbf99e6..e208de5 100644 --- a/README.md +++ b/README.md @@ -26,9 +26,10 @@ System1-Omni models, designed around a Rust frontend, model-owned execution, and high-performance CUDA and Metal backends. -The Rust frontend forwards requests to a separately running model worker. The -Cua-S1 4B 0.2 `text` adapter and Open-Jev-27B-v1.1 have native workers using -shared CUDA kernels in this repository. +The Rust frontend forwards requests to a separately running model worker, and +defines the small engine interface a model implements to serve in-process +instead. The Cua-S1 4B 0.2 `text` adapter and Open-Jev-27B-v1.1 have native +workers using shared CUDA kernels in this repository. ## News diff --git a/src/frontend/README.md b/src/frontend/README.md index 9569594..6ea66e0 100644 --- a/src/frontend/README.md +++ b/src/frontend/README.md @@ -1,8 +1,14 @@ # Rust frontend -An Axum/Tokio server that forwards requests to a separately running model worker -using Reqwest. The worker handles validation, media loading, preprocessing and -inference. +An Axum/Tokio server for the `/v1/systemone` decision API. It serves the same HTTP +surface two ways: + +- **Forwarding** (default). Requests are passed to a separately running model worker + with Reqwest. The worker handles validation, media loading, preprocessing and + inference, and this process owns no model state. +- **In-process.** A model linked into this binary answers the requests through the + [`Engine`](src/engine.rs) trait. Admission, the request budget and readiness belong to + the transport; the model owns the bytes. ## Run @@ -22,6 +28,8 @@ connections bypass system HTTP proxies. ## HTTP interface +### Forwarding mode + - `POST /v1/systemone` forwards the `model`, `state` and `questions` envelope unchanged. Workers return `choice`, `score` or `noul` decisions; see the [Jev API reference](https://docs.typesafe.ai/api). @@ -37,6 +45,33 @@ Text, image, audio, video and mixed payloads pass through as bytes. Actual infer support depends on the worker. The [Laya recipe](../../recipe/laya/README.md) verifies text decisions against a real backend. +### In-process mode + +The same two routes, answered by [`engine::app`](src/engine.rs): + +- `POST /v1/systemone` passes the body to the engine and returns the body it produced. + Nothing is parsed or re-serialized, so a field the frontend has never heard of + survives, and a model's error text cannot produce invalid JSON. +- `GET /health` is the readiness probe. `200` carries `{"status":"ok"}`; before the + engine can take work it is `503` with `{"status":"starting"}`, and an engine that + failed answers `503` with `{"status":"failed","reason":…}`. When the engine reports + its queue, `depth` and `rejected` are included; an engine that reports nothing gets + neither field rather than a zero it did not measure. +- Successful responses carry `x-queue-depth`, measured when the request was dequeued. +- Status codes: `400` invalid request, `413` body over the configured limit, `500` the + engine failed to run an accepted request, `503` not ready / no queue capacity / + request budget expired / the engine stopped, `504` the engine did not answer inside + the budget. + +```rust +let app = omni_jev::engine::app(engine, omni_jev::engine::ServiceConfig::default()); +``` + +An `Engine` reports readiness, accepts one request and hands back a reply carrying the +response body and the queue depth to publish with it. A model implements that and +nothing else; the thread that owns it, the bounded queue and the shutdown sequence +around it arrive in a following change. + ## Checks From the repository root: @@ -51,6 +86,8 @@ cargo clippy --workspace --locked --all-targets -- -D warnings cargo test --workspace --locked ``` -Tests use local mock workers; no model weights or GPU are needed. They cover -multimodal byte preservation, authorization, connection reuse, large uploads, -backend errors, timeouts, health and binary startup/shutdown. +Tests use local mock workers and a fake engine; no model weights and no GPU are needed. +They cover multimodal byte preservation, authorization, connection reuse, large uploads, +backend errors, timeouts, health, and — for the in-process path — readiness, byte +preservation, the body limit, queue-full refusal, the deadline, error mapping and the +queue depth a reply publishes. diff --git a/src/frontend/src/engine.rs b/src/frontend/src/engine.rs new file mode 100644 index 0000000..a478f98 --- /dev/null +++ b/src/frontend/src/engine.rs @@ -0,0 +1,287 @@ +//! In-process serving: the same `/v1/systemone` surface, answered by a linked engine. +//! +//! The forwarding path in [`crate`] stays available for workers that already speak HTTP. +//! This module is for a model in the same binary: the transport owns admission, the +//! request budget and readiness, and the model owns the bytes. Nothing here parses the +//! decision envelope, so a field this crate has never heard of survives. + +use std::{ + fmt, + future::Future, + sync::Arc, + time::{Duration, Instant}, +}; + +use axum::{ + Router, + body::{Body, to_bytes}, + extract::{Request, State}, + http::{HeaderName, HeaderValue, StatusCode, header}, + response::Response, + routing::{get, post}, +}; +use tokio::{net::TcpListener, sync::oneshot}; + +use crate::BoxError; + +/// Requests accepted and not yet completed, as reported when a response was produced. Only +/// sent when the engine reports one, so a missing header means "not observed", not "empty". +pub const QUEUE_DEPTH: HeaderName = HeaderName::from_static("x-queue-depth"); + +/// How far the engine has got. `GET /health` is the readiness probe a load balancer or a +/// benchmark script polls before sending traffic, and `Starting` is what it sees while +/// weights load and the engine warms up. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum Readiness { + Starting, + Ready, + /// Startup failed. The process keeps running and `/health` carries the reason, so the + /// failure is readable by an operator instead of vanishing with the process. + Failed, +} + +/// What an engine tells `/health` about itself. Advisory: readiness never waits on it. +#[derive(Clone, Copy, Debug, Default)] +pub struct Report { + /// Requests accepted and not yet completed, including the one running. + pub depth: usize, + /// Requests the engine refused or dropped since it started. Monotonic. + pub rejected: u64, +} + +/// The reply the transport sends. `depth` is how many requests were outstanding when this +/// one was dequeued, so it includes this one; it becomes the `x-queue-depth` header, which +/// makes backpressure visible before it turns into a refusal. +pub struct Answer { + pub body: Vec, + pub depth: usize, +} + +/// A closed receiver means the client is gone: a queued request must then be skipped +/// rather than executed, and an engine that has already started work must still finish it +/// before reusing GPU buffers. +pub type Reply = oneshot::Receiver>; + +/// Why an engine could not answer. The transport maps each variant to one status code and +/// one message, so the mapping is a property of this crate rather than of a model. +#[derive(Debug)] +pub enum EngineError { + /// The caller's body is not a request this engine accepts. + InvalidRequest(String), + /// Out of capacity right now. The caller may retry. + Busy, + /// Not able to accept work at all: still loading, drained, or dead. + Unavailable, + /// The request was well formed and the engine failed to execute it. + InferenceFailed, +} + +impl fmt::Display for EngineError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::InvalidRequest(detail) => write!(f, "invalid request: {detail}"), + Self::Busy => f.write_str("no queue capacity"), + Self::Unavailable => f.write_str("engine unavailable"), + Self::InferenceFailed => f.write_str("inference failed"), + } + } +} + +impl std::error::Error for EngineError {} + +/// The model side of the transport. Implementations are usually a [`crate::worker::Handle`], +/// which forwards to the thread that owns the engine. +pub trait Engine: Send + Sync + 'static { + /// Non-blocking. Polled on every request, so it must not lock behind queued work. + fn readiness(&self) -> Readiness; + + /// Accepts one request. `Ok` means the engine owns it and the transport waits for the + /// reply until `deadline`; [`EngineError::Busy`] means nothing was accepted. + fn submit(&self, body: Vec, deadline: Instant) -> Result; + + /// What `/health` publishes while ready. `None` leaves those fields out. + fn report(&self) -> Option { + None + } + + /// Set only when [`Engine::readiness`] is [`Readiness::Failed`], and surfaced verbatim + /// so an operator can read the cause without the server log. + fn failure(&self) -> Option { + None + } +} + +/// Admission policy: a bounded queue, and a budget that starts when the headers arrive. +#[derive(Clone, Debug)] +pub struct ServiceConfig { + /// Total request budget: upload, queueing and inference together. + pub timeout: Duration, + pub max_body: usize, +} + +impl Default for ServiceConfig { + fn default() -> Self { + Self { + timeout: Duration::from_secs(30), + max_body: 1024 * 1024, + } + } +} + +#[derive(Clone)] +struct Service { + engine: Arc, + config: ServiceConfig, +} + +/// The native router: `POST /v1/systemone` and `GET /health`. +pub fn app(engine: Arc, config: ServiceConfig) -> Router { + Router::new() + .route("/v1/systemone", post(infer)) + .route("/health", get(health)) + .with_state(Service { engine, config }) +} + +/// Serves until `shutdown` resolves, then drains in-flight connections. +/// +/// Draining here covers the transport only. Stopping admission and joining the engine's +/// own worker happen after this returns, and belong to the caller. +pub async fn serve( + listener: TcpListener, + router: Router, + shutdown: impl Future + Send + 'static, +) -> Result<(), BoxError> { + axum::serve(listener, router) + .with_graceful_shutdown(shutdown) + .await?; + Ok(()) +} + +async fn health(State(service): State) -> Response { + let readiness = service.engine.readiness(); + let report = match readiness { + Readiness::Ready => service.engine.report(), + // A starting engine knows nothing yet, and a failed one has already said why. + _ => None, + }; + let (status, value) = match readiness { + Readiness::Ready => (StatusCode::OK, "ok"), + Readiness::Starting => (StatusCode::SERVICE_UNAVAILABLE, "starting"), + Readiness::Failed => (StatusCode::SERVICE_UNAVAILABLE, "failed"), + }; + let mut fields = vec![format!("\"status\":\"{value}\"")]; + if let (Readiness::Failed, Some(reason)) = (readiness, service.engine.failure()) { + fields.push(format!("\"reason\":{}", quote(&reason))); + } + if let Some(report) = report { + fields.push(format!("\"depth\":{}", report.depth)); + fields.push(format!("\"rejected\":{}", report.rejected)); + } + json(status, &format!("{{{}}}", fields.join(","))) +} + +async fn infer(State(service): State, request: Request) -> Response { + // The budget starts here, so a slow upload is charged to the caller rather than to the + // engine, and work that cannot finish in time is refused before it queues. + let deadline = Instant::now() + service.config.timeout; + if service.engine.readiness() != Readiness::Ready { + return fail(StatusCode::SERVICE_UNAVAILABLE, "model unavailable"); + } + + // Reading the body is charged to the same budget. Too large and too slow are separate + // answers so a caller can tell a big request from a late one. + let body = match tokio::time::timeout_at( + deadline.into(), + to_bytes(request.into_body(), service.config.max_body), + ) + .await + { + Ok(Ok(bytes)) => bytes.to_vec(), + Ok(Err(_)) => return fail(StatusCode::PAYLOAD_TOO_LARGE, "request body too large"), + Err(_) => { + return fail( + StatusCode::SERVICE_UNAVAILABLE, + "request deadline expired while uploading", + ); + } + }; + if Instant::now() >= deadline { + return fail(StatusCode::SERVICE_UNAVAILABLE, "request deadline expired"); + } + + let reply = match service.engine.submit(body, deadline) { + Ok(reply) => reply, + Err(error) => return error_response(error), + }; + match tokio::time::timeout_at(deadline.into(), reply).await { + Ok(Ok(Ok(answer))) => with_depth(answer), + Ok(Ok(Err(error))) => error_response(error), + // The engine dropped a request it had already accepted. Not bad input and not a + // failed inference, so it reads as lost capacity. + Ok(Err(_)) => fail(StatusCode::SERVICE_UNAVAILABLE, "engine stopped"), + Err(_) => fail(StatusCode::GATEWAY_TIMEOUT, "inference timed out"), + } +} + +fn error_response(error: EngineError) -> Response { + let (status, message) = match error { + EngineError::InvalidRequest(detail) => (StatusCode::BAD_REQUEST, detail), + EngineError::Busy => ( + StatusCode::SERVICE_UNAVAILABLE, + "inference queue full".into(), + ), + EngineError::Unavailable => (StatusCode::SERVICE_UNAVAILABLE, "model unavailable".into()), + EngineError::InferenceFailed => { + (StatusCode::INTERNAL_SERVER_ERROR, "inference failed".into()) + } + }; + fail(status, &message) +} + +fn fail(status: StatusCode, message: &str) -> Response { + json(status, &format!("{{\"error\":{}}}", quote(message))) +} + +/// Answers with the engine's body, publishing the queue depth it reported. +fn with_depth(answer: Answer) -> Response { + let mut response = Response::new(Body::from(answer.body)); + let headers = response.headers_mut(); + headers.insert( + header::CONTENT_TYPE, + HeaderValue::from_static("application/json"), + ); + if let Ok(depth) = HeaderValue::from_str(&answer.depth.to_string()) { + headers.insert(QUEUE_DEPTH, depth); + } + response +} + +fn json(status: StatusCode, body: &str) -> Response { + let mut response = Response::new(Body::from(body.to_owned())); + *response.status_mut() = status; + response.headers_mut().insert( + header::CONTENT_TYPE, + HeaderValue::from_static("application/json"), + ); + response +} + +/// Serializes a string as a JSON string literal. Hand-rolled so this crate keeps parsing +/// nothing: a model's error text must not be able to produce invalid JSON. +fn quote(value: &str) -> String { + let mut out = String::with_capacity(value.len() + 2); + out.push('"'); + for c in value.chars() { + match c { + '"' => out.push_str("\\\""), + '\\' => out.push_str("\\\\"), + '\n' => out.push_str("\\n"), + '\r' => out.push_str("\\r"), + '\t' => out.push_str("\\t"), + c if (c as u32) < 0x20 => out.push_str(&format!("\\u{:04x}", c as u32)), + c => out.push(c), + } + } + out.push('"'); + out +} diff --git a/src/frontend/src/lib.rs b/src/frontend/src/lib.rs index 0de9ba5..69afe90 100644 --- a/src/frontend/src/lib.rs +++ b/src/frontend/src/lib.rs @@ -1,4 +1,13 @@ //! Jev HTTP transport. The worker owns request parsing and inference. +//! +//! Two ways to answer `/v1/systemone`: +//! +//! - [`Config`] and [`app`] forward to a worker that already speaks HTTP. The frontend +//! owns no model state, which is how the Jev path works today. +//! - [`engine`] serves a model linked into this binary. Admission, the request budget and +//! readiness are the transport's; the model owns the bytes. + +pub mod engine; use std::{env, error::Error, net::SocketAddr, time::Duration}; diff --git a/src/frontend/tests/engine_service.rs b/src/frontend/tests/engine_service.rs new file mode 100644 index 0000000..6a17954 --- /dev/null +++ b/src/frontend/tests/engine_service.rs @@ -0,0 +1,317 @@ +//! End-to-end tests over real sockets for the in-process path: client -> omni-jev -> +//! a fake engine. The engine is a stand-in, so none of this needs a model or a GPU. +//! +//! The last test runs the compiled binary with `OMNI_SYSTEMONE_ENGINE=passthrough`, which +//! is the only way to cover the whole lifecycle: readiness, signals and process exit. + +use std::{ + sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, + }, + time::{Duration, Instant}, +}; + +use axum::Router; +use omni_jev::engine::{ + self, Answer, Engine, EngineError, Readiness, Reply, Report, ServiceConfig, +}; +use tokio::{net::TcpListener, sync::oneshot}; + +const REQUEST: &str = r#"{"model":"english","state":"refund please","questions":{"q":{"type":"noul","instructions":"Ask?"}}}"#; + +/// How the fake engine answers. Every mode is one of the states the transport has to tell +/// apart, so the status codes below are the contract rather than a convention. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum Mode { + Ready, + Starting, + Failed, + /// Accepts the request and never answers it. + Silent, + /// Refuses with a full queue. + Full, + /// Echoes its own error variants. + Invalid, + Inference, +} + +struct Fake { + mode: Mode, + depth: usize, + rejected: u64, + /// Number of requests that reached `submit`. + submitted: Arc, +} + +impl Default for Fake { + fn default() -> Self { + Self { + mode: Mode::Ready, + depth: 0, + rejected: 0, + submitted: Arc::new(AtomicUsize::new(0)), + } + } +} + +impl Engine for Fake { + fn readiness(&self) -> Readiness { + match self.mode { + Mode::Starting => Readiness::Starting, + Mode::Failed => Readiness::Failed, + _ => Readiness::Ready, + } + } + + fn submit(&self, body: Vec, _deadline: Instant) -> Result { + self.submitted.fetch_add(1, Ordering::AcqRel); + match self.mode { + Mode::Full => Err(EngineError::Busy), + Mode::Invalid => Err(EngineError::InvalidRequest("bad \"envelope\"\n".into())), + Mode::Inference => Err(EngineError::InferenceFailed), + Mode::Silent => { + // Keeping the sender alive means the request neither completes nor fails: + // only the deadline can end it. + let (reply, receiver) = oneshot::channel(); + std::mem::forget(reply); + Ok(receiver) + } + _ => { + let answer = Answer { + body, + depth: self.depth, + }; + // Answered from another task, because a real engine never answers inside + // `submit` and the transport must not depend on getting one immediately. + let (reply, receiver) = oneshot::channel(); + tokio::spawn(async move { + let _ = reply.send(Ok(answer)); + }); + Ok(receiver) + } + } + } + + fn report(&self) -> Option { + Some(Report { + depth: self.depth, + rejected: self.rejected, + }) + } + + fn failure(&self) -> Option { + matches!(self.mode, Mode::Failed).then(|| "no such checkpoint".to_owned()) + } +} + +fn client() -> reqwest::Client { + reqwest::Client::builder() + .no_proxy() + .redirect(reqwest::redirect::Policy::none()) + .timeout(Duration::from_secs(10)) + .build() + .unwrap() +} + +/// Starts the native service and returns its base URL. +async fn start(fake: Arc, config: ServiceConfig) -> String { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = format!("http://{}", listener.local_addr().unwrap()); + let app = engine::app(fake, config); + tokio::spawn(async move { + engine::serve(listener, app, std::future::pending()) + .await + .unwrap(); + }); + url +} + +async fn post(url: &str, body: &str) -> reqwest::Response { + client() + .post(format!("{url}/v1/systemone")) + .body(body.to_owned()) + .send() + .await + .unwrap() +} + +/// The body a response carried. Asserting the whole JSON text rather than a parsed field +/// pins the wire format itself, and keeps the transport from depending on a parser it +/// deliberately does not have. +async fn body(response: reqwest::Response) -> String { + response.text().await.unwrap() +} + +#[tokio::test] +async fn readiness_and_failure_are_distinguishable() { + for (mode, status, expected) in [ + ( + Mode::Ready, + 200, + r#"{"status":"ok","depth":0,"rejected":0}"#, + ), + (Mode::Starting, 503, r#"{"status":"starting"}"#), + ( + Mode::Failed, + 503, + r#"{"status":"failed","reason":"no such checkpoint"}"#, + ), + ] { + let fake = Arc::new(Fake { + mode, + ..Default::default() + }); + let url = start(fake.clone(), ServiceConfig::default()).await; + + let health = client().get(format!("{url}/health")).send().await.unwrap(); + assert_eq!(health.status().as_u16(), status, "mode {mode:?}"); + // A ready engine describes its queue; a failed one says why. A starting one claims + // neither, because a half-loaded engine knows nothing worth reporting. + assert_eq!(body(health).await, expected, "mode {mode:?}"); + + // Inference is refused before the engine is asked, so an unready engine never sees + // the request at all. + let response = post(&url, REQUEST).await; + assert_eq!(response.status().as_u16(), status, "mode {mode:?}"); + if mode != Mode::Ready { + assert_eq!(fake.submitted.load(Ordering::Acquire), 0); + } + } +} + +#[tokio::test] +async fn a_valid_request_keeps_its_bytes() { + let fake = Arc::new(Fake::default()); + let url = start(fake, ServiceConfig::default()).await; + + let response = post(&url, REQUEST).await; + assert_eq!(response.status(), 200); + assert_eq!( + response.headers()["content-type"], + "application/json", + "the transport labels its own responses" + ); + assert_eq!(response.text().await.unwrap(), REQUEST); +} + +#[tokio::test] +async fn a_body_over_the_limit_is_refused_without_touching_the_engine() { + let fake = Arc::new(Fake::default()); + let url = start( + fake.clone(), + ServiceConfig { + max_body: 32, + ..Default::default() + }, + ) + .await; + + let response = post(&url, &"x".repeat(64)).await; + assert_eq!(response.status(), 413); + assert_eq!(fake.submitted.load(Ordering::Acquire), 0); +} + +#[tokio::test] +async fn a_full_queue_is_a_503_and_a_lost_slot() { + let fake = Arc::new(Fake { + mode: Mode::Full, + rejected: 7, + ..Default::default() + }); + let url = start(fake, ServiceConfig::default()).await; + + let response = post(&url, REQUEST).await; + assert_eq!(response.status(), 503); + assert_eq!(body(response).await, r#"{"error":"inference queue full"}"#); +} + +#[tokio::test] +async fn a_silent_engine_is_cut_off_by_the_deadline() { + let fake = Arc::new(Fake { + mode: Mode::Silent, + ..Default::default() + }); + let url = start( + fake, + ServiceConfig { + timeout: Duration::from_millis(150), + ..Default::default() + }, + ) + .await; + + let started = Instant::now(); + let response = post(&url, REQUEST).await; + assert_eq!(response.status(), 504); + assert!( + started.elapsed() < Duration::from_secs(5), + "the deadline, not the client, must end the wait" + ); + assert_eq!(body(response).await, r#"{"error":"inference timed out"}"#); +} + +#[tokio::test] +async fn engine_errors_map_to_one_status_each() { + for (mode, status, message) in [ + ( + Mode::Invalid, + 400, + "{\"error\":\"bad \\\"envelope\\\"\\n\"}", + ), + (Mode::Inference, 500, r#"{"error":"inference failed"}"#), + ] { + let fake = Arc::new(Fake { + mode, + ..Default::default() + }); + let url = start(fake, ServiceConfig::default()).await; + + let response = post(&url, REQUEST).await; + assert_eq!(response.status().as_u16(), status, "mode {mode:?}"); + // The engine's text is passed through as JSON, so it has to survive being quoted: + // the fake's message contains a quote and a newline on purpose. + assert_eq!(body(response).await, message); + } +} + +#[tokio::test] +async fn the_answer_carries_the_depth_the_engine_reported() { + // The engine, not the transport, measures the queue, so the header is whatever the + // answer carries. A depth of zero is included rather than treated as absent. + let fake = Arc::new(Fake { + depth: 3, + ..Default::default() + }); + let url = start(fake, ServiceConfig::default()).await; + + let response = post(&url, REQUEST).await; + assert_eq!(response.status(), 200); + assert_eq!( + response.headers()["x-queue-depth"], + "3", + "backpressure must be visible before it becomes a refusal" + ); +} + +/// The forwarding path is unchanged, but both modes now share one serve loop, so this pins +/// that the shared loop binds and keeps serving rather than returning early. +#[tokio::test] +async fn engine_serve_keeps_serving_until_shutdown() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let app = Router::new(); + let handle = tokio::spawn(async move { + engine::serve(listener, app, std::future::pending()) + .await + .unwrap(); + }); + tokio::time::sleep(Duration::from_millis(20)).await; + assert!(!handle.is_finished(), "serve returned without a shutdown"); + + // And it stops when the shutdown future resolves. + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + engine::serve(listener, Router::new(), async {}) + .await + .unwrap(); + handle.abort(); +}