Skip to main content

LlmExecutor

Struct LlmExecutor 

Source
pub struct LlmExecutor { /* private fields */ }

Implementations§

Source§

impl LlmExecutor

Source

pub fn new(model: Box<dyn DecoderOnlyLLM>, info: ModelInfo) -> Self

Source

pub fn new_with_vnext_model( model: Box<dyn DecoderOnlyLLM>, info: ModelInfo, vnext_model: Arc<PreparedProductionModel>, ) -> Self

Source

pub fn vnext_model(&self) -> Option<&Arc<PreparedProductionModel>>

Source

pub fn vnext_family(&self) -> Option<&PreparedModelFamily>

Source

pub fn truncate_kv_for_cache_id(&self, cache_id: &str, new_len: usize)

Roll the KV cache for cache_id back to new_len positions. Used by speculative decoding on partial rejection. The caller must supply a GenericKvCacheHandle whose seq_len is also updated.

Trait Implementations§

Source§

impl ModelExecutor for LlmExecutor

Source§

fn batch_prefill<'life0, 'life1, 'async_trait>( &'life0 self, inputs: &'life1 [PrefillInput], ) -> Pin<Box<dyn Future<Output = Result<Vec<PrefillOutput>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Batched prefill: combine all prompts into ONE model.unified_forward call so launch / kernel-overhead is amortized across the cohort.

Falls back to the trait default (serial per-item) when the model returns Err(unsupported) from unified_forward — e.g. Qwen3MoeModel today, until Phase 2 adds its native unified path.

Source§

fn batch_decode<'life0, 'life1, 'async_trait>( &'life0 self, inputs: &'life1 [DecodeInput], ) -> Pin<Box<dyn Future<Output = Result<Vec<DecodeOutput>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Override default fallback to acquire the model lock ONCE for the whole batch, avoiding N round-trips through parking_lot. Does not yet do true attention batching (each cache has its own kv_len), but removes mutex churn that was serialising concurrent requests at async level.

Source§

fn unified_decode<'life0, 'life1, 'async_trait>( &'life0 self, batch: &'life1 UnifiedBatch, ) -> Pin<Box<dyn Future<Output = Result<Vec<Option<Vec<f32>>>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Unified mixed-batch dispatch (chunked-prefill API).

This impl is a behavior-preserving FALLBACK over the existing trait methods on DecoderOnlyLLM: prefill items go through model.prefill(seq_id, &q_tokens) (one at a time, mirroring the engine’s current sequential prefill loop), decode items (q_len == 1 && is_final_chunk) are grouped into a single model.decode_batch(...) call. Net behavior is identical to the engine’s pre-Phase-13 path; this just changes WHO orchestrates the prefill/decode split (caller → unified_decode) so the engine can converge on a single call.

The real performance unlock comes in Step 5 when models override this with a true unified-forward (one [M_total, hidden] forward

  • varlen attention) — at that point the kernel-level mix replaces the host-side serial dispatch here.
Source§

fn info(&self) -> &ModelInfo

Get model information and metadata
Source§

fn supports_native_unified_decode(&self) -> bool

Whether this executor’s backend can run the unified mixed prefill+decode forward natively. When false, the engine routes Qwen3-MoE batches through the legacy split path. Reported by the (backend-aware) executor so the engine stays backend-agnostic — replaces a cfg(target_os) branch that previously hard-coded “Metal/CPU lack native unified” in the hot path. Read more
Source§

fn kv_capacity(&self) -> Option<usize>

Per-request KV capacity in tokens when the executor owns a smaller runtime cache window than the model’s declared context length.
Source§

fn reserve_kv_slots( &self, requests: &[KvSlotRequest], ) -> Result<Option<KvSlotReservation>>

Reserve model-owned KV slots before a forward is dispatched. Read more
Source§

fn kv_slot_capacity_snapshot(&self) -> Option<KvSlotCapacitySnapshot>

Snapshot model-owned paged-KV capacity without allocating slots. Read more
Source§

fn recurrent_state_spec( &self, request_id: &RequestId, input_tokens: &[TokenId], ) -> Result<Option<RecurrentStateSpec>>

Recurrent-state allocation spec for this request, when the model has state-space or hybrid layers that need per-request recurrent state. Read more
Source§

fn prefill<'life0, 'life1, 'async_trait>( &'life0 self, input: &'life1 PrefillInput, ) -> Pin<Box<dyn Future<Output = Result<PrefillOutput>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Execute prefill phase (process initial prompt)
Source§

fn truncate_kv<'life0, 'life1, 'async_trait>( &'life0 self, kv_cache: &'life1 Arc<dyn KvCacheHandle>, new_len: usize, ) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Roll the KV cache for this executor’s sequence back to new_len. Used by speculative decoding on partial rejection so the next iteration sees a KV prefix that matches the accepted token stream. Default: Ok(()) — executors that don’t cache per-sequence state (stub, mock) are inherently tolerant; real LLM executors override.
Source§

