Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ Documentation: <https://thinkflowlab.github.io/system1-omni/>

A community-maintained inference engine for prefill-only 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 has a native worker with CUDA kernels in this repository; other in-repository model engines and GPU backends are not implemented yet.
The Rust frontend serves `/v1/systemone` either by forwarding to a separately running model worker, as the Cua-S1 4B 0.2 `text` adapter does through its native worker and CUDA kernels, or — with a model linked in — in-process behind a small engine interface. Model engines and GPU backends beyond that are not implemented yet.

## Run the frontend

Expand Down
49 changes: 43 additions & 6 deletions src/frontend/README.md
Original file line number Diff line number Diff line change
@@ -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

Expand All @@ -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).
Expand All @@ -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:
Expand All @@ -47,6 +82,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.
287 changes: 287 additions & 0 deletions src/frontend/src/engine.rs
Original file line number Diff line number Diff line change
@@ -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<u8>,
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<Result<Answer, EngineError>>;

/// 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<u8>, deadline: Instant) -> Result<Reply, EngineError>;

/// What `/health` publishes while ready. `None` leaves those fields out.
fn report(&self) -> Option<Report> {
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<String> {
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<dyn Engine>,
config: ServiceConfig,
}

/// The native router: `POST /v1/systemone` and `GET /health`.
pub fn app(engine: Arc<dyn Engine>, 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<Output = ()> + Send + 'static,
) -> Result<(), BoxError> {
axum::serve(listener, router)
.with_graceful_shutdown(shutdown)
.await?;
Ok(())
}

async fn health(State(service): State<Service>) -> 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<Service>, 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
}
9 changes: 9 additions & 0 deletions src/frontend/src/lib.rs
Original file line number Diff line number Diff line change
@@ -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};

Expand Down
Loading