Skip to content
Closed
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
4 changes: 3 additions & 1 deletion crates/kv_index/benches/throughput_bench.rs
Original file line number Diff line number Diff line change
Expand Up @@ -667,7 +667,9 @@ async fn run_benchmark(args: &Args, traces: Vec<Vec<TimedEntry>>) -> BenchmarkRe

let num_total_workers = args.num_workers * args.duplication_factor;
for w in 0..num_total_workers {
indexer.intern_worker(&format!("worker-{w}"));
indexer
.intern_worker(&format!("worker-{w}"))
.expect("worker id space exhausted");
}

let traces: Vec<Arc<Vec<TimedEntry>>> = traces.into_iter().map(Arc::new).collect();
Expand Down
533 changes: 424 additions & 109 deletions crates/kv_index/src/event_tree.rs

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion crates/kv_index/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ mod token_tree;
pub use common::{MatchResult, TenantId};
pub use event_tree::{
compute_content_hash, compute_request_content_hashes, ApplyError, ContentHash, OverlapScores,
PositionalIndexer, SequenceHash, StoredBlock, WorkerBlockMap, WorkerId,
PositionalIndexer, SequenceHash, StoredBlock, WorkerBlockMap, WorkerId, WorkerIdExhausted,
};
pub use path_hash::{hash_node_path, hash_token_path, GLOBAL_EVICTION_HASH};
// Re-export under names matching old tree.rs API for easier migration
Expand Down
6 changes: 3 additions & 3 deletions model_gateway/benches/radix_tree_benchmark.rs
Original file line number Diff line number Diff line change
Expand Up @@ -445,7 +445,7 @@ fn build_populated_indexer(
let mut all_worker_chunks = Vec::with_capacity(workers.len());

for worker in workers {
let worker_id = indexer.intern_worker(worker);
let worker_id = indexer.intern_worker(worker).unwrap();
let mut wb = WorkerBlockMap::default();
indexer
.apply_stored(worker_id, &shared_blocks, None, &mut wb)
Expand Down Expand Up @@ -500,7 +500,7 @@ macro_rules! bench_indexer_store {
for _ in 0..iters {
let indexer = PositionalIndexer::new(32);
for worker in &workers {
let worker_id = indexer.intern_worker(worker);
let worker_id = indexer.intern_worker(worker).unwrap();
let mut wb = WorkerBlockMap::default();
let chunks = generate_token_chunks($blocks_per_worker, $block_size);
let blocks = chunks_to_stored_blocks(&chunks);
Expand Down Expand Up @@ -596,7 +596,7 @@ macro_rules! bench_indexer_concurrent {
(0..$num_threads)
.map(|t| {
let chunks = &worker_chunks[t % workers.len()];
let worker_id = indexer.intern_worker(&workers[t % workers.len()]);
let worker_id = indexer.intern_worker(&workers[t % workers.len()]).unwrap();

// Read data: pre-computed content hashes
let query_tokens = flatten_tokens(chunks);
Expand Down
17 changes: 17 additions & 0 deletions model_gateway/src/observability/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -230,6 +230,11 @@ pub(crate) fn init_metrics() {
"smg_worker_errors_total",
"Worker-level errors by worker_type, connection_mode, error_type"
);
describe_counter!(
"smg_kv_event_subscription_failures_total",
"KV event subscription task failures by worker and reason \
(panic, join_error, intern_failed)"
);
describe_gauge!(
"smg_manual_policy_cache_entries",
"Number of routing entries in manual policy cache"
Expand Down Expand Up @@ -944,6 +949,18 @@ impl Metrics {
.set(if healthy { 1.0 } else { 0.0 });
}

/// Record a KV event subscription task failure (panic, join error, or
/// worker-id intern failure)
pub fn record_kv_event_subscription_failure(worker_url: &str, reason: &'static str) {
let worker_interned = intern_string(worker_url);
counter!(
"smg_kv_event_subscription_failures_total",
"worker" => worker_interned,
"reason" => reason
)
.increment(1);
}

// ========================================================================
// Layer 3: Worker resilience metrics (circuit breaker)
// ========================================================================
Expand Down
Loading