Files
DS4Server/src/metrics.rs
Georg Bauer 76a5dd5b26 Name session checkpoints by their session title
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-26 09:25:37 +02:00

940 lines
35 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
use crate::model::ModelChoice;
use std::fs;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, AtomicU8, AtomicU16, AtomicU32, AtomicU64, Ordering};
use std::time::{Duration, Instant, SystemTime};
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
#[repr(u8)]
pub(crate) enum RuntimePhase {
#[default]
Unloaded,
Loading,
Prefilling,
Generating,
Ready,
Failed,
}
impl RuntimePhase {
pub(crate) fn label(self) -> &'static str {
match self {
Self::Unloaded => "Unloaded",
Self::Loading => "Loading",
Self::Prefilling => "Prefilling",
Self::Generating => "Generating",
Self::Ready => "Ready",
Self::Failed => "Failed",
}
}
}
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
#[repr(u8)]
pub(crate) enum WorkSource {
#[default]
None,
LocalChat,
Http,
}
impl WorkSource {
pub(crate) fn label(self) -> &'static str {
match self {
Self::None => "None",
Self::LocalChat => "Local chat",
Self::Http => "HTTP endpoint",
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum KvLookup {
MemoryHit,
DiskHit,
Miss,
Invalid,
}
#[derive(Clone, Debug, Default)]
pub(crate) struct MetricsSnapshot {
pub(crate) uptime_seconds: u64,
pub(crate) phase: RuntimePhase,
pub(crate) source: WorkSource,
pub(crate) model: &'static str,
pub(crate) queue_depth: u32,
pub(crate) runtime_requests: u64,
pub(crate) local_requests: u64,
pub(crate) endpoint_generations: u64,
pub(crate) completed_requests: u64,
pub(crate) failed_requests: u64,
pub(crate) model_loads: u64,
pub(crate) model_unloads: u64,
pub(crate) model_load_ms: u64,
pub(crate) model_bytes: u64,
pub(crate) tensor_count: u64,
pub(crate) vocabulary_size: u64,
pub(crate) context_used: u32,
pub(crate) context_limit: u32,
pub(crate) decode_tokens_per_second: f32,
pub(crate) prefill_tokens_per_second: f32,
pub(crate) last_runtime_ms: u64,
pub(crate) average_runtime_ms: u64,
pub(crate) last_prompt_tokens: u64,
pub(crate) last_cached_tokens: u64,
pub(crate) last_completion_tokens: u64,
pub(crate) prompt_tokens: u64,
pub(crate) cached_tokens: u64,
pub(crate) completion_tokens: u64,
pub(crate) kv_lookups: u64,
pub(crate) kv_hits: u64,
pub(crate) kv_memory_hits: u64,
pub(crate) kv_disk_hits: u64,
pub(crate) kv_misses: u64,
pub(crate) kv_invalid: u64,
pub(crate) kv_prefix_hits: u64,
pub(crate) checkpoint_writes: u64,
pub(crate) kv_read_active: bool,
pub(crate) kv_write_active: bool,
pub(crate) kv_read_operations: u64,
pub(crate) kv_write_operations: u64,
pub(crate) kv_read_errors: u64,
pub(crate) kv_write_errors: u64,
pub(crate) kv_read_bytes: u64,
pub(crate) kv_write_bytes: u64,
pub(crate) last_kv_read_ms: u64,
pub(crate) last_kv_write_ms: u64,
pub(crate) kv_files: u64,
pub(crate) kv_bytes: u64,
pub(crate) local_kv_files: u64,
pub(crate) local_kv_bytes: u64,
pub(crate) http_kv_files: u64,
pub(crate) http_kv_bytes: u64,
pub(crate) server_listening: bool,
pub(crate) server_port: u16,
pub(crate) http_active: u32,
pub(crate) http_requests: u64,
pub(crate) http_completed: u64,
pub(crate) http_errors: u64,
pub(crate) http_chat_requests: u64,
pub(crate) http_model_requests: u64,
pub(crate) http_streaming_requests: u64,
pub(crate) http_bytes_received: u64,
pub(crate) last_http_ms: u64,
pub(crate) average_http_ms: u64,
}
pub(crate) struct Metrics {
started: Instant,
phase: AtomicU8,
source: AtomicU8,
model: AtomicU8,
queue_depth: AtomicU32,
runtime_requests: AtomicU64,
local_requests: AtomicU64,
endpoint_generations: AtomicU64,
completed_requests: AtomicU64,
failed_requests: AtomicU64,
model_loads: AtomicU64,
model_unloads: AtomicU64,
model_load_ms: AtomicU64,
model_bytes: AtomicU64,
tensor_count: AtomicU64,
vocabulary_size: AtomicU64,
context_used: AtomicU32,
context_limit: AtomicU32,
decode_tps: AtomicU32,
prefill_tps: AtomicU32,
prefill_sample: AtomicU32,
last_runtime_ms: AtomicU64,
total_runtime_ms: AtomicU64,
last_prompt_tokens: AtomicU64,
last_cached_tokens: AtomicU64,
last_completion_tokens: AtomicU64,
prompt_tokens: AtomicU64,
cached_tokens: AtomicU64,
completion_tokens: AtomicU64,
kv_lookups: AtomicU64,
kv_hits: AtomicU64,
kv_memory_hits: AtomicU64,
kv_disk_hits: AtomicU64,
kv_misses: AtomicU64,
kv_invalid: AtomicU64,
kv_prefix_hits: AtomicU64,
checkpoint_writes: AtomicU64,
kv_read_active: AtomicBool,
kv_write_active: AtomicBool,
kv_read_operations: AtomicU64,
kv_write_operations: AtomicU64,
kv_read_errors: AtomicU64,
kv_write_errors: AtomicU64,
kv_read_bytes: AtomicU64,
kv_write_bytes: AtomicU64,
kv_read_sample_bytes: AtomicU64,
kv_write_sample_bytes: AtomicU64,
last_kv_read_ms: AtomicU64,
last_kv_write_ms: AtomicU64,
kv_files: AtomicU64,
kv_bytes: AtomicU64,
local_kv_files: AtomicU64,
local_kv_bytes: AtomicU64,
http_kv_files: AtomicU64,
http_kv_bytes: AtomicU64,
server_listening: AtomicBool,
server_port: AtomicU16,
http_active: AtomicU32,
http_requests: AtomicU64,
http_completed: AtomicU64,
http_errors: AtomicU64,
http_chat_requests: AtomicU64,
http_model_requests: AtomicU64,
http_streaming_requests: AtomicU64,
http_bytes_received: AtomicU64,
last_http_ms: AtomicU64,
total_http_ms: AtomicU64,
}
impl Metrics {
pub(crate) fn new(cache_root: &Path) -> Self {
let cache = cache_usage(cache_root);
Self {
started: Instant::now(),
phase: AtomicU8::new(RuntimePhase::Unloaded as u8),
source: AtomicU8::new(WorkSource::None as u8),
model: AtomicU8::new(0),
queue_depth: AtomicU32::new(0),
runtime_requests: AtomicU64::new(0),
local_requests: AtomicU64::new(0),
endpoint_generations: AtomicU64::new(0),
completed_requests: AtomicU64::new(0),
failed_requests: AtomicU64::new(0),
model_loads: AtomicU64::new(0),
model_unloads: AtomicU64::new(0),
model_load_ms: AtomicU64::new(0),
model_bytes: AtomicU64::new(0),
tensor_count: AtomicU64::new(0),
vocabulary_size: AtomicU64::new(0),
context_used: AtomicU32::new(0),
context_limit: AtomicU32::new(0),
decode_tps: AtomicU32::new(0),
prefill_tps: AtomicU32::new(0),
prefill_sample: AtomicU32::new(0),
last_runtime_ms: AtomicU64::new(0),
total_runtime_ms: AtomicU64::new(0),
last_prompt_tokens: AtomicU64::new(0),
last_cached_tokens: AtomicU64::new(0),
last_completion_tokens: AtomicU64::new(0),
prompt_tokens: AtomicU64::new(0),
cached_tokens: AtomicU64::new(0),
completion_tokens: AtomicU64::new(0),
kv_lookups: AtomicU64::new(0),
kv_hits: AtomicU64::new(0),
kv_memory_hits: AtomicU64::new(0),
kv_disk_hits: AtomicU64::new(0),
kv_misses: AtomicU64::new(0),
kv_invalid: AtomicU64::new(0),
kv_prefix_hits: AtomicU64::new(0),
checkpoint_writes: AtomicU64::new(0),
kv_read_active: AtomicBool::new(false),
kv_write_active: AtomicBool::new(false),
kv_read_operations: AtomicU64::new(0),
kv_write_operations: AtomicU64::new(0),
kv_read_errors: AtomicU64::new(0),
kv_write_errors: AtomicU64::new(0),
kv_read_bytes: AtomicU64::new(0),
kv_write_bytes: AtomicU64::new(0),
kv_read_sample_bytes: AtomicU64::new(0),
kv_write_sample_bytes: AtomicU64::new(0),
last_kv_read_ms: AtomicU64::new(0),
last_kv_write_ms: AtomicU64::new(0),
kv_files: AtomicU64::new(cache.files),
kv_bytes: AtomicU64::new(cache.bytes),
local_kv_files: AtomicU64::new(cache.local_files),
local_kv_bytes: AtomicU64::new(cache.local_bytes),
http_kv_files: AtomicU64::new(cache.http_files),
http_kv_bytes: AtomicU64::new(cache.http_bytes),
server_listening: AtomicBool::new(false),
server_port: AtomicU16::new(0),
http_active: AtomicU32::new(0),
http_requests: AtomicU64::new(0),
http_completed: AtomicU64::new(0),
http_errors: AtomicU64::new(0),
http_chat_requests: AtomicU64::new(0),
http_model_requests: AtomicU64::new(0),
http_streaming_requests: AtomicU64::new(0),
http_bytes_received: AtomicU64::new(0),
last_http_ms: AtomicU64::new(0),
total_http_ms: AtomicU64::new(0),
}
}
pub(crate) fn request_queued(&self, source: WorkSource) {
self.runtime_requests.fetch_add(1, Ordering::Relaxed);
match source {
WorkSource::LocalChat => self.local_requests.fetch_add(1, Ordering::Relaxed),
WorkSource::Http => self.endpoint_generations.fetch_add(1, Ordering::Relaxed),
WorkSource::None => 0,
};
self.queue_depth.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn request_rejected(&self) {
self.queue_depth.fetch_sub(1, Ordering::Relaxed);
self.failed_requests.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn request_started(&self, source: WorkSource) {
self.queue_depth.fetch_sub(1, Ordering::Relaxed);
self.source.store(source as u8, Ordering::Relaxed);
self.decode_tps.store(0, Ordering::Relaxed);
self.prefill_tps.store(0, Ordering::Relaxed);
self.prefill_sample.store(0, Ordering::Relaxed);
}
pub(crate) fn loading(&self) {
self.phase
.store(RuntimePhase::Loading as u8, Ordering::Relaxed);
}
pub(crate) fn loaded(
&self,
model: ModelChoice,
elapsed: Duration,
mapped_bytes: u64,
tensors: usize,
vocabulary_size: usize,
) {
self.model.store(model_code(model), Ordering::Relaxed);
self.model_loads.fetch_add(1, Ordering::Relaxed);
self.model_load_ms
.store(milliseconds(elapsed), Ordering::Relaxed);
self.model_bytes.store(mapped_bytes, Ordering::Relaxed);
self.tensor_count.store(tensors as u64, Ordering::Relaxed);
self.vocabulary_size
.store(vocabulary_size as u64, Ordering::Relaxed);
}
pub(crate) fn prefill_progress(&self, used: u32, limit: u32, speed: f32) {
self.phase
.store(RuntimePhase::Prefilling as u8, Ordering::Relaxed);
self.context_used.store(used, Ordering::Relaxed);
self.context_limit.store(limit, Ordering::Relaxed);
store_f32(&self.prefill_tps, speed);
store_f32(&self.prefill_sample, speed);
}
pub(crate) fn generation_progress(&self, used: u32, limit: u32, speed: f32) {
self.phase
.store(RuntimePhase::Generating as u8, Ordering::Relaxed);
self.context_used.store(used, Ordering::Relaxed);
self.context_limit.store(limit, Ordering::Relaxed);
store_f32(&self.decode_tps, speed);
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn request_finished(
&self,
source: WorkSource,
elapsed: Duration,
prompt_tokens: u32,
cached_tokens: u32,
completion_tokens: u32,
previous_checkpoint_bytes: Option<u64>,
checkpoint_bytes: u64,
) {
let elapsed = milliseconds(elapsed);
self.completed_requests.fetch_add(1, Ordering::Relaxed);
self.last_runtime_ms.store(elapsed, Ordering::Relaxed);
self.total_runtime_ms.fetch_add(elapsed, Ordering::Relaxed);
self.last_prompt_tokens
.store(u64::from(prompt_tokens), Ordering::Relaxed);
self.last_cached_tokens
.store(u64::from(cached_tokens), Ordering::Relaxed);
self.last_completion_tokens
.store(u64::from(completion_tokens), Ordering::Relaxed);
self.prompt_tokens
.fetch_add(u64::from(prompt_tokens), Ordering::Relaxed);
self.cached_tokens
.fetch_add(u64::from(cached_tokens), Ordering::Relaxed);
self.completion_tokens
.fetch_add(u64::from(completion_tokens), Ordering::Relaxed);
self.record_checkpoint(source, previous_checkpoint_bytes, checkpoint_bytes);
self.phase
.store(RuntimePhase::Ready as u8, Ordering::Relaxed);
self.source.store(WorkSource::None as u8, Ordering::Relaxed);
}
pub(crate) fn request_failed(&self, elapsed: Duration) {
let elapsed = milliseconds(elapsed);
self.failed_requests.fetch_add(1, Ordering::Relaxed);
self.last_runtime_ms.store(elapsed, Ordering::Relaxed);
self.total_runtime_ms.fetch_add(elapsed, Ordering::Relaxed);
self.phase
.store(RuntimePhase::Failed as u8, Ordering::Relaxed);
self.source.store(WorkSource::None as u8, Ordering::Relaxed);
}
pub(crate) fn unloaded(&self) {
self.phase
.store(RuntimePhase::Unloaded as u8, Ordering::Relaxed);
self.model.store(0, Ordering::Relaxed);
self.model_bytes.store(0, Ordering::Relaxed);
self.tensor_count.store(0, Ordering::Relaxed);
self.vocabulary_size.store(0, Ordering::Relaxed);
self.context_used.store(0, Ordering::Relaxed);
self.decode_tps.store(0, Ordering::Relaxed);
self.prefill_tps.store(0, Ordering::Relaxed);
self.prefill_sample.store(0, Ordering::Relaxed);
self.model_unloads.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn server_listening(&self, port: u16) {
self.server_port.store(port, Ordering::Relaxed);
self.server_listening.store(true, Ordering::Relaxed);
}
pub(crate) fn kv_lookup(&self, result: KvLookup) {
self.kv_lookups.fetch_add(1, Ordering::Relaxed);
match result {
KvLookup::MemoryHit => {
self.kv_hits.fetch_add(1, Ordering::Relaxed);
self.kv_memory_hits.fetch_add(1, Ordering::Relaxed);
}
KvLookup::DiskHit => {
self.kv_hits.fetch_add(1, Ordering::Relaxed);
self.kv_disk_hits.fetch_add(1, Ordering::Relaxed);
}
KvLookup::Miss => {
self.kv_misses.fetch_add(1, Ordering::Relaxed);
}
KvLookup::Invalid => {
self.kv_misses.fetch_add(1, Ordering::Relaxed);
self.kv_invalid.fetch_add(1, Ordering::Relaxed);
}
}
}
pub(crate) fn kv_prefix_reused(&self, tokens: usize) {
if tokens > 0 {
self.kv_prefix_hits.fetch_add(1, Ordering::Relaxed);
}
}
pub(crate) fn kv_read_started(&self) {
self.kv_read_active.store(true, Ordering::Relaxed);
self.kv_read_operations.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn kv_read_bytes(&self, bytes: u64) {
self.kv_read_bytes.fetch_add(bytes, Ordering::Relaxed);
self.kv_read_sample_bytes
.fetch_add(bytes, Ordering::Relaxed);
}
pub(crate) fn kv_read_finished(&self, elapsed: Duration, failed: bool) {
self.kv_read_active.store(false, Ordering::Relaxed);
self.last_kv_read_ms
.store(milliseconds(elapsed), Ordering::Relaxed);
if failed {
self.kv_read_errors.fetch_add(1, Ordering::Relaxed);
}
}
pub(crate) fn kv_write_started(&self) {
self.kv_write_active.store(true, Ordering::Relaxed);
self.kv_write_operations.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn kv_write_bytes(&self, bytes: u64) {
self.kv_write_bytes.fetch_add(bytes, Ordering::Relaxed);
self.kv_write_sample_bytes
.fetch_add(bytes, Ordering::Relaxed);
}
pub(crate) fn kv_write_finished(&self, elapsed: Duration, failed: bool) {
self.kv_write_active.store(false, Ordering::Relaxed);
self.last_kv_write_ms
.store(milliseconds(elapsed), Ordering::Relaxed);
if failed {
self.kv_write_errors.fetch_add(1, Ordering::Relaxed);
}
}
pub(crate) fn server_stopped(&self, port: u16) {
if self.server_port.load(Ordering::Relaxed) == port {
self.server_listening.store(false, Ordering::Relaxed);
self.http_active.store(0, Ordering::Relaxed);
}
}
pub(crate) fn http_started(&self, path: &str, bytes: usize) {
self.http_active.fetch_add(1, Ordering::Relaxed);
self.http_requests.fetch_add(1, Ordering::Relaxed);
self.http_bytes_received
.fetch_add(bytes as u64, Ordering::Relaxed);
match path {
"/v1/chat/completions" => {
self.http_chat_requests.fetch_add(1, Ordering::Relaxed);
}
path if path == "/v1/models" || path.starts_with("/v1/models/") => {
self.http_model_requests.fetch_add(1, Ordering::Relaxed);
}
_ => {}
}
}
pub(crate) fn http_streaming(&self) {
self.http_streaming_requests.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn http_finished(&self, elapsed: Duration, failed: bool) {
self.http_active.fetch_sub(1, Ordering::Relaxed);
self.http_completed.fetch_add(1, Ordering::Relaxed);
if failed {
self.http_errors.fetch_add(1, Ordering::Relaxed);
}
let elapsed = milliseconds(elapsed);
self.last_http_ms.store(elapsed, Ordering::Relaxed);
self.total_http_ms.fetch_add(elapsed, Ordering::Relaxed);
}
pub(crate) fn snapshot(&self) -> MetricsSnapshot {
let completed = self.completed_requests.load(Ordering::Relaxed);
let failed = self.failed_requests.load(Ordering::Relaxed);
let http_completed = self.http_completed.load(Ordering::Relaxed);
MetricsSnapshot {
uptime_seconds: self.started.elapsed().as_secs(),
phase: phase(self.phase.load(Ordering::Relaxed)),
source: source(self.source.load(Ordering::Relaxed)),
model: model_name(self.model.load(Ordering::Relaxed)),
queue_depth: self.queue_depth.load(Ordering::Relaxed),
runtime_requests: self.runtime_requests.load(Ordering::Relaxed),
local_requests: self.local_requests.load(Ordering::Relaxed),
endpoint_generations: self.endpoint_generations.load(Ordering::Relaxed),
completed_requests: completed,
failed_requests: failed,
model_loads: self.model_loads.load(Ordering::Relaxed),
model_unloads: self.model_unloads.load(Ordering::Relaxed),
model_load_ms: self.model_load_ms.load(Ordering::Relaxed),
model_bytes: self.model_bytes.load(Ordering::Relaxed),
tensor_count: self.tensor_count.load(Ordering::Relaxed),
vocabulary_size: self.vocabulary_size.load(Ordering::Relaxed),
context_used: self.context_used.load(Ordering::Relaxed),
context_limit: self.context_limit.load(Ordering::Relaxed),
decode_tokens_per_second: load_f32(&self.decode_tps),
prefill_tokens_per_second: load_f32(&self.prefill_tps),
last_runtime_ms: self.last_runtime_ms.load(Ordering::Relaxed),
average_runtime_ms: average(
self.total_runtime_ms.load(Ordering::Relaxed),
completed + failed,
),
last_prompt_tokens: self.last_prompt_tokens.load(Ordering::Relaxed),
last_cached_tokens: self.last_cached_tokens.load(Ordering::Relaxed),
last_completion_tokens: self.last_completion_tokens.load(Ordering::Relaxed),
prompt_tokens: self.prompt_tokens.load(Ordering::Relaxed),
cached_tokens: self.cached_tokens.load(Ordering::Relaxed),
completion_tokens: self.completion_tokens.load(Ordering::Relaxed),
kv_lookups: self.kv_lookups.load(Ordering::Relaxed),
kv_hits: self.kv_hits.load(Ordering::Relaxed),
kv_memory_hits: self.kv_memory_hits.load(Ordering::Relaxed),
kv_disk_hits: self.kv_disk_hits.load(Ordering::Relaxed),
kv_misses: self.kv_misses.load(Ordering::Relaxed),
kv_invalid: self.kv_invalid.load(Ordering::Relaxed),
kv_prefix_hits: self.kv_prefix_hits.load(Ordering::Relaxed),
checkpoint_writes: self.checkpoint_writes.load(Ordering::Relaxed),
kv_read_active: self.kv_read_active.load(Ordering::Relaxed),
kv_write_active: self.kv_write_active.load(Ordering::Relaxed),
kv_read_operations: self.kv_read_operations.load(Ordering::Relaxed),
kv_write_operations: self.kv_write_operations.load(Ordering::Relaxed),
kv_read_errors: self.kv_read_errors.load(Ordering::Relaxed),
kv_write_errors: self.kv_write_errors.load(Ordering::Relaxed),
kv_read_bytes: self.kv_read_bytes.load(Ordering::Relaxed),
kv_write_bytes: self.kv_write_bytes.load(Ordering::Relaxed),
last_kv_read_ms: self.last_kv_read_ms.load(Ordering::Relaxed),
last_kv_write_ms: self.last_kv_write_ms.load(Ordering::Relaxed),
kv_files: self.kv_files.load(Ordering::Relaxed),
kv_bytes: self.kv_bytes.load(Ordering::Relaxed),
local_kv_files: self.local_kv_files.load(Ordering::Relaxed),
local_kv_bytes: self.local_kv_bytes.load(Ordering::Relaxed),
http_kv_files: self.http_kv_files.load(Ordering::Relaxed),
http_kv_bytes: self.http_kv_bytes.load(Ordering::Relaxed),
server_listening: self.server_listening.load(Ordering::Relaxed),
server_port: self.server_port.load(Ordering::Relaxed),
http_active: self.http_active.load(Ordering::Relaxed),
http_requests: self.http_requests.load(Ordering::Relaxed),
http_completed,
http_errors: self.http_errors.load(Ordering::Relaxed),
http_chat_requests: self.http_chat_requests.load(Ordering::Relaxed),
http_model_requests: self.http_model_requests.load(Ordering::Relaxed),
http_streaming_requests: self.http_streaming_requests.load(Ordering::Relaxed),
http_bytes_received: self.http_bytes_received.load(Ordering::Relaxed),
last_http_ms: self.last_http_ms.load(Ordering::Relaxed),
average_http_ms: average(self.total_http_ms.load(Ordering::Relaxed), http_completed),
}
}
pub(crate) fn take_prefill_sample(&self) -> f32 {
f32::from_bits(self.prefill_sample.swap(0, Ordering::Relaxed))
}
pub(crate) fn take_kv_io_sample(&self) -> (u64, u64) {
(
self.kv_read_sample_bytes.swap(0, Ordering::Relaxed),
self.kv_write_sample_bytes.swap(0, Ordering::Relaxed),
)
}
/// Re-reads the cache directories after files are removed outside the
/// runtime, so the counters do not drift away from the disc.
pub(crate) fn rescan_cache(&self, root: &Path) {
let usage = cache_usage(root);
self.kv_files.store(usage.files, Ordering::Relaxed);
self.kv_bytes.store(usage.bytes, Ordering::Relaxed);
self.local_kv_files
.store(usage.local_files, Ordering::Relaxed);
self.local_kv_bytes
.store(usage.local_bytes, Ordering::Relaxed);
self.http_kv_files
.store(usage.http_files, Ordering::Relaxed);
self.http_kv_bytes
.store(usage.http_bytes, Ordering::Relaxed);
}
fn record_checkpoint(&self, source: WorkSource, previous_bytes: Option<u64>, bytes: u64) {
if bytes == 0 {
return;
}
self.checkpoint_writes.fetch_add(1, Ordering::Relaxed);
let previous = previous_bytes.unwrap_or(0);
update_bytes(&self.kv_bytes, previous, bytes);
let (files, total) = match source {
WorkSource::LocalChat => (&self.local_kv_files, &self.local_kv_bytes),
WorkSource::Http => (&self.http_kv_files, &self.http_kv_bytes),
WorkSource::None => return,
};
update_bytes(total, previous, bytes);
if previous_bytes.is_none() {
self.kv_files.fetch_add(1, Ordering::Relaxed);
files.fetch_add(1, Ordering::Relaxed);
}
}
}
fn store_f32(target: &AtomicU32, value: f32) {
target.store(value.max(0.0).to_bits(), Ordering::Relaxed);
}
fn load_f32(source: &AtomicU32) -> f32 {
f32::from_bits(source.load(Ordering::Relaxed))
}
fn milliseconds(duration: Duration) -> u64 {
duration.as_millis().min(u128::from(u64::MAX)) as u64
}
fn average(total: u64, count: u64) -> u64 {
total.checked_div(count).unwrap_or(0)
}
fn update_bytes(target: &AtomicU64, previous: u64, current: u64) {
let total = target.load(Ordering::Relaxed);
target.store(total.saturating_sub(previous) + current, Ordering::Relaxed);
}
fn phase(value: u8) -> RuntimePhase {
match value {
1 => RuntimePhase::Loading,
2 => RuntimePhase::Prefilling,
3 => RuntimePhase::Generating,
4 => RuntimePhase::Ready,
5 => RuntimePhase::Failed,
_ => RuntimePhase::Unloaded,
}
}
fn source(value: u8) -> WorkSource {
match value {
1 => WorkSource::LocalChat,
2 => WorkSource::Http,
_ => WorkSource::None,
}
}
fn model_code(model: ModelChoice) -> u8 {
match model {
ModelChoice::DeepSeekV4Flash => 1,
ModelChoice::DeepSeekV4Pro => 2,
ModelChoice::Glm52 => 3,
}
}
fn model_name(value: u8) -> &'static str {
match value {
1 => "DeepSeek V4 Flash",
2 => "DeepSeek V4 Pro",
3 => "GLM 5.2",
_ => "No model loaded",
}
}
/// Age buckets for the disc usage explorer. The six hour boundary is one
/// eviction hit half-life, so the distribution shows what is about to lose its
/// reuse credit rather than just what is old.
pub(crate) const CACHE_AGE_BUCKETS: [(&str, u64); 5] = [
("Under 1h", 3_600),
("16h", 6 * 3_600),
("624h", 24 * 3_600),
("17d", 7 * 86_400),
("Older", u64::MAX),
];
/// Rows shown in the explorer, oldest first — the order eviction works through.
const CACHE_ENTRY_ROWS: usize = 12;
#[derive(Clone, Debug)]
pub(crate) struct KvCacheEntry {
pub(crate) path: PathBuf,
pub(crate) name: String,
/// Set for session checkpoints, so the view can name the conversation.
pub(crate) session: Option<i32>,
pub(crate) bytes: u64,
pub(crate) age_seconds: u64,
}
#[derive(Clone, Debug, Default)]
pub(crate) struct KvCacheReport {
pub(crate) budget_bytes: u64,
pub(crate) session_bytes: u64,
pub(crate) session_files: u64,
pub(crate) transient_bytes: u64,
pub(crate) transient_files: u64,
/// Bytes per [`CACHE_AGE_BUCKETS`] entry, both buckets combined.
pub(crate) age_bytes: [u64; CACHE_AGE_BUCKETS.len()],
pub(crate) entries: Vec<KvCacheEntry>,
}
impl KvCacheReport {
pub(crate) fn total_bytes(&self) -> u64 {
self.session_bytes + self.transient_bytes
}
}
/// Scans the KV cache directories for the stats explorer. The budget applies to
/// the transient store only; session checkpoints are user-owned and are freed
/// with their session.
pub(crate) fn kv_cache_report(root: &Path, budget_bytes: u64) -> KvCacheReport {
let mut report = KvCacheReport {
budget_bytes,
..KvCacheReport::default()
};
let now = SystemTime::now();
for (directory, transient) in [(root.to_path_buf(), false), (root.join("http"), true)] {
let Ok(files) = fs::read_dir(&directory) else {
continue;
};
for file in files.flatten() {
let path = file.path();
let Ok(metadata) = file.metadata() else {
continue;
};
if !metadata.is_file() {
continue;
}
let bytes = metadata.len();
if transient {
report.transient_bytes += bytes;
report.transient_files += 1;
} else {
report.session_bytes += bytes;
report.session_files += 1;
}
let age_seconds = metadata
.modified()
.ok()
.and_then(|modified| now.duration_since(modified).ok())
.map_or(0, |age| age.as_secs());
let bucket = CACHE_AGE_BUCKETS
.iter()
.position(|(_, limit)| age_seconds < *limit)
.unwrap_or(CACHE_AGE_BUCKETS.len() - 1);
report.age_bytes[bucket] += bytes;
// Index files are counted in the totals but folded into their
// checkpoint's row.
if path.extension().is_some_and(|value| value == "meta") {
continue;
}
report.entries.push(KvCacheEntry {
name: entry_name(&path, transient),
session: (!transient)
.then(|| path.file_stem().and_then(|value| value.to_str()))
.flatten()
.and_then(|stem| stem.parse().ok()),
bytes: bytes.saturating_add(
fs::metadata(path.with_extension("meta")).map_or(0, |index| index.len()),
),
path,
age_seconds,
});
}
}
report.entries.sort_by(|left, right| {
right
.age_seconds
.cmp(&left.age_seconds)
.then_with(|| right.bytes.cmp(&left.bytes))
.then_with(|| left.name.cmp(&right.name))
});
report.entries.truncate(CACHE_ENTRY_ROWS);
report
}
fn entry_name(path: &Path, transient: bool) -> String {
let stem = path
.file_stem()
.and_then(|value| value.to_str())
.unwrap_or("unknown");
if transient {
format!("Transient {}", &stem[..stem.len().min(12)])
} else {
format!("Session {stem}")
}
}
#[derive(Default)]
struct CacheUsage {
files: u64,
bytes: u64,
local_files: u64,
local_bytes: u64,
http_files: u64,
http_bytes: u64,
}
fn cache_usage(root: &Path) -> CacheUsage {
let mut usage = CacheUsage::default();
for (path, http) in [(root.to_path_buf(), false), (root.join("http"), true)] {
let Ok(entries) = fs::read_dir(path) else {
continue;
};
for entry in entries.flatten() {
let Ok(metadata) = entry.metadata() else {
continue;
};
if !metadata.is_file() {
continue;
}
usage.files += 1;
usage.bytes += metadata.len();
if http {
usage.http_files += 1;
usage.http_bytes += metadata.len();
} else {
usage.local_files += 1;
usage.local_bytes += metadata.len();
}
}
}
usage
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn snapshots_track_runtime_and_server_counters() {
let metrics = Metrics::new(Path::new("/path/that/does/not/exist"));
metrics.request_queued(WorkSource::LocalChat);
metrics.request_started(WorkSource::LocalChat);
metrics.prefill_progress(100, 1_000, 250.0);
metrics.generation_progress(120, 1_000, 20.0);
assert_eq!(metrics.take_prefill_sample(), 250.0);
assert_eq!(metrics.take_prefill_sample(), 0.0);
metrics.kv_lookup(KvLookup::DiskHit);
metrics.kv_lookup(KvLookup::MemoryHit);
metrics.kv_lookup(KvLookup::Miss);
metrics.kv_lookup(KvLookup::Invalid);
metrics.kv_prefix_reused(80);
metrics.kv_prefix_reused(0);
metrics.kv_read_started();
metrics.kv_read_bytes(2_048);
metrics.kv_read_finished(Duration::from_millis(10), false);
metrics.kv_write_started();
metrics.kv_write_bytes(4_096);
metrics.kv_write_finished(Duration::from_millis(20), false);
assert_eq!(metrics.take_kv_io_sample(), (2_048, 4_096));
assert_eq!(metrics.take_kv_io_sample(), (0, 0));
metrics.request_finished(
WorkSource::LocalChat,
Duration::from_millis(250),
100,
80,
20,
None,
4_096,
);
metrics.http_started("/v1/models", 0);
metrics.http_finished(Duration::from_millis(5), false);
let snapshot = metrics.snapshot();
assert_eq!(snapshot.phase, RuntimePhase::Ready);
assert_eq!(snapshot.cached_tokens, 80);
assert_eq!(snapshot.kv_lookups, 4);
assert_eq!(snapshot.kv_hits, 2);
assert_eq!(snapshot.kv_memory_hits, 1);
assert_eq!(snapshot.kv_disk_hits, 1);
assert_eq!(snapshot.kv_misses, 2);
assert_eq!(snapshot.kv_invalid, 1);
assert_eq!(snapshot.kv_prefix_hits, 1);
assert_eq!(snapshot.kv_read_bytes, 2_048);
assert_eq!(snapshot.kv_write_bytes, 4_096);
assert_eq!(snapshot.local_kv_bytes, 4_096);
assert_eq!(snapshot.http_model_requests, 1);
}
#[test]
fn cache_report_splits_buckets_and_folds_index_files_into_their_checkpoint() {
let root = std::env::temp_dir().join(format!(
"ds4-server-cache-report-{}-{:?}",
std::process::id(),
Instant::now()
));
fs::create_dir_all(root.join("http")).unwrap();
fs::write(root.join("7.bin"), vec![0; 300]).unwrap();
fs::write(root.join("http/abc.bin"), vec![0; 100]).unwrap();
fs::write(root.join("http/abc.meta"), vec![0; 40]).unwrap();
// The largest file is also the oldest, so size cannot pass for age.
let aged = fs::File::options()
.write(true)
.open(root.join("http/abc.bin"))
.unwrap();
aged.set_modified(SystemTime::now() - Duration::from_secs(2 * 86_400))
.unwrap();
let report = kv_cache_report(&root, 4 * 1024);
assert_eq!(report.session_bytes, 300);
assert_eq!(report.session_files, 1);
assert_eq!(report.transient_bytes, 140);
assert_eq!(report.transient_files, 2);
// Fresh files in the first bucket, the two-day-old one in "17d".
assert_eq!(report.age_bytes[0], 340);
assert_eq!(report.age_bytes[3], 100);
assert_eq!(
(
report.age_bytes[1],
report.age_bytes[2],
report.age_bytes[4]
),
(0, 0, 0)
);
// Oldest first, and the index file rides along with its checkpoint.
assert_eq!(report.entries.len(), 2);
assert_eq!(report.entries[0].bytes, 140);
assert_eq!(report.entries[1].name, "Session 7");
// The session id survives on the row, so the view can add its title.
assert_eq!(report.entries[1].session, Some(7));
assert_eq!(report.entries[0].session, None);
fs::remove_dir_all(root).unwrap();
}
}