1use 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#[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 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 model_eos_token_ids: Vec<u32>,
313 stop_token_ids: HashSet<u32>,
315 user_stop_token_ids: HashSet<u32>,
316 stop_text_seqs: Vec<String>,
317}
318
319fn 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 ¶ms.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 ¶ms.stop_sequences {
376 text_seqs.push(stop_seq.clone());
377 }
378 }
379 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 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 !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#[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 structured_output_factory: OnceLock<std::result::Result<Arc<StructuredOutputFactory>, String>>,
1411 #[allow(dead_code)]
1412 sampler: Arc<dyn Sampler + Send + Sync>,
1414 resource_composition: EngineResourceComposition,
1415 model_executor: Arc<dyn ModelExecutor + Send + Sync>,
1416 draft_executor: Option<Arc<dyn ModelExecutor + Send + Sync>>,
1419 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 iteration_lock: tokio::sync::Mutex<()>,
1429 work_notify: Arc<Notify>,
1431 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 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 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 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
2625pub 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 #[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 #[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 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 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 pub fn prefix_cache_hits(&self) -> u64 {
2909 self.inner.prefix_cache_hits.load(Ordering::Relaxed)
2910 }
2911
2912 pub fn prefix_cache_stats(&self) -> ferrum_kv::cache::prefix::PrefixCacheStats {
2914 self.inner.prefix_cache.stats()
2915 }
2916
2917 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 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 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 } 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 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 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 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#[cfg(test)]
3488mod tests;