runtime: add the worker-side engine contract - #34
xiaoyu-xyz wants to merge 1 commit into
Conversation
hsliuustc0106
left a comment
There was a problem hiding this comment.
Reviewed commit e8eec5d0aecdfb3efca50413a0cd4bbabae05a59.
No actionable findings. All 21 workspace tests passed. Reviewed the Engine boundary, readiness/status mapping, body limit, shared upload/inference deadline, response-byte preservation and error escaping.
Levius-Fubuki
left a comment
There was a problem hiding this comment.
Reviewed e8eec5d0aecdfb3efca50413a0cd4bbabae05a59 and the full Engine/HTTP-service diff.
No actionable findings. Reviewed readiness/failure reporting, request-body and shared deadline handling, status mapping, response-byte preservation, and JSON escaping. On Linux with Rust 1.98.1, formatting, strict all-target Clippy, release build, and all 21 workspace tests passed.
The branch still conflicts with current main, so this is a COMMENT rather than merge approval. Please resolve the conflicts and rerun the checks on the rebased head; this review does not certify conflict-resolution changes that do not exist yet. The separate lifecycle implementation in #35 is reviewed separately; this PR does not itself provide a production model engine.
cf8e73f to
5dd97f8
Compare
|
@Levius-Fubuki — rebased and rerun as you asked, details below. Rebased onto current The conflict was structural.
No source file changed in the resolution, and the diff is still this PR's five files. Rerun on
The counts move because On your "not a production model engine" point: agreed, and the README does not claim one — the in-repo engines are still described as the native workers that exist, and this PR only adds the interface a model would implement. The lifecycle behind that interface is #35. One thing for whoever merges: this PR and #35 are a stack. #35 is now rebased onto this branch, so they merge in order and the second diff stays incremental rather than repeating this one. |
|
lgtm @hsliuustc0106 👀 |
hsliuustc0106
left a comment
There was a problem hiding this comment.
Independent local review — in-process engine contract
Verdict: changes requested — placement conflicts with the merged architecture baseline.
The engineering is good — I ran it: 8/8 engine-service tests plus the 12 proxy tests pass at head 5dd97f8b, clippy clean; the deadline-carrying submit signature, three-state readiness with failure reasons, engine-measured x-queue-depth, and exact-response-text tests are all well designed.
But the design lands the request lifecycle in the wrong layer. It is written against RFC #14's original placement ("the request lifecycle in src/frontend/"); since then the repo re-baselined (#76) and implemented worker-side admission (#80 omni-runtime). docs/architecture.md now assigns the frontend "HTTP transport, forwarding, and response delivery" only, with queues/admission/back-pressure owned by the worker-side runtime — which is where this trait's shape belongs. Landing it as-is (and #35's worker-thread queue behind it) reintroduces exactly the boundary the baseline removed, and would give the crate two timeout semantics (proxy: 504; in-process budget: 503).
Proposed resolution: rehome the trait to src/runtime as the worker-side engine contract — its submit(body, deadline) -> Reply is a natural fit for what multi-request workers and the future shared runtime need — keeping only the health/status mapping and body-limit handling in the frontend. Also note the base predates the crate rename (omni-jev → frontend), so a rebase is due regardless.
Happy to re-review once the placement is decided; the trait itself is worth keeping.
Local reviewer report per the repo review skill; reflects head 5dd97f8b only.
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.
|
Superseded by #91. Closing rather than redirecting. Your finding was about placement, and it was right: the contract belongs to the worker-side runtime, and the frontend should keep transport. #91 does exactly that — Two things made a fresh branch the cheaper path than redirecting this one. The base here predates 102 commits of Also, one correction for the record: this branch's base does not predate a crate rename. |
5dd97f8 to
1363d47
Compare
|
Reopened with the placement revision as the head, so this thread keeps the review rather than starting a new one. What changed, against the three findings:
The earlier discussion in this thread was about the previous placement and does not describe this head. Two things in it were also wrong about the repository and I would rather correct the record than leave them: the base never carried a crate rename — Tests are under the repository On |
|
Rechecked |
|
A note on where the review state sits relative to the head, since the two no longer line up. The So nothing is outstanding from my side, and the blocking state on the PR describes a head that no longer exists. I am not asking for a status change — just recording the mapping so it is not read as "changes were requested and not addressed":
Rebase state: |
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.mdgives the frontend "HTTP transport, forwarding, and response delivery" and gives admission, queues and back-pressure to the runtime, andomni-runtimealready carries the admission half asSerialScheduler.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)Two required methods: say whether work may be sent, and accept it.
submitreturns 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.Readiness,Report,Answer,EngineError), because they are what the contract is made of.Readinessis three-state: a worker that failed keeps running and reports why, so an operator reads the reason instead of guessing from an exit code.Reportcarries 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/failurehave defaults, so the smallest useful worker implements two methods and nothing else.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.rsis 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.503is reserved for a worker that cannot take work at all: still loading, failed, out of capacity, or gone. The forwarding path answers504for 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_answersholds that line.Not in this PR
Reportis shaped to describe one, andtests/runtime/engine.rsexercises the contract with a worker that owns one, but the implementation stays where the runtime README puts it:SerialScheduleris admission with concurrency one, and queue-length limits are listed there as planned. This PR does not claim them.omni-jevbinary. No production engine exists yet to select.Test Plan
cargo fmt --all --check cargo clippy --workspace --locked --all-targets -- -D warnings cargo test --workspace --lockedPer
CONTRIBUTING, all test bodies are in the repositorytests/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 produces503and is not retried;InvalidRequest/InferenceFailedmap to400/500with the worker's text passed through as valid JSON; the reply's depth is published asx-queue-depth; a worker that never answers is cut off by the budget at504; and504and503do not blur.Test Result
System1-Omni Version / Commit:
7f39ac4(main) as base; head1363d47.cargo fmt --all --checkcleancargo clippy --workspace --locked --all-targets -- -D warningscleancargo test --workspace --locked— 98 passed, 0 failed, 10 ignored, five consecutive runs without a failure. The ignored tests came withmain; none are mine.cargo build --workspace --release --lockedpassesNo new dependency, and one new edge:
src/frontendnow depends onomni-runtime, which is the point of the change. The manifests otherwise only register the two test targets, andCargo.lockchanges 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.mdsays the runtime owns "admission, queues, batch budgets, compatibility grouping, request bookkeeping, and result routing". A trait whosesubmitcarries a deadline and whoseReportdescribes 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.