aprender-serve 0.70.2

Pure Rust ML inference engine built from scratch - model serving for GGUF and safetensors
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
//! HTTP API for model inference
//!
//! Provides REST endpoints for tokenization and text generation using axum.
//!
//! ## Endpoints
//!
//! - `GET /health` - Health check
//! - `GET /metrics` - Prometheus-formatted metrics
//! - `GET /metrics/dispatch` - CPU/GPU dispatch statistics (?format=prometheus|json)
//! - `POST /tokenize` - Tokenize text
//! - `POST /generate` - Generate text from prompt
//! - `POST /batch/tokenize` - Batch tokenize multiple texts
//! - `POST /batch/generate` - Batch generate for multiple prompts
//! - `POST /stream/generate` - Stream generated tokens via SSE
//! - `POST /v1/gpu/warmup` - Warmup GPU cache for batch inference (PARITY-022)
//! - `GET /v1/gpu/status` - Check GPU cache status (PARITY-022)
//! - `POST /v1/batch/completions` - GPU-accelerated batch inference (PARITY-022)
//! - `GET /v1/metrics` - JSON metrics for TUI monitoring (PARITY-107)
//!
//! ## Example
//!
//! ```rust,ignore
//! use realizar::api::{create_router, AppState};
//!
//! let state = AppState::new(model, tokenizer);
//! let app = create_router(state);
//! axum::serve(listener, app).await?;
//! ```

use std::sync::Arc;

use axum::{
    extract::{Query, State},
    http::StatusCode,
    routing::{get, post},
    Json, Router,
};
use serde::{Deserialize, Serialize};

use crate::{
    apr::{AprModel, HEADER_SIZE, MAGIC},
    audit::{AuditLogger, AuditRecord, InMemoryAuditSink},
    cache::{CacheKey, ModelCache},
    error::RealizarError,
    explain::ShapExplanation,
    layers::{Model, ModelConfig},
    metrics::MetricsCollector,
    registry::ModelRegistry,
    tokenizer::BPETokenizer,
};

// aprender#2376(3): request-scoped cancellation for the generate handlers.
mod cancel_scope;
pub(crate) use cancel_scope::cancel_on_disconnect;
pub use cancel_scope::request_cancel_token;

