Skip to content

runtime: add the worker-side engine contract - #91

Closed
xiaoyu-xyz wants to merge 1 commit into
ThinkFlowLab:mainfrom
xiaoyu-xyz:runtime-engine-contract
Closed

xiaoyu-xyz wants to merge 1 commit into
ThinkFlowLab:mainfrom
xiaoyu-xyz:runtime-engine-contract

Conversation

@xiaoyu-xyz

@xiaoyu-xyz xiaoyu-xyz commented Oct 5, 2026 •

Copy link
Copy Markdown
Contributor

Purpose

Review of the in-process serving work asked for its placement to change: the request lifecycle belongs to the worker-side runtime, not to the frontend. docs/architecture.md gives the frontend "HTTP transport, forwarding, and response delivery" and gives admission, queues and back-pressure to the runtime, and omni-runtime already carries the admission half as SerialScheduler.

This adds the other half — the contract a model implements to be served — and the transport that serves it. It replaces the older placement rather than adding to it; the two branches that carried that placement are closed in favour of this one.

The contract, in src/runtime/src/engine.rs (140 lines)

pub trait Engine: Send + Sync + 'static {
    fn readiness(&self) -> Readiness;
    fn submit(&self, body: Vec<u8>, deadline: Instant) -> Reply;
    fn report(&self) -> Option<Report> { None }
    fn failure(&self) -> Option<String> { None }
}

Two required methods: say whether work may be sent, and accept it.

  • submit returns a reply, not a result. A refusal travels the same path as an answer, so a caller has one thing to await and can still tell "there was no room" from "the runtime dropped my receiver" — two channels would make those easy to confuse.
  • The types a worker must name live with the trait (Readiness, Report, Answer, EngineError), because they are what the contract is made of. Readiness is three-state: a worker that failed keeps running and reports why, so an operator reads the reason instead of guessing from an exit code.
  • Queue state is published by whoever owns the queue. Report carries depth, capacity and a monotonic rejected count; both are optional so a worker that measures neither reports nothing rather than a zero it did not measure.
  • report/failure have defaults, so the smallest useful worker implements two methods and nothing else.
  • Nothing parses a decision envelope: bytes in, bytes out.

The transport, in src/frontend/src/engine.rs (212 lines)

The two routes, the request budget measured from the request headers, the body limit, the status mapping, and the queue depth a reply publishes. The forwarding path in src/frontend/src/lib.rs is untouched apart from exporting the module.

One timeout semantic

A request that cannot be answered inside its budget is 504, whatever part of the budget ran out — the upload, the queue, or the inference. 503 is reserved for a worker that cannot take work at all: still loading, failed, out of capacity, or gone. The forwarding path answers 504 for a backend timeout as well, so a caller does not have to know which mode it is talking to in order to read a timeout. tests/frontend/in_process.rs::timeout_and_unavailability_are_different_answers holds that line.

Not in this PR

  • A queue with admission limits. Report is shaped to describe one, and tests/runtime/engine.rs exercises the contract with a worker that owns one, but the implementation stays where the runtime README puts it: SerialScheduler is admission with concurrency one, and queue-length limits are listed there as planned. This PR does not claim them.
  • The lifecycle around a queue — worker-thread ownership, readiness transitions, drain on shutdown. Same reason.
  • A model that implements the contract. Nothing here is a checkpoint reader or an executor; the only implementations in this PR are the two test workers. Wiring LAYA to it is a separate change.
  • A runnable in-process mode in the omni-jev binary. No production engine exists yet to select.

Test Plan

cargo fmt --all --check
cargo clippy --workspace --locked --all-targets -- -D warnings
cargo test --workspace --locked

Per CONTRIBUTING, all test bodies are in the repository tests/ tree, with explicit [[test]] paths registered in the manifest.

  • tests/runtime/engine.rs — the contract through its public API: an answer and a refusal both arrive on the reply; a worker without observability reports nothing and answers no work when asked how it is; a queued worker admits to capacity, reports the depth and capacity it is holding, refuses the next request without counting it, and admits again once a slot frees; a dropped reply does not disturb the queue's accounting.
  • tests/frontend/in_process.rs — the transport over real sockets: readiness decides whether work is accepted and an unready worker never sees the request; bytes and content type survive; a body over the limit is refused before the worker sees it; a busy worker produces 503 and is not retried; InvalidRequest/InferenceFailed map to 400/500 with the worker's text passed through as valid JSON; the reply's depth is published as x-queue-depth; a worker that never answers is cut off by the budget at 504; and 504 and 503 do not blur.

