Skip to content
Merged
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
1 change: 1 addition & 0 deletions crates/kv_index/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ categories = ["data-structures", "caching"]

[dependencies]
bincode = "1.3"
blake3 = { workspace = true }
dashmap = { workspace = true }
once_cell = "1.21.3"
parking_lot = { workspace = true }
Expand Down
2 changes: 2 additions & 0 deletions crates/kv_index/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@

mod common;
mod event_tree;
mod path_hash;
pub mod snapshot;
mod string_tree;
mod token_tree;
Expand All @@ -22,6 +23,7 @@ pub use event_tree::{
compute_content_hash, compute_request_content_hashes, ApplyError, ContentHash, OverlapScores,
PositionalIndexer, SequenceHash, StoredBlock, WorkerBlockMap, WorkerId,
};
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
pub use string_tree::Tree;
pub use string_tree::{
Expand Down
File renamed without changes.
2 changes: 0 additions & 2 deletions crates/mesh/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@
mod crdt_kv;
mod gossip_controller;
mod gossip_service;
mod hash;
pub mod kv;
mod metrics;
mod mtls;
Expand All @@ -26,7 +25,6 @@ pub use crdt_kv::{
decode as decode_epoch_count, encode as encode_epoch_count, CrdtOrMap, EpochCount,
MergeStrategy, OperationLog, EPOCH_MAX_WINS_ENCODED_LEN,
};
pub use hash::{hash_node_path, hash_token_path, GLOBAL_EVICTION_HASH};
pub use kv::{
CrdtNamespace, DrainHandle, MeshKV, StreamConfig, StreamDrainFn, StreamNamespace,
StreamRouting, Subscription,
Expand Down
4 changes: 2 additions & 2 deletions model_gateway/src/mesh/adapters/tree_sync.rs
Original file line number Diff line number Diff line change
Expand Up @@ -663,8 +663,8 @@ impl TreeSyncAdapter {
/// Buffer a local tree insert for the next gossip round. Hot
/// path — keep it cheap; the drain does the serialisation.
/// `delta.node_hash` must be non-zero: 0 is
/// `smg_mesh::tree_ops::GLOBAL_EVICTION_HASH` (producers remap
/// 0→1 to keep the space disjoint).
/// [`kv_index::GLOBAL_EVICTION_HASH`] (producers remap 0→1 to
/// keep the space disjoint).
pub fn on_local_insert(&self, model_id: &str, delta: TreeDelta) {
debug_assert!(
!model_id.is_empty(),
Expand Down
24 changes: 12 additions & 12 deletions model_gateway/src/policies/cache_aware.rs
Original file line number Diff line number Diff line change
Expand Up @@ -398,7 +398,7 @@ impl CacheAwarePolicy {
.entry(model_id.to_string())
.or_default()
.token_tree
.insert(smg_mesh::hash_token_path(tokens), matched_prefix);
.insert(kv_index::hash_token_path(tokens), matched_prefix);
}
} else if let Some(text) = info.request_text {
// HTTP request: update string tree
Expand All @@ -419,7 +419,7 @@ impl CacheAwarePolicy {

tree.insert_text(text, worker_url);

let path_hash = smg_mesh::hash_node_path(text);
let path_hash = kv_index::hash_node_path(text);
self.hash_index
.entry(model_id.to_string())
.or_default()
Expand Down Expand Up @@ -595,7 +595,7 @@ impl TreeHandle for CacheAwarePolicy {
.entry(model_id.to_string())
.or_default()
.string_tree
.insert(smg_mesh::hash_node_path(path), path.clone());
.insert(kv_index::hash_node_path(path), path.clone());
applied += 1;
}
RepairEntry::Token { .. } => {
Expand Down Expand Up @@ -625,7 +625,7 @@ impl TreeHandle for CacheAwarePolicy {
.entry(model_id.to_string())
.or_default()
.token_tree
.insert(smg_mesh::hash_token_path(tokens), tokens.clone());
.insert(kv_index::hash_token_path(tokens), tokens.clone());
applied += 1;
}
RepairEntry::String { .. } => {
Expand Down Expand Up @@ -883,7 +883,7 @@ impl CacheAwarePolicy {
.entry(model_id.to_string())
.or_default()
.token_tree
.insert(smg_mesh::hash_token_path(tokens), matched_prefix);
.insert(kv_index::hash_token_path(tokens), matched_prefix);

workers[idx].increment_processed();
return Some(idx);
Expand Down Expand Up @@ -948,7 +948,7 @@ impl CacheAwarePolicy {
// insert_text(matched_prefix, worker) which routes to the same
// tree node. This keeps the index memory-bounded.
let matched_prefix: String = text.chars().take(result.matched_char_count).collect();
let path_hash = smg_mesh::hash_node_path(text);
let path_hash = kv_index::hash_node_path(text);
self.hash_index
.entry(model_id.to_string())
.or_default()
Expand Down Expand Up @@ -1200,8 +1200,8 @@ mod tests {
};
assert_eq!(policy.apply_repair_page(&token_page), 1);

let text_hash = smg_mesh::hash_node_path(text);
let token_hash = smg_mesh::hash_token_path(&tokens);
let text_hash = kv_index::hash_node_path(text);
let token_hash = kv_index::hash_token_path(&tokens);

// Known hashes apply for the matching kind.
assert!(policy.apply_known_remote_insert(
Expand Down Expand Up @@ -1266,7 +1266,7 @@ mod tests {
assert!(policy.apply_known_remote_insert(
"model1",
TreeKind::String,
smg_mesh::hash_node_path(text),
kv_index::hash_node_path(text),
"http://w2",
));

Expand All @@ -1286,7 +1286,7 @@ mod tests {
assert!(policy.apply_known_remote_insert(
"model1",
TreeKind::Token,
smg_mesh::hash_token_path(&tokens),
kv_index::hash_token_path(&tokens),
"http://w2",
));
}
Expand Down Expand Up @@ -1333,7 +1333,7 @@ mod tests {
},
)
.unwrap();
let text_hash = smg_mesh::hash_node_path(text);
let text_hash = kv_index::hash_node_path(text);

// Drive a token request — populates the token-side
// hash_index. select_worker uses the model_id from the
Expand All @@ -1349,7 +1349,7 @@ mod tests {
},
)
.unwrap();
let token_hash = smg_mesh::hash_token_path(&tokens);
let token_hash = kv_index::hash_token_path(&tokens);

// Both populate sites use UNKNOWN_MODEL_ID for these
// workers (no model_id set on the builder), and the
Expand Down
Loading