// PMAT-802: Extracted handlers
//
// aprender#2465(1): NOT `#[cfg(feature = "cuda")]` on the module. Every
// CUDA-dependent item inside is individually gated; the session config
// (`q4k_generate_config`), the forward over it (`apr_q4k_forward`) and their
// cancellation falsifiers are not, so they compile and run under the default
// feature set. Gating the whole module put the only cancellation-free decode
// loop in the crate outside every CI test job.
pub mod apr_q4k_forward;
pub mod apr_q4k_scheduler;
// PERF-041: NOT `#[cfg(feature = "cuda")]`, on purpose. It holds the admission
// predicate of contracts/batch-admission-v1.yaml and its exhaustive table test,
// which `cargo test -p aprender-serve --lib batch_admission` must be able to
// select in the default feature set โ€” the same dark-target reasoning as the
// comment above `apr_q4k_scheduler`.
pub mod batch_admission;
#[cfg(feature = "cuda")]
pub mod cuda_batch_scheduler;
#[cfg(feature = "cuda")]
pub mod iteration_scheduler;
mod openai_handlers;
pub(crate) use openai_handlers::LiveUtf8Deltas;
pub(crate) use openai_handlers::{
    openai_chat_completions_handler, openai_chat_completions_stream_handler, openai_models_handler,
};
// PMAT-923: Ollama HTTP compat (/api/chat, /api/generate) โ€” delegates to the
// OpenAI chat path so `apr serve` is a drop-in Ollama HTTP replacement.
mod ollama_handlers;
pub(crate) use ollama_handlers::{
    ollama_chat_handler, ollama_embeddings_handler, ollama_generate_handler, ollama_show_handler,
    ollama_tags_handler, ollama_version_handler,
};
// What this server actually measured about the model it loaded. Metadata
// handlers read it instead of substituting plausible-looking constants.
mod model_source;
pub use model_source::{detect_format_from_magic, gguf_qtype_name, ModelSourceInfo};
// PP-LLAMA-001 ยง12 row 6 / PP-2: `GET /v1/effective-config` โ€” what THIS
// process resolved, read from the process. NOT `#[cfg(feature = "cuda")]`: the
// route is unconditional and its JSON shape is identical on every build, or the
// PP-2 must-not-fire case ("the CPU cell reports cpu") is unreachable.
pub mod effective_config;
// Carried on `AppState` and read by the handler; not part of the public API,
// which reaches the same facts through `effective_config(&state)`.
pub use effective_config::{
    admission_ceiling_reason, build_features, compute_class_from_residency, effective_config,
    pp14_holds, AdmissionCounter, EffectiveConfigResponse, InFlightCounter, KvReport, ModelReport,
    OffloadReport, SchedulerReport, ServerClock, ServerReport, ADMISSION_POLICY,
};
pub(crate) use effective_config::{effective_config_handler, EffectiveConfigState};
mod gpu_handlers;
pub(crate) use gpu_handlers::{
    batch_generate_handler, batch_tokenize_handler, generate_handler,
    gpu_batch_completions_handler, gpu_status_handler, gpu_warmup_handler, models_handler,
    stream_generate_handler, tokenize_handler,
};
// Public exports for tests (GPU-only types)
#[cfg(feature = "gpu")]
pub use gpu_handlers::{
    BatchProcessResult, BatchQueueStats, ContinuousBatchRequest, ContinuousBatchResponse,
    GpuBatchRequest, GpuBatchResponse, GpuBatchResult, GpuBatchStats, GpuStatusResponse,
    GpuWarmupResponse,
};
// Public exports for apr-cli CUDA integration (PMAT-GPU-001)
#[cfg(feature = "gpu")]
pub use gpu_handlers::{spawn_batch_processor, BatchConfig};
mod realize_handlers;
pub(crate) use realize_handlers::{
    clean_chat_output, format_chat_messages, format_chat_messages_for_state,
    format_chat_messages_for_state_thinking, format_chat_messages_for_state_thinking_tools,
    format_chat_messages_official, format_chat_messages_official_thinking,
    format_chat_messages_official_thinking_tools, openai_completions_handler,
    openai_embeddings_handler, realize_embed_handler, realize_model_handler,
    realize_reload_handler,
};
#[cfg(feature = "cuda")]
pub(crate) use realize_handlers::{logprobs_handler, perplexity_handler};
// Public exports for tests
pub use realize_handlers::{
    CompletionChoice, CompletionRequest, CompletionResponse, ContextWindowConfig,
    ContextWindowManager, EmbeddingData, EmbeddingInput, EmbeddingRequest, EmbeddingResponse,
    EmbeddingUsage, ModelLineage, ModelMetadataResponse, ReloadRequest, ReloadResponse,
};
mod apr_handlers;
pub(crate) use apr_handlers::{apr_audit_handler, apr_explain_handler, apr_predict_handler};
mod types;
pub use crate::registry::ModelInfo;
pub use types::{default_max_tokens, default_top_k};
#[cfg(test)]
pub(crate) use types::{default_strategy, default_temperature, default_top_p};
pub use types::{
    BatchGenerateRequest, BatchGenerateResponse, BatchTokenizeRequest, BatchTokenizeResponse,
    ChoiceCount, ErrorResponse, FinishReason, GenerateRequest, GenerateResponse, HealthResponse,
    ModelsResponse, StreamDoneEvent, StreamTokenEvent, TokenizeRequest, TokenizeResponse,
};

