use std::{
any::Any,
collections::HashMap,
sync::{
Arc, OnceLock,
atomic::{AtomicU64, Ordering},
},
};
use dashmap::DashMap;
use hyphae::{MapDiff, SelectExt};
use serde::de::DeserializeOwned;
use serde_json::Value;
use uuid::Uuid;
use super::{
super::item::Eventable,
cell::FilteredCellMap,
context::{QueryCellContext, QueryContext},
request::QueryRequest,
traits::{AnyQuery, QueryBuildCellCtx, QueryHandler, QueryParams, QueryTestCtx},
};
use crate::{
common::with_id::WithId, core::item::downcast_any_item_arc, request::RequestContext,
server::CellServerCtx, store::StoreRegistry,
};
pub type QueryParseFn = fn(Value) -> Result<Arc<dyn AnyQuery>, anyhow::Error>;
pub type QueryCellFactory = fn(
Arc<dyn AnyQuery>,
Arc<StoreRegistry>,
Arc<RequestContext>,
Option<Arc<CellServerCtx>>,
) -> Result<FilteredCellMap, String>;
type AnyItemArc = Arc<dyn crate::core::item::AnyItem>;
type AnyItemMap = hyphae::CellMap<Arc<str>, AnyItemArc>;
type WeakAnyItemMap = hyphae::WeakCellMap<Arc<str>, AnyItemArc>;
type BucketEntries = Vec<(Arc<str>, AnyItemArc)>;
type BucketDiff = MapDiff<Arc<str>, AnyItemArc>;
type BucketDiffs = Vec<BucketDiff>;
type CompoundKey = Vec<Arc<str>>;
type CompoundFkExtractor = fn(&dyn std::any::Any) -> Option<Vec<Arc<str>>>;
inventory::collect!(QueryRegistration);
#[derive(Debug, Clone, Copy, Default)]
pub struct QueryRuntimeMetrics {
pub cell_factories_created: u64,
pub per_item_guards_created: u64,
pub per_item_guards_removed: u64,
}
#[derive(Debug, Clone, Default)]
pub struct QueryRuntimePerIdMetrics {
pub query_id: Arc<str>,
pub cell_factories_created: u64,
pub per_item_guards_created: u64,
pub per_item_guards_removed: u64,
}
static QUERY_CELL_FACTORIES_CREATED: AtomicU64 = AtomicU64::new(0);
static QUERY_PER_ITEM_GUARDS_CREATED: AtomicU64 = AtomicU64::new(0);
static QUERY_PER_ITEM_GUARDS_REMOVED: AtomicU64 = AtomicU64::new(0);
static QUERY_FACTORIES_BY_ID: OnceLock<DashMap<Arc<str>, u64>> = OnceLock::new();
static QUERY_GUARDS_CREATED_BY_ID: OnceLock<DashMap<Arc<str>, u64>> = OnceLock::new();
static QUERY_GUARDS_REMOVED_BY_ID: OnceLock<DashMap<Arc<str>, u64>> = OnceLock::new();
static BELONGS_TO_SOURCE_INDEXES: OnceLock<DashMap<String, Arc<BelongsToSourceIndex>>> =
OnceLock::new();
fn query_factories_by_id() -> &'static DashMap<Arc<str>, u64> {
QUERY_FACTORIES_BY_ID.get_or_init(DashMap::new)
}
fn query_guards_created_by_id() -> &'static DashMap<Arc<str>, u64> {
QUERY_GUARDS_CREATED_BY_ID.get_or_init(DashMap::new)
}
fn query_guards_removed_by_id() -> &'static DashMap<Arc<str>, u64> {
QUERY_GUARDS_REMOVED_BY_ID.get_or_init(DashMap::new)
}
fn belongs_to_source_indexes() -> &'static DashMap<String, Arc<BelongsToSourceIndex>> {
BELONGS_TO_SOURCE_INDEXES.get_or_init(DashMap::new)
}
fn increment_counter(map: &DashMap<Arc<str>, u64>, key: Arc<str>) {
if let Some(mut value) = map.get_mut(&key) {
*value = value.saturating_add(1);
} else {
map.insert(key, 1);
}
}
pub fn query_runtime_metrics() -> QueryRuntimeMetrics {
QueryRuntimeMetrics {
cell_factories_created: QUERY_CELL_FACTORIES_CREATED.load(Ordering::Relaxed),
per_item_guards_created: QUERY_PER_ITEM_GUARDS_CREATED.load(Ordering::Relaxed),
per_item_guards_removed: QUERY_PER_ITEM_GUARDS_REMOVED.load(Ordering::Relaxed),
}
}
pub fn query_runtime_metrics_by_id(limit: usize) -> Vec<QueryRuntimePerIdMetrics> {
let mut rows: Vec<QueryRuntimePerIdMetrics> = query_factories_by_id()
.iter()
.map(|entry| {
let query_id = entry.key().clone();
let cell_factories_created = *entry.value();
let per_item_guards_created = query_guards_created_by_id()
.get(&query_id)
.map(|v| *v.value())
.unwrap_or(0);
let per_item_guards_removed = query_guards_removed_by_id()
.get(&query_id)
.map(|v| *v.value())
.unwrap_or(0);
QueryRuntimePerIdMetrics {
query_id,
cell_factories_created,
per_item_guards_created,
per_item_guards_removed,
}
})
.collect();
rows.sort_by(|a, b| {
let a_live = a
.per_item_guards_created
.saturating_sub(a.per_item_guards_removed);
let b_live = b
.per_item_guards_created
.saturating_sub(b.per_item_guards_removed);
b_live
.cmp(&a_live)
.then_with(|| b.cell_factories_created.cmp(&a.cell_factories_created))
});
if rows.len() > limit {
rows.truncate(limit);
}
rows
}
struct BelongsToSourceIndex {
store: Arc<crate::store::EntityStore>,
buckets: DashMap<CompoundKey, WeakAnyItemMap>,
_driver: Arc<AnyItemMap>,
}
impl BelongsToSourceIndex {
fn new(store: Arc<crate::store::EntityStore>, extract_fk: CompoundFkExtractor) -> Arc<Self> {
let driver = Arc::new(AnyItemMap::new());
let index = Arc::new(Self {
store: store.clone(),
buckets: DashMap::new(),
_driver: driver.clone(),
});
let index_for_diffs = index.clone();
let guard = store.subscribe_diffs(move |diff| {
index_for_diffs.apply_diff(diff, extract_fk);
});
driver.own_guard(guard);
index
}
fn route_to_live_bucket(&self, key: &CompoundKey) -> Option<AnyItemMap> {
let entry = self.buckets.get(key)?;
if let Some(map) = entry.upgrade() {
return Some(map);
}
drop(entry);
self.buckets.remove(key);
None
}
fn bucket_for(&self, key: CompoundKey, extract_fk: CompoundFkExtractor) -> AnyItemMap {
match self.buckets.entry(key) {
dashmap::mapref::entry::Entry::Occupied(mut occupied) => {
if let Some(map) = occupied.get().upgrade() {
return map;
}
let map = Self::build_backfilled_bucket(&self.store, occupied.key(), extract_fk);
occupied.insert(map.downgrade());
map
}
dashmap::mapref::entry::Entry::Vacant(vacant) => {
let map = Self::build_backfilled_bucket(&self.store, vacant.key(), extract_fk);
vacant.insert(map.downgrade());
map
}
}
}
fn build_backfilled_bucket(
store: &crate::store::EntityStore,
key: &CompoundKey,
extract_fk: CompoundFkExtractor,
) -> AnyItemMap {
let map = AnyItemMap::new();
let backfill: BucketEntries = store
.snapshot()
.into_iter()
.filter(|(_, item)| extract_fk(item.as_any()).as_ref() == Some(key))
.collect();
if !backfill.is_empty() {
map.apply_batch(vec![MapDiff::Initial { entries: backfill }]);
}
map
}
fn sweep_dead_buckets(&self) {
self.buckets.retain(|_, weak| weak.upgrade().is_some());
}
fn apply_diff(&self, diff: &BucketDiff, extract_fk: CompoundFkExtractor) {
match diff {
MapDiff::Initial { entries } => {
let mut grouped: HashMap<CompoundKey, BucketEntries> = HashMap::new();
for (id, item) in entries {
if let Some(fk) = extract_fk(item.as_any()) {
grouped
.entry(fk)
.or_default()
.push((id.clone(), item.clone()));
}
}
let live_keys: Vec<CompoundKey> = self
.buckets
.iter()
.filter(|entry| entry.value().upgrade().is_some())
.map(|entry| entry.key().clone())
.collect();
for key in &live_keys {
if !grouped.contains_key(key)
&& let Some(bucket) = self.route_to_live_bucket(key)
{
bucket.apply_batch(vec![MapDiff::Initial {
entries: Vec::new(),
}]);
}
}
for (fk, bucket_entries) in grouped {
if let Some(bucket) = self.route_to_live_bucket(&fk) {
bucket.apply_batch(vec![MapDiff::Initial {
entries: bucket_entries,
}]);
}
}
}
MapDiff::Insert { key, value } => {
if let Some(fk) = extract_fk(value.as_any())
&& let Some(bucket) = self.route_to_live_bucket(&fk)
{
bucket.apply_batch(vec![MapDiff::Insert {
key: key.clone(),
value: value.clone(),
}]);
}
}
MapDiff::Remove { key, old_value } => {
if let Some(fk) = extract_fk(old_value.as_any())
&& let Some(bucket) = self.route_to_live_bucket(&fk)
{
bucket.apply_batch(vec![MapDiff::Remove {
key: key.clone(),
old_value: old_value.clone(),
}]);
}
}
MapDiff::Update {
key,
old_value,
new_value,
} => {
let old_fk = extract_fk(old_value.as_any());
let new_fk = extract_fk(new_value.as_any());
match (old_fk, new_fk) {
(Some(old_fk), Some(new_fk)) if old_fk == new_fk => {
if let Some(bucket) = self.route_to_live_bucket(&new_fk) {
bucket.apply_batch(vec![MapDiff::Update {
key: key.clone(),
old_value: old_value.clone(),
new_value: new_value.clone(),
}]);
}
}
(Some(old_fk), Some(new_fk)) => {
if let Some(bucket) = self.route_to_live_bucket(&old_fk) {
bucket.apply_batch(vec![MapDiff::Remove {
key: key.clone(),
old_value: old_value.clone(),
}]);
}
if let Some(bucket) = self.route_to_live_bucket(&new_fk) {
bucket.apply_batch(vec![MapDiff::Insert {
key: key.clone(),
value: new_value.clone(),
}]);
}
}
(Some(old_fk), None) => {
if let Some(bucket) = self.route_to_live_bucket(&old_fk) {
bucket.apply_batch(vec![MapDiff::Remove {
key: key.clone(),
old_value: old_value.clone(),
}]);
}
}
(None, Some(new_fk)) => {
if let Some(bucket) = self.route_to_live_bucket(&new_fk) {
bucket.apply_batch(vec![MapDiff::Insert {
key: key.clone(),
value: new_value.clone(),
}]);
}
}
(None, None) => {}
}
}
MapDiff::Batch { changes } => {
let mut by_fk: HashMap<CompoundKey, BucketDiffs> = HashMap::new();
for change in changes {
match change {
MapDiff::Insert { key, value } => {
if let Some(fk) = extract_fk(value.as_any()) {
by_fk.entry(fk).or_default().push(MapDiff::Insert {
key: key.clone(),
value: value.clone(),
});
}
}
MapDiff::Remove { key, old_value } => {
if let Some(fk) = extract_fk(old_value.as_any()) {
by_fk.entry(fk).or_default().push(MapDiff::Remove {
key: key.clone(),
old_value: old_value.clone(),
});
}
}
MapDiff::Update {
key,
old_value,
new_value,
} => {
let old_fk = extract_fk(old_value.as_any());
let new_fk = extract_fk(new_value.as_any());
match (old_fk, new_fk) {
(Some(old_fk), Some(new_fk)) if old_fk == new_fk => {
by_fk.entry(new_fk).or_default().push(MapDiff::Update {
key: key.clone(),
old_value: old_value.clone(),
new_value: new_value.clone(),
});
}
(Some(old_fk), Some(new_fk)) => {
by_fk.entry(old_fk).or_default().push(MapDiff::Remove {
key: key.clone(),
old_value: old_value.clone(),
});
by_fk.entry(new_fk).or_default().push(MapDiff::Insert {
key: key.clone(),
value: new_value.clone(),
});
}
(Some(old_fk), None) => {
by_fk.entry(old_fk).or_default().push(MapDiff::Remove {
key: key.clone(),
old_value: old_value.clone(),
});
}
(None, Some(new_fk)) => {
by_fk.entry(new_fk).or_default().push(MapDiff::Insert {
key: key.clone(),
value: new_value.clone(),
});
}
(None, None) => {}
}
}
MapDiff::Initial { .. } | MapDiff::Batch { .. } => {
self.apply_diff(change, extract_fk);
}
}
}
for (fk, bucket_changes) in by_fk {
if let Some(bucket) = self.route_to_live_bucket(&fk) {
bucket.apply_batch(bucket_changes);
}
}
}
}
}
}
pub fn build_belongs_to_source_map(
registry: Arc<StoreRegistry>,
host_id: Uuid,
local_type: &'static str,
field_names: &'static [&'static str],
extract_fk: CompoundFkExtractor,
foreign_ids: Vec<Arc<str>>,
) -> FilteredCellMap {
debug_assert_eq!(
field_names.len(),
foreign_ids.len(),
"build_belongs_to_source_map: field_names and foreign_ids must be positionally paired"
);
let key = format!("{host_id}:{local_type}:{}", field_names.join("+"));
let index = belongs_to_source_indexes()
.entry(key)
.or_insert_with(|| {
let store = registry.get_or_create(local_type);
BelongsToSourceIndex::new(store, extract_fk)
})
.clone();
index.bucket_for(foreign_ids, extract_fk).lock()
}
pub fn sweep_all_belongs_to_source_indexes() {
for entry in belongs_to_source_indexes().iter() {
entry.value().sweep_dead_buckets();
}
}
pub fn build_ids_source_map(
store: &Arc<crate::store::EntityStore>,
ids: &[Arc<str>],
) -> FilteredCellMap {
use hyphae::{Signal, Watchable};
let result: hyphae::CellMap<Arc<str>, AnyItemArc> = hyphae::CellMap::new();
for id in ids {
let key_cell = store.get(id);
let result_weak = result.downgrade();
let key_for_cb = id.clone();
let guard = key_cell.subscribe(move |signal| {
let Some(result_for_cb) = result_weak.upgrade() else {
return;
};
if let Signal::Value(arc_opt) = signal {
match arc_opt.as_ref() {
Some(item) => {
result_for_cb.insert(key_for_cb.clone(), item.clone());
}
None => {
result_for_cb.remove(&key_for_cb);
}
}
}
});
result.own(guard);
}
result.lock()
}
pub fn filter_query_over_source<Q>(
source: FilteredCellMap,
query: Arc<Q>,
query_context: Arc<QueryContext>,
) -> FilteredCellMap
where
Q: QueryHandler + QueryParams + Clone + Send + Sync + 'static,
Q::Item:
DeserializeOwned + Eventable + WithId + Clone + std::fmt::Debug + Send + Sync + 'static,
{
hyphae::MapQuery::materialize(source.select(move |item_any: &AnyItemArc| {
let item = downcast_any_item_arc::<Q::Item>(item_any, "filter_query_over_source");
Q::test_entity(QueryTestCtx {
item,
query: query.clone(),
query_context: query_context.clone(),
})
}))
}
pub struct QueryRegistration {
pub query_id: &'static str,
pub query_item_type: &'static str,
pub crate_name: &'static str,
pub parse: QueryParseFn,
pub cell_factory: QueryCellFactory,
pub args: &'static [crate::reflection::OperationArgField],
pub description: Option<&'static str>,
}
pub trait QueryFactory: QueryParams {
fn parse(value: Value) -> Result<Arc<dyn AnyQuery>, anyhow::Error>;
fn cell_factory(
query: Arc<dyn AnyQuery>,
registry: Arc<StoreRegistry>,
request_ctx: Arc<RequestContext>,
server_ctx: Option<Arc<CellServerCtx>>,
) -> Result<FilteredCellMap, String>;
}
impl<Q: QueryParams> QueryFactory for Q
where
Q::Item:
Eventable + WithId + DeserializeOwned + Clone + std::fmt::Debug + Send + Sync + 'static,
{
fn parse(value: Value) -> Result<Arc<dyn AnyQuery>, anyhow::Error> {
let query = serde_json::from_value::<QueryRequest<Q>>(value)?;
Ok(Arc::new(query))
}
fn cell_factory(
any_query: Arc<dyn AnyQuery>,
registry: Arc<StoreRegistry>,
request_ctx: Arc<RequestContext>,
server_ctx: Option<Arc<CellServerCtx>>,
) -> Result<FilteredCellMap, String> {
QUERY_CELL_FACTORIES_CREATED.fetch_add(1, Ordering::Relaxed);
let query_id = Q::query_id_static();
let _span = tracing::trace_span!("myko.query", query = query_id.as_ref()).entered();
crate::server::dispatch_metrics::record_query(query_id.as_ref(), request_ctx.origin());
increment_counter(query_factories_by_id(), query_id);
let any_ref: &dyn Any = any_query.as_ref();
let request: QueryRequest<Q> = any_ref
.downcast_ref::<QueryRequest<Q>>()
.cloned()
.ok_or_else(|| "Failed to downcast query payload".to_string())?;
let query: Arc<Q> = Arc::new(request.query);
let query_ctx = Arc::new(QueryContext {
req: request_ctx.clone(),
});
let query_cell_ctx =
QueryCellContext::new(request_ctx, query_ctx.clone(), registry.clone(), server_ctx);
if let Some(built) = Q::build_view(QueryBuildCellCtx {
query: query.clone(),
query_context: query_cell_ctx.clone(),
}) {
return Ok(hyphae::MapQuery::materialize(built));
}
let store: crate::store::EntityStore =
(*registry.get_or_create(&Q::query_item_type_static())).clone();
Ok(hyphae::MapQuery::materialize(store.select(
move |item_any: &AnyItemArc| {
let item = downcast_any_item_arc::<Q::Item>(item_any, "QueryFactory::cell_factory");
Q::test_entity(QueryTestCtx {
item,
query: query.clone(),
query_context: query_ctx.clone(),
})
},
)))
}
}
#[cfg(test)]
mod belongs_to_source_index_tests {
use std::any::Any;
use serde::Serialize;
use super::*;
use crate::common::with_id::WithId;
#[derive(Debug, Clone, PartialEq, Serialize)]
struct TestChild {
id: Arc<str>,
parent_id: Arc<str>,
}
impl WithId for TestChild {
fn id(&self) -> Arc<str> {
self.id.clone()
}
}
impl crate::core::item::AnyItem for TestChild {
fn as_any(&self) -> &dyn Any {
self
}
fn entity_type(&self) -> &'static str {
"TestChild"
}
fn equals(&self, other: &dyn crate::core::item::AnyItem) -> bool {
other
.as_any()
.downcast_ref::<Self>()
.map(|t| t == self)
.unwrap_or(false)
}
}
fn extract_parent_fk(item: &dyn Any) -> Option<Vec<Arc<str>>> {
item.downcast_ref::<TestChild>()
.map(|c| vec![c.parent_id.clone()])
}
fn child(id: &str, parent: &str) -> (Arc<str>, AnyItemArc) {
(
Arc::from(id),
Arc::new(TestChild {
id: Arc::from(id),
parent_id: Arc::from(parent),
}) as AnyItemArc,
)
}
fn new_store() -> Arc<crate::store::EntityStore> {
Arc::new(hyphae::CellMap::new())
}
#[test]
fn dropped_subscriptions_do_not_leak_across_many_distinct_parents() {
let store = new_store();
let index = BelongsToSourceIndex::new(store, extract_parent_fk);
for i in 0..50 {
let parent: Arc<str> = Arc::from(format!("parent-{i}"));
let bucket = index.bucket_for(vec![parent], extract_parent_fk);
drop(bucket); }
index.sweep_dead_buckets();
assert_eq!(
index.buckets.len(),
0,
"sweep must reap all buckets once every subscriber has dropped"
);
}
#[test]
fn live_subscription_survives_going_empty_then_repopulating() {
let store = new_store();
let (id, item) = child("c1", "parent-x");
store.insert(id.clone(), item);
let index = BelongsToSourceIndex::new(store.clone(), extract_parent_fk);
let bucket = index.bucket_for(vec![Arc::from("parent-x")], extract_parent_fk);
assert_eq!(bucket.snapshot().len(), 1);
store.remove(&id);
assert_eq!(bucket.snapshot().len(), 0);
let (id2, item2) = child("c2", "parent-x");
store.insert(id2, item2);
assert_eq!(
bucket.snapshot().len(),
1,
"a live subscriber must see re-population after its bucket went empty"
);
}
#[test]
fn resubscribing_after_reap_backfills_current_children() {
let store = new_store();
let (id, item) = child("c1", "parent-y");
store.insert(id, item);
let index = BelongsToSourceIndex::new(store.clone(), extract_parent_fk);
{
let bucket = index.bucket_for(vec![Arc::from("parent-y")], extract_parent_fk);
assert_eq!(bucket.snapshot().len(), 1);
}
index.sweep_dead_buckets();
assert!(index.buckets.is_empty());
let bucket = index.bucket_for(vec![Arc::from("parent-y")], extract_parent_fk);
assert_eq!(
bucket.snapshot().len(),
1,
"resubscribing after the bucket was reaped must backfill current children, not start empty"
);
}
#[test]
fn concurrent_bucket_for_calls_for_the_same_key_never_orphan_a_bucket() {
let store = new_store();
let index = Arc::new(BelongsToSourceIndex::new(store.clone(), extract_parent_fk));
let key: CompoundKey = vec![Arc::from("parent-race")];
const N: usize = 16;
let barrier = Arc::new(std::sync::Barrier::new(N));
let handles: Vec<_> = (0..N)
.map(|_| {
let index = index.clone();
let key = key.clone();
let barrier = barrier.clone();
std::thread::spawn(move || {
barrier.wait();
index.bucket_for(key, extract_parent_fk)
})
})
.collect();
let buckets: Vec<AnyItemMap> = handles.into_iter().map(|h| h.join().unwrap()).collect();
assert_eq!(
index.buckets.len(),
1,
"N concurrent creators for the same key must converge on exactly one bucket entry"
);
let (id, item) = child("c-race", "parent-race");
store.insert(id, item);
for (i, bucket) in buckets.iter().enumerate() {
assert_eq!(
bucket.snapshot().len(),
1,
"handle {i} must observe the post-race insert — an orphaned bucket stays empty forever"
);
}
}
#[derive(Debug, Clone, PartialEq, Serialize)]
struct TestCursor {
id: Arc<str>,
node_id: Arc<str>,
session_id: Arc<str>,
}
impl WithId for TestCursor {
fn id(&self) -> Arc<str> {
self.id.clone()
}
}
impl crate::core::item::AnyItem for TestCursor {
fn as_any(&self) -> &dyn Any {
self
}
fn entity_type(&self) -> &'static str {
"TestCursor"
}
fn equals(&self, other: &dyn crate::core::item::AnyItem) -> bool {
other
.as_any()
.downcast_ref::<Self>()
.map(|t| t == self)
.unwrap_or(false)
}
}
fn cursor(id: &str, node: &str, session: &str) -> (Arc<str>, AnyItemArc) {
(
Arc::from(id),
Arc::new(TestCursor {
id: Arc::from(id),
node_id: Arc::from(node),
session_id: Arc::from(session),
}) as AnyItemArc,
)
}
fn extract_node_and_session_fk(item: &dyn Any) -> Option<Vec<Arc<str>>> {
item.downcast_ref::<TestCursor>()
.map(|c| vec![c.node_id.clone(), c.session_id.clone()])
}
#[test]
fn compound_key_separates_watchers_sharing_one_field_but_not_the_other() {
let store = new_store();
let (id_a, item_a) = cursor("cursor-a", "node-A", "session-PROD");
let (id_b, item_b) = cursor("cursor-b", "node-B", "session-PROD");
store.insert(id_a.clone(), item_a);
store.insert(id_b.clone(), item_b);
let index = BelongsToSourceIndex::new(store.clone(), extract_node_and_session_fk);
let key_a: CompoundKey = vec![Arc::from("node-A"), Arc::from("session-PROD")];
let key_b: CompoundKey = vec![Arc::from("node-B"), Arc::from("session-PROD")];
let bucket_a = index.bucket_for(key_a.clone(), extract_node_and_session_fk);
let bucket_b = index.bucket_for(key_b.clone(), extract_node_and_session_fk);
assert_eq!(
bucket_a.snapshot().len(),
1,
"node-A's bucket sees only its own cursor"
);
assert_eq!(
bucket_b.snapshot().len(),
1,
"node-B's bucket sees only its own cursor"
);
assert_eq!(
index.buckets.len(),
2,
"distinct (node, session) pairs get distinct buckets"
);
let (id_a2, item_a2) = cursor("cursor-a-tick2", "node-A", "session-PROD");
store.insert(id_a2, item_a2);
assert_eq!(
bucket_a.snapshot().len(),
2,
"node-A's bucket must see its own new entry"
);
assert_eq!(
bucket_b.snapshot().len(),
1,
"node-B's bucket must be unaffected by node-A's insert"
);
let (id_b2, item_b2) = cursor("cursor-b-tick2", "node-B", "session-PROD");
store.insert(id_b2, item_b2);
assert_eq!(
bucket_b.snapshot().len(),
2,
"node-B's bucket must independently see its own new entry — this is the \
regression rship-qtu hit: the second-registered watcher never received \
a diff under single-field session-only routing"
);
}
#[test]
fn compound_and_single_field_routing_never_share_a_bucket() {
let store = new_store();
let index = BelongsToSourceIndex::new(store, extract_node_and_session_fk);
let key: CompoundKey = vec![Arc::from("node-A"), Arc::from("session-PROD")];
let bucket = index.bucket_for(key.clone(), extract_node_and_session_fk);
assert_eq!(bucket.snapshot().len(), 0);
assert!(index.buckets.contains_key(&key));
}
}