Skip to content

feat(frontend): own the engine in a worker thread with a bounded queue - #35

Draft
xiaoyu-xyz wants to merge 2 commits into
ThinkFlowLab:mainfrom
xiaoyu-xyz:frontend-engine-lifecycle
Draft

xiaoyu-xyz wants to merge 2 commits into
ThinkFlowLab:mainfrom
xiaoyu-xyz:frontend-engine-lifecycle

Conversation

@xiaoyu-xyz

@xiaoyu-xyz xiaoyu-xyz commented Sep 29, 2026 •

Copy link
Copy Markdown

Stacked on #34. This branch is built on the in-process engine contract and HTTP service, which is not merged yet, so the diff below still contains that PR's files (engine.rs, engine_service.rs). Only the eight files listed here are this change. Once #34 lands I will rebase and reopen against main.

Purpose

Stacked on #. That PR added the Engine contract and engine::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 to engine_service.rs.

Model code is deliberately absent. A model supplies load, warmup and one request function; everything below comes from the transport.

How the lifecycle is defined

Event Behaviour
Before load/warmup finishes /health is 503 starting; inference is refused without touching the engine.
Ready /health is 200 ok with queue depth and rejected.
Queue full 503 immediately. The queue refuses; it does not grow.
Budget expired (upload, queueing or inference) 503/504. The budget starts at the request headers, so a slow upload is charged to the caller.
Client gone before the worker dequeued it Dropped without reaching the engine. No GPU work for an answer nobody will read.
Engine failed an accepted request 500, the engine is retired, /health becomes 503 failed with the reason, and the process keeps running so the failure is readable. A malformed request retires nothing.
SIGTERM / SIGINT Admission closes at the signal, responses in flight get the same budget to finish, accepted work drains within it, then the process exits.
Drain budget exhausted The process exits and says so. A wedged kernel is not interruptible and the code does not pretend otherwise.

spawn returns only after loading and warmup, so a caller holding a handle cannot observe a half-loaded engine, and the first request after /health is 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:

  1. /health is observable. Add native Laya inference with Rust and CUDA #16's ready() is a bool, so an operator cannot tell a full queue from a dead engine. Here depth and a cumulative rejection count are published, x-queue-depth makes backpressure visible before it becomes a refusal, and a retired engine reports why without the server log.
  2. SIGTERM actually drains. Add native Laya inference with Rust and CUDA #16's main returns 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's engine.rs warns about: GPU allocations outliving the process.
  3. A total order between admission and shutdown. Without one lock spanning the readiness check and the send, a request can pass the check and then enqueue after the queue closed — accepted, with no worker left to answer it. submit holds that lock across both, and a test covers a drained worker refusing further work.
  4. Failure isolation is stated, not implied. A panic in model code becomes a 500 for that caller and retires the engine instead of unwinding the worker thread; bad input does not retire anything.
  5. Model code is one function, not its own worker loop. The queue, readiness and drain are written once, which is what [RFC] Support emerging System 1 models, with CLM as the second model after LAYA #9's second model needs.

Review follow-up