/// Application state shared across handlers
#[derive(Clone)]
pub struct AppState {
    /// Model for inference (single model mode)
    model: Option<Arc<Model>>,
    /// Tokenizer for encoding/decoding (single model mode)
    tokenizer: Option<Arc<BPETokenizer>>,
    /// Model cache for multi-model support
    #[allow(dead_code)]
    cache: Option<Arc<ModelCache>>,
    /// Default cache key for single model mode
    #[allow(dead_code)]
    cache_key: Option<CacheKey>,
    /// Metrics collector for monitoring
    metrics: Arc<MetricsCollector>,
    /// Model registry for multi-model serving
    registry: Option<Arc<ModelRegistry>>,
    /// Default model ID for multi-model mode
    default_model_id: Option<String>,
    /// APR model for /v1/predict endpoint (real inference, not mock)
    apr_model: Option<Arc<AprModel>>,
    /// Audit logger for /v1/audit endpoint (real records, not mock)
    audit_logger: Arc<AuditLogger>,
    /// In-memory audit sink for record retrieval
    audit_sink: Arc<InMemoryAuditSink>,
    /// GPU model for GGUF inference (M33: IMP-084)
    #[cfg(feature = "gpu")]
    gpu_model: Option<Arc<std::sync::RwLock<crate::gpu::GpuModel>>>,
    /// Quantized model for fused Q4_K inference (IMP-100)
    /// Faster than the dequantized GpuModel because it reads less memory per token
    /// (IMP-100; the measured factor has no receipt under evidence/, so it is not stated here)
    quantized_model: Option<Arc<crate::gguf::OwnedQuantizedModel>>,
    /// Thread-safe cached model for HTTP serving (IMP-116)
    /// Uses Mutex-based scheduler caching (IMP-116; the measured speedup has no receipt under evidence/)
    #[cfg(feature = "gpu")]
    cached_model: Option<Arc<crate::gguf::OwnedQuantizedModelCachedSync>>,
    /// Dispatch metrics for adaptive CPU/GPU tracking (IMP-126)
    #[cfg(feature = "gpu")]
    dispatch_metrics: Option<Arc<crate::gguf::DispatchMetrics>>,
    /// Batch request channel for continuous batching (PARITY-052)
    /// Requests sent here are queued and processed in batches
    #[cfg(feature = "gpu")]
    batch_request_tx: Option<tokio::sync::mpsc::Sender<ContinuousBatchRequest>>,
    /// Batch configuration for window timing and size thresholds (PARITY-052)
    #[cfg(feature = "gpu")]
    batch_config: Option<BatchConfig>,
    /// CUDA-optimized model for GPU inference (PAR-111).
    /// Uses pre-uploaded weights and batched workspaces. The throughput and
    /// Ollama-ratio this comment used to assert were withdrawn: they were taken
    /// on the batched path while it emitted garbage tokens (aprender#2753).
    #[cfg(feature = "cuda")]
    cuda_model: Option<Arc<std::sync::RwLock<crate::gguf::OwnedQuantizedModelCuda>>>,
    /// PMAT-044: CUDA batch scheduler for continuous batching on /v1/chat/completions
    #[cfg(feature = "cuda")]
    cuda_batch_tx: Option<tokio::sync::mpsc::Sender<cuda_batch_scheduler::CudaBatchRequest>>,
    /// ALB-095: APR Q4K GPU inference channel (dedicated thread owns CudaExecutor)
    #[cfg(feature = "cuda")]
    apr_q4k_tx: Option<tokio::sync::mpsc::Sender<apr_q4k_scheduler::AprQ4kRequest>>,
    /// APR Transformer for SafeTensors/APR inference (PMAT-SERVE-FIX-001)
    /// Supports F32 weights from SafeTensors or APR format
    apr_transformer: Option<Arc<crate::apr_transformer::AprTransformer>>,
    /// #169: SafeTensors CUDA model for GPU-accelerated F16/F32 inference
    #[cfg(feature = "cuda")]
    safetensors_cuda_model:
        Option<Arc<std::sync::Mutex<crate::safetensors_cuda::SafeTensorsCudaModel>>>,
    /// GH-319: Cached model architecture string (avoids RwLock in hot path)
    cached_architecture: Option<String>,
    /// aprender#1789 Option B: retained MappedGGUFModel for MoE-aware HTTP
    /// dispatch. `run_qwen3_moe_generate` borrows per-expert tensors
    /// directly from the mmap, so the mapped model must outlive any
    /// inference call. Held in an `Arc` to share between the chat handler
    /// and any future streaming/batch backends.
    /// See `contracts/qwen3-moe-serve-dispatch-v1.yaml` (V1_001, V1_003).
    mapped_gguf_model: Option<Arc<crate::gguf::MappedGGUFModel>>,
    /// #3987: whether qwen3moe generation must stay on the CPU. `true` for every
    /// constructor, which is exactly the pre-#3987 behaviour (the serve MoE backend
    /// called the CPU-only generator). Only a CUDA server opts in, via
    /// `with_moe_gpu()`, and then the MoE backend goes through the ONE dispatch
    /// `apr run` uses (`run_qwen3_moe_generate_dispatch`), proven by #3714's parity tests.
    moe_no_gpu: bool,
    /// #3571: the Qwen3.5 hybrid, resident for the server's lifetime. Its
    /// Gated-DeltaNet and attention layers live here, not in
    /// `quantized_model` โ€” the hybrid has no dense layers โ€” so a Qwen3.5
    /// request is served from this session or refused, never decoded through
    /// the base (embeddings, norm, `lm_head`) alone.
    qwen35_session: Option<Arc<Qwen35Served>>,
    /// GH-330: Cached EOS token ID (avoids RwLock in hot path)
    cached_eos_token_id: Option<u32>,
    /// D5 (ruling-3715-2010): positions a dense CUDA serve turn can reach, the
    /// model's context capped by the device KV cache. Cached at construction so
    /// the pre-flight length check never waits on the scheduler's write lock.
    cached_serving_context: Option<usize>,
    /// GH-152: Enable verbose request/response logging
    verbose: bool,
    /// GH-103: Enable inference tracing (propagates into QuantizedGenerateConfig.trace)
    trace: bool,
    /// What the loader measured about the served model (path, size, format,
    /// quantization, context length). `None` means this server was built
    /// without that knowledge โ€” the metadata handlers then report the fields
    /// as ABSENT rather than inventing values.
    model_source: Option<Arc<ModelSourceInfo>>,
    /// PP-LLAMA-001 ยง5.2: the process clock, the resolved offload, the
    /// scheduler identity and the in-flight counter that `GET
    /// /v1/effective-config` reports. One field rather than four so that adding
    /// a reported fact does not mean editing sixteen struct literals.
    effective: EffectiveConfigState,
}