fn forward_verify<'life0, 'life1, 'async_trait>( &'life0 self, inputs: &'life1 [DecodeInput], ) -> Pin<Box<dyn Future<Output = Result<Vec<DecodeOutput>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Multi-position decode-verify: one forward over N+1 tokens, producing one logits row per position. Used by speculative decoding’s target path so we don’t pay N+1 sequential forwards. Read more
Source§

fn decode<'life0, 'life1, 'async_trait>( &'life0 self, input: &'life1 DecodeInput, ) -> Pin<Box<dyn Future<Output = Result<DecodeOutput>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Execute decode phase (generate next token)
Source§

fn release_cache(&self, cache_id: &str)

Release KV cache and state without asserting successful completion. Read more
Source§

fn capabilities(&self) -> ExecutorCapabilities

Get executor capabilities
Source§

fn status(&self) -> ExecutorStatus

Get current executor status
Source§

fn cache_metrics_snapshot(&self) -> Option<Value>

Optional model/executor cache metrics. Read more
Source§

fn lora_metrics_snapshot(&self) -> Option<Value>

Optional LoRA runtime metrics.
Source§

fn execution_resource_authority(&self) -> ExecutionResourceAuthority

Selects the single authority for request-lifetime model resources. Existing executors remain on the transitional legacy-engine path by default. A runtime that returns PlanRuntime must return the shared runtime’s opaque cache handle from prefill/decode and delegate release of that authority from release_cache.
Source§

fn admission_limits( &self, ) -> Result<Option<ExecutorAdmissionLimits>, FerrumError>

Immutable admission limits compiled into this executor’s runtime plan. A PlanRuntime executor must return Some; legacy executors may defer to the engine-owned scheduler and recurrent-state limits.
Source§

fn resolved_model_plan(&self) -> Option<&ResolvedModelPlan>

Returns the immutable product plan that owns planning, provider selection, and resource authority for a plan-runtime executor. Legacy executors return None; PlanRuntime executors must expose the exact plan used for provisioning and dispatch.
Source§

fn plan_runtime_resource_snapshot( &self, ) -> Result<Option<PlanRuntimeResourceSnapshot>, FerrumError>

Returns the authoritative memory breakdown for a shared plan runtime. LegacyEngine executors return None; PlanRuntime executors must return Some while they are ready.
Source§

fn attach_execution_event_sink(&self, _sink: Arc<dyn ExecutionEventSink>)

Installs the product-owned execution event sink before requests start. Read more
Source§

fn execution_capacity_epochs( &self, ) -> Result<Option<ExecutorAdmissionEpochs>, FerrumError>

Current plan-local capacity evidence for scheduler wake suppression. Legacy-engine executors return None; an executor declaring ExecutionResourceAuthority::PlanRuntime must return Some.
Source§

fn write_execution_capacity_snapshot( &self, availability: &mut Vec<CapacityAvailabilityEpoch>, ) -> Result<Option<ExecutorAdmissionEpochs>, FerrumError>

Writes the canonical per-source availability generations into caller-owned storage and returns the matching global audit epochs. Executors with typed dynamic admission override this to avoid allocating on steady scheduler ticks.
Source§

fn register_execution_capacity_waiter( &self, _observed: &CapacityWaitCondition, ) -> Result<Option<ExecutorCapacityWaitRegistration>, FerrumError>

Synchronously subscribes to every source named by one passive capacity wait. The returned registration must remain alive until it is awaited or deliberately cancelled by being dropped. Read more
Source§