Both findings from the review are fixed, and each has a regression test that fails against
the code that was reviewed:

  • [P2] Queue depth is now reserved before the job is published (worker.rs). The slot is
    taken before try_send and rolled back if the send fails, because the worker does not take
    the 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_published submits 20,000 requests against a worker
    that 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 26 with 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.
  • [P2] The drain budget now starts when shutdown is signalled (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_drained takes that deadline instead of measuring a fresh budget of its
    own.
    shutdown_is_bounded_by_the_drain_budget_not_the_request_budget opens a request whose body
    never arrives and sends SIGTERM. Against the reviewed commit the process took 3.009 s
    to 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

  • CUDA Graph capture, RoPE, weights, the encoder and the scorer. The model half.
  • Metrics beyond queue depth. /metrics is a separate decision about an exposition format; Engine::report is the extension point it would use.
  • TLS, auth and ingress limits. Unchanged from the forwarding mode.

Test Plan

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

tests/worker_lifecycle.rs covers 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; head 1b63ed7, 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:

$ OMNI_SYSTEMONE_ENGINE=passthrough OMNI_JEV_BIND=127.0.0.1:18096 ./target/release/omni-jev
omni-jev listening on 127.0.0.1:18096 (engine passthrough, queue 32)

$ curl -s localhost:18096/health   → {"status":"ok","depth":0,"rejected":0}
$ curl -si -X POST --data '{"a":1}' localhost:18096/v1/systemone
                                   → 200, x-queue-depth: 1, body {"a":1}
$ kill -TERM <pid>                 → exit 0, within the drain budget

No GPU and no model weights were needed for any of it. Cargo.lock is unchanged.

Related: #14 (this is its HTTP integration step), #16 (closed reference implementation), #3, #9, #1.

@hsliuustc0106 hsliuustc0106 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@xiaoyu-xyz
xiaoyu-xyz force-pushed the frontend-engine-lifecycle branch from 1b63ed7 to 3219451 Compare September 30, 2026 00:55
@xiaoyu-xyz xiaoyu-xyz changed the title feat(frontend): own the engine in a worker thread with a bounded queue [stacked on #34] feat(frontend): own the engine in a worker thread with a bounded queue Sep 30, 2026
@xiaoyu-xyz

xiaoyu-xyz commented Sep 30, 2026 •

Copy link
Copy Markdown
Author

All three are fixed in a9f7fb83, rebased onto main (the merge conflict is resolved; it was the README sentence both sides had rewritten). Each fix has a test that fails against the code you reviewed.

[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 serve future meant a healthy process retired once the budget elapsed with no signal sent, so the default was a 10-second service lifetime. One serve call now runs for the life of the process; its shutdown future waits for the signal, closes admission, and starts the budget there. the_service_stays_up_beyond_the_drain_budget_without_a_signal serves past five times the budget and then answers a request; against 3219451 it fails with the process exited on its own 200 after starting.

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. mpsc::try_send only bounds the channel's buffer, and the worker frees its slot the moment it dequeues, so a capacity-1 engine served two requests at once. Admission now takes the slot and checks it under the admission lock, before the handover; both failed-send paths roll it back. a_running_request_counts_against_capacity holds the first inference behind a gate so the channel is provably empty, then asserts the second is refused; against 3219451 it fails with a second request was accepted while one was already running.

[P3] An unknown engine name is a configuration error. OMNI_SYSTEMONE_ENGINE is matched explicitly now; unset and empty still forward, anything else is a startup error. an_unknown_engine_name_is_a_startup_error runs the binary with typo-engine and asserts a non-zero exit naming the value; against 3219451 it fails because the process starts a forwarding server instead.

A correction to my previous comment

I claimed the depth test "fires within the first hundred submissions" against 1b63ed7. That was wrong, and I should not have written it. When I re-checked the test I had by then, it passed against the buggy ordering — the observation I quoted came from an earlier version of the test, before I changed the fake to spin, and I reported it as if it still held.

What is true now: the depth fix itself is in place, and queue_depth_returns_to_zero_after_the_storm covers the consequence — a counter released before it was reserved stays wrapped and never drains — but it does not reproduce your interleaving, and it passes against the released-before-reserved ordering too. The test's own doc comment says so. I could not make that window deterministic here: it is a few instructions wide, and forcing it would mean injecting a delay into submit. The ordering is covered by your probe and by inspection, not by my test.

Verified on a clean checkout of a9f7fb83, with the changed files hashed against the pushed branch to confirm they are identical: cargo fmt --all --check clean, cargo clippy --workspace --locked --all-targets -- -D warnings clean, cargo test --workspace --locked 34 passed across repeated runs, release build passes.

@Levius-Fubuki Levius-Fubuki left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread src/frontend/src/main.rs Outdated
// 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();

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[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.

Comment thread src/frontend/src/worker.rs Outdated
// 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);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[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.

Comment thread src/frontend/src/main.rs Outdated

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") {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[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.
@xiaoyu-xyz
xiaoyu-xyz force-pushed the frontend-engine-lifecycle branch from 3219451 to a9f7b83 Compare October 3, 2026 12:34

This branch has not been deployed

No deployments
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.

3 participants