impl AppState {
    /// Attach measured model provenance/metadata.
    ///
    /// Call this from whatever loaded the model; it is what makes
    /// `/realize/model`, `/api/tags` and `/api/show` report the truth instead
    /// of constants.
    #[must_use]
    pub fn with_model_source(mut self, source: ModelSourceInfo) -> Self {
        self.model_source = Some(Arc::new(source));
        self
    }

    /// Measured model provenance/metadata, if the loader supplied any.
    #[must_use]
    pub fn model_source(&self) -> Option<&ModelSourceInfo> {
        self.model_source.as_deref()
    }

    /// Attach what the loader resolved for `--gpu-layers` (PP-14, PP-15).
    ///
    /// Called by whatever placed the layers, next to the line that prints them,
    /// so the endpoint reports the same resolution the operator saw and not a
    /// second derivation of it.
    #[must_use]
    pub fn with_offload_report(mut self, report: effective_config::OffloadReport) -> Self {
        self.effective.offload = Some(Arc::new(report));
        self
    }

    /// Attach the scheduler identity and its live in-flight counter (PP-13/PP-24).
    #[must_use]
    pub fn with_scheduler_report(
        mut self,
        report: effective_config::SchedulerReport,
        counter: Option<Arc<effective_config::InFlightCounter>>,
    ) -> Self {
        self.effective.scheduler = Some(Arc::new(report));
        self.effective.in_flight = counter;
        self
    }

    /// The `/v1/effective-config` state carried on this server.
    #[must_use]
    pub(crate) fn effective_config_state(&self) -> &EffectiveConfigState {
        &self.effective
    }

    /// This server's process clock (PP-30).
    #[must_use]
    pub fn clock(&self) -> &Arc<effective_config::ServerClock> {
        &self.effective.clock
    }

    /// Record one request REFUSED admission (ยง5.2 `kv.admission_rejected`).
    ///
    /// Called from the `503` return itself, not from a wrapper around it: every
    /// site that turns a client away because the scheduler could not take the
    /// request is an admission rejection, and a counter incremented anywhere
    /// else would drift from the response the client actually got.
    pub fn record_admission_rejected(&self) -> u64 {
        self.effective.admission.record_rejected()
    }

