Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 7 additions & 2 deletions lib/kv-router/src/indexer/local.rs
Original file line number Diff line number Diff line change
Expand Up @@ -213,6 +213,9 @@ pub struct LocalKvIndexer {
/// This stays separate from `event_buffer` so dump wait/build state can be
/// managed on the async path without holding the buffer lock across `.await`.
recovery_cache: Arc<RecoverySnapshotCache>,
/// Shared metrics handle, also wired into lazily created lower-tier
/// indexers so HostPinned/Disk/External traffic is counted too.
metrics: Arc<KvIndexerMetrics>,
/// Maximum number of events to keep in buffer
max_buffer_size: usize, // Router sets this to WORKER_KV_INDEXER_BUFFER_SIZE
#[cfg(test)]
Expand All @@ -230,7 +233,8 @@ impl LocalKvIndexer {
max_buffer_size: usize,
) -> Self {
Self {
indexer: KvIndexer::new(token, kv_block_size, metrics),
indexer: KvIndexer::new(token, kv_block_size, metrics.clone()),
metrics,
lower_tier_indexers: Arc::new(Mutex::new(HashMap::new())),
event_buffer: Mutex::new(VecDeque::with_capacity(max_buffer_size)),
recovery_cache: Arc::new(RecoverySnapshotCache::new()),
Expand Down Expand Up @@ -669,10 +673,11 @@ impl LocalKvIndexer {
indexers
.entry(storage_tier)
.or_insert_with(|| {
Arc::new(ThreadPoolIndexer::new(
Arc::new(ThreadPoolIndexer::new_with_metrics(
LowerTierIndexer::new(),
1,
self.block_size(),
Some(self.metrics.clone()),
))
})
.clone()
Expand Down
18 changes: 14 additions & 4 deletions lib/kv-router/src/indexer/lower_tier.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ use std::sync::Arc;
use dashmap::DashMap;
use rustc_hash::{FxBuildHasher, FxHashMap, FxHashSet};

use super::{KvIndexerMetrics, SyncIndexer, WorkerLookupStats, WorkerTask};
use super::{EventKind, KvIndexerMetrics, SyncIndexer, WorkerLookupStats, WorkerTask};
use crate::protocols::{
ExternalSequenceBlockHash, KvCacheEvent, KvCacheEventData, KvCacheEventError, KvCacheStoreData,
KvCacheStoredBlockData, LocalBlockHash, OverlapScores, RouterEvent, WorkerWithDpRank,
Expand Down Expand Up @@ -516,23 +516,33 @@ impl SyncIndexer for LowerTierIndexer {
fn worker(
&self,
event_receiver: flume::Receiver<WorkerTask>,
_metrics: Option<Arc<KvIndexerMetrics>>,
metrics: Option<Arc<KvIndexerMetrics>>,
) -> anyhow::Result<()> {
let mut worker_blocks = WorkerBlockIndex::default();
let counters = metrics.as_ref().map(|m| m.prebind());

while let Ok(task) = event_receiver.recv() {
match task {
WorkerTask::Event(event) => {
if let Err(error) = self.apply_event(&mut worker_blocks, event) {
let kind = EventKind::of(&event.event.data);
let result = self.apply_event(&mut worker_blocks, event);
if let Err(ref error) = result {
tracing::warn!(%error, "Failed to apply lower-tier event");
}
if let Some(ref c) = counters {
c.inc(kind, result);
}
}
WorkerTask::EventWithAck { event, resp } => {
let kind = EventKind::of(&event.event.data);
let result = self.apply_event(&mut worker_blocks, event);
let applied = result.is_ok();
if let Err(error) = result {
if let Err(ref error) = result {
tracing::warn!(%error, "Failed to apply lower-tier event");
}
if let Some(ref c) = counters {
c.inc(kind, result);
}
let _ = resp.send(applied);
}
WorkerTask::Anchor { worker, anchor } => {
Expand Down
21 changes: 19 additions & 2 deletions lib/kv-router/src/indexer/lower_tier_indexers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ use std::collections::HashMap;
use std::sync::{Arc, RwLock};

use crate::indexer::{
LowerTierContinuation, LowerTierIndexer, LowerTierMatchDetails, MatchDetails,
KvIndexerMetrics, LowerTierContinuation, LowerTierIndexer, LowerTierMatchDetails, MatchDetails,
ThreadPoolIndexer, WireTieredMatchDetails,
};
use crate::protocols::{LocalBlockHash, StorageTier};
Expand All @@ -27,20 +27,36 @@ use crate::protocols::{LocalBlockHash, StorageTier};
/// non-device [`StorageTier`] that has received at least one event.
#[derive(Clone)]
pub struct LowerTierIndexers {
metrics: Option<Arc<KvIndexerMetrics>>,
num_threads: usize,
block_size: u32,
indexers: Arc<RwLock<HashMap<StorageTier, Arc<ThreadPoolIndexer<LowerTierIndexer>>>>>,
}

impl LowerTierIndexers {
/// Metrics-less constructor for call sites without a `KvIndexerMetrics` handle.
/// Router production assembly should use [`new_with_metrics`](Self::new_with_metrics)
/// so lower-tier traffic is included in `kv_cache_events_applied`.
pub fn new(num_threads: usize, block_size: u32) -> Self {
Self::new_with_metrics(num_threads, block_size, None)
}

/// Same as [`new`](Self::new) but wires `kv_cache_events_applied`
/// counters into every lazily created per-tier indexer, matching the
/// observability of the device-tier path.
pub fn new_with_metrics(
num_threads: usize,
block_size: u32,
metrics: Option<Arc<KvIndexerMetrics>>,
) -> Self {
assert!(
num_threads > 0,
"lower-tier indexer threads must be non-zero"
);
Self {
num_threads,
block_size,
metrics,
indexers: Arc::new(RwLock::new(HashMap::new())),
}
}
Expand All @@ -60,10 +76,11 @@ impl LowerTierIndexers {
.unwrap()
.entry(storage_tier)
.or_insert_with(|| {
Arc::new(ThreadPoolIndexer::new(
Arc::new(ThreadPoolIndexer::new_with_metrics(
LowerTierIndexer::new(),
self.num_threads,
self.block_size,
self.metrics.clone(),
))
})
.clone()
Expand Down
2 changes: 1 addition & 1 deletion lib/kv-router/src/protocols.rs
Original file line number Diff line number Diff line change
Expand Up @@ -330,7 +330,7 @@ impl StorageTier {
pub fn from_kv_medium(medium: &str) -> Option<Self> {
match medium {
"GPU" | "DEVICE" => Some(Self::Device),
"CPU_PINNED" | "CPU_TIER1" => Some(Self::HostPinned),
"CPU" | "CPU_PINNED" | "CPU_TIER1" => Some(Self::HostPinned),
"CPU_TIER2" | "DISK" | "NVME" => Some(Self::Disk),
"EXTERNAL" | "NETWORK" | "REMOTE" | "SHARED" => Some(Self::External),
_ => None,
Expand Down
16 changes: 15 additions & 1 deletion lib/kv-router/src/standalone_indexer/listener.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ use crate::zmq_wire::{ZmqEventNormalizer, decode_event_batch};
use super::evictions::PendingEvictions;
use super::indexer::Indexer;
use super::registry::ListenerRecord;
use super::tier_bridge::TierBridge;
use super::zmq::{MultipartMessage, SharedSocket, connect_sub_socket, recv_multipart};

const WATERMARK_UNSET: u64 = u64::MAX;
Expand Down Expand Up @@ -129,6 +130,7 @@ struct ListenerLoop {
/// drains it through there. Untouched unless `--keep-evictions` is set.
pending_evictions: Arc<Mutex<PendingEvictions>>,
normalizer: ZmqEventNormalizer,
tier_bridge: TierBridge,
messages_processed: u64,
}

Expand Down Expand Up @@ -157,6 +159,7 @@ impl ListenerLoop {
watermark,
pending_evictions,
normalizer: ZmqEventNormalizer::new(block_size),
tier_bridge: TierBridge::new(),
messages_processed: 0,
}
}
Expand Down Expand Up @@ -330,6 +333,7 @@ impl ListenerLoop {
// them too — replaying them against the snapshot could remove
// blocks the dump says are live.
self.pending_evictions.lock().clear();
self.tier_bridge.reset();
self.indexer
.remove_worker_dp_rank(self.worker_id, self.dp_rank)
.await;
Expand Down Expand Up @@ -388,7 +392,7 @@ impl ListenerLoop {
/// events under the worker that was queried), so they are rewritten before
/// applying. Returns the count actually applied (after the
/// `keep_evictions` measurement filter).
async fn apply_recovered_events(&self, events: Vec<RouterEvent>) -> u64 {
async fn apply_recovered_events(&mut self, events: Vec<RouterEvent>) -> u64 {
let mut applied = 0;
for mut event in events {
event.worker_id = self.worker_id;
Expand All @@ -397,6 +401,11 @@ impl ListenerLoop {
// Audit-log the recovered event before the measurement filter.
audit_log_event(&event, event.event.event_id, "recover");

if let Some(promoted) = self.tier_bridge.observe(&event)
&& !self.keep_evictions_intercept(&promoted)
{
self.indexer.apply_event_routed(promoted).await;
}
// Feed-layer measurement filter (same as apply_live_batch).
if self.keep_evictions_intercept(&event) {
continue;
Expand Down Expand Up @@ -481,6 +490,11 @@ impl ListenerLoop {
// Audit-log the event as published by the engine, before the
// measurement filter below can drop it.
audit_log_event(&router_event, seq, "live");
if let Some(promoted) = self.tier_bridge.observe(&router_event)
&& !self.keep_evictions_intercept(&promoted)
{
self.indexer.apply_event_routed(promoted).await;
}
// Feed-layer measurement filter.
if self.keep_evictions_intercept(&router_event) {
continue;
Expand Down
1 change: 1 addition & 0 deletions lib/kv-router/src/standalone_indexer/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ pub mod pod_watcher;
pub mod recovery;
pub mod registry;
pub mod server;
mod tier_bridge;
mod zmq;

use std::sync::{Arc, OnceLock};
Expand Down
Loading
Loading