#![cfg(feature = "bench")]
use std::sync::Arc;
use myko::{
bench_entities::{BenchItem, BenchItemQuery, GetBenchItemsByQuery, SwitchMapReport},
hyphae::{Cell, Gettable, Materialize, Mutable, SwitchMapExt},
query::StringFilter,
search::SearchIndex,
server::{HandlerRegistry, MykoServerContext, RelationshipManager, persister::PersisterRouter},
store::StoreRegistry,
wire::{MEvent, MEventType},
};
use uuid::Uuid;
fn scheduler_test_serial() -> std::sync::MutexGuard<'static, ()> {
static LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
LOCK.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
fn make_ctx() -> MykoServerContext {
MykoServerContext::new(
Uuid::new_v4(),
Arc::new(StoreRegistry::new()),
Arc::new(HandlerRegistry::new()),
Arc::new(RelationshipManager::new()),
Arc::new(PersisterRouter::default()),
Arc::new(SearchIndex::new()),
myko::server::MykoServerRuntime {
peer_clients: Arc::new(dashmap::DashMap::new()),
event_sink: None,
history_replay: None,
},
)
}
fn insert_bench_item(ctx: &MykoServerContext, id: &str, category: &str, value: i64) {
let item = BenchItem {
id: id.into(),
name: format!("item-{id}"),
category: category.to_string(),
value,
};
let event = MEvent::from_item(&item, MEventType::SET, &format!("tx-{id}"));
assert!(ctx.apply_event_batch(vec![event]).is_ok());
}
#[test]
fn query_map_inside_switch_map_cache_entries_become_reclaimable() {
let _serial = scheduler_test_serial();
let ctx = make_ctx();
for i in 0..10 {
insert_bench_item(&ctx, &format!("a-{i}"), "alpha", i);
insert_bench_item(&ctx, &format!("b-{i}"), "beta", i);
insert_bench_item(&ctx, &format!("g-{i}"), "gamma", i);
}
let categories = ["alpha", "beta", "gamma"];
let selector = Cell::new(0usize);
let ctx_clone = ctx.clone();
let switched = selector
.clone()
.switch_map(move |idx| {
let category = idx
.checked_rem(categories.len())
.and_then(|index| categories.get(index))
.copied()
.unwrap_or_default()
.to_string();
let request = ctx_clone.new_server_transaction();
let query_result = ctx_clone.query_map(
GetBenchItemsByQuery(BenchItemQuery {
category: Some(StringFilter::Eq(category.into())),
..Default::default()
}),
request,
);
query_result.items().materialize()
})
.materialize();
assert_eq!(switched.get().len(), 10);
let cache_before = ctx.query_cache_len();
for i in 1..=30 {
selector.set(i);
}
assert_eq!(switched.get().len(), 10);
let cache_after = ctx.query_cache_len();
assert!(
cache_after >= cache_before,
"cache should have grown or stayed the same"
);
drop(switched);
drop(selector);
ctx.sweep_dead_cache_entries();
let cache_final = ctx.query_cache_len();
assert!(
cache_final < cache_after,
"sweep should have removed dead query cache entries \
(cache before={cache_before}, after switch={cache_after}, final={cache_final})",
);
}
#[test]
fn query_map_same_params_inside_switch_map_reuses_cache() {
let _serial = scheduler_test_serial();
let ctx = make_ctx();
for i in 0..5 {
insert_bench_item(&ctx, &format!("item-{i}"), "tools", i);
}
let trigger = Cell::new(0u64);
let ctx_clone = ctx.clone();
let switched = trigger
.clone()
.switch_map(move |_| {
let request = ctx_clone.new_server_transaction();
ctx_clone
.query_map(
GetBenchItemsByQuery(BenchItemQuery {
category: Some(StringFilter::Eq("tools".into())),
..Default::default()
}),
request,
)
.items()
.materialize()
})
.materialize();
assert_eq!(switched.get().len(), 5);
let cache_after_init = ctx.query_cache_len();
for i in 1..=50 {
trigger.set(i);
}
assert_eq!(switched.get().len(), 5);
let cache_after_switches = ctx.query_cache_len();
assert!(
cache_after_switches <= cache_after_init.saturating_add(5),
"cache grew from {cache_after_init} to {cache_after_switches} after 50 same-param switches — \
expected cache reuse, got accumulation",
);
}
#[test]
fn query_map_inside_switch_map_with_active_store_mutations() {
let _serial = scheduler_test_serial();
let ctx = make_ctx();
for i in 0..10 {
insert_bench_item(&ctx, &format!("a-{i}"), "alpha", i);
insert_bench_item(&ctx, &format!("b-{i}"), "beta", i);
}
let categories = ["alpha", "beta"];
let selector = Cell::new(0usize);
let ctx_clone = ctx.clone();
let switched = selector
.clone()
.switch_map(move |idx| {
let category = idx
.checked_rem(categories.len())
.and_then(|index| categories.get(index))
.copied()
.unwrap_or_default()
.to_string();
let request = ctx_clone.new_server_transaction();
ctx_clone
.query_map(
GetBenchItemsByQuery(BenchItemQuery {
category: Some(StringFilter::Eq(category.into())),
..Default::default()
}),
request,
)
.items()
.materialize()
})
.materialize();
assert_eq!(switched.get().len(), 10);
for round in 0usize..20 {
selector.set(round.saturating_add(1));
let cat = round
.checked_rem(categories.len())
.and_then(|index| categories.get(index))
.copied()
.unwrap_or_default();
for i in 0usize..5 {
insert_bench_item(
&ctx,
&format!("{}-{i}", cat.chars().next().unwrap_or('?')),
cat,
i64::try_from(round.saturating_mul(10).saturating_add(i)).unwrap_or(i64::MAX),
);
}
}
let live_before_drop = ctx.query_cache_live_count();
let total_before_drop = ctx.query_cache_len();
drop(switched);
drop(selector);
ctx.sweep_dead_cache_entries();
let live_after = ctx.query_cache_live_count();
assert_eq!(
live_after, 0,
"expected 0 live cache entries after dropping switch_map, found {live_after} \
(total before drop={total_before_drop}, live before drop={live_before_drop})",
);
}
#[test]
fn query_cache_live_count_tracks_reachable_entries() {
let _serial = scheduler_test_serial();
let ctx = make_ctx();
for i in 0..3 {
insert_bench_item(&ctx, &format!("item-{i}"), "test", i);
}
assert_eq!(ctx.query_cache_live_count(), 0);
let request = ctx.new_server_transaction();
let map = ctx.query_map(
GetBenchItemsByQuery(BenchItemQuery {
category: Some(StringFilter::Eq("test".into())),
..Default::default()
}),
request,
);
assert_eq!(ctx.query_cache_live_count(), 1);
let request2 = ctx.new_server_transaction();
let map2 = ctx.query_map(
GetBenchItemsByQuery(BenchItemQuery {
category: Some(StringFilter::Eq("other".into())),
..Default::default()
}),
request2,
);
assert_eq!(ctx.query_cache_live_count(), 2);
drop(map);
ctx.sweep_dead_cache_entries();
assert_eq!(ctx.query_cache_live_count(), 1);
drop(map2);
ctx.sweep_dead_cache_entries();
assert_eq!(ctx.query_cache_live_count(), 0);
}
#[test]
fn report_with_switch_map_query_map_cleans_up_cache() {
let _serial = scheduler_test_serial();
let ctx = make_ctx();
for i in 0..10 {
insert_bench_item(&ctx, &format!("item-{i}"), "alpha", i);
}
let cache_before = ctx.query_cache_len();
let report_cache_before = ctx.report_cache_len();
let request = ctx.new_server_transaction();
let report_cell = ctx.report(
SwitchMapReport {
category: "alpha".to_string(),
},
request,
);
assert_eq!(report_cell.get().len(), 10);
let cache_after_report = ctx.query_cache_len();
let report_cache_after_report = ctx.report_cache_len();
assert!(cache_after_report > cache_before);
assert!(report_cache_after_report > report_cache_before);
for round in 0..20 {
insert_bench_item(
&ctx,
&format!("new-{round}"),
"alpha",
100 + i64::from(round),
);
}
assert_eq!(report_cell.get().len(), 30);
let cache_during = ctx.query_cache_len();
drop(report_cell);
ctx.sweep_dead_cache_entries();
let cache_final = ctx.query_cache_live_count();
let report_cache_final = ctx.report_cache_live_count();
assert_eq!(
cache_final, 0,
"expected 0 live query cache entries after dropping report, found {cache_final} \
(cache grew to {cache_during} during mutations)",
);
assert_eq!(
report_cache_final, 0,
"expected 0 live report cache entries after dropping report, found {report_cache_final}",
);
}
#[test]
fn report_switch_map_cache_bounded_during_active_mutations() {
let _serial = scheduler_test_serial();
let ctx = make_ctx();
for i in 0..5 {
insert_bench_item(&ctx, &format!("item-{i}"), "beta", i);
}
let request = ctx.new_server_transaction();
let report_cell = ctx.report(
SwitchMapReport {
category: "beta".to_string(),
},
request,
);
assert_eq!(report_cell.get().len(), 5);
for round in 0..100 {
insert_bench_item(
&ctx,
&format!("stress-{round}"),
"beta",
1000_i64.saturating_add(i64::from(round)),
);
}
assert_eq!(report_cell.get().len(), 105);
ctx.sweep_dead_cache_entries();
let live_count = ctx.query_cache_live_count();
let total_count = ctx.query_cache_len();
assert!(
live_count <= 10,
"live query cache entries ({live_count}) should be bounded during active use, \
but found {total_count} total cache entries — old entries are being kept alive",
);
drop(report_cell);
ctx.sweep_dead_cache_entries();
assert_eq!(
ctx.query_cache_live_count(),
0,
"all entries should be dead after drop"
);
}