Request Engine#

class SubmitError : public std::runtime_error#

Thrown by submit() when the request is refused outright rather than queued.

Refusal is not an error condition of the engine: a bounded queue is the backpressure signal, and the caller is expected to translate this into whatever its transport calls “busy”.

Public Functions

inline explicit SubmitError(std::string const &message)#
class RequestHandle#

A caller’s claim on one in-flight request. Move-only: a channel has exactly one consumer, and copying the handle would quietly create a second one.

Public Functions

RequestHandle() = default#
RequestHandle(
std::shared_ptr<EngineChannel> channel,
std::shared_ptr<RequestRecord> record
) noexcept#
RequestHandle(RequestHandle&&) noexcept#
RequestHandle &operator=(RequestHandle&&) noexcept#
RequestHandle(RequestHandle const&) = delete#
RequestHandle &operator=(RequestHandle const&) = delete#
~RequestHandle() = default#
RequestId id() const noexcept#
StreamChannel &stream()#

The token channel for this request, as the runtime writes it. Valid while the handle is.

Tokens appear here as they are produced. To cancel, use RequestHandle::cancel() or RequestEngine::cancel(): cancelling the channel directly also stops generation at the next step boundary, but a get() waiting on this request is only released once the actor reports the outcome.

std::shared_ptr<StreamChannel> streamShared() const noexcept#

The same channel with shared ownership, for callers that outlive the handle.

The pybind layer holds StreamChannel through a shared_ptr holder, so handing it the reference above would manufacture a second owner of the same object. Anything else should prefer stream().

LLMGenerationResponse get()#

Block until the request reaches a terminal state, then hand back its response.

Throws:
  • std::runtime_error – if the request ended in kExecutionError

  • std::runtime_error – if the request was cancelled

bool ready() const noexcept#

True once the request has an outcome, so get() would return without blocking.

Distinct from the channel being finished, which reports only that the token stream ended.

void cancel() noexcept#

Ask the engine to stop this request. Non-blocking and idempotent; a request that has already retired simply does not match anything the actor still holds.

Releases this caller immediately, but does not interrupt a request the actor has already started: that request runs to completion and is then reported as cancelled.

inline bool valid() const noexcept#
class RequestEngine#

Owns the runtime and the one thread allowed to touch it.

submit() and cancel() are callable from any thread: each posts a command to the queue and returns without waiting for the actor. shutdown() is the exception — it joins, so it blocks and must not be called from the actor itself. The actor starts in the constructor and is joined by shutdown(), which the destructor implies, so there is no start() to forget.

What the actor serialises is feeding the GPU, which has to be serial anyway. Callers never block on each other, and no caller ever holds a pointer to the runtime, so a data race on it cannot be written rather than being prevented by a lock the caller has to remember to take.

Public Types

using SteppedFactory = std::function<std::unique_ptr<SteppedExecution>(LLMGenerationRequest&, cudaStream_t)>#

How the actor opens one founding request under the stepped control plane.

In production this is bound to LLMInferenceRuntime::beginStepped. It is a seam rather than a direct call so the actor’s admission, cancellation and shutdown paths can be tested without a GPU: LLMInferenceRuntime is a concrete class and needs a real CUDA stream to assemble, which would put every one of those tests behind a device. A null return means the runtime refused the request; a throw carries the reason.

Public Functions

RequestEngine(
std::unique_ptr<LLMInferenceRuntime> runtime,
cudaStream_t stream,
EngineConfig config = {}
)#
Parameters:
  • runtime – Taken over outright. Nothing else may hold or call it afterwards.

  • stream – The CUDA stream the actor drives the runtime on.

Throws:

std::runtime_error – if runtime spans more than one rank: the stepped control plane runs on a single rank in this release.

RequestEngine(SteppedFactory beginStepped, EngineConfig config)#

Test seam: an engine with no runtime, whose founding batches are whatever beginStepped hands back.

~RequestEngine()#
RequestEngine(RequestEngine const&) = delete#
RequestEngine &operator=(RequestEngine const&) = delete#
RequestEngine(RequestEngine&&) = delete#
RequestEngine &operator=(RequestEngine&&) = delete#
RequestHandle submit(LLMGenerationRequest request)#

Any thread: hand a request to the actor. Does not wait for it to start.

Throws:

SubmitError – if the engine is shutting down or the command queue is full

void cancel(RequestId id)#

