#![allow(dead_code)]
use std::collections::BTreeMap;
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct ShardKey(pub String);
impl ShardKey {
#[must_use]
pub fn new(key: impl Into<String>) -> Self {
Self(key.into())
}
#[must_use]
pub fn hash64(&self) -> u64 {
const FNV_OFFSET: u64 = 14_695_981_039_346_656_037;
const FNV_PRIME: u64 = 1_099_511_628_211;
let mut h = FNV_OFFSET;
for &b in self.0.as_bytes() {
h ^= u64::from(b);
h = h.wrapping_mul(FNV_PRIME);
}
h ^= h >> 30;
h = h.wrapping_mul(0xbf58476d1ce4e5b9);
h ^= h >> 27;
h = h.wrapping_mul(0x94d049bb133111eb);
h ^= h >> 31;
h
}
}
#[derive(Debug, Clone)]
pub struct VirtualNode {
pub position: u64,
pub node_id: String,
pub replica: u32,
}
#[derive(Debug, Default)]
pub struct ConsistentHashRing {
ring: BTreeMap<u64, VirtualNode>,
replicas_per_node: u32,
}
impl ConsistentHashRing {
#[must_use]
pub fn new(replicas_per_node: u32) -> Self {
Self {
ring: BTreeMap::new(),
replicas_per_node: replicas_per_node.max(1),
}
}
pub fn add_node(&mut self, node_id: impl Into<String>) {
let node_id = node_id.into();
for replica in 0..self.replicas_per_node {
let key = ShardKey::new(format!("{node_id}:{replica}"));
let position = key.hash64();
self.ring.insert(
position,
VirtualNode {
position,
node_id: node_id.clone(),
replica,
},
);
}
}
pub fn remove_node(&mut self, node_id: &str) {
for replica in 0..self.replicas_per_node {
let key = ShardKey::new(format!("{node_id}:{replica}"));
let position = key.hash64();
self.ring.remove(&position);
}
}
#[must_use]
pub fn get_node(&self, key: &ShardKey) -> Option<&str> {
if self.ring.is_empty() {
return None;
}
let hash = key.hash64();
let node = self
.ring
.range(hash..)
.next()
.or_else(|| self.ring.iter().next())
.map(|(_, v)| v.node_id.as_str());
node
}
#[must_use]
pub fn virtual_node_count(&self) -> usize {
self.ring.len()
}
#[must_use]
pub fn physical_node_count(&self) -> usize {
let mut nodes: Vec<&str> = self.ring.values().map(|v| v.node_id.as_str()).collect();
nodes.sort_unstable();
nodes.dedup();
nodes.len()
}
#[must_use]
pub fn nodes(&self) -> Vec<String> {
let mut ids: Vec<String> = self.ring.values().map(|v| v.node_id.clone()).collect();
ids.sort();
ids.dedup();
ids
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ShardState {
Active,
Migrating,
Archived,
}
#[derive(Debug, Clone)]
pub struct ShardMetadata {
pub key: ShardKey,
pub owner_node: String,
pub size_bytes: u64,
pub state: ShardState,
pub updated_at_ms: u64,
}
impl ShardMetadata {
#[must_use]
pub fn new(key: ShardKey, owner_node: impl Into<String>, size_bytes: u64, now_ms: u64) -> Self {
Self {
key,
owner_node: owner_node.into(),
size_bytes,
state: ShardState::Active,
updated_at_ms: now_ms,
}
}
pub fn begin_migration(&mut self, target_node: impl Into<String>, now_ms: u64) {
self.owner_node = target_node.into();
self.state = ShardState::Migrating;
self.updated_at_ms = now_ms;
}
pub fn complete_migration(&mut self, now_ms: u64) {
self.state = ShardState::Active;
self.updated_at_ms = now_ms;
}
#[must_use]
pub fn is_active(&self) -> bool {
self.state == ShardState::Active
}
}
#[derive(Debug, Default)]
pub struct ShardCatalog {
shards: Vec<ShardMetadata>,
}
impl ShardCatalog {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn upsert(&mut self, shard: ShardMetadata) {
if let Some(existing) = self.shards.iter_mut().find(|s| s.key == shard.key) {
*existing = shard;
} else {
self.shards.push(shard);
}
}
#[must_use]
pub fn get(&self, key: &ShardKey) -> Option<&ShardMetadata> {
self.shards.iter().find(|s| &s.key == key)
}
#[must_use]
pub fn shards_for_node(&self, node_id: &str) -> Vec<&ShardMetadata> {
self.shards
.iter()
.filter(|s| s.owner_node == node_id)
.collect()
}
#[must_use]
pub fn total_size_bytes(&self) -> u64 {
self.shards.iter().map(|s| s.size_bytes).sum()
}
#[must_use]
pub fn len(&self) -> usize {
self.shards.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.shards.is_empty()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_shard_key_hash64_deterministic() {
let k = ShardKey::new("video-chunk-0042");
assert_eq!(k.hash64(), k.hash64());
}
#[test]
fn test_shard_key_different_keys_different_hashes() {
let k1 = ShardKey::new("a");
let k2 = ShardKey::new("b");
assert_ne!(k1.hash64(), k2.hash64());
}
#[test]
fn test_ring_empty_returns_none() {
let ring = ConsistentHashRing::new(3);
assert!(ring.get_node(&ShardKey::new("key")).is_none());
}
#[test]
fn test_ring_single_node_always_assigned() {
let mut ring = ConsistentHashRing::new(10);
ring.add_node("node-0");
for i in 0..20_u32 {
let key = ShardKey::new(format!("key-{i}"));
assert_eq!(ring.get_node(&key), Some("node-0"));
}
}
#[test]
fn test_ring_virtual_node_count() {
let mut ring = ConsistentHashRing::new(5);
ring.add_node("n0");
ring.add_node("n1");
assert_eq!(ring.virtual_node_count(), 10);
}
#[test]
fn test_ring_physical_node_count() {
let mut ring = ConsistentHashRing::new(5);
ring.add_node("n0");
ring.add_node("n1");
ring.add_node("n2");
assert_eq!(ring.physical_node_count(), 3);
}
#[test]
fn test_ring_remove_node() {
let mut ring = ConsistentHashRing::new(5);
ring.add_node("n0");
ring.add_node("n1");
ring.remove_node("n0");
assert_eq!(ring.physical_node_count(), 1);
}
#[test]
fn test_ring_nodes_list() {
let mut ring = ConsistentHashRing::new(3);
ring.add_node("alpha");
ring.add_node("beta");
let nodes = ring.nodes();
assert!(nodes.contains(&"alpha".to_string()));
assert!(nodes.contains(&"beta".to_string()));
assert_eq!(nodes.len(), 2);
}
#[test]
fn test_ring_distribution_two_nodes() {
let mut ring = ConsistentHashRing::new(50);
ring.add_node("node-A");
ring.add_node("node-B");
let mut a_count = 0_u32;
let mut b_count = 0_u32;
for i in 0..100_u32 {
let k = ShardKey::new(format!("shard-{i}"));
match ring.get_node(&k) {
Some("node-A") => a_count += 1,
Some("node-B") => b_count += 1,
_ => {}
}
}
assert!(a_count > 20 && b_count > 20, "a={a_count}, b={b_count}");
}
#[test]
fn test_shard_metadata_initial_state() {
let meta = ShardMetadata::new(ShardKey::new("k"), "node-0", 1024, 1000);
assert!(meta.is_active());
assert_eq!(meta.state, ShardState::Active);
}
#[test]
fn test_shard_metadata_begin_migration() {
let mut meta = ShardMetadata::new(ShardKey::new("k"), "node-0", 1024, 1000);
meta.begin_migration("node-1", 2000);
assert_eq!(meta.state, ShardState::Migrating);
assert_eq!(meta.owner_node, "node-1");
}
#[test]
fn test_shard_metadata_complete_migration() {
let mut meta = ShardMetadata::new(ShardKey::new("k"), "node-0", 1024, 1000);
meta.begin_migration("node-1", 2000);
meta.complete_migration(3000);
assert!(meta.is_active());
}
#[test]
fn test_catalog_empty() {
let catalog = ShardCatalog::new();
assert!(catalog.is_empty());
assert_eq!(catalog.len(), 0);
}
#[test]
fn test_catalog_upsert_and_get() {
let mut catalog = ShardCatalog::new();
let key = ShardKey::new("shard-0");
catalog.upsert(ShardMetadata::new(key.clone(), "n0", 512, 100));
let meta = catalog.get(&key).expect("get should return a value");
assert_eq!(meta.size_bytes, 512);
}
#[test]
fn test_catalog_upsert_updates_existing() {
let mut catalog = ShardCatalog::new();
let key = ShardKey::new("shard-0");
catalog.upsert(ShardMetadata::new(key.clone(), "n0", 512, 100));
catalog.upsert(ShardMetadata::new(key.clone(), "n1", 1024, 200));
assert_eq!(catalog.len(), 1);
assert_eq!(
catalog
.get(&key)
.expect("get should return a value")
.size_bytes,
1024
);
}
#[test]
fn test_catalog_shards_for_node() {
let mut catalog = ShardCatalog::new();
catalog.upsert(ShardMetadata::new(ShardKey::new("s0"), "n0", 100, 1));
catalog.upsert(ShardMetadata::new(ShardKey::new("s1"), "n0", 200, 1));
catalog.upsert(ShardMetadata::new(ShardKey::new("s2"), "n1", 300, 1));
let shards = catalog.shards_for_node("n0");
assert_eq!(shards.len(), 2);
}
#[test]
fn test_catalog_total_size_bytes() {
let mut catalog = ShardCatalog::new();
catalog.upsert(ShardMetadata::new(ShardKey::new("s0"), "n0", 100, 1));
catalog.upsert(ShardMetadata::new(ShardKey::new("s1"), "n0", 200, 1));
assert_eq!(catalog.total_size_bytes(), 300);
}
}