    /// Requests refused admission since start.
    #[must_use]
    pub fn admission_rejected(&self) -> u64 {
        self.effective.admission.rejected()
    }
}

/// Helper to create default audit infrastructure
fn create_audit_state() -> (Arc<AuditLogger>, Arc<InMemoryAuditSink>) {
    let sink = Arc::new(InMemoryAuditSink::new());
    let logger = AuditLogger::new(Box::new(InMemorySinkWrapper(sink.clone())))
        .with_model_hash("demo-model-hash");
    (Arc::new(logger), sink)
}

/// Wrapper to make Arc<InMemoryAuditSink> implement AuditSink
struct InMemorySinkWrapper(Arc<InMemoryAuditSink>);

impl crate::audit::AuditSink for InMemorySinkWrapper {
    fn write_batch(&self, records: &[AuditRecord]) -> Result<(), crate::audit::AuditError> {
        self.0.write_batch(records)
    }

    fn flush(&self) -> Result<(), crate::audit::AuditError> {
        self.0.flush()
    }
}

/// HTTP status for a model/tokenizer resolution failure.
///
/// One server-side condition must map to one status code. Before this existed the
/// identical `"Model registry error: No model available"` came back as 404 from
/// `/tokenize`, `/stream/generate` and `/realize/embed` but as 500 from
/// `/batch/tokenize` and `/batch/generate`, so a client retry policy keyed on
/// status treated the same failure as permanent on one route and as a server bug
/// on the next (aprender#2376 finding 5).
///
/// * [`RealizarError::ModelNotFound`] โ€” the client named a model this server does
///   not have: 404, the route and request are fine.
/// * [`RealizarError::RegistryError`] โ€” the server has no usable model at all.
///   That is a server-side condition, so 503 (the shape `/metrics/dispatch`
///   already uses), never 404: the resource exists, the server cannot serve it.
pub(crate) fn model_resolution_status(err: &RealizarError) -> StatusCode {
    match err {
        RealizarError::ModelNotFound(_) => StatusCode::NOT_FOUND,
        RealizarError::RegistryError(_) => StatusCode::SERVICE_UNAVAILABLE,
        _ => StatusCode::INTERNAL_SERVER_ERROR,
    }
}

/// HTTP status for a generation failure.
///
/// A prompt or `max_tokens` that does not fit the model's context window is fully
/// determined by the request, so it is a client error. Reporting it as 500 tells
/// the caller the server broke and invites a retry of the identical request
/// (aprender#2376 findings 9 and 11).
pub(crate) fn generation_error_status(err: &RealizarError) -> StatusCode {
    match err {
        RealizarError::ContextLimitExceeded { .. } => StatusCode::BAD_REQUEST,
        _ => StatusCode::INTERNAL_SERVER_ERROR,
    }
}

/// D5 (ruling-3715-2010): the refusal for a prompt that cannot fit a serve turn,
/// or `None` when it fits with room to answer.
///
/// `serving_context` is the model's context capped by the device KV cache
/// ([`AppState::serving_context`]). A prompt at or past it used to reach the GPU
/// forward and come back as a 500 ("the device KV cache holds 4096"), which
/// tells the client the server broke. It is decided by the request alone, so it
/// is a 400 and says which limit it hit.
pub(crate) fn serve_context_refusal(
    prompt_tokens: usize,
    serving_context: usize,
) -> Option<String> {
    (prompt_tokens >= serving_context).then(|| {
        format!(
            "context exceeds {serving_context}: the prompt is {prompt_tokens} tokens and this \
             server holds {serving_context} positions per turn (the model's context, capped by \
             the GPU KV cache). It was refused whole, not truncated."
        )
    })
}

include!("mod_app_state_gpu.rs");
include!("mod_app_state_qwen35.rs");
include!("mod_create_demo.rs");
include!("router.rs");
include!("dispatch_metrics.rs");

#[cfg(test)]
#[path = "tests_engine_identity.rs"]
mod tests_engine_identity;