use std::sync::Arc;
use myko::{
entities::{
client::{Client, ClientQuery},
server::ServerId,
},
hyphae::{Cell, CellMutable, MapExt, Materialize, Mutable},
query::IdFilter,
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(myko::search::SearchIndex::new()),
myko::server::MykoServerRuntime {
peer_clients: Arc::new(dashmap::DashMap::new()),
event_sink: None,
history_replay: None,
},
)
}
fn insert_client(ctx: &MykoServerContext, id: &str, server_id: &str) {
let client = Client {
id: id.into(),
server_id: ServerId::from(Arc::<str>::from(server_id)),
address: None,
windback: None,
};
let event = MEvent::from_item(&client, MEventType::SET, &format!("tx-{id}"));
assert!(ctx.apply_event_batch(vec![event]).is_ok());
}
fn eq_filter(server_id: &str) -> ClientQuery {
ClientQuery {
server_id: Some(IdFilter::Eq(ServerId::from(Arc::<str>::from(server_id)))),
..Default::default()
}
}
fn in_filter(server_ids: &[&str]) -> ClientQuery {
ClientQuery {
server_id: Some(IdFilter::In(
server_ids
.iter()
.map(|s| ServerId::from(Arc::<str>::from(*s)))
.collect(),
)),
..Default::default()
}
}
#[test]
fn query_live_initial_population_matches_the_starting_filter() {
let _serial = scheduler_test_serial();
let ctx = make_ctx();
insert_client(&ctx, "c1", "server-A");
insert_client(&ctx, "c2", "server-B");
let filter_cell: Cell<ClientQuery, CellMutable> = Cell::new(eq_filter("server-A"));
let result = ctx.query_live(filter_cell);
assert_eq!(result.snapshot().len(), 1);
}
#[test]
fn query_live_updates_when_the_in_set_grows() {
let _serial = scheduler_test_serial();
let ctx = make_ctx();
insert_client(&ctx, "c1", "server-A");
insert_client(&ctx, "c2", "server-B");
insert_client(&ctx, "c3", "server-C");
let filter_cell: Cell<ClientQuery, CellMutable> = Cell::new(in_filter(&["server-A"]));
let result = ctx.query_live(filter_cell.clone());
assert_eq!(result.snapshot().len(), 1);
filter_cell.set(in_filter(&["server-A", "server-B"]));
assert_eq!(
result.snapshot().len(),
2,
"growing the In set must add the newly-included server's clients"
);
filter_cell.set(in_filter(&["server-A", "server-B", "server-C"]));
assert_eq!(result.snapshot().len(), 3);
}
#[test]
fn query_live_updates_when_the_in_set_shrinks() {
let _serial = scheduler_test_serial();
let ctx = make_ctx();
insert_client(&ctx, "c1", "server-A");
insert_client(&ctx, "c2", "server-B");
let filter_cell: Cell<ClientQuery, CellMutable> =
Cell::new(in_filter(&["server-A", "server-B"]));
let result = ctx.query_live(filter_cell.clone());
assert_eq!(result.snapshot().len(), 2);
filter_cell.set(in_filter(&["server-A"]));
assert_eq!(
result.snapshot().len(),
1,
"shrinking the In set must retract the dropped server's clients"
);
}
#[test]
fn query_live_still_tracks_store_writes_after_a_filter_change() {
let _serial = scheduler_test_serial();
let ctx = make_ctx();
insert_client(&ctx, "c1", "server-A");
let filter_cell: Cell<ClientQuery, CellMutable> =
Cell::new(in_filter(&["server-A", "server-B"]));
let result = ctx.query_live(filter_cell.clone());
assert_eq!(result.snapshot().len(), 1);
filter_cell.set(in_filter(&["server-B"]));
assert_eq!(result.snapshot().len(), 0);
insert_client(&ctx, "c2", "server-B");
assert_eq!(
result.snapshot().len(),
1,
"a write to a bucket added by the LATEST filter tick must still be tracked reactively"
);
}
#[test]
fn query_live_range_filter_change_reevaluates_correctly() {
let _serial = scheduler_test_serial();
let ctx = make_ctx();
let client_a = Client {
id: "c1".into(),
server_id: ServerId::from(Arc::<str>::from("server-A")),
address: None,
windback: Some(Arc::from("2026-01-01T00:00:00Z")),
};
assert!(
ctx.apply_event_batch(vec![MEvent::from_item(&client_a, MEventType::SET, "tx-1")])
.is_ok()
);
let filter_cell: Cell<ClientQuery, CellMutable> = Cell::new(ClientQuery {
windback: Some(myko::query::StringFilter::Eq(Arc::from(
"2026-01-01T00:00:00Z",
))),
..Default::default()
});
let result = ctx.query_live(filter_cell.clone());
assert_eq!(result.snapshot().len(), 1);
filter_cell.set(ClientQuery {
windback: Some(myko::query::StringFilter::Eq(Arc::from("nope"))),
..Default::default()
});
assert_eq!(
result.snapshot().len(),
0,
"a non-indexed filter change must re-evaluate the scan-mode scope"
);
filter_cell.set(ClientQuery {
windback: Some(myko::query::StringFilter::Eq(Arc::from(
"2026-01-01T00:00:00Z",
))),
..Default::default()
});
assert_eq!(result.snapshot().len(), 1);
}
#[test]
fn query_live_downstream_state_survives_a_filter_change() {
use std::sync::atomic::{AtomicUsize, Ordering};
let _serial = scheduler_test_serial();
let ctx = make_ctx();
insert_client(&ctx, "c1", "server-A");
let filter_cell: Cell<ClientQuery, CellMutable> = Cell::new(in_filter(&["server-A"]));
let result = ctx.query_live(filter_cell.clone());
let build_count = Arc::new(AtomicUsize::new(0));
let build_count_for_map = build_count.clone();
let downstream = result
.entries()
.map(move |_entries| {
build_count_for_map.fetch_add(1, Ordering::SeqCst);
build_count_for_map.load(Ordering::SeqCst)
})
.materialize();
let _ = myko::hyphae::Gettable::get(&downstream);
let builds_after_first_read = build_count.load(Ordering::SeqCst);
assert!(builds_after_first_read >= 1);
filter_cell.set(in_filter(&["server-A", "server-B"]));
insert_client(&ctx, "c2", "server-B");
let _ = myko::hyphae::Gettable::get(&downstream);
assert_eq!(
result.snapshot().len(),
2,
"result must reflect both the filter change and the subsequent write"
);
}
#[test]
fn query_live_transitions_between_indexed_and_scan_mode() {
let _serial = scheduler_test_serial();
let ctx = make_ctx();
let client_a = Client {
id: "c1".into(),
server_id: ServerId::from(Arc::<str>::from("server-A")),
address: None,
windback: Some(Arc::from("2026-01-01T00:00:00Z")),
};
let client_b = Client {
id: "c2".into(),
server_id: ServerId::from(Arc::<str>::from("server-B")),
address: None,
windback: None,
};
assert!(
ctx.apply_event_batch(vec![
MEvent::from_item(&client_a, MEventType::SET, "tx-1"),
MEvent::from_item(&client_b, MEventType::SET, "tx-2"),
])
.is_ok()
);
let filter_cell: Cell<ClientQuery, CellMutable> = Cell::new(eq_filter("server-A"));
let result = ctx.query_live(filter_cell.clone());
assert_eq!(result.snapshot().len(), 1);
filter_cell.set(ClientQuery {
windback: Some(myko::query::StringFilter::Eq(Arc::from(
"2026-01-01T00:00:00Z",
))),
..Default::default()
});
assert_eq!(
result.snapshot().len(),
1,
"switching to scan mode must re-evaluate correctly (client_a matches via windback)"
);
filter_cell.set(eq_filter("server-B"));
assert_eq!(
result.snapshot().len(),
1,
"switching back to indexed mode must re-evaluate correctly (client_b matches server-B)"
);
insert_client(&ctx, "c3", "server-B");
assert_eq!(result.snapshot().len(), 2);
}
fn id_in_filter(ids: &[&str]) -> ClientQuery {
use myko::entities::client::ClientId;
ClientQuery {
id: Some(IdFilter::In(
ids.iter()
.map(|s| ClientId::from(Arc::<str>::from(*s)))
.collect(),
)),
..Default::default()
}
}
#[test]
fn query_live_id_filter_routes_and_tracks_disjoint_swaps() {
let _serial = scheduler_test_serial();
let ctx = make_ctx();
insert_client(&ctx, "c1", "server-A");
insert_client(&ctx, "c2", "server-B");
let filter_cell: Cell<ClientQuery, CellMutable> = Cell::new(id_in_filter(&["c1", "c3"]));
let result = ctx.query_live(filter_cell.clone());
assert_eq!(result.snapshot().len(), 1);
insert_client(&ctx, "c3", "server-C");
assert_eq!(result.snapshot().len(), 2);
filter_cell.set(id_in_filter(&["c2"]));
let snap = result.snapshot();
assert_eq!(snap.len(), 1);
assert_eq!(snap.first().map(|entry| entry.0.as_ref()), Some("c2"));
}
#[test]
fn query_live_switches_between_id_and_belongs_to_modes() {
let _serial = scheduler_test_serial();
let ctx = make_ctx();
insert_client(&ctx, "c1", "server-A");
insert_client(&ctx, "c2", "server-A");
insert_client(&ctx, "c3", "server-B");
let filter_cell: Cell<ClientQuery, CellMutable> = Cell::new(id_in_filter(&["c3"]));
let result = ctx.query_live(filter_cell.clone());
assert_eq!(result.snapshot().len(), 1);
filter_cell.set(eq_filter("server-A"));
assert_eq!(result.snapshot().len(), 2);
filter_cell.set(id_in_filter(&["c1"]));
let snap = result.snapshot();
assert_eq!(snap.len(), 1);
assert_eq!(snap.first().map(|entry| entry.0.as_ref()), Some("c1"));
}
#[test]
fn query_live_id_route_narrows_by_secondary_fields() {
let _serial = scheduler_test_serial();
let ctx = make_ctx();
insert_client(&ctx, "c1", "server-A");
insert_client(&ctx, "c2", "server-B");
let mut filter = id_in_filter(&["c1", "c2"]);
filter.server_id = Some(IdFilter::Eq(ServerId::from(Arc::<str>::from("server-B"))));
let filter_cell: Cell<ClientQuery, CellMutable> = Cell::new(filter);
let result = ctx.query_live(filter_cell);
let snap = result.snapshot();
assert_eq!(snap.len(), 1);
assert_eq!(snap.first().map(|entry| entry.0.as_ref()), Some("c2"));
}