diff --git a/Cargo.lock b/Cargo.lock index 0605e28f6..b73f81300 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7541,6 +7541,7 @@ dependencies = [ "clap", "futures-util", "libc", + "mesh-llm-events", "mesh-native-serving-plugin-api", "model-artifact", "openai-frontend", diff --git a/crates/skippy-cache/src/disk_tier.rs b/crates/skippy-cache/src/disk_tier.rs index 06900734a..022425c29 100644 --- a/crates/skippy-cache/src/disk_tier.rs +++ b/crates/skippy-cache/src/disk_tier.rs @@ -565,6 +565,17 @@ impl PrefixDiskTier { })) } + /// Mark an entry as recently used without mapping or verifying its payload. + pub fn touch(&mut self, page_id: &str) -> bool { + let Some(entry) = self.entries.get_mut(page_id) else { + return false; + }; + self.use_clock = self.use_clock.saturating_add(1); + entry.last_used_secs = now_secs(); + entry.use_sequence = self.use_clock; + true + } + pub fn contains(&self, page_id: &str) -> bool { self.entries.contains_key(page_id) } diff --git a/crates/skippy-cache/src/exact_state.rs b/crates/skippy-cache/src/exact_state.rs index 5107b01a8..e86c775df 100644 --- a/crates/skippy-cache/src/exact_state.rs +++ b/crates/skippy-cache/src/exact_state.rs @@ -206,6 +206,21 @@ impl ExactStateCache { self.misses.note_identity_mismatch(); } + /// Mark an existing page as recently used without reconstructing or + /// replacing its payload. + /// + /// Page identities include the complete prefix identity, so an existing + /// entry is already the checkpoint the caller intends to record. This + /// lets record paths avoid exporting and re-hashing the same state. + pub fn touch(&mut self, page_id: &str) -> bool { + self.clock = self.clock.saturating_add(1); + if let Some(entry) = self.entries.get_mut(page_id) { + entry.last_used = self.clock; + return true; + } + self.disk.as_mut().is_some_and(|disk| disk.touch(page_id)) + } + pub fn disk_contains(&self, page_id: &str) -> bool { self.disk .as_ref() @@ -515,6 +530,36 @@ mod tests { assert_eq!(cache.stats().entries, 1); } + #[test] + fn touching_existing_page_refreshes_lru_without_replacing_payload() { + let mut cache = ExactStateCache::new(2, 0); + cache.record( + "first".to_string(), + 2, + ExactStatePayload::full_state(vec![1, 2]), + (), + ); + cache.record( + "second".to_string(), + 2, + ExactStatePayload::full_state(vec![3, 4]), + (), + ); + + assert!(cache.touch("first")); + assert!(!cache.touch("missing")); + cache.record( + "third".to_string(), + 2, + ExactStatePayload::full_state(vec![5, 6]), + (), + ); + + assert!(cache.lookup("first").is_some()); + assert!(cache.lookup("second").is_none()); + assert!(cache.lookup("third").is_some()); + } + #[test] fn cached_token_counts_are_bounded_sorted_and_deduplicated() { let mut cache = ExactStateCache::new(4, 0); diff --git a/crates/skippy-server/Cargo.toml b/crates/skippy-server/Cargo.toml index 0a868decd..aac69516b 100644 --- a/crates/skippy-server/Cargo.toml +++ b/crates/skippy-server/Cargo.toml @@ -28,6 +28,7 @@ futures-util = "0.3" skippy-runtime = { path = "../skippy-runtime", version = "0.76.0-rc4" } skippy-tokenizer = { path = "../skippy-tokenizer", version = "0.76.0-rc4" } model-artifact = { path = "../model-artifact", version = "0.76.0-rc4" } +mesh-llm-events = { path = "../mesh-llm-events", version = "0.76.0-rc4" } mesh-native-serving-plugin-api = { path = "../mesh-native-serving-plugin-api", version = "0.76.0-rc4" } skippy-protocol = { path = "../skippy-protocol", version = "0.76.0-rc4" } skippy-cache = { path = "../skippy-cache", version = "0.76.0-rc4" } diff --git a/crates/skippy-server/src/frontend/local_generation/token_generation.rs b/crates/skippy-server/src/frontend/local_generation/token_generation.rs index 943807ecc..5ca75c2ee 100644 --- a/crates/skippy-server/src/frontend/local_generation/token_generation.rs +++ b/crates/skippy-server/src/frontend/local_generation/token_generation.rs @@ -600,6 +600,18 @@ impl StageOpenAiBackend { "skippy.exact_cache.reconstruct_blocks".to_string(), json!(restored.reconstruct_blocks), ); + attrs.insert( + "skippy.exact_cache.lookup_ms".to_string(), + json!(restored.lookup_ms), + ); + attrs.insert( + "skippy.exact_cache.kv_import_ms".to_string(), + json!(restored.kv_import_ms), + ); + attrs.insert( + "skippy.exact_cache.recurrent_import_ms".to_string(), + json!(restored.recurrent_import_ms), + ); self.telemetry .emit("stage.openai_kv_lookup_decision", attrs); } @@ -1088,13 +1100,10 @@ impl StageOpenAiBackend { "skippy.exact_cache.recorded_tokens".to_string(), json!(record.token_count), ); - attrs.insert( - "skippy.exact_cache.stored".to_string(), - json!(record.stored), - ); + attrs.insert("skippy.exact_cache.queued".to_string(), json!(true)); self.telemetry .emit("stage.openai_kv_record_decision", attrs); - record.stored + true } Ok(None) => false, Err(error) => { diff --git a/crates/skippy-server/src/kv_integration/config.rs b/crates/skippy-server/src/kv_integration/config.rs index 2cccee8f1..8fca1e21e 100644 --- a/crates/skippy-server/src/kv_integration/config.rs +++ b/crates/skippy-server/src/kv_integration/config.rs @@ -27,6 +27,7 @@ pub fn configure_kv_disk_cache(config: KvDiskCacheConfig) -> Result<(), KvDiskCa } use anyhow::Result; +use mesh_llm_events::{OutputEvent, emit_event}; use skippy_cache::{ ExactStateCache, PrefixCandidatePolicy, PrefixDiskTier, ResidentActivationCache, ResidentCacheConfig, ResidentPrefixCache, @@ -38,8 +39,8 @@ use skippy_runtime::ModelInfo; use skippy_topology::{STAGE_RUNTIME_LLAMA_FAMILY_EXPECTATIONS, infer_family_capability}; use super::{ - ExactStateExtra, KvStageIntegration, StageKvMode, StagePrefixCachePayload, disk_budget, - disk_budget::NodeBudget, + EXACT_STATE_RECORD_CAPACITY, ExactStateExtra, KvStageIntegration, PendingExactStateRecord, + StageKvMode, StagePrefixCachePayload, disk_budget, disk_budget::NodeBudget, }; impl KvStageIntegration { @@ -78,16 +79,54 @@ impl KvStageIntegration { if let Some(opened) = disk { exact_states = exact_states.with_disk_tier(opened.tier); } + let exact_states = Arc::new(Mutex::new(exact_states)); + let (exact_state_record_tx, exact_state_record_rx) = + std::sync::mpsc::sync_channel::(EXACT_STATE_RECORD_CAPACITY); + let worker_exact_states = exact_states.clone(); + let inflight_records = Arc::new(Mutex::new(BTreeSet::new())); + let worker_inflight_records = inflight_records.clone(); + let exact_state_records_queued = Arc::new(std::sync::atomic::AtomicU64::new(0)); + let exact_state_records_dropped = Arc::new(std::sync::atomic::AtomicU64::new(0)); + let exact_state_records_pending = Arc::new(std::sync::atomic::AtomicUsize::new(0)); + let worker_exact_state_records_pending = exact_state_records_pending.clone(); + let worker_disk_budget_reservation = disk_budget_reservation.clone(); + std::thread::Builder::new() + .name(format!("skippy-exact-cache-{}", config.stage_id)) + .spawn(move || { + let _disk_budget_reservation = worker_disk_budget_reservation; + while let Ok(pending) = exact_state_record_rx.recv() { + let page_id = pending.page_id.clone(); + worker_exact_states + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .record( + pending.page_id, + pending.token_count, + pending.payload, + pending.extra, + ); + worker_inflight_records + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .remove(&page_id); + worker_exact_state_records_pending + .fetch_sub(1, std::sync::atomic::Ordering::Release); + } + })?; Ok(Some(Self { mode, payload, correctness_mode: false, trust_local_writes: true, candidate_policy, - inflight_records: Arc::new(Mutex::new(BTreeSet::new())), + inflight_records, resident: Arc::new(Mutex::new(ResidentPrefixCache::new(resident_config))), activations: Arc::new(Mutex::new(ResidentActivationCache::new(resident_config))), - exact_states: Arc::new(Mutex::new(exact_states)), + exact_states, + exact_state_record_tx, + exact_state_records_queued, + exact_state_records_dropped, + exact_state_records_pending, first_tokens: Arc::new(Mutex::new(BTreeMap::new())), replay_tokens: Arc::new(Mutex::new(BTreeMap::new())), split_prefill_tokens: Arc::new(Mutex::new(BTreeMap::new())), @@ -97,6 +136,13 @@ impl KvStageIntegration { } } +fn emit_warning(message: String) { + let _ = emit_event(OutputEvent::Warning { + message, + context: None, + }); +} + /// Open the KV disk tier for this stage. Host configuration wins; legacy /// environment variables remain compatibility input when no host configured it. struct OpenedDiskTier { @@ -107,17 +153,19 @@ struct OpenedDiskTier { fn open_disk_tier(config: &StageConfig) -> Option { let root = disk_tier_root(config); if !has_valid_content_digest(config) { - eprintln!( + emit_warning(format!( "skippy: KV disk tier disabled for stage {}: no valid content digest", config.stage_id - ); + )); return None; } let reservation = stage_disk_budget(&root, config)?; match PrefixDiskTier::open(&root, reservation.bytes()) { Ok(tier) => Some(OpenedDiskTier { tier, reservation }), Err(error) => { - eprintln!("skippy: KV disk tier unavailable, continuing without it: {error}"); + emit_warning(format!( + "skippy: KV disk tier unavailable, continuing without it: {error}" + )); None } } @@ -149,12 +197,12 @@ fn stage_disk_budget(root: &Path, config: &StageConfig) -> Option { let (full_state, stats) = lookup @@ -60,11 +67,13 @@ impl KvStageIntegration { &mut reconstruct_blocks, stats, ); + let import_started = Instant::now(); runtime.import_full_state_for_token_count( session_id, full_state.as_ref(), lookup.token_count, )?; + kv_import_ms = import_started.elapsed().as_secs_f64() * 1000.0; } StagePrefixCachePayload::KvRecurrent => { if let Some((kv, stats)) = lookup @@ -79,7 +88,9 @@ impl KvStageIntegration { stats, ); if let Some(desc) = lookup.extra.kv_desc { + let import_started = Instant::now(); runtime.import_kv_page(session_id, &desc, kv.as_ref())?; + kv_import_ms = import_started.elapsed().as_secs_f64() * 1000.0; } else if !kv.is_empty() { continue; } @@ -94,11 +105,13 @@ impl KvStageIntegration { &mut reconstruct_blocks, stats, ); + let import_started = Instant::now(); runtime.import_recurrent_state_for_token_count( session_id, recurrent.as_ref(), lookup.token_count, )?; + recurrent_import_ms = import_started.elapsed().as_secs_f64() * 1000.0; } _ => continue, } @@ -116,6 +129,9 @@ impl KvStageIntegration { reconstruct_ms, reconstruct_bytes, reconstruct_blocks, + lookup_ms, + kv_import_ms, + recurrent_import_ms, })); } Ok(None) @@ -137,51 +153,116 @@ impl KvStageIntegration { if !self.try_begin_record(&identity.page_id) { return Ok(None); } - let result = (|| { - let (payload, extra) = match self.payload { - StagePrefixCachePayload::FullState => ( - ExactStatePayload::full_state(runtime.export_full_state(session_id)?), - ExactStateExtra::default(), - ), - StagePrefixCachePayload::KvRecurrent => { - let kv = match runtime.export_kv_page(session_id, 0, token_count) { - Ok(kv) => Some(kv), - Err(error) if is_native_kv_unavailable(&error) => None, - Err(error) => return Err(error), - }; - let recurrent = runtime.export_recurrent_state(session_id)?; + let already_recorded = match try_touch_exact_state(&self.exact_states, &identity.page_id) { + Ok(Some(already_recorded)) => already_recorded, + Ok(None) => { + // Recording is optional. A background worker may hold this lock + // while hashing hundreds of MiB; never make inference wait for it. + self.finish_record(&identity.page_id); + return Ok(None); + } + Err(error) => { + self.finish_record(&identity.page_id); + return Err(error); + } + }; + if already_recorded { + self.finish_record(&identity.page_id); + return Ok(None); + } + // Avoid paying a potentially multi-hundred-MiB runtime export when the + // bounded worker queue is already occupied. Admission remains best-effort: + // another producer may win the race before `try_send` below. + if !self.has_exact_state_record_capacity() { + self.finish_record(&identity.page_id); + return Ok(None); + } + let exported = match self.payload { + StagePrefixCachePayload::FullState => { + runtime.export_full_state(session_id).map(|state| { ( - ExactStatePayload::kv_recurrent( - kv.as_ref().map(|kv| kv.payload.clone()).unwrap_or_default(), - recurrent, - ), - ExactStateExtra { - kv_desc: kv.as_ref().map(|kv| kv.desc.clone()), - }, + ExactStatePayload::full_state(state), + ExactStateExtra::default(), ) - } - _ => return Ok(None), - }; - let outcome = self - .exact_states - .lock() - .expect("exact state cache lock poisoned") - .record(identity.page_id.clone(), token_count, payload, extra); - Ok(Some(ExactStateRecord { - page_id: outcome.page_id, - token_count: outcome.token_count as usize, - payload_kind: self.payload.into(), - stored: outcome.stored, - logical_bytes: outcome.logical_bytes, - physical_bytes: outcome.physical_bytes, - entries: outcome.entries, - evicted_entries: outcome.evicted_entries, - evicted_logical_bytes: outcome.evicted_logical_bytes, - dedupe: outcome.dedupe, - })) - })(); - self.finish_record(&identity.page_id); - result + }) + } + StagePrefixCachePayload::KvRecurrent => (|| { + let kv = match runtime.export_kv_page(session_id, 0, token_count) { + Ok(kv) => Some(kv), + Err(error) if is_native_kv_unavailable(&error) => None, + Err(error) => return Err(error), + }; + let recurrent = runtime.export_recurrent_state(session_id)?; + Ok(( + ExactStatePayload::kv_recurrent( + kv.as_ref().map(|kv| kv.payload.clone()).unwrap_or_default(), + recurrent, + ), + ExactStateExtra { + kv_desc: kv.as_ref().map(|kv| kv.desc.clone()), + }, + )) + })(), + StagePrefixCachePayload::Disabled | StagePrefixCachePayload::ResidentKv => { + self.finish_record(&identity.page_id); + return Ok(None); + } + }; + let (payload, extra) = match exported { + Ok(exported) => exported, + Err(error) => { + self.finish_record(&identity.page_id); + return Err(error); + } + }; + let payload_kind = payload.kind(); + let logical_bytes = payload.byte_len(); + match self.enqueue_exact_state_record(PendingExactStateRecord { + page_id: identity.page_id.clone(), + token_count, + payload, + extra, + }) { + ExactStateRecordAdmission::Queued => { + // Recording owns the exact-state mutex while it hashes a potentially + // multi-hundred-MiB payload. Telemetry must not turn that background + // work back into request latency by waiting for cache stats here. + let stats = self + .exact_states + .try_lock() + .ok() + .map(|states| states.stats()) + .unwrap_or_default(); + Ok(Some(ExactStateRecord { + page_id: identity.page_id.clone(), + token_count: token_count as usize, + payload_kind, + stored: false, + logical_bytes, + physical_bytes: stats.physical_bytes, + entries: stats.entries, + evicted_entries: 0, + evicted_logical_bytes: 0, + dedupe: Default::default(), + })) + } + ExactStateRecordAdmission::DroppedFull | ExactStateRecordAdmission::WorkerStopped => { + Ok(None) + } + } + } +} + +fn try_touch_exact_state( + exact_states: &std::sync::Mutex>, + page_id: &str, +) -> Result> { + match exact_states.try_lock() { + Ok(mut exact_states) => Ok(Some(exact_states.touch(page_id))), + Err(std::sync::TryLockError::WouldBlock) => Ok(None), + Err(std::sync::TryLockError::Poisoned(poisoned)) => { + Ok(Some(poisoned.into_inner().touch(page_id))) + } } } @@ -224,3 +305,55 @@ impl From for skippy_cache::ExactStatePayloadKind { } } } + +#[cfg(test)] +mod tests { + use std::{ + sync::{Arc, Mutex}, + time::{Duration, Instant}, + }; + + use skippy_cache::ExactStateCache; + + use super::{ExactStateExtra, try_touch_exact_state}; + + #[test] + fn busy_exact_state_lock_skips_touch_without_waiting() { + let cache = Arc::new(Mutex::new(ExactStateCache::::new(1, 1024))); + let locked = cache.clone(); + let (locked_tx, locked_rx) = std::sync::mpsc::channel(); + let (release_tx, release_rx) = std::sync::mpsc::channel(); + let holder = std::thread::spawn(move || { + let _guard = locked.lock().unwrap(); + locked_tx.send(()).unwrap(); + release_rx.recv().unwrap(); + }); + locked_rx.recv().unwrap(); + + let started = Instant::now(); + assert_eq!(try_touch_exact_state(&cache, "busy").unwrap(), None); + assert!(started.elapsed() < Duration::from_millis(100)); + + release_tx.send(()).unwrap(); + holder.join().unwrap(); + } + + #[test] + fn poisoned_exact_state_lock_recovers_without_panicking() { + let cache = Arc::new(Mutex::new(ExactStateCache::::new(1, 1024))); + let poisoned = cache.clone(); + assert!( + std::thread::spawn(move || { + let _guard = poisoned.lock().unwrap(); + panic!("poison exact-state cache for test"); + }) + .join() + .is_err() + ); + + assert_eq!( + try_touch_exact_state(&cache, "poisoned").unwrap(), + Some(false) + ); + } +} diff --git a/crates/skippy-server/src/kv_integration/mod.rs b/crates/skippy-server/src/kv_integration/mod.rs index 27dd41ec4..c5e7a48cf 100644 --- a/crates/skippy-server/src/kv_integration/mod.rs +++ b/crates/skippy-server/src/kv_integration/mod.rs @@ -1,13 +1,18 @@ use std::{ collections::{BTreeMap, BTreeSet}, - sync::{Arc, Mutex, atomic::AtomicBool}, + sync::{ + Arc, Mutex, + atomic::{AtomicBool, AtomicU64, AtomicUsize}, + mpsc::{SyncSender, TrySendError}, + }, }; use anyhow::{Result, bail}; use serde_json::{Value, json}; use sha2::{Digest, Sha256}; use skippy_cache::{ - ExactStateCache, PrefixCandidatePolicy, ResidentActivationCache, ResidentPrefixCache, + ExactStateCache, ExactStatePayload, PrefixCandidatePolicy, ResidentActivationCache, + ResidentPrefixCache, }; use skippy_metrics::attr as attr_key; use skippy_runtime::{ActivationFrame, RuntimeKvPageDesc}; @@ -139,6 +144,10 @@ pub struct KvStageIntegration { pub(crate) resident: Arc>, pub(crate) activations: Arc>>, pub(crate) exact_states: Arc>>, + pub(crate) exact_state_record_tx: SyncSender, + pub(crate) exact_state_records_queued: Arc, + pub(crate) exact_state_records_dropped: Arc, + pub(crate) exact_state_records_pending: Arc, pub(crate) first_tokens: Arc>>, pub(crate) replay_tokens: Arc>>>, pub(crate) split_prefill_tokens: Arc>>>, @@ -158,6 +167,62 @@ pub enum StagePrefixCachePayload { FullState, } +pub(crate) const EXACT_STATE_RECORD_CAPACITY: usize = 1; + +#[derive(Debug)] +pub(crate) struct PendingExactStateRecord { + pub(crate) page_id: String, + pub(crate) token_count: u64, + pub(crate) payload: ExactStatePayload, + pub(crate) extra: ExactStateExtra, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum ExactStateRecordAdmission { + Queued, + DroppedFull, + WorkerStopped, +} + +fn has_exact_state_record_capacity(pending_count: &AtomicUsize) -> bool { + pending_count.load(std::sync::atomic::Ordering::Acquire) < EXACT_STATE_RECORD_CAPACITY +} + +fn enqueue_exact_state_record( + sender: &SyncSender, + inflight_records: &Mutex>, + queued: &AtomicU64, + dropped: &AtomicU64, + pending_count: &AtomicUsize, + pending: PendingExactStateRecord, +) -> ExactStateRecordAdmission { + pending_count.fetch_add(1, std::sync::atomic::Ordering::Release); + match sender.try_send(pending) { + Ok(()) => { + queued.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + ExactStateRecordAdmission::Queued + } + Err(TrySendError::Full(pending)) => { + pending_count.fetch_sub(1, std::sync::atomic::Ordering::Release); + inflight_records + .lock() + .expect("kv inflight record lock poisoned") + .remove(&pending.page_id); + dropped.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + ExactStateRecordAdmission::DroppedFull + } + Err(TrySendError::Disconnected(pending)) => { + pending_count.fetch_sub(1, std::sync::atomic::Ordering::Release); + inflight_records + .lock() + .expect("kv inflight record lock poisoned") + .remove(&pending.page_id); + dropped.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + ExactStateRecordAdmission::WorkerStopped + } + } +} + #[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)] pub(crate) struct ExactStateExtra { pub(crate) kv_desc: Option, @@ -200,6 +265,24 @@ impl KvStageIntegration { .remove(page_id); } + pub(crate) fn has_exact_state_record_capacity(&self) -> bool { + has_exact_state_record_capacity(&self.exact_state_records_pending) + } + + pub(crate) fn enqueue_exact_state_record( + &self, + pending: PendingExactStateRecord, + ) -> ExactStateRecordAdmission { + enqueue_exact_state_record( + &self.exact_state_record_tx, + &self.inflight_records, + &self.exact_state_records_queued, + &self.exact_state_records_dropped, + &self.exact_state_records_pending, + pending, + ) + } + pub async fn hello(&self) -> Result<()> { Ok(()) } @@ -283,6 +366,13 @@ impl KvStageIntegration { .lock() .expect("resident activation cache lock poisoned"); let activations = activations.stats(); + let exact_state_stats = self + .exact_states + .try_lock() + .ok() + .map(|states| states.stats()); + let exact_state_stats_busy = exact_state_stats.is_none(); + let exact_state_stats = exact_state_stats.unwrap_or_default(); vec![ ("skippy.kv.mode", json!(format!("{:?}", self.mode))), ("skippy.kv.payload", json!(format!("{:?}", self.payload))), @@ -308,19 +398,44 @@ impl KvStageIntegration { ), ( "skippy.exact_cache.entries", - json!(self.exact_state_stats().entries), + json!(exact_state_stats.entries), ), ( "skippy.exact_cache.logical_bytes", - json!(self.exact_state_stats().logical_bytes), + json!(exact_state_stats.logical_bytes), ), ( "skippy.exact_cache.physical_bytes", - json!(self.exact_state_stats().physical_bytes), + json!(exact_state_stats.physical_bytes), + ), + ( + "skippy.exact_cache.stats_busy", + json!(exact_state_stats_busy), + ), + ( + "skippy.exact_cache.records_queued", + json!( + self.exact_state_records_queued + .load(std::sync::atomic::Ordering::Relaxed) + ), + ), + ( + "skippy.exact_cache.records_dropped", + json!( + self.exact_state_records_dropped + .load(std::sync::atomic::Ordering::Relaxed) + ), + ), + ( + "skippy.exact_cache.records_pending", + json!( + self.exact_state_records_pending + .load(std::sync::atomic::Ordering::Relaxed) + ), ), ( "skippy.exact_cache.max_bytes", - json!(self.exact_state_stats().max_bytes), + json!(exact_state_stats.max_bytes), ), ("skippy.kv.correctness_mode", json!(self.correctness_mode)), ( @@ -353,40 +468,7 @@ impl KvStageIntegration { /// (full budget, every write failing, every entry quarantined) is /// indistinguishable from one that is simply never probed. pub fn disk_tier_attrs(&self) -> Vec<(&'static str, Value)> { - let Some(stats) = self - .exact_states - .lock() - .expect("exact state cache lock poisoned") - .disk_stats() - else { - return vec![("skippy.kv.disk_tier_enabled", json!(false))]; - }; - vec![ - ("skippy.kv.disk_tier_enabled", json!(true)), - ("skippy.kv.disk_entries", json!(stats.entries)), - ("skippy.kv.disk_bytes", json!(stats.bytes)), - ("skippy.kv.disk_max_bytes", json!(stats.max_bytes)), - ("skippy.kv.disk_demotions", json!(stats.demotions)), - ("skippy.kv.disk_promotions", json!(stats.promotions)), - ("skippy.kv.disk_evictions", json!(stats.evictions)), - ( - "skippy.kv.disk_pages_rejected_too_large", - json!(stats.pages_rejected_too_large), - ), - ( - "skippy.kv.disk_last_rejected_page_bytes", - json!(stats.last_rejected_page_bytes), - ), - ( - "skippy.kv.disk_corrupt_entries", - json!(stats.corrupt_entries), - ), - ("skippy.kv.disk_verifications", json!(stats.verifications)), - ( - "skippy.kv.disk_verifications_skipped", - json!(stats.verifications_skipped), - ), - ] + disk_tier_attrs(&self.exact_states) } /// Attach the disk-tier counters to a decision event. @@ -401,13 +483,6 @@ impl KvStageIntegration { } } - fn exact_state_stats(&self) -> skippy_cache::ExactStateCacheStats { - self.exact_states - .lock() - .expect("exact state cache lock poisoned") - .stats() - } - fn record_candidate_token_counts(&self, token_count: u64) -> Vec { self.candidate_policy .record_candidate_token_counts(token_count) @@ -483,6 +558,47 @@ impl KvStageIntegration { } } +fn disk_tier_attrs( + exact_states: &Mutex>, +) -> Vec<(&'static str, Value)> { + let exact_states = match exact_states.try_lock() { + Ok(exact_states) => exact_states, + Err(std::sync::TryLockError::WouldBlock) => { + return vec![("skippy.kv.disk_tier_stats_busy", json!(true))]; + } + Err(std::sync::TryLockError::Poisoned(poisoned)) => poisoned.into_inner(), + }; + let Some(stats) = exact_states.disk_stats() else { + return vec![("skippy.kv.disk_tier_enabled", json!(false))]; + }; + vec![ + ("skippy.kv.disk_tier_enabled", json!(true)), + ("skippy.kv.disk_entries", json!(stats.entries)), + ("skippy.kv.disk_bytes", json!(stats.bytes)), + ("skippy.kv.disk_max_bytes", json!(stats.max_bytes)), + ("skippy.kv.disk_demotions", json!(stats.demotions)), + ("skippy.kv.disk_promotions", json!(stats.promotions)), + ("skippy.kv.disk_evictions", json!(stats.evictions)), + ( + "skippy.kv.disk_pages_rejected_too_large", + json!(stats.pages_rejected_too_large), + ), + ( + "skippy.kv.disk_last_rejected_page_bytes", + json!(stats.last_rejected_page_bytes), + ), + ( + "skippy.kv.disk_corrupt_entries", + json!(stats.corrupt_entries), + ), + ("skippy.kv.disk_verifications", json!(stats.verifications)), + ( + "skippy.kv.disk_verifications_skipped", + json!(stats.verifications_skipped), + ), + ] +} + fn local_trust_checksum(page_id: &str, byte_size: u64) -> Checksum { let mut digest = Sha256::new(); digest.update(b"skippy-local-trust-v1"); @@ -494,6 +610,154 @@ fn local_trust_checksum(page_id: &str, byte_size: u64) -> Checksum { } } +#[cfg(test)] +mod exact_state_record_queue_tests { + use std::sync::{ + Arc, Mutex, + atomic::{AtomicU64, AtomicUsize, Ordering}, + mpsc::sync_channel, + }; + use std::time::{Duration, Instant}; + + use skippy_cache::{ExactStateCache, ExactStatePayload}; + + use super::{ + BTreeSet, EXACT_STATE_RECORD_CAPACITY, ExactStateExtra, ExactStateRecordAdmission, + PendingExactStateRecord, disk_tier_attrs, enqueue_exact_state_record, + has_exact_state_record_capacity, + }; + + fn pending(page_id: &str) -> PendingExactStateRecord { + PendingExactStateRecord { + page_id: page_id.to_string(), + token_count: 1, + payload: ExactStatePayload::full_state(vec![1]), + extra: ExactStateExtra::default(), + } + } + + #[test] + fn busy_exact_state_lock_skips_disk_stats_without_waiting() { + let cache = Arc::new(Mutex::new(ExactStateCache::::new(1, 1024))); + let locked = cache.clone(); + let (locked_tx, locked_rx) = std::sync::mpsc::channel(); + let (release_tx, release_rx) = std::sync::mpsc::channel(); + let holder = std::thread::spawn(move || { + let _guard = locked.lock().unwrap(); + locked_tx.send(()).unwrap(); + release_rx.recv().unwrap(); + }); + locked_rx.recv().unwrap(); + + let started = Instant::now(); + assert_eq!( + disk_tier_attrs(&cache), + vec![("skippy.kv.disk_tier_stats_busy", serde_json::json!(true))] + ); + assert!(started.elapsed() < Duration::from_millis(100)); + + release_tx.send(()).unwrap(); + holder.join().unwrap(); + } + + #[test] + fn pending_capacity_signal_rejects_work_before_export() { + let pending_count = AtomicUsize::new(0); + assert!(has_exact_state_record_capacity(&pending_count)); + + pending_count.store(EXACT_STATE_RECORD_CAPACITY, Ordering::Release); + assert!(!has_exact_state_record_capacity(&pending_count)); + } + + #[test] + fn full_queue_drops_optional_record_and_releases_inflight_page() { + let (sender, _receiver) = sync_channel(1); + sender.send(pending("queued")).unwrap(); + let inflight = Mutex::new(BTreeSet::from(["dropped".to_string()])); + let queued = AtomicU64::new(0); + let dropped = AtomicU64::new(0); + let pending_count = AtomicUsize::new(0); + + assert_eq!( + enqueue_exact_state_record( + &sender, + &inflight, + &queued, + &dropped, + &pending_count, + pending("dropped"), + ), + ExactStateRecordAdmission::DroppedFull + ); + assert!(!inflight.lock().unwrap().contains("dropped")); + assert_eq!(queued.load(Ordering::Relaxed), 0); + assert_eq!(dropped.load(Ordering::Relaxed), 1); + assert_eq!(pending_count.load(Ordering::Relaxed), 0); + } + + #[test] + fn disconnected_worker_releases_inflight_page() { + let (sender, receiver) = sync_channel(1); + drop(receiver); + let inflight = Mutex::new(BTreeSet::from(["orphaned".to_string()])); + let queued = AtomicU64::new(0); + let dropped = AtomicU64::new(0); + let pending_count = AtomicUsize::new(0); + + assert_eq!( + enqueue_exact_state_record( + &sender, + &inflight, + &queued, + &dropped, + &pending_count, + pending("orphaned"), + ), + ExactStateRecordAdmission::WorkerStopped + ); + assert!(!inflight.lock().unwrap().contains("orphaned")); + assert_eq!(dropped.load(Ordering::Relaxed), 1); + assert_eq!(pending_count.load(Ordering::Relaxed), 0); + } + + #[test] + fn worker_keeps_page_inflight_until_record_finishes() { + let (sender, receiver) = sync_channel(1); + let inflight = Arc::new(Mutex::new(BTreeSet::from(["page".to_string()]))); + let worker_inflight = inflight.clone(); + let queued = AtomicU64::new(0); + let dropped = AtomicU64::new(0); + let pending_count = Arc::new(AtomicUsize::new(0)); + let worker_pending_count = pending_count.clone(); + + assert_eq!( + enqueue_exact_state_record( + &sender, + &inflight, + &queued, + &dropped, + &pending_count, + pending("page"), + ), + ExactStateRecordAdmission::Queued + ); + assert!(inflight.lock().unwrap().contains("page")); + + let worker = std::thread::spawn(move || { + let pending = receiver.recv().unwrap(); + assert!(worker_inflight.lock().unwrap().contains(&pending.page_id)); + worker_inflight.lock().unwrap().remove(&pending.page_id); + worker_pending_count.fetch_sub(1, Ordering::Relaxed); + }); + worker.join().unwrap(); + + assert!(!inflight.lock().unwrap().contains("page")); + assert_eq!(queued.load(Ordering::Relaxed), 1); + assert_eq!(dropped.load(Ordering::Relaxed), 0); + assert_eq!(pending_count.load(Ordering::Relaxed), 0); + } +} + #[cfg(test)] mod telemetry_error_class_tests { use super::telemetry_error_class_from_message; diff --git a/crates/skippy-server/src/kv_integration/records.rs b/crates/skippy-server/src/kv_integration/records.rs index 56bc9fa66..b632468dd 100644 --- a/crates/skippy-server/src/kv_integration/records.rs +++ b/crates/skippy-server/src/kv_integration/records.rs @@ -103,6 +103,9 @@ pub struct ExactStateRestore { pub reconstruct_ms: f64, pub reconstruct_bytes: u64, pub reconstruct_blocks: usize, + pub lookup_ms: f64, + pub kv_import_ms: f64, + pub recurrent_import_ms: f64, } #[derive(Debug, Clone)] diff --git a/tools/xtask/data/console_print_allowlist.json b/tools/xtask/data/console_print_allowlist.json index ffdf69a36..f1b4f0159 100644 --- a/tools/xtask/data/console_print_allowlist.json +++ b/tools/xtask/data/console_print_allowlist.json @@ -4347,20 +4347,6 @@ "macro_name": "println!" } ], - "crates/skippy-server/src/kv_integration/config.rs": [ - { - "line": 110, - "macro_name": "eprintln!" - }, - { - "line": 120, - "macro_name": "eprintln!" - }, - { - "line": 152, - "macro_name": "eprintln!" - } - ], "crates/skippy-server/src/main.rs": [ { "line": 16,