feat(frontend): own the engine in a worker thread with a bounded queue - #35
xiaoyu-xyz wants to merge 2 commits into
Conversation
hsliuustc0106
left a comment
There was a problem hiding this comment.
Reviewed commit 1b63ed7431ea4107852600d07799a826773e3bf4.
- [P2] Reserve queue depth before publishing the job. src/frontend/src/worker.rs:217 increments depth after try_send, but the consumer never takes the admission mutex. It can dequeue/process/decrement first, yielding depth=0 in a response and transient usize underflow in health. A standalone CPU probe against the built crate observed four zero-depth answers in 100,000 sequential immediate requests. Reserve before sending and roll back on send failure, or synchronize dequeue/accounting.
- [P2] Start the drain budget when shutdown is signaled. src/frontend/src/main.rs:68 awaits graceful HTTP shutdown before stop_accepting/run_until_drained. In-flight upload/response handlers can consume their entire request budget before the drain timer even begins, and uploads can still submit during this interval. A real binary with DRAIN_MS=50 and TIMEOUT_MS=3000 took 2.917 seconds to exit after SIGTERM with one partial upload. Stop admission on the signal and bound the complete shutdown sequence with the drain deadline.
All 31 ordinary tests plus 1 doctest passed; these probes expose gaps in the existing suite.
1b63ed7 to
3219451
Compare
|
All three are fixed in [P1] The budget no longer starts at startup. You were right, and it was worse than the finding says: arming the timer around the whole Two related things the fix had to get right, both of which I got wrong first: draining the worker and draining the transport now run concurrently (sequentially, an idle engine paid the full budget on every shutdown), and an unknown engine name is a startup error rather than a silent fall back to forwarding — that is [P2] below. [P2] The limit covers outstanding requests, running or queued. [P3] An unknown engine name is a configuration error. A correction to my previous commentI claimed the depth test "fires within the first hundred submissions" against What is true now: the depth fix itself is in place, and Verified on a clean checkout of |
Levius-Fubuki
left a comment
There was a problem hiding this comment.
Reviewed 3219451e57ab7ad96fc92ae47d4636fe404f9c7b, including the incremental changes on #34 and the two previous findings.
The depth reservation now precedes publication and rolls back on both failed-send paths, addressing the original counter-underflow mechanism. However, the shutdown follow-up introduces a reproducible service-lifetime regression. Changes are requested for the three issues below.
Formatting passed. Clippy, workspace tests and release build passed on Linux/Rust 1.98.1 after rebuilding in a dedicated target directory (an initial shared-target run had stale cross-checkout artifacts and is not treated as a PR failure). Independent probes reproduce the no-signal exit and invalid-mode fallback on Linux and macOS; a gated CPU worker reproduces the capacity breach. No GPU is needed for these cases. This remains a draft and also needs its existing merge conflicts resolved.
| // Responses in flight get the drain budget to finish. A handler still stuck on a slow | ||
| // upload at the deadline is not going to finish, and holding the process open for it | ||
| // would mean the budget bounded nothing. | ||
| let deadline = Instant::now() + running.drain_budget(); |
There was a problem hiding this comment.
[P1] Start the drain deadline on the shutdown signal, not at startup
This deadline is created as soon as the listener starts and wraps the entire engine::serve future. Consequently every native process exits after the drain duration even if no SIGTERM/SIGINT is sent (10 seconds by default). With OMNI_SYSTEMONE_ENGINE=passthrough, OMNI_SYSTEMONE_DRAIN_MS=200 and OMNI_JEV_BIND=127.0.0.1:0, the release binary exited 0 after 0.206 seconds on Linux, logging shutdown exceeded its drain budget; the probe sent no signal. Wait for the signal while serving normally, then close admission and compute one absolute deadline for graceful HTTP shutdown plus worker drain. Add a test that remains healthy beyond the drain duration before sending any signal.
| // answered and released — counting afterwards would release a slot that had not | ||
| // been taken yet, which reads as a wrong depth and wraps the counter. The window is | ||
| // small but real, and `depth_is_reserved_before_the_job_is_published` caught it. | ||
| self.depth.fetch_add(1, Ordering::AcqRel); |
There was a problem hiding this comment.
[P2] Enforce the documented limit on all outstanding requests
This increment reserves accounting but never checks the configured outstanding-request limit; only mpsc::try_send enforces capacity. Once the worker dequeues a job, the channel slot is free even while inference still runs. A deterministic probe with capacity 1 holds the first inference behind a gate, waits until it has started, then submits a second job: it is accepted and report().depth becomes 2. Both Options::queue_capacity and the README promise refusal once that many accepted requests are unfinished. Check/reserve the outstanding count under the admission lock (with rollback on send failure), and cover a running-plus-queued case instead of relying on whether the consumer has happened to dequeue.
|
|
||
| async fn run() -> Result<(), BoxError> { | ||
| let bind = setting("OMNI_JEV_BIND", Config::DEFAULT_BIND.to_owned())?; | ||
| if env::var("OMNI_SYSTEMONE_ENGINE").as_deref() != Ok("passthrough") { |
There was a problem hiding this comment.
[P2] Reject unsupported engine names instead of silently selecting forwarding
Every value other than passthrough enters forwarding mode here. With OMNI_SYSTEMONE_ENGINE=typo-engine, the built binary starts and logs forwarding to http://127.0.0.1:8000/, although the README says unknown modes are startup errors. A configuration typo can therefore route inference to a different backend instead of failing visibly. Match supported values explicitly, preserve the documented unset/default behavior, and return a configuration error for unknown values.
RFC ThinkFlowLab#14 puts the request lifecycle in src/frontend/ and the model behind "a small engine interface", and #3's acceptance criteria depend on that interface existing. Neither half is on main today. This adds the interface and the HTTP service behind it: the Engine trait, its readiness/answer/error types, and engine::app serving /v1/systemone and /health with a request budget measured from the headers, a body limit, and one status code per failure mode. Nothing parses the decision envelope, so a field this crate has never heard of survives and a model's error text cannot produce invalid JSON. submit takes owned bytes and returns a channel, which is the only shape that works for a worker thread with a queue behind it; the reply carries the queue depth to publish, so x-queue-depth is measured by the engine rather than guessed by the transport. No model, no GPU, and no change to main.rs or the forwarding path: the thread, the bounded queue and the drain that plug into this boundary come next.
The previous change defined the Engine contract and the HTTP service behind it. This makes them usable, and is where #3's acceptance criteria for persistent serving live: readiness that follows loading and warmup, a queue that refuses rather than grows, a request budget that can expire before work starts, and a shutdown that drains accepted work. - worker.rs: the thread that owns the engine, admission, readiness, the budget and drain. spawn returns only after load and warmup, so a handle cannot observe a half-loaded engine. A panic in model code answers its caller and retires the engine instead of unwinding the worker; bad input retires nothing. - main.rs: mode selection and shutdown wiring. The forwarding path is unchanged. - passthrough.rs: a stand-in engine that echoes requests, so the native path can be run and tested with no model and no GPU. Rebased onto main, which had moved on; the only conflict was the README sentence both sides had rewritten, and it now carries both facts. Three findings from review, each with a test that fails against the code that was reviewed: - Queue depth is reserved before the job is published, and rolled back when the send fails. The worker does not take the admission lock, so a job in the queue can be dequeued, answered and released before the producer counted it. - The configured limit now covers outstanding requests, running or queued, not only the channel's buffer. The worker frees its channel slot on dequeue, so a capacity-1 engine was serving two requests at once while the README promised a refusal for the second. - The drain budget starts at the signal and one deadline bounds the whole sequence. Arming it at startup retired a healthy process once it elapsed, and starting it only after the listener stopped let a request whose body never arrived hold the process open for its entire request budget. - An unknown OMNI_SYSTEMONE_ENGINE value is a startup error instead of a silent fall back to forwarding, which would send inference to a different backend. Shutdown drains the worker and the transport concurrently and exits when the worker is done plus a short bounded grace period, rather than paying the whole budget on every shutdown. Verified on a clean checkout of this commit: fmt and clippy -D warnings clean, 34 tests pass repeatedly, release build passes, and the built binary serves, refuses over-size bodies, reports readiness and exits 0 on SIGTERM inside its budget. Cargo.lock unchanged.
3219451 to
a9f7b83
Compare
Purpose
Stacked on #. That PR added the
Enginecontract andengine::app; this one makes them usable, and is where #3's acceptance criteria for persistent serving actually live: readiness that follows loading and warmup, a bounded queue, a request budget, and a shutdown that drains accepted work.src/frontend/src/worker.rs(450 lines) — the thread that owns the engine, the bounded queue, readiness, the request budget and drain.src/frontend/src/main.rs(+107) — mode selection and shutdown wiring. The forwarding path is unchanged.src/frontend/src/passthrough.rs(20 lines) — a stand-in engine that echoes requests, so the whole path can be run and tested with no model and no GPU.src/frontend/tests/worker_lifecycle.rs(449 lines) and the binary test added toengine_service.rs.Model code is deliberately absent. A model supplies
load,warmupand one request function; everything below comes from the transport.How the lifecycle is defined
/healthis503 starting; inference is refused without touching the engine./healthis200 okwith queuedepthandrejected.503immediately. The queue refuses; it does not grow.503/504. The budget starts at the request headers, so a slow upload is charged to the caller.500, the engine is retired,/healthbecomes503 failedwith the reason, and the process keeps running so the failure is readable. A malformed request retires nothing.SIGTERM/SIGINTspawnreturns only after loading and warmup, so a caller holding a handle cannot observe a half-loaded engine, and the first request after/healthis ready is not the one that pays for compilation.Deliberate differences from #16
#16 already prototyped an in-process path (
engine.rs,serve.rs) and was closed. This was written against the RFC text rather than derived from it, and the differences are the reason it is worth landing:/healthis observable. Add native Laya inference with Rust and CUDA #16'sready()is abool, so an operator cannot tell a full queue from a dead engine. Here depth and a cumulative rejection count are published,x-queue-depthmakes backpressure visible before it becomes a refusal, and a retired engine reports why without the server log.SIGTERMactually drains. Add native Laya inference with Rust and CUDA #16'smainreturns on the signal without joining the engine thread, while its README claims to drain accepted work. Here shutdown takes the queue sender under the admission lock, waits out the drain budget, and reports whether it finished. That is the case the doc comment at the top of Add native Laya inference with Rust and CUDA #16'sengine.rswarns about: GPU allocations outliving the process.submitholds that lock across both, and a test covers a drained worker refusing further work.500for that caller and retires the engine instead of unwinding the worker thread; bad input does not retire anything.Review follow-up
Both findings from the review are fixed, and each has a regression test that fails against
the code that was reviewed:
worker.rs). The slot istaken before
try_sendand rolled back if the send fails, because the worker does not takethe admission lock: the moment a job is in the queue it can be dequeued, answered and
released, so counting afterwards could release a slot that had not been taken. Both
rollback paths (
Full,Closed) restore the count.depth_is_reserved_before_the_job_is_publishedsubmits 20,000 requests against a workerthat is still draining and asserts that the depth never exceeds the number of requests the
caller has accepted. Against the reviewed commit that assertion fires within the first
hundred submissions (observed
depth 26with 27 accepted); against the fix it holds.The test asserts the one direction that is a hard invariant — a slot cannot be counted that
was never taken — rather than equality, because a job may legitimately complete between the
send and the read. I could not make the specific underflow branch deterministic: the
window is a few instructions wide and forcing it would mean injecting a delay into
submit,which is worse than the thing it tests.
main.rs,worker.rs).Admission closes inside the shutdown future itself, and one absolute deadline — computed
from the signal — bounds every step: the graceful listener shutdown, and the drain.
Running::run_until_drainedtakes that deadline instead of measuring a fresh budget of itsown.
shutdown_is_bounded_by_the_drain_budget_not_the_request_budgetopens a request whose bodynever arrives and sends
SIGTERM. Against the reviewed commit the process took 3.009 sto exit with a 200 ms drain budget and a 3000 ms request budget — the request budget, not
the drain budget, was bounding the shutdown. Against the fix it exits in well under a
second.
Not in this PR
/metricsis a separate decision about an exposition format;Engine::reportis the extension point it would use.Test Plan
cargo fmt --all --check cargo clippy --workspace --locked --all-targets -- -D warnings cargo test --workspace --lockedtests/worker_lifecycle.rscovers readiness following warmup, startup failure, queue refusal with depth accounting, cancelled and expired requests never reaching the engine, an inference failure retiring the worker, panic isolation, drain reporting that it cannot finish, a drained worker refusing work, a zero queue capacity rejected at startup, and the depth reservation above.Those tests synchronise on a gate the fake inference holds and a handshake signalling that it has started, so they assert on queue state rather than sleeping and hoping.
Test Result
System1-Omni Version / Commit:
3062243(main) as base; head1b63ed7, stacked on the in-process engine contract PR.cargo fmt --all --check: clean.cargo clippy --workspace --locked --all-targets -- -D warnings: clean.cargo test --workspace --locked: 34 passed, 0 failed, run six times without a flake.Verified against the built binary, not only in-process:
No GPU and no model weights were needed for any of it.
Cargo.lockis unchanged.Related: #14 (this is its HTTP integration step), #16 (closed reference implementation), #3, #9, #1.