fn try_admit_prefill( &self, _input: ExecutorPrefillAdmission<'_>, ) -> Result<ExecutorPrefillAdmissionDecision, FerrumError>

Probe and retain the exact request/sequence authority needed by a future prefill. No provider encode, kernel launch, or device submit may occur in this method.
Source§

fn cancel_prefill_admission(&self, _request_id: &RequestId) -> bool

Release an admitted but not yet active prefill authority. Read more
Source§

fn write_execution_capacity_release_sources( &self, _preemption: &ExecutorExecutionCapacityPreemption, sources: &mut Vec<CapacityAvailabilitySource>, ) -> Result<bool, FerrumError>

Writes the exact availability sources advanced when this request authority is preempted for recompute. Read more
Source§

fn preempt_execution_capacity<'life0, 'async_trait>( &'life0 self, _preemption: ExecutorExecutionCapacityPreemption, ) -> Pin<Box<dyn Future<Output = Result<ExecutorExecutionCapacityPreemptionReceipt, FerrumError>> + Send + 'async_trait>>
where 'life0: 'async_trait, Self: 'async_trait,

Retire one exact request-scoped runtime authority for recompute. Read more
Source§

fn maintain_prefill_backing( &self, _request_id: &RequestId, ) -> Result<ExecutorPrefillMaintenanceOutcome, FerrumError>

Consume one retained logical/physical backing deferral after the scheduler waiting lock has been released. Implementations must perform at most one bounded maintenance attempt and publish capacity epochs only after real backing is installed.
Source§

fn prefill_with_capacity<'life0, 'life1, 'async_trait>( &'life0 self, input: &'life1 PrefillInput, ) -> Pin<Box<dyn Future<Output = Result<ExecutorPrefillOutcome, FerrumError>> + Send + 'async_trait>>
where 'life0: 'async_trait, 'life1: 'async_trait, Self: 'async_trait,

Execute one exact prefill chunk with an explicit pre-submit capacity deferral edge. Legacy executors inherit full-prefill behavior.
Source§

fn batch_prefill_with_capacity<'life0, 'life1, 'async_trait>( &'life0 self, _inputs: &'life1 [PrefillInput], ) -> Pin<Box<dyn Future<Output = Result<ExecutorBatchPrefillOutcome, FerrumError>> + Send + 'async_trait>>
where 'life0: 'async_trait, 'life1: 'async_trait, Self: 'async_trait,

Attempt one physical, capacity-aware prefill batch. Read more
Source§

fn plan_runtime_prefill_with_capacity<'life0, 'life1, 'async_trait>( &'life0 self, _input: &'life1 PlanRuntimePrefillInput, ) -> Pin<Box<dyn Future<Output = Result<PlanRuntimePrefillOutcome, FerrumError>> + Send + 'async_trait>>
where 'life0: 'async_trait, 'life1: 'async_trait, Self: 'async_trait,

Tensor-free prefill for executors that declare ExecutionResourceAuthority::PlanRuntime. Read more
Source§

fn plan_runtime_batch_prefill_with_capacity<'life0, 'life1, 'async_trait>( &'life0 self, _inputs: &'life1 [PlanRuntimePrefillInput], ) -> Pin<Box<dyn Future<Output = Result<PlanRuntimeBatchPrefillOutcome, FerrumError>> + Send + 'async_trait>>
where 'life0: 'async_trait, 'life1: 'async_trait, Self: 'async_trait,

Attempt one physical tensor-free prefill batch. Read more
Source§

fn discard_plan_runtime_prefill( &self, authority: PlanRuntimePrefillAuthority, ) -> Result<(), FerrumError>

Discard an exact prefill authority after engine-side validation, sampling, or scheduler commit fails. Read more
Source§

fn batch_decode_with_capacity<'life0, 'life1, 'async_trait>( &'life0 self, inputs: &'life1 [DecodeInput], ) -> Pin<Box<dyn Future<Output = Result<ExecutorBatchDecodeOutcome, FerrumError>> + Send + 'async_trait>>
where 'life0: 'async_trait, 'life1: 'async_trait, Self: 'async_trait,

Batch decode with an explicit pre-submit capacity deferral edge. Read more
Source§

fn plan_runtime_batch_decode_with_capacity<'life0, 'life1, 'async_trait>( &'life0 self, _inputs: &'life1 [PlanRuntimeDecodeInput], ) -> Pin<Box<dyn Future<Output = Result<PlanRuntimeBatchDecodeOutcome, FerrumError>> + Send + 'async_trait>>
where 'life0: 'async_trait, 'life1: 'async_trait, Self: 'async_trait,

Tensor-free batch decode for executors that declare ExecutionResourceAuthority::PlanRuntime. Read more
Source§

fn forward<'life0, 'life1, 'async_trait>( &'life0 self, _input: &'life1 Arc<dyn TensorLike>, ) -> Pin<Box<dyn Future<Output = Result<Arc<dyn TensorLike>, FerrumError>> + Send + 'async_trait>>
where 'life0: 'async_trait, 'life1: 'async_trait, Self: 'async_trait,

Optional: full forward pass (for non-autoregressive use cases)
Source§

fn prepare_startup<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<(), FerrumError>> + Send + 'async_trait>>
where 'life0: 'async_trait, Self: 'async_trait,

Complete executor-owned startup preparation before either product entrypoint can accept a request. Read more
Source§

fn warmup<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), FerrumError>> + Send + 'async_trait>>
where 'life0: 'async_trait, Self: 'async_trait,

Warm up executor (load model, allocate memory, etc.)
Source§

fn shutdown<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), FerrumError>> + Send + 'async_trait>>
where 'life0: 'async_trait, Self: 'async_trait,

Shutdown executor gracefully
Source§

fn complete_cache( &self, completion: ExecutorSequenceCompletion, ) -> Result<(), FerrumError>

Complete and release one cache authority with product-authoritative terminal token counts. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> ErasedDestructor for T
where T: 'static,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<F, T> IntoSample<T> for F
where T: FromSample<F>,

Source§

fn into_sample(self) -> T

Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more