use std::collections::HashSet;
use crate::engine::graph::edge_store::parse_versioned_edge_key;
use crate::types::TenantDataSnapshot;
pub fn extract_db_tenant_scoped_collection(key: &str, tenant_id: u64) -> Option<&str> {
let mut it = key.splitn(3, ':');
let _db = it.next()?;
let tid = it.next()?;
if tid.parse::<u64>().ok()? != tenant_id {
return None;
}
let rest = it.next()?;
let coll = rest.split([':', '\u{0}']).next()?;
if coll.is_empty() { None } else { Some(coll) }
}
pub fn extract_db_scoped_collection(key: &str, tenant_id: u64) -> Option<&str> {
let mut it = key.splitn(3, ':');
let _db = it.next()?;
let tid = it.next()?;
let coll = it.next()?;
if tid.parse::<u64>().ok()? != tenant_id || coll.is_empty() {
return None;
}
Some(coll)
}
pub fn retain_tenant_data_for_vshards(
snap: &mut TenantDataSnapshot,
tenant_id: u64,
source_vshards: &HashSet<u32>,
vshard_of: impl Fn(&str) -> u32,
) {
let in_group_db_tenant_scoped = |key: &str| {
extract_db_tenant_scoped_collection(key, tenant_id)
.map(|c| source_vshards.contains(&vshard_of(c)))
.unwrap_or(false)
};
let in_group_db_scoped = |key: &str| {
extract_db_scoped_collection(key, tenant_id)
.map(|c| source_vshards.contains(&vshard_of(c)))
.unwrap_or(false)
};
snap.documents.retain(|(k, _)| in_group_db_tenant_scoped(k));
snap.indexes.retain(|(k, _)| in_group_db_tenant_scoped(k));
snap.vectors.retain(|(k, _)| in_group_db_tenant_scoped(k));
snap.timeseries
.retain(|(k, _)| in_group_db_tenant_scoped(k));
snap.flushed_ts_segments
.retain(|b| in_group_db_scoped(&b.collection_key));
snap.columnar_engines.retain(|(k, _)| in_group_db_scoped(k));
snap.kv_tables
.retain(|(k, _)| source_vshards.contains(&vshard_of(k)));
snap.surrogate_pk
.retain(|e| source_vshards.contains(&vshard_of(&e.collection)));
snap.edges.retain(|(k, _)| {
parse_versioned_edge_key(k)
.map(|(collection, ..)| source_vshards.contains(&vshard_of(collection)))
.unwrap_or(false)
});
snap.crdt_state
.retain(|(_, _, collection, _)| source_vshards.contains(&vshard_of(collection)));
}
#[cfg(test)]
mod tests {
use super::{extract_db_scoped_collection, extract_db_tenant_scoped_collection};
#[test]
fn extract_db_tenant_scoped_collection_parses_key() {
assert_eq!(
extract_db_tenant_scoped_collection("0:1:snap_rt_docs:abcd1234", 1),
Some("snap_rt_docs")
);
assert_eq!(
extract_db_tenant_scoped_collection("0:1:users\u{0}doc1", 1),
Some("users")
);
assert_eq!(
extract_db_tenant_scoped_collection("0:1:metrics", 1),
Some("metrics")
);
assert_eq!(extract_db_tenant_scoped_collection("0:2:x:y", 1), None);
assert_eq!(extract_db_tenant_scoped_collection("0:1:", 1), None);
assert_eq!(extract_db_tenant_scoped_collection("0:1", 1), None);
}
#[test]
fn extract_db_scoped_collection_parses_db_prefixed_key() {
assert_eq!(
extract_db_scoped_collection("0:7:metrics", 7),
Some("metrics")
);
assert_eq!(
extract_db_scoped_collection("0:7:a:b", 7),
Some("a:b"),
"collection retains embedded ':'"
);
assert_eq!(extract_db_scoped_collection("0:8:metrics", 7), None);
assert_eq!(extract_db_scoped_collection("0:7", 7), None);
assert_eq!(extract_db_scoped_collection("0:7:", 7), None);
}
#[test]
fn retain_filters_append_sections_to_owning_vshard_only() {
use super::retain_tenant_data_for_vshards;
use crate::types::{TenantDataSnapshot, TsFlushedCollectionBlob};
use std::collections::HashSet;
const TID: u64 = 1;
let vshard_of = |c: &str| c.bytes().next().map(u32::from).unwrap_or(0);
let va = vshard_of("alpha"); let vb = vshard_of("beta");
let template = || TenantDataSnapshot {
timeseries: vec![
(format!("0:{TID}:alpha"), b"a".to_vec()),
(format!("0:{TID}:beta"), b"b".to_vec()),
],
columnar_engines: vec![
(format!("0:{TID}:alpha"), b"a".to_vec()),
(format!("0:{TID}:beta"), b"b".to_vec()),
],
flushed_ts_segments: vec![
TsFlushedCollectionBlob {
collection_key: format!("0:{TID}:alpha"),
partitions: vec![],
},
TsFlushedCollectionBlob {
collection_key: format!("0:{TID}:beta"),
partitions: vec![],
},
],
kv_tables: vec![("alpha".into(), b"a".to_vec())],
..Default::default()
};
let mut node_a = template();
let only_a: HashSet<u32> = [va].into_iter().collect();
retain_tenant_data_for_vshards(&mut node_a, TID, &only_a, vshard_of);
assert_eq!(node_a.timeseries.len(), 1);
assert_eq!(node_a.timeseries[0].0, format!("0:{TID}:alpha"));
assert_eq!(node_a.columnar_engines.len(), 1);
assert_eq!(node_a.flushed_ts_segments.len(), 1);
assert_eq!(node_a.kv_tables.len(), 1);
let mut node_b = template();
let only_b: HashSet<u32> = [vb].into_iter().collect();
retain_tenant_data_for_vshards(&mut node_b, TID, &only_b, vshard_of);
assert_eq!(node_b.timeseries.len(), 1);
assert_eq!(node_b.timeseries[0].0, format!("0:{TID}:beta"));
assert_eq!(node_b.kv_tables.len(), 0, "alpha kv not owned by beta node");
let mut node_c = template();
let none: HashSet<u32> = HashSet::new();
retain_tenant_data_for_vshards(&mut node_c, TID, &none, vshard_of);
assert!(node_c.timeseries.is_empty());
assert!(node_c.columnar_engines.is_empty());
let union_ts = node_a.timeseries.len() + node_b.timeseries.len() + node_c.timeseries.len();
assert_eq!(
union_ts, 2,
"each timeseries collection captured exactly once"
);
}
#[test]
fn retain_is_noop_when_node_owns_all_vshards() {
use super::retain_tenant_data_for_vshards;
use crate::types::TenantDataSnapshot;
use std::collections::HashSet;
let vshard_of = |c: &str| c.bytes().next().map(u32::from).unwrap_or(0);
let mut snap = TenantDataSnapshot {
timeseries: vec![("0:1:alpha".into(), b"a".to_vec())],
columnar_engines: vec![("0:1:beta".into(), b"b".to_vec())],
kv_tables: vec![("gamma".into(), b"g".to_vec())],
..Default::default()
};
let all: HashSet<u32> = ["alpha", "beta", "gamma"]
.iter()
.map(|c| vshard_of(c))
.collect();
retain_tenant_data_for_vshards(&mut snap, 1, &all, vshard_of);
assert_eq!(snap.timeseries.len(), 1);
assert_eq!(snap.columnar_engines.len(), 1);
assert_eq!(snap.kv_tables.len(), 1);
}
}