Any thread: cancel a request by id. Same effect and same guarantees as RequestHandle::cancel(): the flag is planted immediately on the caller’s thread and nothing rides the command queue, so it cannot be lost under overload. Idempotent; unknown or retired ids are ignored (ids are never reused).

bool handleRequest(
LLMGenerationRequest const &request,
LLMGenerationResponse &response
) noexcept#

Submit one request and block until it finishes, in the shape of the old blocking call.

Exists so a caller written against LLMInferenceRuntime::handleRequest can move onto the engine without being restructured first: it reports failure by returning false rather than throwing, and leaves response untouched on failure.

This is a convenience over submit() plus RequestHandle::get(), not a second execution path. It gives up everything the engine exists to provide &#8212; the caller occupies a thread for the whole generation and cannot see tokens as they arrive &#8212; so new code should submit instead.

Returns:

false if the request was refused, cancelled, or failed during execution.

void shutdown(ShutdownMode mode)#

Any thread except the actor’s: stop the actor and join it. Idempotent.

Unlike the other operations this one blocks, because it joins. kDrain lets every queued request run; kCancel abandons what is still queued. Neither can interrupt a request already executing &#8212; the actor is inside it and cannot reach a boundary &#8212; so both wait at least that long. Both leave every live stream terminal, so no caller stays parked in waitPop().

The first caller’s mode wins; a later kCancel cannot escalate a drain already under way.

int32_t queued() const noexcept#

Requests the actor has taken off the queue but not yet admitted. Observability only.

Not a count of everything submitted: a request stays invisible here until the actor reaches a boundary and drains it, which it cannot do while it is inside a request. Do not use this to wait for a submission to become visible.

int32_t resident() const noexcept#

Requests currently holding runtime resources. Observability only.

bool running() const noexcept#

False once shutdown() has been called.

EngineMetrics metrics() const noexcept#

Counters for observability; see EngineMetrics for the consistency contract.

struct EngineConfig#

Public Members

int32_t commandQueueCapacity = {64}#

Commands the queue holds before submit() starts refusing. Rounded up to a power of two and to CommandQueue::kMinCapacity.

int32_t maxPendingRequests = {256}#

Requests accepted but not yet terminal, across the queue and the actor’s own backlog. Bounding the command queue alone does not bound this: each boundary moves a queue’s worth of commands into actor-owned state, so without this the backlog grows without limit and the backpressure the bounded queue is supposed to provide never materialises.

int32_t maxBatchSize = {1}#

Requests allowed to be resident at once.

Above 1, arriving requests join the running batch at step boundaries (in-flight batching), which needs the batch-aware executor; an engine built on the plain per-request executor has nowhere to admit into and rejects values above 1 at construction. This is also the KV capacity bound: page tables partition their pool per slot, so a request beyond this count has no pages to run on and waits in the queue &#8212; which is the V0 self-exit story, no eviction of others required.

std::chrono::milliseconds idleWait = {50}#

How long an idle actor sleeps before waking to look around again. It does not affect a busy engine, which never reaches the park at all.

Every operation that gives the actor something to do also wakes it, so this timeout is a backstop rather than the mechanism: it bounds the damage if a future caller changes actor state without waking it, turning a permanent hang into a delay of this length.

struct EngineMetrics#

A point-in-time copy of the engine’s counters. Each field is read with relaxed ordering, so the snapshot is per-field accurate but not a single atomic cut &#8212; fine for observability, not for control flow.

Public Members

uint64_t submitted = {}#

Requests submit() accepted / refused (queue full, backlog full, or shutting down).

uint64_t refused = {}#
uint64_t stallsIncompatible = {}#

Boundaries where the queue head could not join the running batch, by reason. High stall counts quantify the cost of FCFS head blocking under a mixed workload.

uint64_t stallsGuided = {}#
uint64_t stallsFounderOnly = {}#
uint64_t stallsNoCapacity = {}#
uint64_t admittedMidFlight = {}#

Requests that joined a running batch mid-flight (excludes founders).

uint64_t completed = {}#

Running totals of terminal outcomes across all requests, one per TerminalStatus value: countOutcome() bumps exactly one of them each time a request retires, so they always sum to the outcomes callers saw. A single request’s own state is the TerminalStatus on its record, not anything here.

uint64_t cancelled = {}#
uint64_t failed = {}#
uint64_t queueLatencyTotalUs = {}#

Time from submit() to the actor starting to serve the request (founding or admission).

uint64_t queueLatencyMaxUs = {}#
uint64_t queueLatencyCount = {}#