Test Result

System1-Omni Version / Commit: 7f39ac4 (main) as base; head 1363d47.

  • cargo fmt --all --check clean
  • cargo clippy --workspace --locked --all-targets -- -D warnings clean
  • cargo test --workspace --locked — 98 passed, 0 failed, 10 ignored, five consecutive runs without a failure. The ignored tests came with main; none are mine.
  • cargo build --workspace --release --locked passes
  • The existing forwarding suites are unchanged and still green

No new dependency, and one new edge: src/frontend now depends on omni-runtime, which is the point of the change. The manifests otherwise only register the two test targets, and Cargo.lock changes by a single line for that edge.

Not verified: Linux (this was macOS with Rust 1.98.1), and any model implementation, since none exists yet.

Review note

docs/architecture.md says the runtime owns "admission, queues, batch budgets, compatibility grouping, request bookkeeping, and result routing". A trait whose submit carries a deadline and whose Report describes a bounded queue is a small part of that, and it is the part a worker needs in order to be served without writing its own HTTP. If the intended home for the trait is narrower — say, admission only, with the frontend still owning any deadline arithmetic — that is a smaller change than this one and I would rather make it now than after a model depends on the shape.

Review of the in-process serving work asked for its placement to change: the
request lifecycle belongs to the worker-side runtime, not to the frontend.
`docs/architecture.md` gives the frontend HTTP transport, forwarding and response
delivery, and gives admission, queues and back-pressure to the runtime, and
`omni-runtime` already carries the admission half as `SerialScheduler`.

This adds the other half: the contract a model implements to be served. It is
two required methods — say whether work may be sent, and accept it — plus two
optional ones for observability:

    fn readiness(&self) -> Readiness;
    fn submit(&self, body: Vec<u8>, deadline: Instant) -> Reply;
    fn report(&self) -> Option<Report> { None }
    fn failure(&self) -> Option<String> { None }

`submit` returns a reply rather than a result, so a refusal travels the same path
as an answer. A caller then has one thing to await and can still tell "there was
no room" from "the runtime dropped my receiver", which two channels would make
easy to confuse. The types a worker has to name — `Readiness`, `Report`,
`Answer`, `EngineError` — live here with the trait, because they are what the
contract is made of. Queue state is published by the worker that owns the queue:
`Report` carries depth, capacity and a rejected count.

Nothing here parses a decision envelope. Request bytes go in and response bytes
come out, so a field the runtime has never heard of survives and a worker's error
text cannot produce invalid JSON.

The frontend gains the transport that serves such a worker: the two routes, the
request budget measured from the headers, the body limit, the status mapping and
the queue depth a reply publishes. Nothing else changes for the forwarding path.

One timeout semantic, which the same review asked for. A request that cannot be
answered inside its budget is `504` whatever part of the budget ran out — upload,
queue or inference. `503` is reserved for a worker that cannot take work at all:
still loading, failed, out of capacity, or gone. The forwarding path answers
504 for a backend timeout too, so a caller does not have to learn which mode it
is talking to in order to read one.

Tests are under the repository `tests/` tree per CONTRIBUTING: `tests/runtime/engine.rs`
exercises the contract through its public API with two workers (an inline one and
one that owns a queue), and `tests/frontend/in_process.rs` runs the transport over
real sockets for readiness, byte preservation, the body limit, refusal without
retry, error mapping, queue depth, and a worker that never answers.

Verified on this commit: fmt and clippy -D warnings clean over the workspace,
70 tests passing across five consecutive runs, release build passes, and the
existing forwarding suites are untouched and still green.
@xiaoyu-xyz

Copy link
Copy Markdown
Contributor Author

Closing as a duplicate: this change now lives in #34, reopened with the same commit 1363d47 so that the review already on that thread carries over instead of restarting here.

Nothing is lost by closing this one — the head is the same commit, and #34's body and note describe the revision.

@xiaoyu-xyz xiaoyu-xyz closed this Oct 5, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant