Skip to main content

ferrum_engine/
continuous_engine.rs

1//! Continuous Batching Engine
2//!
3//! Iteration-level continuous batching: each step processes a mixed batch of
4//! prefill and decode requests selected by the scheduler.  Multiple callers
5//! can submit requests concurrently. An `iteration_lock` makes publication
6//! atomic with planning, while the single background driver owns device-wave
7//! execution.
8
9use crate::resource_lifecycle::{
10    ResourceLedgerTransition, ResourceLifecycleLedger, ResourceOwnerCloseSummary,
11};
12use async_trait::async_trait;
13use ferrum_bench_core::{
14    global_profile, profile_fields_from_json, JsonlJournal, JsonlJournalError,
15};
16use ferrum_interfaces::{
17    engine::{InferenceEngine, LlmInferenceEngine},
18    kv_cache::AllocationRequest,
19    model_executor::{
20        ExecutionResourceAuthority, ExecutorAdmissionEpochs, ExecutorCapacityWaitRegistration,
21        ExecutorExecutionCapacityDeferral, ExecutorExecutionCapacityPreemption,
22        ExecutorExecutionCapacityStage, ExecutorExecutionDeferral, ExecutorPrefillAdmission,
23        ExecutorPrefillAdmissionDecision, ExecutorPrefillAdmissionReceipt,
24        ExecutorPrefillMaintenanceDeferral, ExecutorPrefillMaintenanceOutcome,
25        ExecutorRequestOrigin, ExecutorRequestStateDeferral, ExecutorSamplingOutput,
26        ExecutorSequenceCompletion, GreedyRepetitionPenalty, KvSlotRequest, LogitsReturnPolicy,
27        PlanRuntimeBatchDecodeOutcome, PlanRuntimeBatchPrefillOutcome, PlanRuntimeDecodeInput,
28        PlanRuntimePrefillInput, PlanRuntimePrefillOutcome, PlanRuntimePrefillProduct,
29        TokenSelectionMask,
30    },
31    sampler::{SamplingConfig as TokenSamplingPlan, SamplingRng},
32    vnext::{
33        AdmissionDeferred, AdmissionRejected, BoundDeviceSubmissionAttribution,
34        BoundExecutionResourceMaintenance, CapacityAvailabilityEpoch, DeferredAction,
35        DeviceCapacityPressureScope, DeviceExecutionSpanKind, DeviceExecutionSpanMeasurement,
36        DeviceSubmissionExecutionSpan, DeviceSubmissionExecutionTiming, DeviceTimingMeasurement,
37        DeviceTimingUnavailableReason, EventBatchEmissionPermit, EventEmissionPermit,
38        ExecutionEvent, ExecutionEventCapturePolicy, ExecutionEventDetail,
39        ExecutionEventKind as VNextExecutionEventKind, ExecutionEventSink, ExecutionEventSinkError,
40        OperationCompletionReceipt,
41    },
42    KvCacheHandle, KvCacheManager, ModelExecutor, RecurrentStateHandle, RecurrentStateManager,
43    Sampler, SchedulerInterface as Scheduler, TensorFactory, TensorRef, Tokenizer,
44};
45use ferrum_kv::cache::prefix::PrefixCache;
46use ferrum_sampler::structured_output::{StructuredOutputFactory, StructuredOutputProcessor};
47use ferrum_scheduler::implementations::{
48    ContinuousBatchScheduler, ExecutionCapacityAction, ExecutionCapacityReleaseSnapshot,
49    ExecutionReadinessWake, ExecutorAdmissionProbeOutcome, ExecutorAdmissionQueueObservation,
50    PressureYieldTransaction, RequestPhase,
51};
52#[cfg(test)]
53use ferrum_scheduler::implementations::{PressureTransitionKind, PressureYieldKind};
54use ferrum_scheduler::vnext::{
55    AdmissionDeferral, AdmissionProbeOutcome, AdmissionWakeEpochs, AdmissionWakeSnapshot,
56};
57use ferrum_types::{
58    DataType, Device, EngineConfig, EngineDecodeStage, EngineDecodeStageInterval, EngineStatus,
59    EngineTokenTimingEvidence, FerrumError, FerrumProfileEvent, FinishReason,
60    InferenceExecutionEvidence, InferenceRequest, InferenceResponse, ObservabilityProfileDetail,
61    Priority, ProfileEntrypoint, ProfileError, ProfileEventKind, ProfileStatus, RequestId,
62    ResourceAction, ResourceTraceEvent, ResponseCompletionBoundary, Result, SamplingParams,
63    StreamChunk, TokenId, TokenUsage, DEFAULT_MAX_TOKENS_METADATA_KEY,
64    ENGINE_RUNTIME_TRACE_PRESET_HASH, OBSERVABILITY_PROFILE_SCHEMA_VERSION,
65    PROMPT_TOKENS_METADATA_KEY,
66};
67use futures::{stream::Stream, FutureExt};
68use metrics::{counter, gauge, histogram};
69use parking_lot::{Mutex, RwLock};
70use serde::Serialize;
71use sha2::{Digest, Sha256};
72use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
73use std::future::Future;
74use std::io::Write;
75use std::path::{Path, PathBuf};
76use std::pin::Pin;
77use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
78use std::sync::OnceLock;
79use std::sync::{Arc, Weak};
80use std::task::{Context, Poll};
81use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
82use tokio::sync::{mpsc, Notify, Semaphore};
83use tracing::{debug, info, warn};
84
85// Env-name constants + `from_env_vars` are retained as test-only parse
86// helpers: production resolves these knobs via EngineConfig.runtime
87// (apply_runtime_config_snapshot), not env. The unit tests still exercise the
88// env-name → field mapping.
89#[cfg(test)]
90const BATCH_DECODE_PROF_ENV: &str = "FERRUM_BATCH_DECODE_PROF";
91#[cfg(test)]
92const CHUNKED_PREFILL_ENV: &str = "FERRUM_CHUNKED_PREFILL";
93#[cfg(test)]
94const KV_CAPACITY_ENV: &str = "FERRUM_KV_CAPACITY";
95#[cfg(test)]
96const MAX_MODEL_LEN_ENV: &str = "FERRUM_MAX_MODEL_LEN";
97#[cfg(test)]
98const NEXT_BATCH_PROF_ENV: &str = "FERRUM_NEXT_BATCH_PROF";
99#[cfg(test)]
100const WHOLE_PROMPT_PREFIX_CACHE_ENV: &str = "FERRUM_WHOLE_PROMPT_PREFIX_CACHE";
101#[cfg(test)]
102const RBD_PROF_ENV: &str = "FERRUM_RBD_PROF";
103#[cfg(test)]
104const UNIFIED_POST_PROF_ENV: &str = "FERRUM_UNIFIED_POST_PROF";
105const GENERATION_POLICY_SCAN_LIMIT: usize = 262_144;
106const FORBIDDEN_DECODE_RESAMPLE_LIMIT: usize = 64;
107const TOKEN_TRACE_PROMPT_PREFIX_LIMIT: usize = 64;
108const TOKEN_TRACE_PROMPT_TAIL_LIMIT: usize = 128;
109const TOKEN_TRACE_GENERATED_PREFIX_LIMIT: usize = 256;
110const TOKEN_TRACE_GENERATED_TAIL_LIMIT: usize = 32;
111const KV_ADMISSION_TARGET_LEN_METADATA_KEY: &str = "ferrum_kv_admission_target_len";
112const GENERATED_CONTROL_TOKEN_TEXTS: &[&str] = &[
113    "<think>",
114    "</think>",
115    "<|im_end|>",
116    "<|endoftext|>",
117    "<|eot_id|>",
118    "<|eom_id|>",
119    "</s>",
120];
121
122struct TokenPolicyCacheEntry {
123    tokenizer: Weak<dyn Tokenizer + Send + Sync>,
124    forbidden: HashSet<u32>,
125    model_greedy_forbidden: HashSet<u32>,
126}
127
128static TOKEN_POLICY_CACHE: OnceLock<
129    std::sync::Mutex<HashMap<(usize, usize), TokenPolicyCacheEntry>>,
130> = OnceLock::new();
131
132#[derive(Debug, Clone, Default, PartialEq, Eq)]
133struct ContinuousEngineRuntimeConfig {
134    active_decode_prefill_chunk: Option<usize>,
135    batch_decode_prof: bool,
136    chunked_prefill_present: bool,
137    chunked_prefill_size: Option<usize>,
138    kv_capacity: Option<usize>,
139    max_model_len: Option<usize>,
140    next_batch_prof: bool,
141    profile_entrypoint: Option<ProfileEntrypoint>,
142    profile_jsonl: Option<PathBuf>,
143    prefix_cache_enabled: bool,
144    rbd_prof: bool,
145    scheduler_trace_jsonl: Option<PathBuf>,
146    legacy_scheduler_trace_jsonl: Option<PathBuf>,
147    unified_post_prof: bool,
148}
149
150impl ContinuousEngineRuntimeConfig {
151    /// Build from the typed `EngineConfig.runtime` knobs (resolved by the CLI/
152    /// autosizer via the runtime-config snapshot). Reads no environment — the
153    /// env bridge stays at the composition root.
154    fn from_engine_config(config: &EngineConfig) -> Self {
155        let r = &config.runtime;
156        Self {
157            active_decode_prefill_chunk: config.scheduler.active_decode_prefill_chunk,
158            batch_decode_prof: r.batch_decode_prof,
159            chunked_prefill_present: r.chunked_prefill_size.is_some(),
160            chunked_prefill_size: r.chunked_prefill_size,
161            kv_capacity: r.kv_capacity,
162            max_model_len: r.max_model_len,
163            next_batch_prof: r.next_batch_prof,
164            profile_entrypoint: r.profile_entrypoint,
165            profile_jsonl: r.profile_jsonl.clone(),
166            prefix_cache_enabled: r.prefix_cache_enabled,
167            rbd_prof: r.rbd_prof,
168            scheduler_trace_jsonl: r.scheduler_trace_jsonl.clone(),
169            legacy_scheduler_trace_jsonl: r.legacy_scheduler_trace_jsonl.clone(),
170            unified_post_prof: r.unified_post_prof,
171        }
172    }
173
174    #[cfg(test)]
175    fn from_env_vars<I, K, V>(active_decode_prefill_chunk: Option<usize>, vars: I) -> Self
176    where
177        I: IntoIterator<Item = (K, V)>,
178        K: Into<String>,
179        V: Into<String>,
180    {
181        let vars: HashMap<String, String> = vars
182            .into_iter()
183            .map(|(key, value)| (key.into(), value.into()))
184            .collect();
185        Self {
186            active_decode_prefill_chunk,
187            batch_decode_prof: vars.contains_key(BATCH_DECODE_PROF_ENV),
188            chunked_prefill_present: vars.contains_key(CHUNKED_PREFILL_ENV),
189            chunked_prefill_size: parse_positive_usize_env(&vars, CHUNKED_PREFILL_ENV),
190            kv_capacity: parse_positive_usize_env(&vars, KV_CAPACITY_ENV),
191            max_model_len: parse_positive_usize_env(&vars, MAX_MODEL_LEN_ENV),
192            next_batch_prof: vars.contains_key(NEXT_BATCH_PROF_ENV),
193            profile_entrypoint: vars
194                .get("FERRUM_PROFILE_ENTRYPOINT")
195                .and_then(|value| ProfileEntrypoint::parse(value)),
196            profile_jsonl: vars
197                .get("FERRUM_PROFILE_JSONL")
198                .and_then(|value| ferrum_types::parse_path_env_value(value).ok()),
199            prefix_cache_enabled: vars
200                .get(WHOLE_PROMPT_PREFIX_CACHE_ENV)
201                .is_some_and(|v| v == "1"),
202            rbd_prof: vars.contains_key(RBD_PROF_ENV),
203            scheduler_trace_jsonl: vars
204                .get("FERRUM_SCHEDULER_TRACE_JSONL")
205                .and_then(|value| ferrum_types::parse_path_env_value(value).ok()),
206            legacy_scheduler_trace_jsonl: vars
207                .get("FERRUM_LEGACY_SCHEDULER_TRACE_JSONL")
208                .and_then(|value| ferrum_types::parse_path_env_value(value).ok()),
209            unified_post_prof: vars.contains_key(UNIFIED_POST_PROF_ENV),
210        }
211    }
212
213    fn chunked_prefill_size_for(&self, num_tokens: usize) -> Option<usize> {
214        self.chunked_prefill_size.filter(|&n| n < num_tokens)
215    }
216}
217
218#[cfg(test)]
219fn parse_positive_usize_env(vars: &HashMap<String, String>, name: &str) -> Option<usize> {
220    vars.get(name)
221        .and_then(|v| v.parse::<usize>().ok())
222        .filter(|&v| v > 0)
223}
224
225fn effective_request_context_capacity(
226    config: &EngineConfig,
227    runtime_config: &ContinuousEngineRuntimeConfig,
228    executor_kv_capacity: Option<usize>,
229) -> Option<usize> {
230    let kv_capacity = runtime_config
231        .kv_capacity
232        .or(executor_kv_capacity)
233        .or_else(|| (config.kv_cache.max_blocks > 0).then_some(config.kv_cache.max_blocks));
234    let max_model_len = runtime_config.max_model_len.or_else(|| {
235        config
236            .model
237            .model_info
238            .as_ref()
239            .map(|info| info.max_sequence_length)
240            .filter(|&len| len > 0)
241    });
242
243    match (kv_capacity, max_model_len) {
244        (Some(kv), Some(model)) => Some(kv.min(model)),
245        (Some(kv), None) => Some(kv),
246        (None, Some(model)) => Some(model),
247        (None, None) => None,
248    }
249}
250
251fn validate_request_context_budget(
252    request: &InferenceRequest,
253    input_tokens: usize,
254    config: &EngineConfig,
255    runtime_config: &ContinuousEngineRuntimeConfig,
256    executor_kv_capacity: Option<usize>,
257) -> Result<()> {
258    let Some(capacity) =
259        effective_request_context_capacity(config, runtime_config, executor_kv_capacity)
260    else {
261        return Ok(());
262    };
263    let output_tokens = request.sampling_params.max_tokens;
264    if input_tokens.saturating_add(output_tokens) <= capacity {
265        return Ok(());
266    }
267
268    Err(FerrumError::ContextLengthExceeded {
269        capacity,
270        input_tokens,
271        output_tokens,
272    })
273}
274
275fn clamp_default_max_tokens_to_context(
276    request: &mut InferenceRequest,
277    input_tokens: usize,
278    config: &EngineConfig,
279    runtime_config: &ContinuousEngineRuntimeConfig,
280    executor_kv_capacity: Option<usize>,
281) {
282    let default_max_tokens = request
283        .metadata
284        .get(DEFAULT_MAX_TOKENS_METADATA_KEY)
285        .and_then(|value| value.as_bool())
286        .unwrap_or(false);
287    if !default_max_tokens {
288        return;
289    }
290    let Some(capacity) =
291        effective_request_context_capacity(config, runtime_config, executor_kv_capacity)
292    else {
293        return;
294    };
295    let available_output_tokens = capacity.saturating_sub(input_tokens);
296    if available_output_tokens == 0 {
297        return;
298    }
299    let current = request.sampling_params.max_tokens;
300    let clamped = current.min(available_output_tokens);
301    if clamped < current {
302        warn!(
303            "Clamping default max_tokens from {} to {} for context budget: input_tokens={}, capacity={}",
304            current, clamped, input_tokens, capacity
305        );
306        request.sampling_params.max_tokens = clamped;
307    }
308}
309
310struct StopConditions {
311    /// Automatic model termination, excluding explicit user-stop collisions.
312    model_eos_token_ids: Vec<u32>,
313    /// All terminal IDs, including model EOS and explicit user stops.
314    stop_token_ids: HashSet<u32>,
315    user_stop_token_ids: HashSet<u32>,
316    stop_text_seqs: Vec<String>,
317}
318
319/// Resolve per-request automatic and user stop conditions.
320///
321/// Combines:
322/// 1. Model EOS reported by the tokenizer (`special_tokens().eos_token`).
323/// 2. Common chat-EOS literal names looked up in the tokenizer's vocab —
324///    `<|im_end|>`, `<|endoftext|>`, `<|eot_id|>`, `</s>`. Each lookup is
325///    model-specific (only IDs that actually exist in this vocab get added),
326///    so there's no risk of inserting an unrelated token id from a hard-coded
327///    fallback list (e.g. `2` is `</s>` for LLaMA but `!` for Qwen3).
328/// 3. User-supplied `stop_sequences` — each is encoded with `add_special=false`;
329///    one-token results land in `stop_token_ids` for the fast path, and all
330///    user stop strings remain in `stop_text_seqs` so tokens that contain the
331///    stop text as a substring still stop.
332fn resolve_stop_conditions(
333    params: &SamplingParams,
334    tokenizer: Option<&(dyn Tokenizer + Send + Sync)>,
335    ignore_eos: bool,
336) -> StopConditions {
337    let mut model_eos_token_ids = Vec::new();
338    let mut text_seqs: Vec<String> = Vec::new();
339
340    if let Some(tok) = tokenizer {
341        if !ignore_eos {
342            if let Some(eos) = tok.special_tokens().eos_token {
343                model_eos_token_ids.push(eos.get());
344            }
345            for extra in &tok.special_tokens().extra_eos_tokens {
346                model_eos_token_ids.push(extra.get());
347            }
348            for name in ["<|im_end|>", "<|endoftext|>", "<|eot_id|>", "</s>"] {
349                if let Some(t) = tok.token_id(name) {
350                    model_eos_token_ids.push(t.get());
351                }
352            }
353        }
354    }
355    model_eos_token_ids.sort_unstable();
356    model_eos_token_ids.dedup();
357    let mut ids: HashSet<u32> = model_eos_token_ids.iter().copied().collect();
358    let mut user_stop_token_ids = HashSet::new();
359
360    if let Some(tok) = tokenizer {
361        for stop_seq in &params.stop_sequences {
362            if stop_seq.is_empty() {
363                continue;
364            }
365            text_seqs.push(stop_seq.clone());
366            match tok.encode(stop_seq, false) {
367                Ok(toks) if toks.len() == 1 => {
368                    ids.insert(toks[0].get());
369                    user_stop_token_ids.insert(toks[0].get());
370                }
371                _ => {}
372            }
373        }
374    } else {
375        for stop_seq in &params.stop_sequences {
376            text_seqs.push(stop_seq.clone());
377        }
378    }
379    // Explicit stops may interrupt an unfinished response even when their ID
380    // is also model EOS. Keep the union intact for hard grammar validation and
381    // output stripping; only automatic completion gates exclude collisions.
382    model_eos_token_ids.retain(|id| !user_stop_token_ids.contains(id));
383    StopConditions {
384        model_eos_token_ids,
385        stop_token_ids: ids,
386        user_stop_token_ids,
387        stop_text_seqs: text_seqs,
388    }
389}
390
391#[derive(Debug)]
392struct TokenSequenceMatcher {
393    tokens: Vec<u32>,
394    failure: Vec<usize>,
395    matched: usize,
396}
397
398impl TokenSequenceMatcher {
399    fn new(tokens: Vec<u32>, label: &str) -> Result<Self> {
400        if tokens.is_empty() {
401            return Err(FerrumError::config(format!(
402                "{label} requires at least one token"
403            )));
404        }
405        let failure = delimiter_failure_table(&tokens);
406        Ok(Self {
407            tokens,
408            failure,
409            matched: 0,
410        })
411    }
412
413    fn observe(&mut self, token_id: u32) -> bool {
414        if self.matched == self.tokens.len() {
415            self.matched = self.failure[self.matched - 1];
416        }
417        while self.matched > 0 && self.tokens[self.matched] != token_id {
418            self.matched = self.failure[self.matched - 1];
419        }
420        if self.tokens[self.matched] == token_id {
421            self.matched += 1;
422        }
423        let completed = self.matched == self.tokens.len();
424        if completed {
425            self.matched = self.failure[self.matched - 1];
426        }
427        completed
428    }
429
430    fn reset(&mut self) {
431        self.matched = 0;
432    }
433
434    fn is_at_boundary(&self) -> bool {
435        self.matched == 0
436    }
437
438    fn is_partial(&self) -> bool {
439        self.matched > 0
440    }
441}
442
443#[derive(Debug)]
444enum DelimitedPayloadCompletionState {
445    AwaitingDelimiter(TokenSequenceMatcher),
446    AwaitingPayload,
447}
448
449#[derive(Debug, Clone, Copy, PartialEq, Eq)]
450enum EnvelopePrefixEffect {
451    Clear,
452    Pending,
453    Rejected,
454}
455
456impl DelimitedPayloadCompletionState {
457    fn observe(
458        &mut self,
459        previous_tokens: &[TokenId],
460        token: TokenId,
461        tokenizer: Option<&(dyn Tokenizer + Send + Sync)>,
462        envelope_prefix: EnvelopePrefixEffect,
463    ) -> Result<bool> {
464        match self {
465            Self::AwaitingDelimiter(matcher) => {
466                if matcher.observe(token.get()) {
467                    *self = Self::AwaitingPayload;
468                }
469                Ok(false)
470            }
471            Self::AwaitingPayload => {
472                match envelope_prefix {
473                    EnvelopePrefixEffect::Pending => return Ok(false),
474                    EnvelopePrefixEffect::Rejected => return Ok(true),
475                    EnvelopePrefixEffect::Clear => {}
476                }
477                let tokenizer = tokenizer.ok_or_else(|| {
478                    FerrumError::config("response completion boundary lost its tokenizer")
479                })?;
480                let delta = tokenizer.decode_incremental(previous_tokens, token)?;
481                Ok(delta
482                    .chars()
483                    .any(|character| !character.is_whitespace() && character != '\u{FFFD}'))
484            }
485        }
486    }
487}
488
489#[derive(Debug, Clone, Copy, PartialEq, Eq)]
490enum ResponseEnvelopePhase {
491    AwaitingOpen,
492    AwaitingClose,
493}
494
495#[derive(Debug)]
496struct ResponseEnvelopeCompletionState {
497    open: TokenSequenceMatcher,
498    close: TokenSequenceMatcher,
499    phase: ResponseEnvelopePhase,
500    completed_envelopes: usize,
501    max_envelopes: usize,
502}
503
504impl ResponseEnvelopeCompletionState {
505    fn observe(&mut self, token_id: u32) -> Result<()> {
506        match self.phase {
507            ResponseEnvelopePhase::AwaitingOpen => {
508                let matched_before = self.open.matched;
509                if self.open.observe(token_id) {
510                    if self.completed_envelopes == self.max_envelopes {
511                        self.open.matched = matched_before;
512                        return Err(FerrumError::invalid_format(format!(
513                            "generated response exceeded its {}-envelope protocol limit",
514                            self.max_envelopes
515                        )));
516                    }
517                    self.phase = ResponseEnvelopePhase::AwaitingClose;
518                    self.close.reset();
519                }
520            }
521            ResponseEnvelopePhase::AwaitingClose => {
522                if self.close.observe(token_id) {
523                    self.completed_envelopes += 1;
524                    self.phase = ResponseEnvelopePhase::AwaitingOpen;
525                    self.open.reset();
526                    self.close.reset();
527                }
528            }
529        }
530        Ok(())
531    }
532
533    fn allows_model_eos(&self) -> bool {
534        self.completed_envelopes > 0
535            && self.phase == ResponseEnvelopePhase::AwaitingOpen
536            && self.open.is_at_boundary()
537    }
538
539    fn has_committed_to_envelope_path(&self) -> bool {
540        self.phase == ResponseEnvelopePhase::AwaitingClose || self.completed_envelopes > 0
541    }
542
543    fn opener_is_partial(&self) -> bool {
544        self.phase == ResponseEnvelopePhase::AwaitingOpen && self.open.is_partial()
545    }
546}
547
548#[derive(Debug)]
549enum ResponseCompletionState {
550    Satisfied,
551    Pending {
552        delimited_payload: DelimitedPayloadCompletionState,
553        alternate_envelope: Option<ResponseEnvelopeCompletionState>,
554    },
555}
556
557impl ResponseCompletionState {
558    fn compile_token_sequence(
559        text: &str,
560        label: &str,
561        tokenizer: &(dyn Tokenizer + Send + Sync),
562        model_eos_token_ids: &[u32],
563    ) -> Result<Vec<u32>> {
564        let tokens = if let Some(token) = tokenizer.token_id(text) {
565            vec![token.get()]
566        } else {
567            tokenizer
568                .encode(text, false)?
569                .into_iter()
570                .map(TokenId::get)
571                .collect::<Vec<_>>()
572        };
573        if tokens.is_empty() {
574            return Err(FerrumError::invalid_request(format!(
575                "{label} {text:?} did not tokenize"
576            )));
577        }
578        if let Some(token) = tokens
579            .iter()
580            .find(|token| model_eos_token_ids.contains(token))
581        {
582            return Err(FerrumError::invalid_request(format!(
583                "{label} token {token} conflicts with model EOS"
584            )));
585        }
586        Ok(tokens)
587    }
588
589    fn compile(
590        boundary: &ResponseCompletionBoundary,
591        tokenizer: Option<&(dyn Tokenizer + Send + Sync)>,
592        model_eos_token_ids: &[u32],
593        max_tokens: usize,
594    ) -> Result<Self> {
595        let ResponseCompletionBoundary::AfterDelimiterAndPayload {
596            delimiter,
597            alternate_envelope,
598        } = boundary
599        else {
600            return Ok(Self::Satisfied);
601        };
602        if model_eos_token_ids.is_empty() {
603            return Ok(Self::Satisfied);
604        }
605        let tokenizer = tokenizer.ok_or_else(|| {
606            FerrumError::config("response completion boundary requires a tokenizer")
607        })?;
608        let delimiter_tokens = Self::compile_token_sequence(
609            delimiter,
610            "response completion delimiter",
611            tokenizer,
612            model_eos_token_ids,
613        )?;
614        if max_tokens <= delimiter_tokens.len() {
615            return Err(FerrumError::invalid_request(format!(
616                "response completion requires max_tokens greater than its {}-token delimiter",
617                delimiter_tokens.len()
618            )));
619        }
620        let alternate_envelope = alternate_envelope
621            .as_ref()
622            .map(|envelope| -> Result<ResponseEnvelopeCompletionState> {
623                if envelope.max_envelopes == 0 {
624                    return Err(FerrumError::invalid_request(
625                        "response completion envelope limit must be greater than zero",
626                    ));
627                }
628                Ok(ResponseEnvelopeCompletionState {
629                    open: TokenSequenceMatcher::new(
630                        Self::compile_token_sequence(
631                            &envelope.open_token_text,
632                            "response completion envelope opener",
633                            tokenizer,
634                            model_eos_token_ids,
635                        )?,
636                        "response completion envelope opener",
637                    )?,
638                    close: TokenSequenceMatcher::new(
639                        Self::compile_token_sequence(
640                            &envelope.close_token_text,
641                            "response completion envelope closer",
642                            tokenizer,
643                            model_eos_token_ids,
644                        )?,
645                        "response completion envelope closer",
646                    )?,
647                    phase: ResponseEnvelopePhase::AwaitingOpen,
648                    completed_envelopes: 0,
649                    max_envelopes: envelope.max_envelopes,
650                })
651            })
652            .transpose()?;
653        Ok(Self::Pending {
654            delimited_payload: DelimitedPayloadCompletionState::AwaitingDelimiter(
655                TokenSequenceMatcher::new(delimiter_tokens, "response completion delimiter")?,
656            ),
657            alternate_envelope,
658        })
659    }
660
661    fn allows_model_eos(&self) -> bool {
662        match self {
663            Self::Satisfied => true,
664            Self::Pending {
665                alternate_envelope, ..
666            } => alternate_envelope
667                .as_ref()
668                .is_some_and(ResponseEnvelopeCompletionState::allows_model_eos),
669        }
670    }
671
672    /// Advance the completion protocol with one accepted token. Matchers are
673    /// allocation-free after request construction.
674    fn observe(
675        &mut self,
676        previous_tokens: &[TokenId],
677        token: TokenId,
678        tokenizer: Option<&(dyn Tokenizer + Send + Sync)>,
679    ) -> Result<Option<bool>> {
680        let allowed_before = self.allows_model_eos();
681        let payload_completed = match self {
682            Self::Satisfied => return Ok(None),
683            Self::Pending {
684                delimited_payload,
685                alternate_envelope,
686            } => {
687                if let Some(envelope) = alternate_envelope {
688                    let committed_before = envelope.has_committed_to_envelope_path();
689                    let opener_was_partial = envelope.opener_is_partial();
690                    envelope.observe(token.get())?;
691                    let committed_after = envelope.has_committed_to_envelope_path();
692                    let opener_is_partial = envelope.opener_is_partial();
693
694                    if committed_before || committed_after {
695                        false
696                    } else {
697                        let envelope_prefix = if opener_is_partial {
698                            EnvelopePrefixEffect::Pending
699                        } else if opener_was_partial {
700                            EnvelopePrefixEffect::Rejected
701                        } else {
702                            EnvelopePrefixEffect::Clear
703                        };
704                        delimited_payload.observe(
705                            previous_tokens,
706                            token,
707                            tokenizer,
708                            envelope_prefix,
709                        )?
710                    }
711                } else {
712                    delimited_payload.observe(
713                        previous_tokens,
714                        token,
715                        tokenizer,
716                        EnvelopePrefixEffect::Clear,
717                    )?
718                }
719            }
720        };
721        if payload_completed {
722            *self = Self::Satisfied;
723        }
724        let allowed_after = self.allows_model_eos();
725        Ok((allowed_before != allowed_after).then_some(allowed_after))
726    }
727}
728
729fn delimiter_failure_table(tokens: &[u32]) -> Vec<usize> {
730    let mut failure = vec![0usize; tokens.len()];
731    let mut matched = 0usize;
732    for index in 1..tokens.len() {
733        while matched > 0 && tokens[matched] != tokens[index] {
734            matched = failure[matched - 1];
735        }
736        if tokens[matched] == tokens[index] {
737            matched += 1;
738        }
739        failure[index] = matched;
740    }
741    failure
742}
743
744fn resolve_sampling_token_constraints(
745    tokenizer: Option<&Arc<dyn Tokenizer + Send + Sync>>,
746    stop_token_ids: &HashSet<u32>,
747    request_generated_control_token_texts: &[&str],
748) -> (HashSet<u32>, HashSet<u32>, Option<usize>, HashSet<u32>) {
749    let mut allowed_extended = stop_token_ids.clone();
750    let Some(tok) = tokenizer else {
751        return (HashSet::new(), HashSet::new(), None, allowed_extended);
752    };
753
754    if let Some(eos) = tok.special_tokens().eos_token {
755        allowed_extended.insert(eos.get());
756    }
757    for extra in &tok.special_tokens().extra_eos_tokens {
758        allowed_extended.insert(extra.get());
759    }
760    for text in GENERATED_CONTROL_TOKEN_TEXTS {
761        if let Some(token) = tok.token_id(text) {
762            allowed_extended.insert(token.get());
763        }
764    }
765    for text in request_generated_control_token_texts {
766        if let Some(token) = tok.token_id(text) {
767            allowed_extended.insert(token.get());
768        }
769    }
770
771    let forbidden = cached_forbidden_generation_tokens(tok, &allowed_extended);
772    let model_greedy_forbidden = cached_model_greedy_forbidden_tokens(tok);
773
774    (
775        forbidden,
776        model_greedy_forbidden,
777        Some(tok.vocab_size()),
778        allowed_extended,
779    )
780}
781
782fn build_argmax_token_mask(
783    tok: &(dyn Tokenizer + Send + Sync),
784    model_vocab_size: Option<usize>,
785    forbidden_token_ids: &HashSet<u32>,
786    model_greedy_forbidden_token_ids: &HashSet<u32>,
787    initial_forbidden_token_ids: &HashSet<u32>,
788    stop_token_ids: &HashSet<u32>,
789    allowed_extended_token_ids: &HashSet<u32>,
790) -> TokenSelectionMask {
791    let tokenizer_vocab_size = tok.vocab_size();
792    let max_allowed_id = allowed_extended_token_ids
793        .iter()
794        .chain(stop_token_ids.iter())
795        .copied()
796        .max()
797        .map(|id| id as usize + 1)
798        .unwrap_or(0);
799    let mask_len = model_vocab_size
800        .unwrap_or(tokenizer_vocab_size)
801        .max(tokenizer_vocab_size)
802        .max(max_allowed_id);
803    let mut valid = vec![1i8; mask_len];
804    for &token_id in forbidden_token_ids
805        .iter()
806        .chain(model_greedy_forbidden_token_ids.iter())
807        .chain(initial_forbidden_token_ids.iter())
808    {
809        if let Some(slot) = valid.get_mut(token_id as usize) {
810            *slot = 0;
811        }
812    }
813    for token_id in tokenizer_vocab_size..mask_len {
814        if !allowed_extended_token_ids.contains(&(token_id as u32)) {
815            valid[token_id] = 0;
816        }
817    }
818    for &token_id in allowed_extended_token_ids {
819        if stop_token_ids.contains(&token_id) {
820            continue;
821        }
822        let token = TokenId::new(token_id);
823        let should_mask = tok
824            .decode(&[token], true)
825            .map(|text| decoded_delta_has_forbidden_quality(&text, 0, false, true))
826            .unwrap_or(true);
827        if should_mask {
828            if let Some(slot) = valid.get_mut(token_id as usize) {
829                *slot = 0;
830            }
831        }
832    }
833    TokenSelectionMask::new(valid)
834}
835
836fn cached_forbidden_generation_tokens(
837    tok: &Arc<dyn Tokenizer + Send + Sync>,
838    allowed_generated_controls: &HashSet<u32>,
839) -> HashSet<u32> {
840    let key = (tokenizer_cache_key(tok), tok.vocab_size());
841    let cache = TOKEN_POLICY_CACHE.get_or_init(|| std::sync::Mutex::new(HashMap::new()));
842    let tokenizer_identity = Arc::downgrade(tok);
843    {
844        let mut cache = cache.lock().expect("token policy cache poisoned");
845        if let Some(cached) = cache.get(&key) {
846            if cached.tokenizer.ptr_eq(&tokenizer_identity) && cached.tokenizer.strong_count() > 0 {
847                let mut forbidden = cached.forbidden.clone();
848                forbidden.retain(|token_id| !allowed_generated_controls.contains(token_id));
849                return forbidden;
850            }
851        }
852        cache.remove(&key);
853    }
854
855    let mut forbidden = HashSet::new();
856    let mut model_greedy_forbidden = HashSet::new();
857    let scan_limit = tok.vocab_size().min(GENERATION_POLICY_SCAN_LIMIT);
858    let has_reverse_vocab =
859        (0..scan_limit).any(|token_id| tok.token_text(TokenId::new(token_id as u32)).is_some());
860    let special = tok.special_tokens();
861    for token in [
862        special.bos_token,
863        special.unk_token,
864        special.pad_token,
865        special.sep_token,
866        special.cls_token,
867        special.mask_token,
868    ]
869    .into_iter()
870    .flatten()
871    {
872        forbidden.insert(token.get());
873    }
874
875    for text in [
876        "<unk",
877        "<unk>",
878        "[UNK]",
879        "<pad>",
880        "[PAD]",
881        "<|pad|>",
882        "<mask>",
883        "[MASK]",
884        "\u{00ef}\u{00bf}\u{00bd}",
885    ] {
886        if let Some(token) = tok.token_id(text) {
887            forbidden.insert(token.get());
888        }
889    }
890    for token_id in 0..scan_limit {
891        let id = token_id as u32;
892        let token = TokenId::new(id);
893        let raw_text = tok.token_text(token);
894        let missing_token_text = has_reverse_vocab && raw_text.is_none();
895        let raw_text_forbidden = raw_text.is_some_and(is_forbidden_generation_token_text);
896        let decoded_text_forbidden = tok
897            .decode(&[token], true)
898            .map(|text| decoded_token_is_statically_forbidden(tok.as_ref(), token, &text))
899            .unwrap_or(true);
900        let context_free_bytes_invalid = tok
901            .token_bytes(token)
902            .is_none_or(|bytes| advance_pending_utf8_fragment(&[], &bytes).is_err());
903        if context_free_bytes_invalid {
904            model_greedy_forbidden.insert(id);
905        }
906        if missing_token_text || raw_text_forbidden || decoded_text_forbidden {
907            forbidden.insert(id);
908        }
909    }
910
911    let mut cache = cache.lock().expect("token policy cache poisoned");
912    cache.retain(|_, entry| entry.tokenizer.strong_count() > 0);
913    cache.insert(
914        key,
915        TokenPolicyCacheEntry {
916            tokenizer: tokenizer_identity,
917            forbidden: forbidden.clone(),
918            model_greedy_forbidden,
919        },
920    );
921    drop(cache);
922    forbidden.retain(|token_id| !allowed_generated_controls.contains(token_id));
923    forbidden
924}
925
926fn cached_model_greedy_forbidden_tokens(tok: &Arc<dyn Tokenizer + Send + Sync>) -> HashSet<u32> {
927    let key = (tokenizer_cache_key(tok), tok.vocab_size());
928    let cache = TOKEN_POLICY_CACHE.get_or_init(|| std::sync::Mutex::new(HashMap::new()));
929    cache
930        .lock()
931        .expect("token policy cache poisoned")
932        .get(&key)
933        .filter(|entry| {
934            entry.tokenizer.ptr_eq(&Arc::downgrade(tok)) && entry.tokenizer.strong_count() > 0
935        })
936        .map(|entry| entry.model_greedy_forbidden.clone())
937        .unwrap_or_default()
938}
939
940fn tokenizer_cache_key(tok: &Arc<dyn Tokenizer + Send + Sync>) -> usize {
941    Arc::as_ptr(tok).cast::<()>() as usize
942}
943
944fn is_forbidden_generation_token_text(text: &str) -> bool {
945    let text = text.trim();
946    if text.is_empty() {
947        return false;
948    }
949    if text.contains('\u{FFFD}') {
950        return true;
951    }
952    if contains_replacement_char_mojibake(text) {
953        return true;
954    }
955
956    let lower = text.to_ascii_lowercase();
957    let lower = lower.as_str();
958    if matches!(
959        lower,
960        "<unk" | "<unk>" | "[unk]" | "<pad>" | "[pad]" | "<|pad|>" | "<mask>" | "[mask]"
961    ) {
962        return true;
963    }
964
965    let looks_like_special = (lower.starts_with('<') && lower.ends_with('>'))
966        || (lower.starts_with('[') && lower.ends_with(']'));
967    if !looks_like_special {
968        return false;
969    }
970
971    lower.contains("unk")
972        || lower.contains("pad")
973        || lower.contains("mask")
974        || lower.contains("reserved")
975        || lower.contains("unused")
976}
977
978fn decoded_token_is_statically_forbidden(
979    tokenizer: &(dyn Tokenizer + Send + Sync),
980    token: TokenId,
981    decoded: &str,
982) -> bool {
983    if !decoded.contains('\u{FFFD}') {
984        return is_forbidden_generation_token_text(decoded);
985    }
986    if contains_replacement_char_mojibake(decoded) {
987        return true;
988    }
989
990    let Some(bytes) = tokenizer.token_bytes(token) else {
991        return true;
992    };
993    if bytes.is_empty() || std::str::from_utf8(&bytes).is_ok() {
994        return true;
995    }
996
997    // Byte-level vocabularies can split one UTF-8 scalar across tokens. A
998    // one-token string decode must render such a fragment as U+FFFD, but the
999    // raw bytes remain legal in context and are owned by candidate decoding.
1000    !is_potential_utf8_fragment(&bytes)
1001}
1002
1003fn is_potential_utf8_fragment(bytes: &[u8]) -> bool {
1004    if bytes.is_empty() {
1005        return false;
1006    }
1007
1008    let mut offset = 0usize;
1009    while offset < bytes.len() && is_utf8_continuation(bytes[offset]) {
1010        offset += 1;
1011    }
1012    if offset > 3 {
1013        return false;
1014    }
1015
1016    while offset < bytes.len() {
1017        let lead = bytes[offset];
1018        if lead.is_ascii() {
1019            offset += 1;
1020            continue;
1021        }
1022
1023        let continuation_count = match lead {
1024            0xC2..=0xDF => 1,
1025            0xE0..=0xEF => 2,
1026            0xF0..=0xF4 => 3,
1027            _ => return false,
1028        };
1029        let available = bytes.len() - offset - 1;
1030        let present = available.min(continuation_count);
1031        for index in 0..present {
1032            let byte = bytes[offset + index + 1];
1033            if !is_utf8_continuation(byte) {
1034                return false;
1035            }
1036            if index == 0
1037                && matches!(
1038                    (lead, byte),
1039                    (0xE0, 0x80..=0x9F)
1040                        | (0xED, 0xA0..=0xBF)
1041                        | (0xF0, 0x80..=0x8F)
1042                        | (0xF4, 0x90..=0xBF)
1043                )
1044            {
1045                return false;
1046            }
1047        }
1048        if present < continuation_count {
1049            return true;
1050        }
1051        offset += continuation_count + 1;
1052    }
1053
1054    true
1055}
1056
1057fn is_utf8_continuation(byte: u8) -> bool {
1058    matches!(byte, 0x80..=0xBF)
1059}
1060
1061fn advance_pending_utf8_fragment(pending: &[u8], next: &[u8]) -> std::result::Result<Vec<u8>, ()> {
1062    let mut combined = Vec::with_capacity(pending.len().saturating_add(next.len()));
1063    combined.extend_from_slice(pending);
1064    combined.extend_from_slice(next);
1065    match std::str::from_utf8(&combined) {
1066        Ok(text) => {
1067            if text.contains('\u{FFFD}') || contains_replacement_char_mojibake(text) {
1068                Err(())
1069            } else {
1070                Ok(Vec::new())
1071            }
1072        }
1073        Err(error) if error.error_len().is_none() => {
1074            let valid_prefix =
1075                std::str::from_utf8(&combined[..error.valid_up_to()]).map_err(|_| ())?;
1076            if valid_prefix.contains('\u{FFFD}') || contains_replacement_char_mojibake(valid_prefix)
1077            {
1078                return Err(());
1079            }
1080            let fragment = &combined[error.valid_up_to()..];
1081            if fragment.is_empty() || fragment.len() > 3 {
1082                return Err(());
1083            }
1084            Ok(fragment.to_vec())
1085        }
1086        Err(_) => Err(()),
1087    }
1088}
1089
1090fn decoded_delta_has_forbidden_quality(
1091    full_text: &str,
1092    previous_text_len: usize,
1093    candidate_is_stop: bool,
1094    candidate_is_non_stop_control: bool,
1095) -> bool {
1096    if previous_text_len > full_text.len() || !full_text.is_char_boundary(previous_text_len) {
1097        return true;
1098    }
1099    let delta = &full_text[previous_text_len..];
1100    if delta.is_empty() {
1101        return candidate_is_non_stop_control;
1102    }
1103    if contains_replacement_char_mojibake(delta) {
1104        return true;
1105    }
1106    if delta.contains('\u{FFFD}') && (candidate_is_stop || !full_text.ends_with('\u{FFFD}')) {
1107        return true;
1108    }
1109    false
1110}
1111
1112fn contains_replacement_char_mojibake(text: &str) -> bool {
1113    let mut chars = text.chars();
1114    let mut a = chars.next();
1115    let mut b = chars.next();
1116    let mut c = chars.next();
1117    loop {
1118        if matches!(
1119            (a, b, c),
1120            (Some('\u{00ef}'), Some('\u{00bf}'), Some('\u{00bd}'))
1121        ) {
1122            return true;
1123        }
1124        if c.is_none() {
1125            return false;
1126        }
1127        a = b;
1128        b = c;
1129        c = chars.next();
1130    }
1131}
1132
1133// ────────────────────────────────────────────────────────────────────────────
1134// Sequence state
1135// ────────────────────────────────────────────────────────────────────────────
1136
1137#[derive(Debug, Clone, PartialEq, Serialize)]
1138struct SequenceTokenTraceEvidence {
1139    schema_version: u32,
1140    token_encoding: &'static str,
1141    prompt_token_count: usize,
1142    prompt_token_sha256: String,
1143    prompt_token_prefix: Vec<u32>,
1144    prompt_token_tail: Vec<u32>,
1145    generated_token_count: usize,
1146    generated_token_sha256: String,
1147    generated_token_prefix: Vec<u32>,
1148    generated_token_tail: Vec<u32>,
1149    sampling_rng_algorithm: &'static str,
1150    sampling_seed: Option<u64>,
1151    sampler: String,
1152    processors: Vec<String>,
1153    max_tokens: usize,
1154    temperature: f32,
1155    top_p: f32,
1156    top_k: Option<usize>,
1157    repetition_penalty: f32,
1158    presence_penalty: f32,
1159    frequency_penalty: f32,
1160}
1161
1162impl SequenceTokenTraceEvidence {
1163    fn capture(sequence: &SequenceState) -> Self {
1164        Self {
1165            schema_version: OBSERVABILITY_PROFILE_SCHEMA_VERSION,
1166            token_encoding: "u32-le-v1",
1167            prompt_token_count: sequence.input_tokens.len(),
1168            prompt_token_sha256: token_ids_sha256(&sequence.input_tokens),
1169            prompt_token_prefix: token_id_prefix(
1170                &sequence.input_tokens,
1171                TOKEN_TRACE_PROMPT_PREFIX_LIMIT,
1172            ),
1173            prompt_token_tail: token_id_tail(&sequence.input_tokens, TOKEN_TRACE_PROMPT_TAIL_LIMIT),
1174            generated_token_count: sequence.generated_tokens.len(),
1175            generated_token_sha256: token_ids_sha256(&sequence.generated_tokens),
1176            generated_token_prefix: token_id_prefix(
1177                &sequence.generated_tokens,
1178                TOKEN_TRACE_GENERATED_PREFIX_LIMIT,
1179            ),
1180            generated_token_tail: token_id_tail(
1181                &sequence.generated_tokens,
1182                TOKEN_TRACE_GENERATED_TAIL_LIMIT,
1183            ),
1184            sampling_rng_algorithm: SamplingRng::algorithm_id(),
1185            sampling_seed: sequence.sampling_params.seed,
1186            sampler: sequence.sampling_plan.sampler.name().to_string(),
1187            processors: sequence
1188                .sampling_plan
1189                .processor_chain
1190                .processor_names()
1191                .into_iter()
1192                .map(str::to_string)
1193                .collect(),
1194            max_tokens: sequence.sampling_params.max_tokens,
1195            temperature: sequence.sampling_params.temperature,
1196            top_p: sequence.sampling_params.top_p,
1197            top_k: sequence.sampling_params.top_k,
1198            repetition_penalty: sequence.sampling_params.repetition_penalty,
1199            presence_penalty: sequence.sampling_params.presence_penalty,
1200            frequency_penalty: sequence.sampling_params.frequency_penalty,
1201        }
1202    }
1203}
1204
1205fn token_ids_sha256(tokens: &[TokenId]) -> String {
1206    let mut digest = Sha256::new();
1207    digest.update(b"ferrum-token-ids:u32-le-v1\0");
1208    for token in tokens {
1209        digest.update(token.get().to_le_bytes());
1210    }
1211    format!("sha256:{:x}", digest.finalize())
1212}
1213
1214fn token_id_prefix(tokens: &[TokenId], limit: usize) -> Vec<u32> {
1215    tokens.iter().take(limit).map(|token| token.get()).collect()
1216}
1217
1218fn token_id_tail(tokens: &[TokenId], limit: usize) -> Vec<u32> {
1219    let start = tokens.len().saturating_sub(limit);
1220    tokens[start..].iter().map(|token| token.get()).collect()
1221}
1222
1223mod sequence;
1224pub use sequence::SequenceState;
1225use sequence::*;
1226
1227enum EngineIterationOutcome {
1228    Progressed,
1229    Idle,
1230    CapacityBlocked(ExecutorCapacityWaitRegistration),
1231}
1232
1233enum EngineResourceComposition {
1234    LegacyEngine {
1235        kv_cache: Arc<dyn KvCacheManager + Send + Sync>,
1236        recurrent_state_manager: Option<Arc<dyn RecurrentStateManager + Send + Sync>>,
1237    },
1238    PlanRuntime,
1239}
1240
1241const MAX_EXECUTION_READINESS_WAITERS: usize = 64;
1242
1243struct ExecutionReadinessWaitFailure {
1244    ticket_id: u64,
1245    request_ids: Vec<RequestId>,
1246    message: String,
1247}
1248
1249struct ExecutionReadinessWaitTask {
1250    wake: ExecutionReadinessWake,
1251    handle: tokio::task::JoinHandle<()>,
1252}
1253
1254struct ExecutionReadinessWaitRegistry {
1255    slots: Arc<Semaphore>,
1256    tasks: Mutex<HashMap<u64, ExecutionReadinessWaitTask>>,
1257    failures: Arc<Mutex<VecDeque<ExecutionReadinessWaitFailure>>>,
1258}
1259
1260impl ExecutionReadinessWaitRegistry {
1261    fn new() -> Self {
1262        Self {
1263            slots: Arc::new(Semaphore::new(MAX_EXECUTION_READINESS_WAITERS)),
1264            tasks: Mutex::new(HashMap::new()),
1265            failures: Arc::new(Mutex::new(VecDeque::new())),
1266        }
1267    }
1268
1269    async fn reserve(&self) -> Result<tokio::sync::OwnedSemaphorePermit> {
1270        Arc::clone(&self.slots)
1271            .acquire_owned()
1272            .await
1273            .map_err(|_| FerrumError::internal("execution readiness waiter registry is closed"))
1274    }
1275
1276    fn track<F>(
1277        &self,
1278        wake: ExecutionReadinessWake,
1279        wait: F,
1280        request_ids: Vec<RequestId>,
1281        work_notify: Arc<Notify>,
1282        slot: tokio::sync::OwnedSemaphorePermit,
1283    ) -> Result<()>
1284    where
1285        F: Future<Output = std::result::Result<(), String>> + Send + 'static,
1286    {
1287        let ticket_id = wake.ticket_id().get();
1288        let failures = Arc::clone(&self.failures);
1289        let retained_wake = wake.clone();
1290        let mut tasks = self.tasks.lock();
1291        if tasks.contains_key(&ticket_id) {
1292            failures.lock().push_back(ExecutionReadinessWaitFailure {
1293                ticket_id,
1294                request_ids,
1295                message: format!("execution readiness ticket {ticket_id} was registered twice"),
1296            });
1297            wake.mark_failed();
1298            work_notify.notify_one();
1299            drop(slot);
1300            return Err(FerrumError::internal(format!(
1301                "execution readiness ticket {ticket_id} was registered twice"
1302            )));
1303        }
1304        let task = tokio::spawn(async move {
1305            let result = std::panic::AssertUnwindSafe(wait).catch_unwind().await;
1306            match result {
1307                Ok(Ok(_)) => {
1308                    wake.mark_ready();
1309                }
1310                Ok(Err(error)) => {
1311                    failures.lock().push_back(ExecutionReadinessWaitFailure {
1312                        ticket_id,
1313                        request_ids,
1314                        message: error.to_string(),
1315                    });
1316                    wake.mark_failed();
1317                }
1318                Err(_) => {
1319                    failures.lock().push_back(ExecutionReadinessWaitFailure {
1320                        ticket_id,
1321                        request_ids,
1322                        message: "Request-state readiness waiter panicked".to_string(),
1323                    });
1324                    wake.mark_failed();
1325                }
1326            }
1327            drop(slot);
1328            work_notify.notify_one();
1329        });
1330        tasks.insert(
1331            ticket_id,
1332            ExecutionReadinessWaitTask {
1333                wake: retained_wake,
1334                handle: task,
1335            },
1336        );
1337        Ok(())
1338    }
1339
1340    fn reap_finished(&self) {
1341        self.tasks
1342            .lock()
1343            .retain(|_, task| !task.handle.is_finished());
1344    }
1345
1346    fn take_failures(&self) -> Vec<ExecutionReadinessWaitFailure> {
1347        self.failures.lock().drain(..).collect()
1348    }
1349
1350    fn pending_count(&self) -> usize {
1351        MAX_EXECUTION_READINESS_WAITERS - self.slots.available_permits()
1352    }
1353
1354    async fn abort_and_join(&self) -> Result<()> {
1355        let tasks = self
1356            .tasks
1357            .lock()
1358            .drain()
1359            .map(|(_, task)| task)
1360            .collect::<Vec<_>>();
1361        for task in &tasks {
1362            task.wake.cancel();
1363            task.handle.abort();
1364        }
1365        for task in tasks {
1366            if let Err(error) = task.handle.await {
1367                if !error.is_cancelled() {
1368                    return Err(FerrumError::internal(format!(
1369                        "execution readiness waiter join failed: {error}"
1370                    )));
1371                }
1372            }
1373        }
1374        Ok(())
1375    }
1376}
1377
1378impl EngineResourceComposition {
1379    const fn authority(&self) -> ExecutionResourceAuthority {
1380        match self {
1381            Self::LegacyEngine { .. } => ExecutionResourceAuthority::LegacyEngine,
1382            Self::PlanRuntime => ExecutionResourceAuthority::PlanRuntime,
1383        }
1384    }
1385
1386    fn kv_cache(&self) -> Option<&Arc<dyn KvCacheManager + Send + Sync>> {
1387        match self {
1388            Self::LegacyEngine { kv_cache, .. } => Some(kv_cache),
1389            Self::PlanRuntime => None,
1390        }
1391    }
1392
1393    fn recurrent_state_manager(&self) -> Option<&Arc<dyn RecurrentStateManager + Send + Sync>> {
1394        match self {
1395            Self::LegacyEngine {
1396                recurrent_state_manager,
1397                ..
1398            } => recurrent_state_manager.as_ref(),
1399            Self::PlanRuntime => None,
1400        }
1401    }
1402}
1403
1404struct EngineInner {
1405    config: EngineConfig,
1406    scheduler: Arc<ContinuousBatchScheduler>,
1407    tokenizer: Arc<dyn Tokenizer + Send + Sync>,
1408    /// Lazily built once because most requests are plain text. Structured
1409    /// requests reuse its tokenizer trie and compiled grammar templates.
1410    structured_output_factory: OnceLock<std::result::Result<Arc<StructuredOutputFactory>, String>>,
1411    #[allow(dead_code)]
1412    // Retained for constructor API; sampling now uses per-request SamplingConfig
1413    sampler: Arc<dyn Sampler + Send + Sync>,
1414    resource_composition: EngineResourceComposition,
1415    model_executor: Arc<dyn ModelExecutor + Send + Sync>,
1416    /// Optional draft executor for speculative decoding. When set alongside
1417    /// `spec_config`, `run_single_decode` routes through `SpeculativeRunner`.
1418    draft_executor: Option<Arc<dyn ModelExecutor + Send + Sync>>,
1419    /// Speculative decoding parameters (N, temperature). `None` = disabled.
1420    spec_config: Option<crate::speculative::SpeculativeDecodingConfig>,
1421    tensor_factory: Arc<dyn TensorFactory>,
1422    sequences: RwLock<HashMap<RequestId, SequenceState>>,
1423    is_running: AtomicBool,
1424    shutdown_notify: Arc<Notify>,
1425    /// Serializes request publication with cancellation and BatchPlan
1426    /// construction. Device execution deliberately runs after this guard is
1427    /// released so new user requests can enter while a wave is in flight.
1428    iteration_lock: tokio::sync::Mutex<()>,
1429    /// Wakes callers or a background loop when new work is submitted.
1430    work_notify: Arc<Notify>,
1431    /// Prefix cache: shares KV blocks across requests with common prompts.
1432    prefix_cache: PrefixCache,
1433    runtime_config: ContinuousEngineRuntimeConfig,
1434    profile_trace_jsonl: Option<SchedulerTraceJournal>,
1435    scheduler_trace_jsonl: Option<SchedulerTraceJournal>,
1436    legacy_scheduler_trace_jsonl: Option<Arc<Mutex<std::fs::File>>>,
1437    scheduler_trace_none_streak: AtomicU64,
1438    resource_lifecycle: Mutex<ResourceLifecycleLedger>,
1439    resource_trace_event_counter: AtomicU64,
1440    dynamic_admission_availability: Mutex<Vec<CapacityAvailabilityEpoch>>,
1441    execution_readiness_waiters: ExecutionReadinessWaitRegistry,
1442    prefix_rendezvous: Mutex<Vec<inner::prefix_rendezvous::PrefixRendezvous>>,
1443    prefix_restore_pending: Mutex<HashMap<RequestId, inner::prefix_restore::PendingPrefixRestore>>,
1444    // stats
1445    iteration_count: AtomicU64,
1446    total_prefill_tokens: AtomicU64,
1447    total_decode_tokens: AtomicU64,
1448    total_preemptions: AtomicU64,
1449    prefix_cache_hits: AtomicU64,
1450    total_iteration_lock_wait_us: AtomicU64,
1451    iteration_lock_wait_samples: AtomicU64,
1452    total_scheduling_time_us: AtomicU64,
1453    scheduling_time_samples: AtomicU64,
1454    total_model_execution_time_us: AtomicU64,
1455    model_execution_time_samples: AtomicU64,
1456    /// Set true the first time `ensure_bg_loop` runs, so per-request
1457    /// `infer_stream` callers don't each spawn their own competing
1458    /// driver task (16 streaming requests = 16 drivers thrashing on
1459    /// `iteration_lock`, ~5ms/iter of tokio scheduling overhead).
1460    bg_loop_spawned: AtomicBool,
1461    shutdown_started: AtomicBool,
1462    shutdown_lock: tokio::sync::Mutex<()>,
1463    background_loop: Mutex<Option<tokio::task::JoinHandle<()>>>,
1464}
1465
1466struct ClientReceiverDropWake {
1467    work_notify: Arc<Notify>,
1468    armed: bool,
1469}
1470
1471impl EngineInner {
1472    fn signal_shutdown(&self) {
1473        self.shutdown_started.store(true, Ordering::Release);
1474        self.is_running.store(false, Ordering::SeqCst);
1475
1476        // The background loop can be between its state check and registering
1477        // the async wait. `notify_one` retains a permit across that window;
1478        // `notify_waiters` would lose the shutdown signal when no waiter is
1479        // registered yet.
1480        self.shutdown_notify.notify_one();
1481        self.work_notify.notify_one();
1482    }
1483
1484    fn structured_output_factory(&self) -> Result<Arc<StructuredOutputFactory>> {
1485        match self.structured_output_factory.get_or_init(|| {
1486            StructuredOutputFactory::new_with_model_vocab_size(
1487                Arc::clone(&self.tokenizer),
1488                Some(self.model_executor.info().vocab_size),
1489            )
1490            .map(Arc::new)
1491            .map_err(|error| error.to_string())
1492        }) {
1493            Ok(factory) => Ok(Arc::clone(factory)),
1494            Err(message) => Err(FerrumError::config(format!(
1495                "structured-output runtime unavailable: {message}"
1496            ))),
1497        }
1498    }
1499}
1500
1501impl ClientReceiverDropWake {
1502    fn new(work_notify: Arc<Notify>) -> Self {
1503        Self {
1504            work_notify,
1505            armed: true,
1506        }
1507    }
1508
1509    fn disarm(&mut self) {
1510        self.armed = false;
1511    }
1512}
1513
1514impl Drop for ClientReceiverDropWake {
1515    fn drop(&mut self) {
1516        if self.armed {
1517            self.work_notify.notify_one();
1518        }
1519    }
1520}
1521
1522struct CancellationAwareResponseStream {
1523    receiver: tokio_stream::wrappers::ReceiverStream<Result<StreamChunk>>,
1524    receiver_drop_wake: ClientReceiverDropWake,
1525}
1526
1527impl Stream for CancellationAwareResponseStream {
1528    type Item = Result<StreamChunk>;
1529
1530    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
1531        let result = Pin::new(&mut self.receiver).poll_next(cx);
1532        if matches!(&result, Poll::Ready(None)) {
1533            self.receiver_drop_wake.disarm();
1534        }
1535        result
1536    }
1537}
1538
1539impl EngineInner {
1540    fn engine_managed_kv_cache(&self) -> Result<&Arc<dyn KvCacheManager + Send + Sync>> {
1541        self.resource_composition.kv_cache().ok_or_else(|| {
1542            FerrumError::internal(
1543                "plan runtime attempted to use the legacy engine KV-cache manager",
1544            )
1545        })
1546    }
1547
1548    fn recurrent_state_manager(&self) -> Option<&Arc<dyn RecurrentStateManager + Send + Sync>> {
1549        self.resource_composition.recurrent_state_manager()
1550    }
1551
1552    fn record_iteration_lock_wait(&self, duration: Duration) {
1553        self.total_iteration_lock_wait_us
1554            .fetch_add(duration_to_us(duration), Ordering::Relaxed);
1555        self.iteration_lock_wait_samples
1556            .fetch_add(1, Ordering::Relaxed);
1557    }
1558
1559    fn record_scheduling_time(&self, duration: Duration) {
1560        self.total_scheduling_time_us
1561            .fetch_add(duration_to_us(duration), Ordering::Relaxed);
1562        self.scheduling_time_samples.fetch_add(1, Ordering::Relaxed);
1563    }
1564
1565    fn record_model_execution_time(&self, duration: Duration) {
1566        self.total_model_execution_time_us
1567            .fetch_add(duration_to_us(duration), Ordering::Relaxed);
1568        self.model_execution_time_samples
1569            .fetch_add(1, Ordering::Relaxed);
1570    }
1571
1572    fn trace_entrypoint(&self) -> ProfileEntrypoint {
1573        self.runtime_config
1574            .profile_entrypoint
1575            .unwrap_or(ProfileEntrypoint::Synthetic)
1576    }
1577
1578    fn extend_scheduler_timeline_attributes(
1579        &self,
1580        attributes: &mut BTreeMap<String, serde_json::Value>,
1581    ) {
1582        let snapshot = self.scheduler.trace_snapshot();
1583        attributes.extend([
1584            (
1585                "active_sequence_count".to_string(),
1586                serde_json::json!(snapshot.active_len),
1587            ),
1588            (
1589                "monotonic_nanos".to_string(),
1590                serde_json::json!(inner::scheduler_trace_monotonic_nanos()),
1591            ),
1592            (
1593                "scheduler_snapshot".to_string(),
1594                serde_json::to_value(snapshot).unwrap_or(serde_json::Value::Null),
1595            ),
1596        ]);
1597    }
1598
1599    #[allow(clippy::too_many_arguments)]
1600    fn trace_resource_event(
1601        &self,
1602        request_id: &RequestId,
1603        owner_kind: &str,
1604        owner_id: &str,
1605        resource_kind: &str,
1606        phase: &str,
1607        action: ResourceAction,
1608        amount: Option<i64>,
1609        before: Option<i64>,
1610        after: Option<i64>,
1611        capacity: Option<i64>,
1612        reason: Option<String>,
1613    ) {
1614        let Some(sink) = &self.scheduler_trace_jsonl else {
1615            return;
1616        };
1617        let entrypoint = self.trace_entrypoint();
1618        let event_num = self
1619            .resource_trace_event_counter
1620            .fetch_add(1, Ordering::Relaxed);
1621        let mut attributes = BTreeMap::from([
1622            (
1623                "actual_model_smoke".to_string(),
1624                serde_json::json!(matches!(
1625                    entrypoint,
1626                    ProfileEntrypoint::Run | ProfileEntrypoint::Serve
1627                )),
1628            ),
1629            (
1630                "backend_device".to_string(),
1631                serde_json::json!(format!("{:?}", self.config.backend.device)),
1632            ),
1633            (
1634                "backend_type".to_string(),
1635                serde_json::json!(format!("{:?}", self.config.backend.backend_type)),
1636            ),
1637            (
1638                "diagnostic_only".to_string(),
1639                serde_json::json!(self.config.runtime.profile_detail.diagnostic_only()),
1640            ),
1641            ("l0_only".to_string(), serde_json::json!(false)),
1642            (
1643                "profile_detail".to_string(),
1644                serde_json::json!(self.config.runtime.profile_detail.as_str()),
1645            ),
1646            (
1647                "resource_trace_source".to_string(),
1648                serde_json::json!("engine"),
1649            ),
1650        ]);
1651        if let Some(reason) = reason.as_deref() {
1652            attributes.insert("resource_reason".to_string(), serde_json::json!(reason));
1653        }
1654        let underflow_amount = match (action, amount, before) {
1655            (ResourceAction::Release | ResourceAction::Rollback, Some(amount), Some(before))
1656                if amount > before =>
1657            {
1658                Some(amount.saturating_sub(before))
1659            }
1660            _ => None,
1661        };
1662        if let Some(underflow_amount) = underflow_amount {
1663            attributes.insert(
1664                "resource_underflow_amount".to_string(),
1665                serde_json::json!(underflow_amount),
1666            );
1667        }
1668        if matches!(
1669            action,
1670            ResourceAction::RequestOpen
1671                | ResourceAction::RequestClose
1672                | ResourceAction::Defer
1673                | ResourceAction::Reject
1674        ) {
1675            self.extend_scheduler_timeline_attributes(&mut attributes);
1676        }
1677        let timestamp = chrono::Utc::now();
1678        let mut shape =
1679            BTreeMap::from([("resource_amount".to_string(), serde_json::json!(amount))]);
1680        if let Some(capacity) = capacity {
1681            shape.insert("resource_capacity".to_string(), serde_json::json!(capacity));
1682        }
1683        let event = FerrumProfileEvent {
1684            schema_version: OBSERVABILITY_PROFILE_SCHEMA_VERSION,
1685            ts_unix_nanos: timestamp
1686                .timestamp_nanos_opt()
1687                .unwrap_or_else(|| timestamp.timestamp_micros() * 1_000),
1688            event_id: format!("evt-engine-resource-{event_num}"),
1689            request_id: request_id.to_string(),
1690            correlation_id: Some(request_id.to_string()),
1691            entrypoint,
1692            backend: "actual".to_string(),
1693            runtime_preset_hash: ENGINE_RUNTIME_TRACE_PRESET_HASH.to_string(),
1694            phase: phase.to_string(),
1695            event_kind: ProfileEventKind::Resource,
1696            timestamp,
1697            status: ProfileStatus::Ok,
1698            model: Some(self.config.model.model_id.to_string()),
1699            duration_us: None,
1700            memory: None,
1701            resource: Some(ResourceTraceEvent {
1702                owner_kind: owner_kind.to_string(),
1703                owner_id: owner_id.to_string(),
1704                resource_kind: resource_kind.to_string(),
1705                action,
1706                amount,
1707                before,
1708                after,
1709                capacity,
1710                underflow_amount,
1711                reason,
1712                error_kind: None,
1713                message: None,
1714                resource_error_kind: None,
1715            }),
1716            error: None,
1717            replay: None,
1718            shape,
1719            backend_detail: Some(BTreeMap::from([
1720                (
1721                    "backend_device".to_string(),
1722                    serde_json::json!(format!("{:?}", self.config.backend.device)),
1723                ),
1724                (
1725                    "backend_type".to_string(),
1726                    serde_json::json!(format!("{:?}", self.config.backend.backend_type)),
1727                ),
1728            ])),
1729            attributes,
1730        };
1731        if let Err(error) = event.validate() {
1732            warn!("Skipping invalid engine resource trace event: {}", error);
1733            return;
1734        }
1735        if let Err(error) = sink.enqueue(event) {
1736            warn!("Failed to enqueue engine resource trace event: {}", error);
1737        }
1738    }
1739
1740    #[allow(clippy::too_many_arguments)]
1741    fn trace_resource_event_with_close_summary(
1742        &self,
1743        request_id: &RequestId,
1744        owner_kind: &str,
1745        owner_id: &str,
1746        resource_kind: &str,
1747        phase: &str,
1748        action: ResourceAction,
1749        close_summary: &[ResourceOwnerCloseSummary],
1750        status: ProfileStatus,
1751        message: Option<String>,
1752    ) {
1753        let Some(sink) = &self.scheduler_trace_jsonl else {
1754            return;
1755        };
1756        let entrypoint = self.trace_entrypoint();
1757        let event_num = self
1758            .resource_trace_event_counter
1759            .fetch_add(1, Ordering::Relaxed);
1760        let outstanding: Vec<_> = close_summary
1761            .iter()
1762            .filter(|item| item.outstanding_reserved > 0 || item.outstanding_committed > 0)
1763            .collect();
1764        let close_summary_json: Vec<_> = close_summary
1765            .iter()
1766            .map(|item| {
1767                serde_json::json!({
1768                    "resource_kind": item.resource_kind,
1769                    "reserved": item.reserved,
1770                    "committed": item.committed,
1771                    "released": item.released,
1772                    "rolled_back": item.rolled_back,
1773                    "outstanding_reserved": item.outstanding_reserved,
1774                    "outstanding_committed": item.outstanding_committed,
1775                    "capacity": item.capacity,
1776                })
1777            })
1778            .collect();
1779        let outstanding_kinds: Vec<_> = outstanding
1780            .iter()
1781            .map(|item| item.resource_kind.clone())
1782            .collect();
1783        let mut attributes = BTreeMap::from([
1784            (
1785                "actual_model_smoke".to_string(),
1786                serde_json::json!(matches!(
1787                    entrypoint,
1788                    ProfileEntrypoint::Run | ProfileEntrypoint::Serve
1789                )),
1790            ),
1791            (
1792                "backend_device".to_string(),
1793                serde_json::json!(format!("{:?}", self.config.backend.device)),
1794            ),
1795            (
1796                "backend_type".to_string(),
1797                serde_json::json!(format!("{:?}", self.config.backend.backend_type)),
1798            ),
1799            (
1800                "diagnostic_only".to_string(),
1801                serde_json::json!(self.config.runtime.profile_detail.diagnostic_only()),
1802            ),
1803            ("l0_only".to_string(), serde_json::json!(false)),
1804            (
1805                "profile_detail".to_string(),
1806                serde_json::json!(self.config.runtime.profile_detail.as_str()),
1807            ),
1808            (
1809                "resource_owner_close_summary".to_string(),
1810                serde_json::Value::Array(close_summary_json),
1811            ),
1812            (
1813                "resource_owner_outstanding_count".to_string(),
1814                serde_json::json!(outstanding.len()),
1815            ),
1816            (
1817                "resource_owner_outstanding_kinds".to_string(),
1818                serde_json::json!(outstanding_kinds),
1819            ),
1820            (
1821                "resource_trace_source".to_string(),
1822                serde_json::json!("engine"),
1823            ),
1824        ]);
1825        if let Some(message) = message.as_deref() {
1826            attributes.insert(
1827                "resource_close_error".to_string(),
1828                serde_json::json!(message),
1829            );
1830        }
1831        self.extend_scheduler_timeline_attributes(&mut attributes);
1832        let timestamp = chrono::Utc::now();
1833        let error = message.as_ref().map(|message| ProfileError {
1834            kind: "resource_owner_close_outstanding".to_string(),
1835            message: message.clone(),
1836            blocking: true,
1837        });
1838        let resource_error_kind = error.as_ref().map(|_| "resource_leak".to_string());
1839        let mut shape = BTreeMap::from([("resource_amount".to_string(), serde_json::Value::Null)]);
1840        shape.insert(
1841            "resource_owner_outstanding_count".to_string(),
1842            serde_json::json!(outstanding.len()),
1843        );
1844        let event = FerrumProfileEvent {
1845            schema_version: OBSERVABILITY_PROFILE_SCHEMA_VERSION,
1846            ts_unix_nanos: timestamp
1847                .timestamp_nanos_opt()
1848                .unwrap_or_else(|| timestamp.timestamp_micros() * 1_000),
1849            event_id: format!("evt-engine-resource-{event_num}"),
1850            request_id: request_id.to_string(),
1851            correlation_id: Some(request_id.to_string()),
1852            entrypoint,
1853            backend: "actual".to_string(),
1854            runtime_preset_hash: ENGINE_RUNTIME_TRACE_PRESET_HASH.to_string(),
1855            phase: phase.to_string(),
1856            event_kind: ProfileEventKind::Resource,
1857            timestamp,
1858            status,
1859            model: Some(self.config.model.model_id.to_string()),
1860            duration_us: None,
1861            memory: None,
1862            resource: Some(ResourceTraceEvent {
1863                owner_kind: owner_kind.to_string(),
1864                owner_id: owner_id.to_string(),
1865                resource_kind: resource_kind.to_string(),
1866                action,
1867                amount: None,
1868                before: None,
1869                after: None,
1870                capacity: None,
1871                underflow_amount: None,
1872                reason: None,
1873                error_kind: error.as_ref().map(|error| error.kind.clone()),
1874                message: error.as_ref().map(|error| error.message.clone()),
1875                resource_error_kind,
1876            }),
1877            error,
1878            replay: None,
1879            shape,
1880            backend_detail: Some(BTreeMap::from([
1881                (
1882                    "backend_device".to_string(),
1883                    serde_json::json!(format!("{:?}", self.config.backend.device)),
1884                ),
1885                (
1886                    "backend_type".to_string(),
1887                    serde_json::json!(format!("{:?}", self.config.backend.backend_type)),
1888                ),
1889            ])),
1890            attributes,
1891        };
1892        if let Err(error) = event.validate() {
1893            warn!(
1894                "Skipping invalid engine resource close trace event: {}",
1895                error
1896            );
1897            return;
1898        }
1899        if let Err(error) = sink.enqueue(event) {
1900            warn!(
1901                "Failed to enqueue engine resource close trace event: {}",
1902                error
1903            );
1904        }
1905    }
1906
1907    fn resource_amount_i64(amount: usize) -> i64 {
1908        amount.min(i64::MAX as usize) as i64
1909    }
1910
1911    fn trace_lifecycle_resource_event(
1912        &self,
1913        request_id: &RequestId,
1914        owner_kind: &str,
1915        owner_id: &str,
1916        resource_kind: &str,
1917        phase: &str,
1918        action: ResourceAction,
1919        amount: i64,
1920        transition: ResourceLedgerTransition,
1921    ) {
1922        self.trace_resource_event(
1923            request_id,
1924            owner_kind,
1925            owner_id,
1926            resource_kind,
1927            phase,
1928            action,
1929            Some(amount),
1930            Some(transition.before),
1931            Some(transition.after),
1932            transition.capacity,
1933            None,
1934        );
1935    }
1936
1937    fn trace_request_open(&self, request_id: &RequestId) {
1938        self.trace_resource_event(
1939            request_id,
1940            "request",
1941            &request_id.to_string(),
1942            "request_slot",
1943            "engine_request_open",
1944            ResourceAction::RequestOpen,
1945            None,
1946            None,
1947            None,
1948            None,
1949            None,
1950        );
1951    }
1952
1953    fn trace_request_admitted(&self, request_id: &RequestId) {
1954        self.trace_resource_reserve_commit(
1955            request_id,
1956            "request",
1957            &request_id.to_string(),
1958            "request_slot",
1959            "engine_request_slot",
1960            1,
1961            None,
1962        );
1963    }
1964
1965    fn trace_request_rejected(&self, request_id: &RequestId, reason: String) {
1966        self.trace_resource_event(
1967            request_id,
1968            "request",
1969            &request_id.to_string(),
1970            "request_slot",
1971            "engine_request_reject",
1972            ResourceAction::Reject,
1973            None,
1974            None,
1975            None,
1976            Some(Self::resource_amount_i64(
1977                self.config.scheduler.max_waiting_requests,
1978            )),
1979            Some(reason),
1980        );
1981        self.trace_request_owner_close(request_id);
1982    }
1983
1984    fn trace_request_close(&self, request_id: &RequestId) {
1985        self.trace_resource_release(
1986            request_id,
1987            "request",
1988            &request_id.to_string(),
1989            "request_slot",
1990            "engine_request_slot_release",
1991            1,
1992            None,
1993        );
1994        self.trace_request_owner_close(request_id);
1995    }
1996
1997    fn trace_request_owner_close(&self, request_id: &RequestId) {
1998        let owner_id = request_id.to_string();
1999        if self.scheduler_trace_jsonl.is_none() {
2000            self.trace_resource_event(
2001                request_id,
2002                "request",
2003                &owner_id,
2004                "request_slot",
2005                "engine_request_close",
2006                ResourceAction::RequestClose,
2007                None,
2008                None,
2009                None,
2010                None,
2011                None,
2012            );
2013            return;
2014        }
2015
2016        let summary = {
2017            let mut lifecycle = self.resource_lifecycle.lock();
2018            let summary = lifecycle.owner_close_summary("request", &owner_id);
2019            lifecycle.close_owner("request", &owner_id);
2020            summary
2021        };
2022        self.trace_request_owner_close_with_summary(request_id, &summary);
2023    }
2024
2025    fn trace_request_owner_close_with_summary(
2026        &self,
2027        request_id: &RequestId,
2028        summary: &[ResourceOwnerCloseSummary],
2029    ) {
2030        let outstanding: Vec<_> = summary
2031            .iter()
2032            .filter(|item| item.outstanding_reserved > 0 || item.outstanding_committed > 0)
2033            .collect();
2034        let close_status = if outstanding.is_empty() {
2035            ProfileStatus::Ok
2036        } else {
2037            ProfileStatus::Failure
2038        };
2039        let message = if outstanding.is_empty() {
2040            None
2041        } else {
2042            Some(format!(
2043                "request closed with outstanding resources: {}",
2044                outstanding
2045                    .iter()
2046                    .map(|item| format!(
2047                        "{} reserved={} committed={}",
2048                        item.resource_kind, item.outstanding_reserved, item.outstanding_committed
2049                    ))
2050                    .collect::<Vec<_>>()
2051                    .join(", ")
2052            ))
2053        };
2054        self.trace_resource_event_with_close_summary(
2055            request_id,
2056            "request",
2057            &request_id.to_string(),
2058            "request_slot",
2059            "engine_request_close",
2060            ResourceAction::RequestClose,
2061            summary,
2062            close_status,
2063            message,
2064        );
2065    }
2066
2067    fn trace_scheduler_defer(&self, request_id: &RequestId, phase: &str, reason: &str) {
2068        self.trace_resource_event(
2069            request_id,
2070            "request",
2071            &request_id.to_string(),
2072            "scheduler_capacity",
2073            phase,
2074            ResourceAction::Defer,
2075            None,
2076            None,
2077            None,
2078            Some(Self::resource_amount_i64(
2079                self.config.scheduler.max_running_requests.max(1),
2080            )),
2081            Some(reason.to_string()),
2082        );
2083    }
2084
2085    fn trace_resource_reserve_commit(
2086        &self,
2087        request_id: &RequestId,
2088        owner_kind: &str,
2089        owner_id: &str,
2090        resource_kind: &str,
2091        phase_prefix: &str,
2092        amount: usize,
2093        capacity: Option<usize>,
2094    ) {
2095        if self.scheduler_trace_jsonl.is_none() {
2096            return;
2097        }
2098        let amount = Self::resource_amount_i64(amount.max(1));
2099        let capacity_i64 = capacity.map(Self::resource_amount_i64);
2100        let (reserve, commit) = {
2101            let mut lifecycle = self.resource_lifecycle.lock();
2102            let reserve =
2103                lifecycle.reserve(owner_kind, owner_id, resource_kind, amount, capacity_i64);
2104            let commit =
2105                lifecycle.commit(owner_kind, owner_id, resource_kind, amount, capacity_i64);
2106            (reserve, commit)
2107        };
2108        self.trace_lifecycle_resource_event(
2109            request_id,
2110            owner_kind,
2111            owner_id,
2112            resource_kind,
2113            &format!("{phase_prefix}_reserve"),
2114            ResourceAction::Reserve,
2115            amount,
2116            reserve,
2117        );
2118        self.trace_lifecycle_resource_event(
2119            request_id,
2120            owner_kind,
2121            owner_id,
2122            resource_kind,
2123            &format!("{phase_prefix}_commit"),
2124            ResourceAction::Commit,
2125            amount,
2126            commit,
2127        );
2128    }
2129
2130    fn trace_resource_release(
2131        &self,
2132        request_id: &RequestId,
2133        owner_kind: &str,
2134        owner_id: &str,
2135        resource_kind: &str,
2136        phase: &str,
2137        amount: usize,
2138        capacity: Option<usize>,
2139    ) {
2140        if self.scheduler_trace_jsonl.is_none() {
2141            return;
2142        }
2143        let amount = Self::resource_amount_i64(amount.max(1));
2144        let transition = self.resource_lifecycle.lock().release(
2145            owner_kind,
2146            owner_id,
2147            resource_kind,
2148            amount,
2149            capacity.map(Self::resource_amount_i64),
2150        );
2151        self.trace_lifecycle_resource_event(
2152            request_id,
2153            owner_kind,
2154            owner_id,
2155            resource_kind,
2156            phase,
2157            ResourceAction::Release,
2158            amount,
2159            transition,
2160        );
2161    }
2162
2163    fn trace_resource_release_failure(
2164        &self,
2165        request_id: &RequestId,
2166        resource_kind: &str,
2167        phase: &str,
2168        capacity: Option<usize>,
2169        reason: String,
2170    ) {
2171        self.trace_resource_event(
2172            request_id,
2173            "request",
2174            &request_id.to_string(),
2175            resource_kind,
2176            phase,
2177            ResourceAction::Reject,
2178            None,
2179            None,
2180            None,
2181            capacity.map(Self::resource_amount_i64),
2182            Some(reason),
2183        );
2184    }
2185
2186    fn kv_resource_blocks_for_tokens(&self, tokens: usize) -> usize {
2187        tokens
2188            .div_ceil(self.config.kv_cache.block_size.max(1))
2189            .max(1)
2190    }
2191
2192    fn trace_kv_allocate(&self, request_id: &RequestId, blocks: usize) {
2193        self.trace_resource_reserve_commit(
2194            request_id,
2195            "request",
2196            &request_id.to_string(),
2197            "kv_block",
2198            "engine_kv_block",
2199            blocks,
2200            Some(self.config.kv_cache.max_blocks),
2201        );
2202    }
2203
2204    async fn allocate_kv_lease(
2205        &self,
2206        owner_request_id: &RequestId,
2207        allocation_request_id: RequestId,
2208        request: &AllocationRequest,
2209        tokens: usize,
2210    ) -> Result<KvAllocationLease> {
2211        debug_assert_eq!(allocation_request_id, request.request_id);
2212        let handle = self.engine_managed_kv_cache()?.allocate(request).await?;
2213        let blocks = self.kv_resource_blocks_for_tokens(tokens);
2214        self.trace_kv_allocate(owner_request_id, blocks);
2215        Ok(KvAllocationLease::new(
2216            owner_request_id.clone(),
2217            allocation_request_id,
2218            handle,
2219            blocks,
2220        ))
2221    }
2222
2223    fn trace_kv_release(&self, request_id: &RequestId, blocks: usize) {
2224        self.trace_resource_release(
2225            request_id,
2226            "request",
2227            &request_id.to_string(),
2228            "kv_block",
2229            "engine_kv_block_release",
2230            blocks,
2231            Some(self.config.kv_cache.max_blocks),
2232        );
2233    }
2234
2235    fn trace_model_cache_ref_acquire(&self, request_id: &RequestId) {
2236        self.trace_resource_reserve_commit(
2237            request_id,
2238            "request",
2239            &request_id.to_string(),
2240            "model_cache_ref",
2241            "engine_model_cache_ref",
2242            1,
2243            None,
2244        );
2245    }
2246
2247    fn trace_model_cache_ref_release(&self, request_id: &RequestId) {
2248        self.trace_resource_release(
2249            request_id,
2250            "request",
2251            &request_id.to_string(),
2252            "model_cache_ref",
2253            "engine_model_cache_ref_release",
2254            1,
2255            None,
2256        );
2257    }
2258
2259    fn legacy_backend_workspace_trace_capacity(&self) -> Option<usize> {
2260        Some(self.config.scheduler.max_running_requests.max(1))
2261    }
2262
2263    fn trace_legacy_backend_workspace_acquire(&self, request_id: &RequestId, phase_prefix: &str) {
2264        self.trace_resource_reserve_commit(
2265            request_id,
2266            "request",
2267            &request_id.to_string(),
2268            "backend_workspace",
2269            phase_prefix,
2270            1,
2271            self.legacy_backend_workspace_trace_capacity(),
2272        );
2273    }
2274
2275    fn trace_legacy_backend_workspace_release(&self, request_id: &RequestId, phase: &str) {
2276        self.trace_resource_release(
2277            request_id,
2278            "request",
2279            &request_id.to_string(),
2280            "backend_workspace",
2281            phase,
2282            1,
2283            self.legacy_backend_workspace_trace_capacity(),
2284        );
2285    }
2286
2287    fn trace_legacy_backend_workspace_acquire_many(
2288        &self,
2289        request_ids: &[RequestId],
2290        phase_prefix: &str,
2291    ) {
2292        for request_id in request_ids {
2293            self.trace_legacy_backend_workspace_acquire(request_id, phase_prefix);
2294        }
2295    }
2296
2297    fn trace_legacy_backend_workspace_release_many(&self, request_ids: &[RequestId], phase: &str) {
2298        for request_id in request_ids {
2299            self.trace_legacy_backend_workspace_release(request_id, phase);
2300        }
2301    }
2302
2303    fn acquire_legacy_backend_workspace_trace_lease(
2304        &self,
2305        request_ids: Vec<RequestId>,
2306        phase_prefix: &'static str,
2307        release_phase: &'static str,
2308    ) -> Result<LegacyBackendWorkspaceTraceLease<'_>> {
2309        if self.resource_composition.authority() != ExecutionResourceAuthority::LegacyEngine {
2310            return Err(FerrumError::internal(
2311                "synthetic backend workspace tracing is forbidden for PlanRuntime authority",
2312            ));
2313        }
2314        Ok(LegacyBackendWorkspaceTraceLease::new(
2315            self,
2316            request_ids,
2317            phase_prefix,
2318            release_phase,
2319        ))
2320    }
2321
2322    fn apply_model_cache_ref_update(&self, request_id: &RequestId, update: ModelCacheRefUpdate) {
2323        if let Some(cache_id) = update.released {
2324            self.model_executor.release_cache(&cache_id);
2325            self.trace_model_cache_ref_release(request_id);
2326        }
2327        if update.acquired.is_some() {
2328            self.trace_model_cache_ref_acquire(request_id);
2329        }
2330    }
2331
2332    fn release_model_cache_ref(&self, request_id: &RequestId, cache_id: &str) {
2333        self.model_executor.release_cache(cache_id);
2334        self.trace_model_cache_ref_release(request_id);
2335    }
2336
2337    async fn release_kv_allocation(
2338        &self,
2339        owner_request_id: &RequestId,
2340        allocation_request_id: RequestId,
2341        blocks: usize,
2342    ) {
2343        let kv_cache = match self.engine_managed_kv_cache() {
2344            Ok(kv_cache) => kv_cache,
2345            Err(error) => {
2346                warn!(
2347                    owner_request_id = %owner_request_id,
2348                    allocation_request_id = %allocation_request_id,
2349                    error = %error,
2350                    "Legacy engine KV allocation reached a plan-runtime composition"
2351                );
2352                return;
2353            }
2354        };
2355        match kv_cache.deallocate(allocation_request_id.clone()).await {
2356            Ok(()) => {
2357                self.trace_kv_release(owner_request_id, blocks);
2358            }
2359            Err(error) => {
2360                warn!(
2361                    owner_request_id = %owner_request_id,
2362                    allocation_request_id = %allocation_request_id,
2363                    error = %error,
2364                    "KV allocation release failed"
2365                );
2366                self.trace_resource_release_failure(
2367                    owner_request_id,
2368                    "kv_block",
2369                    "engine_kv_block_release_failed",
2370                    Some(self.config.kv_cache.max_blocks),
2371                    format!("kv release failed for {allocation_request_id}: {error}"),
2372                );
2373            }
2374        }
2375    }
2376
2377    async fn release_sequence_physical_resources(
2378        &self,
2379        request_id: &RequestId,
2380        resources: SequencePhysicalResources,
2381    ) {
2382        if let Some(cache_id) = resources.model_cache_id {
2383            self.release_model_cache_ref(request_id, &cache_id);
2384        }
2385        if let Some(kv_allocation) = resources.legacy_kv_allocation {
2386            self.release_kv_allocation(request_id, kv_allocation.request_id, kv_allocation.blocks)
2387                .await;
2388        }
2389        if let Some(draft_kv_allocation) = resources.legacy_draft_kv_allocation {
2390            self.release_kv_allocation(
2391                request_id,
2392                draft_kv_allocation.request_id,
2393                draft_kv_allocation.blocks,
2394            )
2395            .await;
2396        }
2397        if let Some(recurrent_allocation) = resources.recurrent_state_allocation {
2398            self.release_recurrent_allocation(request_id, recurrent_allocation.slots)
2399                .await;
2400        }
2401    }
2402
2403    async fn complete_sequence_physical_resources(
2404        &self,
2405        request_id: &RequestId,
2406        mut resources: SequencePhysicalResources,
2407        usage: &TokenUsage,
2408    ) -> Result<()> {
2409        let completion_result = if let Some(cache_id) = resources.model_cache_id.take() {
2410            let completion = ExecutorSequenceCompletion::new(
2411                request_id.clone(),
2412                cache_id.clone(),
2413                usage.prompt_tokens,
2414                usage.completion_tokens,
2415            );
2416            let result = match completion {
2417                Ok(completion) => self.model_executor.complete_cache(completion).await,
2418                Err(error) => {
2419                    self.model_executor.release_cache(&cache_id);
2420                    Err(error)
2421                }
2422            };
2423            self.trace_model_cache_ref_release(request_id);
2424            result
2425        } else {
2426            Ok(())
2427        };
2428
2429        self.release_sequence_physical_resources(request_id, resources)
2430            .await;
2431        completion_result
2432    }
2433
2434    fn trace_recurrent_allocate(
2435        &self,
2436        request_id: &RequestId,
2437        slots: usize,
2438        capacity: Option<usize>,
2439    ) {
2440        self.trace_resource_reserve_commit(
2441            request_id,
2442            "request",
2443            &request_id.to_string(),
2444            "recurrent_state_slot",
2445            "engine_recurrent_state_slot",
2446            slots,
2447            capacity,
2448        );
2449    }
2450
2451    fn trace_recurrent_release(
2452        &self,
2453        request_id: &RequestId,
2454        slots: usize,
2455        capacity: Option<usize>,
2456    ) {
2457        self.trace_resource_release(
2458            request_id,
2459            "request",
2460            &request_id.to_string(),
2461            "recurrent_state_slot",
2462            "engine_recurrent_state_slot_release",
2463            slots,
2464            capacity,
2465        );
2466    }
2467
2468    async fn release_recurrent_allocation(&self, request_id: &RequestId, slots: Option<usize>) {
2469        if let Some(manager) = self.recurrent_state_manager() {
2470            let capacity = manager.stats().total_batch_slots;
2471            match manager.deallocate(request_id.clone()).await {
2472                Ok(()) => {
2473                    if let Some(slots) = slots {
2474                        self.trace_recurrent_release(request_id, slots, Some(capacity));
2475                    }
2476                }
2477                Err(error) => {
2478                    warn!(
2479                        request_id = %request_id,
2480                        error = %error,
2481                        "Recurrent-state release failed"
2482                    );
2483                    if slots.is_some() {
2484                        self.trace_resource_release_failure(
2485                            request_id,
2486                            "recurrent_state_slot",
2487                            "engine_recurrent_state_slot_release_failed",
2488                            Some(capacity),
2489                            format!("recurrent-state release failed for {request_id}: {error}"),
2490                        );
2491                    }
2492                }
2493            }
2494        }
2495    }
2496
2497    async fn prepare_recurrent_state(
2498        &self,
2499        request_id: &RequestId,
2500        spec: Option<ferrum_interfaces::RecurrentStateSpec>,
2501    ) -> Result<RecurrentStateAdmission> {
2502        if let Some(existing) = self
2503            .sequences
2504            .read()
2505            .get(request_id)
2506            .and_then(SequenceState::recurrent_state_handle)
2507        {
2508            return Ok(RecurrentStateAdmission::existing(existing));
2509        }
2510
2511        let Some(spec) = spec else {
2512            return Ok(RecurrentStateAdmission::none());
2513        };
2514
2515        debug_assert_eq!(&spec.request_id, request_id);
2516        let Some(manager) = self.recurrent_state_manager() else {
2517            return Err(FerrumError::config(format!(
2518                "model '{}' requires recurrent state for request {}, but no recurrent-state manager is configured",
2519                self.model_executor.info().model_id, request_id
2520            )));
2521        };
2522
2523        let before_stats = manager.stats();
2524        let slots = spec.max_batch_slots.max(1);
2525        let handle = match manager.allocate(&spec).await {
2526            Ok(handle) => handle,
2527            Err(error) => {
2528                self.trace_resource_event(
2529                    request_id,
2530                    "request",
2531                    &request_id.to_string(),
2532                    "recurrent_state_slot",
2533                    "engine_recurrent_state_slot_reject",
2534                    ResourceAction::Reject,
2535                    None,
2536                    None,
2537                    None,
2538                    Some(Self::resource_amount_i64(before_stats.total_batch_slots)),
2539                    Some(error.to_string()),
2540                );
2541                return Err(error);
2542            }
2543        };
2544        let after_stats = manager.stats();
2545        self.trace_recurrent_allocate(request_id, slots, Some(after_stats.total_batch_slots));
2546        Ok(RecurrentStateAdmission::fresh(RecurrentStateLease::new(
2547            request_id.clone(),
2548            handle,
2549            slots,
2550            Some(after_stats.total_batch_slots),
2551        )))
2552    }
2553
2554    async fn ensure_recurrent_state(
2555        &self,
2556        request_id: &RequestId,
2557        spec: Option<ferrum_interfaces::RecurrentStateSpec>,
2558    ) -> Result<Option<Arc<dyn RecurrentStateHandle>>> {
2559        let mut admission = self.prepare_recurrent_state(request_id, spec).await?;
2560        let handle = admission.handle();
2561        if let Some(slots) = admission.fresh_slots() {
2562            let Some(handle) = handle.clone() else {
2563                admission.release_fresh(self).await;
2564                return Err(FerrumError::internal(format!(
2565                    "missing recurrent state handle while committing recurrent slots for {request_id}"
2566                )));
2567            };
2568            let mut found = false;
2569            {
2570                let mut sequences = self.sequences.write();
2571                if let Some(seq) = sequences.get_mut(request_id) {
2572                    seq.commit_recurrent_state_admission(handle, slots);
2573                    found = true;
2574                }
2575            }
2576            if found {
2577                admission.commit_fresh();
2578            } else {
2579                admission.release_fresh(self).await;
2580                return Err(FerrumError::internal(format!(
2581                    "sequence not found while committing recurrent state for {request_id}"
2582                )));
2583            }
2584        }
2585
2586        Ok(handle)
2587    }
2588
2589    fn performance_breakdown(&self) -> ferrum_types::PerformanceBreakdown {
2590        ferrum_types::PerformanceBreakdown {
2591            scheduling_time_ms: avg_duration_ms(
2592                self.total_scheduling_time_us.load(Ordering::Relaxed),
2593                self.scheduling_time_samples.load(Ordering::Relaxed),
2594            ),
2595            model_execution_time_ms: avg_duration_ms(
2596                self.total_model_execution_time_us.load(Ordering::Relaxed),
2597                self.model_execution_time_samples.load(Ordering::Relaxed),
2598            ),
2599            other_overhead_time_ms: avg_duration_ms(
2600                self.total_iteration_lock_wait_us.load(Ordering::Relaxed),
2601                self.iteration_lock_wait_samples.load(Ordering::Relaxed),
2602            ),
2603            ..Default::default()
2604        }
2605    }
2606}
2607
2608fn duration_to_us(duration: Duration) -> u64 {
2609    duration.as_micros().min(u64::MAX as u128) as u64
2610}
2611
2612fn avg_duration_ms(total_us: u64, samples: u64) -> f64 {
2613    if samples == 0 {
2614        0.0
2615    } else {
2616        total_us as f64 / samples as f64 / 1000.0
2617    }
2618}
2619
2620mod profile;
2621use profile::*;
2622
2623mod inner;
2624
2625// ────────────────────────────────────────────────────────────────────────────
2626// Public engine wrapper
2627// ────────────────────────────────────────────────────────────────────────────
2628
2629/// Continuous batching inference engine.
2630///
2631/// Wraps an `Arc<EngineInner>` so it can be cloned and shared freely.
2632/// Multiple concurrent `infer()` / `infer_stream()` calls are safe —
2633/// an internal `iteration_lock` serializes engine steps while allowing
2634/// all pending requests to be processed in each iteration's batch.
2635pub struct ContinuousBatchEngine {
2636    inner: Arc<EngineInner>,
2637}
2638
2639impl ContinuousBatchEngine {
2640    pub fn new(
2641        config: EngineConfig,
2642        scheduler: Arc<ContinuousBatchScheduler>,
2643        tokenizer: Arc<dyn Tokenizer + Send + Sync>,
2644        sampler: Arc<dyn Sampler + Send + Sync>,
2645        kv_cache: Arc<dyn KvCacheManager + Send + Sync>,
2646        model_executor: Arc<dyn ModelExecutor + Send + Sync>,
2647        tensor_factory: Arc<dyn TensorFactory>,
2648    ) -> Result<Self> {
2649        Self::new_with_speculation(
2650            config,
2651            scheduler,
2652            tokenizer,
2653            sampler,
2654            kv_cache,
2655            model_executor,
2656            tensor_factory,
2657            None,
2658            None,
2659        )
2660    }
2661
2662    /// Build an engine with optional speculative decoding. Pass both the
2663    /// draft executor AND the config together — either both or neither.
2664    #[allow(clippy::too_many_arguments)]
2665    pub fn new_with_speculation(
2666        config: EngineConfig,
2667        scheduler: Arc<ContinuousBatchScheduler>,
2668        tokenizer: Arc<dyn Tokenizer + Send + Sync>,
2669        sampler: Arc<dyn Sampler + Send + Sync>,
2670        kv_cache: Arc<dyn KvCacheManager + Send + Sync>,
2671        model_executor: Arc<dyn ModelExecutor + Send + Sync>,
2672        tensor_factory: Arc<dyn TensorFactory>,
2673        draft_executor: Option<Arc<dyn ModelExecutor + Send + Sync>>,
2674        spec_config: Option<crate::speculative::SpeculativeDecodingConfig>,
2675    ) -> Result<Self> {
2676        Self::new_with_speculation_and_recurrent_state_manager(
2677            config,
2678            scheduler,
2679            tokenizer,
2680            sampler,
2681            kv_cache,
2682            model_executor,
2683            tensor_factory,
2684            draft_executor,
2685            spec_config,
2686            None,
2687        )
2688    }
2689
2690    /// Build an engine with optional speculative decoding and an optional
2691    /// recurrent-state manager for state-space / hybrid models.
2692    #[allow(clippy::too_many_arguments)]
2693    pub fn new_with_speculation_and_recurrent_state_manager(
2694        config: EngineConfig,
2695        scheduler: Arc<ContinuousBatchScheduler>,
2696        tokenizer: Arc<dyn Tokenizer + Send + Sync>,
2697        sampler: Arc<dyn Sampler + Send + Sync>,
2698        kv_cache: Arc<dyn KvCacheManager + Send + Sync>,
2699        model_executor: Arc<dyn ModelExecutor + Send + Sync>,
2700        tensor_factory: Arc<dyn TensorFactory>,
2701        draft_executor: Option<Arc<dyn ModelExecutor + Send + Sync>>,
2702        spec_config: Option<crate::speculative::SpeculativeDecodingConfig>,
2703        recurrent_state_manager: Option<Arc<dyn RecurrentStateManager + Send + Sync>>,
2704    ) -> Result<Self> {
2705        Self::new_with_resource_composition(
2706            config,
2707            scheduler,
2708            tokenizer,
2709            sampler,
2710            EngineResourceComposition::LegacyEngine {
2711                kv_cache,
2712                recurrent_state_manager,
2713            },
2714            model_executor,
2715            tensor_factory,
2716            draft_executor,
2717            spec_config,
2718        )
2719    }
2720
2721    /// Build an engine bound to the shared plan runtime, which is the sole
2722    /// owner of request-lifetime KV, recurrent state, and backing capacity.
2723    /// The model executor adapts that runtime but does not own a second
2724    /// resource manager; no legacy engine manager is created or retained.
2725    pub fn new_plan_runtime(
2726        config: EngineConfig,
2727        scheduler: Arc<ContinuousBatchScheduler>,
2728        tokenizer: Arc<dyn Tokenizer + Send + Sync>,
2729        sampler: Arc<dyn Sampler + Send + Sync>,
2730        model_executor: Arc<dyn ModelExecutor + Send + Sync>,
2731        tensor_factory: Arc<dyn TensorFactory>,
2732    ) -> Result<Self> {
2733        Self::new_with_resource_composition(
2734            config,
2735            scheduler,
2736            tokenizer,
2737            sampler,
2738            EngineResourceComposition::PlanRuntime,
2739            model_executor,
2740            tensor_factory,
2741            None,
2742            None,
2743        )
2744    }
2745
2746    #[allow(clippy::too_many_arguments)]
2747    fn new_with_resource_composition(
2748        config: EngineConfig,
2749        scheduler: Arc<ContinuousBatchScheduler>,
2750        tokenizer: Arc<dyn Tokenizer + Send + Sync>,
2751        sampler: Arc<dyn Sampler + Send + Sync>,
2752        resource_composition: EngineResourceComposition,
2753        model_executor: Arc<dyn ModelExecutor + Send + Sync>,
2754        tensor_factory: Arc<dyn TensorFactory>,
2755        draft_executor: Option<Arc<dyn ModelExecutor + Send + Sync>>,
2756        spec_config: Option<crate::speculative::SpeculativeDecodingConfig>,
2757    ) -> Result<Self> {
2758        let executor_authority = model_executor.execution_resource_authority();
2759        if draft_executor.is_some() != spec_config.is_some() {
2760            return Err(FerrumError::config(
2761                "speculative decoding requires both a draft executor and its configuration",
2762            ));
2763        }
2764        if let Some(draft_executor) = draft_executor.as_ref() {
2765            let draft_authority = draft_executor.execution_resource_authority();
2766            if draft_authority != executor_authority {
2767                return Err(FerrumError::config(format!(
2768                    "draft executor authority {draft_authority:?} does not match target authority {executor_authority:?}"
2769                )));
2770            }
2771        }
2772        if resource_composition.authority() != executor_authority {
2773            return Err(FerrumError::config(format!(
2774                "engine resource composition {:?} does not match executor authority {:?}",
2775                resource_composition.authority(),
2776                executor_authority
2777            )));
2778        }
2779        let recurrent_state_manager = resource_composition.recurrent_state_manager().is_some();
2780        info!(
2781            ?executor_authority,
2782            "Creating ContinuousBatchEngine (speculative_decoding={}, recurrent_state_manager={})",
2783            draft_executor.is_some() && spec_config.is_some(),
2784            recurrent_state_manager
2785        );
2786        let runtime_config = ContinuousEngineRuntimeConfig::from_engine_config(&config);
2787        let profile_trace_jsonl = runtime_config
2788            .profile_jsonl
2789            .as_ref()
2790            .map(|path| {
2791                SchedulerTraceJournal::create(path.clone()).map_err(|error| {
2792                    FerrumError::io(format!(
2793                        "open product profile JSONL {}: {error}",
2794                        path.display()
2795                    ))
2796                })
2797            })
2798            .transpose()?;
2799        let scheduler_trace_jsonl = match runtime_config.scheduler_trace_jsonl.as_deref() {
2800            Some(path)
2801                if profile_trace_jsonl
2802                    .as_ref()
2803                    .is_some_and(|journal| journal.path() == path) =>
2804            {
2805                profile_trace_jsonl.clone()
2806            }
2807            path => create_scheduler_trace_sink(path),
2808        };
2809        let legacy_scheduler_trace_jsonl = create_legacy_scheduler_trace_sink(
2810            runtime_config.legacy_scheduler_trace_jsonl.as_deref(),
2811        );
2812        let mut execution_profile_journals = Vec::with_capacity(2);
2813        if let Some(journal) = profile_trace_jsonl.as_ref() {
2814            execution_profile_journals.push(journal.clone());
2815        }
2816        if let Some(journal) = scheduler_trace_jsonl.as_ref() {
2817            if execution_profile_journals
2818                .iter()
2819                .all(|existing| existing.path() != journal.path())
2820            {
2821                execution_profile_journals.push(journal.clone());
2822            }
2823        }
2824        if !execution_profile_journals.is_empty() {
2825            let sink: Arc<dyn ExecutionEventSink> =
2826                Arc::new(VNextProfileExecutionEventSink::with_journals(
2827                    execution_profile_journals,
2828                    runtime_config
2829                        .profile_entrypoint
2830                        .unwrap_or(ProfileEntrypoint::Synthetic),
2831                    &config,
2832                ));
2833            model_executor.attach_execution_event_sink(Arc::clone(&sink));
2834            if let Some(draft_executor) = draft_executor.as_ref() {
2835                draft_executor.attach_execution_event_sink(sink);
2836            }
2837        }
2838
2839        Ok(Self {
2840            inner: Arc::new(EngineInner {
2841                config,
2842                scheduler,
2843                tokenizer,
2844                structured_output_factory: OnceLock::new(),
2845                sampler,
2846                resource_composition,
2847                model_executor,
2848                draft_executor,
2849                spec_config,
2850                tensor_factory,
2851                sequences: RwLock::new(HashMap::new()),
2852                is_running: AtomicBool::new(false),
2853                shutdown_notify: Arc::new(Notify::new()),
2854                iteration_lock: tokio::sync::Mutex::new(()),
2855                work_notify: Arc::new(Notify::new()),
2856                iteration_count: AtomicU64::new(0),
2857                prefix_cache: PrefixCache::new(256, 2),
2858                runtime_config,
2859                profile_trace_jsonl,
2860                scheduler_trace_jsonl,
2861                legacy_scheduler_trace_jsonl,
2862                scheduler_trace_none_streak: AtomicU64::new(0),
2863                resource_lifecycle: Mutex::new(ResourceLifecycleLedger::default()),
2864                resource_trace_event_counter: AtomicU64::new(0),
2865                dynamic_admission_availability: Mutex::new(Vec::with_capacity(16)),
2866                execution_readiness_waiters: ExecutionReadinessWaitRegistry::new(),
2867                prefix_rendezvous: Mutex::new(Vec::new()),
2868                prefix_restore_pending: Mutex::new(HashMap::new()),
2869                total_prefill_tokens: AtomicU64::new(0),
2870                total_decode_tokens: AtomicU64::new(0),
2871                total_preemptions: AtomicU64::new(0),
2872                prefix_cache_hits: AtomicU64::new(0),
2873                total_iteration_lock_wait_us: AtomicU64::new(0),
2874                iteration_lock_wait_samples: AtomicU64::new(0),
2875                total_scheduling_time_us: AtomicU64::new(0),
2876                scheduling_time_samples: AtomicU64::new(0),
2877                total_model_execution_time_us: AtomicU64::new(0),
2878                model_execution_time_samples: AtomicU64::new(0),
2879                bg_loop_spawned: AtomicBool::new(false),
2880                shutdown_started: AtomicBool::new(false),
2881                shutdown_lock: tokio::sync::Mutex::new(()),
2882                background_loop: Mutex::new(None),
2883            }),
2884        })
2885    }
2886
2887    /// Spawn the background iteration loop on first request. Without this,
2888    /// every concurrent infer/infer_stream call spawned its own
2889    /// drive_to_completion task → 16 streaming requests = 16 tasks all
2890    /// racing for `iteration_lock` (thundering herd, observed as ~5ms of
2891    /// per-iter tokio scheduling overhead at c=16). With one bg loop +
2892    /// per-request tasks just consuming their channel, lock is uncontested.
2893    fn ensure_bg_loop(&self) {
2894        if self.inner.bg_loop_spawned.load(Ordering::Acquire) {
2895            return;
2896        }
2897        let mut background_loop = self.inner.background_loop.lock();
2898        if self.inner.shutdown_started.load(Ordering::Acquire) {
2899            return;
2900        }
2901        if !self.inner.bg_loop_spawned.swap(true, Ordering::AcqRel) {
2902            *background_loop = Some(self.start_loop());
2903        }
2904    }
2905
2906    /// Hit count since engine construction (prefix cache). Exposed for
2907    /// tests + /metrics endpoint; monotonic, Relaxed-ordered.
2908    pub fn prefix_cache_hits(&self) -> u64 {
2909        self.inner.prefix_cache_hits.load(Ordering::Relaxed)
2910    }
2911
2912    /// Snapshot of prefix cache stats (hits/misses/evictions/active entries).
2913    pub fn prefix_cache_stats(&self) -> ferrum_kv::cache::prefix::PrefixCacheStats {
2914        self.inner.prefix_cache.stats()
2915    }
2916
2917    /// Start a background iteration loop.  Returns a `JoinHandle` that
2918    /// runs until `shutdown()` is called.  When a background loop is
2919    /// active, `infer()` / `infer_stream()` simply submit and wait.
2920    pub fn start_loop(&self) -> tokio::task::JoinHandle<()> {
2921        let inner = self.inner.clone();
2922        inner.is_running.store(true, Ordering::SeqCst);
2923        tokio::spawn(async move {
2924            info!("Background iteration loop started");
2925            let prof = inner.runtime_config.batch_decode_prof;
2926            let mut last_iter_end: Option<std::time::Instant> = None;
2927            static GAP_PROF_CALLS: std::sync::atomic::AtomicU64 =
2928                std::sync::atomic::AtomicU64::new(0);
2929            loop {
2930                if !inner.is_running.load(Ordering::SeqCst) {
2931                    break;
2932                }
2933                let inter_iter_us = if let Some(prev) = last_iter_end {
2934                    Some(prev.elapsed().as_micros() as u64)
2935                } else {
2936                    None
2937                };
2938                let outcome = match inner.run_iteration().await {
2939                    Ok(outcome) => outcome,
2940                    Err(error) => {
2941                        warn!("Iteration error: {}", error);
2942                        EngineIterationOutcome::Progressed
2943                    }
2944                };
2945                if prof {
2946                    let n = GAP_PROF_CALLS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
2947                    if n.is_multiple_of(8) {
2948                        if let Some(gap_us) = inter_iter_us {
2949                            eprintln!("[bg-loop-gap] call#{} inter_iter={}us", n, gap_us);
2950                        }
2951                    }
2952                }
2953                last_iter_end = Some(std::time::Instant::now());
2954                if !inner.is_running.load(Ordering::SeqCst) {
2955                    break;
2956                }
2957                match outcome {
2958                    EngineIterationOutcome::Progressed => tokio::task::yield_now().await,
2959                    EngineIterationOutcome::Idle => {
2960                        tokio::select! {
2961                            _ = inner.wait_for_prefix_deadline() => {}
2962                            _ = inner.shutdown_notify.notified() => {}
2963                            _ = inner.work_notify.notified() => {}
2964                        }
2965                    }
2966                    EngineIterationOutcome::CapacityBlocked(registration) => {
2967                        tokio::select! {
2968                            _ = inner.wait_for_prefix_deadline() => {}
2969                            _ = inner.shutdown_notify.notified() => {}
2970                            _ = inner.work_notify.notified() => {}
2971                            result = registration.wait_for_change() => {
2972                                if let Err(error) = result {
2973                                    warn!("Executor capacity wait error: {}", error);
2974                                }
2975                            }
2976                        }
2977                    }
2978                }
2979            }
2980            info!("Background iteration loop stopped");
2981        })
2982    }
2983}
2984
2985#[async_trait]
2986impl LlmInferenceEngine for ContinuousBatchEngine {
2987    fn context_capacity(&self) -> Option<usize> {
2988        effective_request_context_capacity(
2989            &self.inner.config,
2990            &self.inner.runtime_config,
2991            self.inner.model_executor.kv_capacity(),
2992        )
2993    }
2994
2995    async fn infer(&self, mut request: InferenceRequest) -> Result<InferenceResponse> {
2996        let request_id = request.id.clone();
2997        let infer_start = Instant::now();
2998        counter!("ferrum.engine.requests_total").increment(1);
2999
3000        let input_tokens = self.inner.tokenizer.encode(&request.prompt, true)?;
3001        clamp_default_max_tokens_to_context(
3002            &mut request,
3003            input_tokens.len(),
3004            &self.inner.config,
3005            &self.inner.runtime_config,
3006            self.inner.model_executor.kv_capacity(),
3007        );
3008        validate_request_context_budget(
3009            &request,
3010            input_tokens.len(),
3011            &self.inner.config,
3012            &self.inner.runtime_config,
3013            self.inner.model_executor.kv_capacity(),
3014        )?;
3015        request.metadata.insert(
3016            PROMPT_TOKENS_METADATA_KEY.to_string(),
3017            serde_json::Value::from(input_tokens.len() as u64),
3018        );
3019
3020        // Publish the tokenized sequence and scheduler item atomically with
3021        // respect to the iteration driver. Typed admission must never observe
3022        // one without the other.
3023        let (resp_tx, resp_rx) = tokio::sync::oneshot::channel();
3024        let mut receiver_drop_wake =
3025            ClientReceiverDropWake::new(Arc::clone(&self.inner.work_notify));
3026        let structured_factory = if request.requires_structured_output() {
3027            Some(self.inner.structured_output_factory()?)
3028        } else {
3029            None
3030        };
3031        let mut seq_state =
3032            SequenceState::try_new_with_tokenizer_model_vocab_and_structured_factory(
3033                request.clone(),
3034                input_tokens,
3035                Some(self.inner.tokenizer.clone()),
3036                Some(self.inner.model_executor.info().vocab_size),
3037                structured_factory.as_deref(),
3038            )?;
3039        gauge!("ferrum.engine.active_requests").increment(1.0);
3040        let request_slot = RequestSlotLease::open(&self.inner, request_id.clone());
3041        seq_state.response_sender = Some(resp_tx);
3042        seq_state.request_slot = Some(request_slot);
3043        {
3044            let _iteration = self.inner.iteration_lock.lock().await;
3045            {
3046                let mut sequences = self.inner.sequences.write();
3047                if sequences.contains_key(&request_id) {
3048                    let error = FerrumError::already_exists(format!(
3049                        "request {} is already active",
3050                        request_id
3051                    ));
3052                    if let Some(request_slot) = seq_state.request_slot.take() {
3053                        request_slot.reject(&self.inner, error.to_string());
3054                    }
3055                    gauge!("ferrum.engine.active_requests").decrement(1.0);
3056                    return Err(error);
3057                }
3058                sequences.insert(request_id.clone(), seq_state);
3059            }
3060            if let Err(error) = self.inner.scheduler.submit(request).await {
3061                let mut sequence = self
3062                    .inner
3063                    .sequences
3064                    .write()
3065                    .remove(&request_id)
3066                    .expect("just-published sequence remains present after submit failure");
3067                if let Some(request_slot) = sequence.request_slot.take() {
3068                    request_slot.reject(&self.inner, error.to_string());
3069                }
3070                gauge!("ferrum.engine.active_requests").decrement(1.0);
3071                return Err(error);
3072            }
3073            self.inner
3074                .sequences
3075                .write()
3076                .get_mut(&request_id)
3077                .and_then(|sequence| sequence.request_slot.as_mut())
3078                .expect("submitted sequence retains its request slot")
3079                .admit(&self.inner);
3080        }
3081
3082        // Make sure the single shared bg loop is running, then just wait
3083        // for our oneshot to fire. Avoids per-request drive_to_completion
3084        // contention on iteration_lock.
3085        self.ensure_bg_loop();
3086        self.inner.work_notify.notify_one();
3087
3088        let result = resp_rx.await.unwrap_or_else(|_| {
3089            Err(FerrumError::internal(
3090                "Response channel closed before response was sent",
3091            ))
3092        });
3093        receiver_drop_wake.disarm();
3094
3095        gauge!("ferrum.engine.active_requests").decrement(1.0);
3096        let elapsed_ms = infer_start.elapsed().as_secs_f64() * 1000.0;
3097        histogram!("ferrum.engine.request_duration_ms").record(elapsed_ms);
3098
3099        if let Ok(ref resp) = result {
3100            counter!("ferrum.engine.requests_completed").increment(1);
3101            counter!("ferrum.engine.tokens_generated_total").increment(resp.tokens.len() as u64);
3102            // NOTE: real TTFT lives in `send_stream_update` —
3103            // emitted as `ferrum.engine.ttft_seconds`. The sync `infer`
3104            // path returns the whole response at once, so there's no
3105            // observable first-token moment to record here.
3106        } else {
3107            counter!("ferrum.engine.requests_failed").increment(1);
3108        }
3109
3110        result
3111    }
3112
3113    async fn infer_stream(
3114        &self,
3115        mut request: InferenceRequest,
3116    ) -> Result<Pin<Box<dyn Stream<Item = Result<StreamChunk>> + Send>>> {
3117        let (tx, rx) = mpsc::channel(100);
3118        let receiver_drop_wake = ClientReceiverDropWake::new(Arc::clone(&self.inner.work_notify));
3119        let request_id = request.id.clone();
3120
3121        let input_tokens = self.inner.tokenizer.encode(&request.prompt, true)?;
3122        clamp_default_max_tokens_to_context(
3123            &mut request,
3124            input_tokens.len(),
3125            &self.inner.config,
3126            &self.inner.runtime_config,
3127            self.inner.model_executor.kv_capacity(),
3128        );
3129        validate_request_context_budget(
3130            &request,
3131            input_tokens.len(),
3132            &self.inner.config,
3133            &self.inner.runtime_config,
3134            self.inner.model_executor.kv_capacity(),
3135        )?;
3136        request.metadata.insert(
3137            PROMPT_TOKENS_METADATA_KEY.to_string(),
3138            serde_json::Value::from(input_tokens.len() as u64),
3139        );
3140
3141        // Publish tokenized state and the scheduler item under the same
3142        // iteration boundary; see the non-streaming path above.
3143        let structured_factory = if request.requires_structured_output() {
3144            Some(self.inner.structured_output_factory()?)
3145        } else {
3146            None
3147        };
3148        let mut seq_state =
3149            SequenceState::try_new_with_tokenizer_model_vocab_and_structured_factory(
3150                request.clone(),
3151                input_tokens,
3152                Some(self.inner.tokenizer.clone()),
3153                Some(self.inner.model_executor.info().vocab_size),
3154                structured_factory.as_deref(),
3155            )?;
3156        let request_slot = RequestSlotLease::open(&self.inner, request_id.clone());
3157        seq_state.stream_sender = Some(tx);
3158        seq_state.request_slot = Some(request_slot);
3159        {
3160            let _iteration = self.inner.iteration_lock.lock().await;
3161            {
3162                let mut sequences = self.inner.sequences.write();
3163                if sequences.contains_key(&request_id) {
3164                    let error = FerrumError::already_exists(format!(
3165                        "request {} is already active",
3166                        request_id
3167                    ));
3168                    if let Some(request_slot) = seq_state.request_slot.take() {
3169                        request_slot.reject(&self.inner, error.to_string());
3170                    }
3171                    return Err(error);
3172                }
3173                sequences.insert(request_id.clone(), seq_state);
3174            }
3175            if let Err(error) = self.inner.scheduler.submit(request).await {
3176                let mut sequence = self
3177                    .inner
3178                    .sequences
3179                    .write()
3180                    .remove(&request_id)
3181                    .expect("just-published sequence remains present after submit failure");
3182                if let Some(request_slot) = sequence.request_slot.take() {
3183                    request_slot.reject(&self.inner, error.to_string());
3184                }
3185                return Err(error);
3186            }
3187            self.inner
3188                .sequences
3189                .write()
3190                .get_mut(&request_id)
3191                .and_then(|sequence| sequence.request_slot.as_mut())
3192                .expect("submitted sequence retains its request slot")
3193                .admit(&self.inner);
3194        }
3195
3196        // Single shared bg loop drives iters; per-request stream just
3197        // consumes from `rx`. Used to spawn a per-request drive_to_completion
3198        // task here, but with c=N concurrent streams that produced N
3199        // tasks all racing for `iteration_lock` — measured ~5ms/iter of
3200        // tokio thundering-herd overhead at c=16.
3201        let _ = request_id;
3202        self.ensure_bg_loop();
3203        self.inner.work_notify.notify_one();
3204
3205        Ok(Box::pin(CancellationAwareResponseStream {
3206            receiver: tokio_stream::wrappers::ReceiverStream::new(rx),
3207            receiver_drop_wake,
3208        }))
3209    }
3210}
3211
3212#[async_trait]
3213impl InferenceEngine for ContinuousBatchEngine {
3214    async fn status(&self) -> EngineStatus {
3215        let metrics = self.inner.scheduler.metrics();
3216        let (total_bytes, used_bytes, cache_memory_bytes, resource_status_ready) =
3217            match &self.inner.resource_composition {
3218                EngineResourceComposition::LegacyEngine { kv_cache, .. } => {
3219                    let kv_stats = kv_cache.stats();
3220                    (
3221                        kv_stats.total_memory_bytes,
3222                        kv_stats.used_memory_bytes,
3223                        kv_stats.used_memory_bytes,
3224                        true,
3225                    )
3226                }
3227                EngineResourceComposition::PlanRuntime => {
3228                    match self.inner.model_executor.plan_runtime_resource_snapshot() {
3229                        Ok(Some(snapshot)) => {
3230                            let total_bytes = usize::try_from(snapshot.usable_capacity_bytes())
3231                                .unwrap_or(usize::MAX);
3232                            let used_bytes = snapshot
3233                                .used_bytes()
3234                                .ok()
3235                                .and_then(|bytes| usize::try_from(bytes).ok())
3236                                .unwrap_or(usize::MAX);
3237                            let dynamic_used_bytes = usize::try_from(snapshot.dynamic_used_bytes())
3238                                .unwrap_or(usize::MAX);
3239                            (total_bytes, used_bytes, dynamic_used_bytes, true)
3240                        }
3241                        Ok(None) => {
3242                            warn!("Plan runtime did not expose its required resource snapshot");
3243                            (0, 0, 0, false)
3244                        }
3245                        Err(error) => {
3246                            warn!(error = %error, "Plan-runtime resource snapshot failed");
3247                            (0, 0, 0, false)
3248                        }
3249                    }
3250                }
3251            };
3252        let free_bytes = total_bytes.saturating_sub(used_bytes);
3253        let mut memory_usage = ferrum_types::MemoryUsage {
3254            total_bytes,
3255            used_bytes,
3256            free_bytes,
3257            gpu_memory_bytes: self
3258                .inner
3259                .config
3260                .backend
3261                .device
3262                .is_gpu()
3263                .then_some(used_bytes),
3264            cpu_memory_bytes: matches!(self.inner.config.backend.device, Device::CPU)
3265                .then_some(used_bytes),
3266            cache_memory_bytes,
3267            utilization_percent: 0.0,
3268        };
3269        memory_usage.calculate_utilization();
3270        EngineStatus {
3271            is_ready: resource_status_ready && self.inner.is_running.load(Ordering::SeqCst),
3272            loaded_models: vec![self.inner.config.model.model_id.clone()],
3273            active_requests: metrics.running_requests,
3274            queued_requests: metrics.waiting_requests,
3275            memory_usage,
3276            uptime_seconds: 0,
3277            last_heartbeat: chrono::Utc::now(),
3278            version: env!("CARGO_PKG_VERSION").to_string(),
3279        }
3280    }
3281
3282    async fn shutdown(&self) -> Result<()> {
3283        let _shutdown_guard = self.inner.shutdown_lock.lock().await;
3284        info!("Shutting down continuous batch engine");
3285        let background_loop = {
3286            let mut background_loop = self.inner.background_loop.lock();
3287            self.inner.signal_shutdown();
3288            background_loop.take()
3289        };
3290
3291        let loop_result = match background_loop {
3292            Some(background_loop) => background_loop.await.map_err(|error| {
3293                FerrumError::internal(format!("background iteration loop failed: {error}"))
3294            }),
3295            None => Ok(()),
3296        };
3297        let readiness_result = self
3298            .inner
3299            .execution_readiness_waiters
3300            .abort_and_join()
3301            .await;
3302        // Drop immutable checkpoint pins outside the cohort mutex before native
3303        // resource shutdown. No pending dependency survives engine shutdown.
3304        let prefix_cohorts = std::mem::take(&mut *self.inner.prefix_rendezvous.lock());
3305        drop(prefix_cohorts);
3306        let prefix_restores = std::mem::take(&mut *self.inner.prefix_restore_pending.lock());
3307        drop(prefix_restores);
3308
3309        let mut trace_journals = Vec::with_capacity(2);
3310        if let Some(journal) = self.inner.profile_trace_jsonl.clone() {
3311            trace_journals.push(journal);
3312        }
3313        if let Some(journal) = self.inner.scheduler_trace_jsonl.clone() {
3314            if trace_journals
3315                .iter()
3316                .all(|existing| existing.path() != journal.path())
3317            {
3318                trace_journals.push(journal);
3319            }
3320        }
3321        let trace_result = if trace_journals.is_empty() {
3322            Ok(())
3323        } else {
3324            tokio::task::spawn_blocking(move || {
3325                for journal in trace_journals {
3326                    journal.close()?;
3327                }
3328                Ok::<(), JsonlJournalError>(())
3329            })
3330            .await
3331            .map_err(|error| {
3332                FerrumError::internal(format!("scheduler trace close task failed: {error}"))
3333            })?
3334            .map_err(|error| {
3335                FerrumError::internal(format!("scheduler trace close failed: {error}"))
3336            })
3337        };
3338
3339        loop_result?;
3340        readiness_result?;
3341        trace_result
3342    }
3343
3344    fn config(&self) -> &EngineConfig {
3345        &self.inner.config
3346    }
3347
3348    fn metrics(&self) -> ferrum_types::EngineMetrics {
3349        let sm = self.inner.scheduler.metrics();
3350        ferrum_types::EngineMetrics {
3351            total_requests: sm.completed_requests + sm.failed_requests,
3352            successful_requests: sm.completed_requests,
3353            failed_requests: sm.failed_requests,
3354            avg_request_latency_ms: 0.0,
3355            p95_request_latency_ms: 0.0,
3356            p99_request_latency_ms: 0.0,
3357            throughput_rps: sm.throughput_rps as f32,
3358            tokens_per_second: 0.0,
3359            queue_metrics: ferrum_types::QueueMetrics {
3360                current_queue_length: sm.waiting_requests,
3361                avg_queue_wait_time_ms: sm.avg_wait_time_ms,
3362                queue_throughput_rps: sm.throughput_rps as f32,
3363                queue_rejection_rate: 0.0,
3364            },
3365            resource_utilization: Default::default(),
3366            error_stats: Default::default(),
3367            performance_breakdown: self.inner.performance_breakdown(),
3368        }
3369    }
3370
3371    fn cache_metrics_snapshot(&self) -> Option<serde_json::Value> {
3372        if let Some(snapshot) = self.inner.model_executor.cache_metrics_snapshot() {
3373            return Some(snapshot);
3374        }
3375
3376        let stats = self.inner.prefix_cache.stats();
3377        Some(serde_json::json!({
3378            "position": "engine-whole-prompt-debug-cache",
3379            "source": "continuous-engine-whole-prompt-prefix-cache",
3380            "enabled": self.inner.runtime_config.prefix_cache_enabled,
3381            "hits": stats.hits as u64,
3382            "misses": stats.misses as u64,
3383            "evictions": stats.evictions as u64,
3384            "saved_prefill_tokens": self.inner.prefix_cache_hits.load(Ordering::Relaxed),
3385            "entries": stats.active_prefixes as u64,
3386            "bytes": 0u64,
3387            "cached_tokens": stats.total_cached_tokens as u64,
3388            "hit_rate": stats.hit_rate,
3389        }))
3390    }
3391
3392    fn execution_attribution_snapshot(&self) -> Option<serde_json::Value> {
3393        self.inner.model_executor.execution_attribution_snapshot()
3394    }
3395
3396    fn admission_snapshot(
3397        &self,
3398    ) -> ferrum_types::Result<Option<ferrum_types::ExecutorAdmissionSnapshot>> {
3399        let scheduler = self.inner.scheduler.admission_phase_counts();
3400        let authority = self.inner.resource_composition.authority();
3401        let limits = match authority {
3402            ExecutionResourceAuthority::PlanRuntime => self
3403                .inner
3404                .model_executor
3405                .admission_limits()?
3406                .ok_or_else(|| {
3407                    FerrumError::internal(
3408                        "PlanRuntime executor did not expose its resolved admission limits",
3409                    )
3410                })?,
3411            ExecutionResourceAuthority::LegacyEngine => {
3412                let scheduler_limit = self.inner.config.scheduler.max_running_requests;
3413                let maximum_active_sequences = self
3414                    .inner
3415                    .recurrent_state_manager()
3416                    .map(|manager| scheduler_limit.min(manager.stats().total_batch_slots))
3417                    .unwrap_or(scheduler_limit);
3418                ferrum_types::ExecutorAdmissionLimits::new(
3419                    u32::try_from(maximum_active_sequences).map_err(|_| {
3420                        FerrumError::internal("legacy admission sequence limit exceeds u32")
3421                    })?,
3422                    u64::try_from(self.inner.config.batching.max_num_batched_tokens).map_err(
3423                        |_| FerrumError::internal("legacy scheduled-token limit exceeds u64"),
3424                    )?,
3425                )
3426                .map_err(|reason| {
3427                    FerrumError::internal(format!(
3428                        "legacy admission limits violated their typed contract: {reason}"
3429                    ))
3430                })?
3431            }
3432        };
3433        ferrum_types::ExecutorAdmissionSnapshot::new(
3434            authority,
3435            limits,
3436            u32::try_from(scheduler.waiting_requests)
3437                .map_err(|_| FerrumError::internal("waiting request count exceeds u32"))?,
3438            u32::try_from(scheduler.active_prefill_sequences)
3439                .map_err(|_| FerrumError::internal("active prefill count exceeds u32"))?,
3440            u32::try_from(scheduler.active_decode_sequences)
3441                .map_err(|_| FerrumError::internal("active decode count exceeds u32"))?,
3442            None,
3443            None,
3444        )
3445        .map(Some)
3446        .map_err(|reason| {
3447            FerrumError::internal(format!(
3448                "runtime admission snapshot violated its typed contract: {reason}"
3449            ))
3450        })
3451    }
3452
3453    fn lora_metrics_snapshot(&self) -> Option<serde_json::Value> {
3454        self.inner.model_executor.lora_metrics_snapshot()
3455    }
3456
3457    async fn health_check(&self) -> ferrum_types::HealthStatus {
3458        if self.inner.is_running.load(Ordering::SeqCst) {
3459            ferrum_types::HealthStatus::healthy()
3460        } else {
3461            ferrum_types::HealthStatus {
3462                status: ferrum_types::HealthStatusType::Unhealthy,
3463                component_status: ferrum_types::ComponentStatus::healthy(),
3464                last_check: chrono::Utc::now(),
3465            }
3466        }
3467    }
3468}
3469
3470impl std::fmt::Debug for ContinuousBatchEngine {
3471    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
3472        f.debug_struct("ContinuousBatchEngine")
3473            .field("is_running", &self.inner.is_running.load(Ordering::SeqCst))
3474            .field(
3475                "iteration_count",
3476                &self.inner.iteration_count.load(Ordering::SeqCst),
3477            )
3478            .field("active_sequences", &self.inner.sequences.read().len())
3479            .finish()
3480    }
3481}
3482
3483// ────────────────────────────────────────────────────────────────────────────
3484// Unit tests
3485// ────────────────────────────────────────────────────────────────────────────
3486
3487#[cfg(test)]
3488mod tests;