1use std::{
9 any::Any,
10 collections::HashMap,
11 sync::{Arc, Mutex},
12 thread,
13 time::Duration,
14};
15
16use dashmap::DashMap;
17use hyphae::{
18 Cell, CellImmutable, CellMap, CellMutable, Gettable, IdFor, MaterializeDefinite, Mutable,
19 WeakCellMap,
20};
21use serde::de::DeserializeOwned;
22use uuid::Uuid;
23
24use super::{
25 HandlerRegistry, RelationshipManager,
26 persister::{PersistError, PersistHealth, PersisterRouter},
27};
28use crate::{
29 cache::CacheKey,
30 client::{ConnectionStatus, MykoClient},
31 common::{
32 to_value::ToValue,
33 with_id::{WithId, WithTypedId},
34 },
35 core::item::{
36 AnyItem, Eventable, IngestBufferPolicy, downcast_any_item_arc, typed_map_arc_from_any_item,
37 typed_map_from_any_item_with_typed_id,
38 },
39 query::{
40 FilteredCellMap, QueryContext, QueryFactory, QueryHandler, QueryParams, QueryRequest,
41 QueryTestCtx,
42 },
43 report::{ReportContext, ReportHandler, ReportId},
44 request::RequestContext,
45 search::SearchIndex,
46 store::StoreRegistry,
47 view::{FilteredViewCellMap, TypedViewCellMap, ViewFactory},
48 wire::{EventOptions, MEvent, MEventType},
49};
50
51type AnyItemArc = Arc<dyn AnyItem>;
52
53#[derive(Debug, Clone, Copy, PartialEq, Eq)]
61pub(crate) enum Origin {
62 Local,
65 Cascade,
67 #[allow(dead_code)]
75 Remote,
76}
77
78impl Origin {
79 pub(crate) fn from_options(options: &EventOptions) -> Origin {
82 if options.prevent_relationship_updates {
83 Origin::Cascade
84 } else {
85 Origin::Local
86 }
87 }
88
89 fn should_cascade(self, change: MEventType) -> bool {
105 match self {
106 Origin::Local => true,
107 Origin::Cascade => change == MEventType::DEL,
108 Origin::Remote => false,
109 }
110 }
111
112 fn should_produce(self) -> bool {
119 self != Origin::Remote
120 }
121}
122
123trait ReportCacheEntryDyn: Any + Send + Sync {
127 fn as_any(&self) -> &dyn Any;
128 fn is_alive(&self) -> bool;
129}
130
131struct ReportCacheEntry<T> {
132 weak: hyphae::cell::WeakCell<T, CellImmutable>,
133}
134
135impl<T> ReportCacheEntry<T>
136where
137 T: Clone + Send + Sync + 'static,
138{
139 fn new(cell: &Cell<T, CellImmutable>) -> Self {
140 Self {
141 weak: cell.downgrade(),
142 }
143 }
144
145 fn get(&self) -> Option<Cell<T, CellImmutable>> {
146 self.weak.upgrade()
147 }
148}
149
150impl<T> ReportCacheEntryDyn for ReportCacheEntry<T>
151where
152 T: Clone + Send + Sync + 'static,
153{
154 fn as_any(&self) -> &dyn Any {
155 self
156 }
157
158 fn is_alive(&self) -> bool {
159 self.weak.upgrade().is_some()
160 }
161}
162
163struct MapCacheEntry {
164 weak: WeakCellMap<Arc<str>, AnyItemArc>,
165 typed: Mutex<HashMap<std::any::TypeId, Box<dyn Any + Send + Sync>>>,
169}
170
171#[derive(Default)]
172struct BufferedIngestState {
173 events: Vec<MEvent>,
174 flush_scheduled: bool,
175}
176
177struct BufferedIngestType {
178 state: Mutex<BufferedIngestState>,
179}
180
181impl BufferedIngestType {
182 fn new() -> Self {
183 Self {
184 state: Mutex::new(BufferedIngestState::default()),
185 }
186 }
187}
188
189impl MapCacheEntry {
190 fn new(map: &FilteredCellMap) -> Self {
191 Self {
192 weak: map.downgrade(),
193 typed: Mutex::new(HashMap::new()),
194 }
195 }
196
197 fn get(&self) -> Option<FilteredCellMap> {
198 self.weak.upgrade().map(|map| map.lock())
199 }
200
201 fn get_or_create_typed<K, V, F>(&self, create: F) -> Option<CellMap<K, V, CellImmutable>>
206 where
207 K: std::hash::Hash + Eq + hyphae::traits::CellValue + 'static,
208 V: hyphae::traits::CellValue + 'static,
209 F: FnOnce(FilteredCellMap) -> CellMap<K, V, CellImmutable>,
210 {
211 let type_key = std::any::TypeId::of::<WeakCellMap<K, V>>();
212 let mut typed = self.typed.lock().unwrap();
213
214 if let Some(entry) = typed.get(&type_key) {
216 if let Some(weak) = entry.downcast_ref::<WeakCellMap<K, V>>()
217 && let Some(strong) = weak.upgrade()
218 {
219 return Some(strong.lock());
220 }
221 typed.remove(&type_key);
223 }
224
225 let source = self.weak.upgrade()?.lock();
227 let built = create(source);
228 typed.insert(type_key, Box::new(built.downgrade()));
229 Some(built)
230 }
231}
232
233#[derive(Clone)]
240pub struct CellServerCtx {
241 pub host_id: Uuid,
243 pub registry: Arc<StoreRegistry>,
245 pub handler_registry: Arc<HandlerRegistry>,
247 relationship_manager: Arc<RelationshipManager>,
249 persisters: Arc<PersisterRouter>,
251 search_index: Arc<SearchIndex>,
253 peer_clients: Arc<DashMap<Arc<str>, Arc<MykoClient>>>,
255 peer_clients_tick: Cell<u64, CellMutable>,
257 event_sink: Option<flume::Sender<MEvent>>,
259 query_cache: Arc<DashMap<String, MapCacheEntry, ahash::RandomState>>,
264 view_cache: Arc<DashMap<String, MapCacheEntry, ahash::RandomState>>,
266 report_cache: Arc<DashMap<String, Arc<dyn ReportCacheEntryDyn>, ahash::RandomState>>,
268 compute_gates: Arc<DashMap<String, Arc<std::sync::Mutex<()>>, ahash::RandomState>>,
271 ingest_buffers: Arc<DashMap<Arc<str>, Arc<BufferedIngestType>, ahash::RandomState>>,
273 history_replay: Option<Arc<dyn crate::server::HistoryReplayProvider>>,
275}
276
277impl CellServerCtx {
278 #[allow(clippy::too_many_arguments)]
280 pub fn new(
281 host_id: Uuid,
282 registry: Arc<StoreRegistry>,
283 handler_registry: Arc<HandlerRegistry>,
284 relationship_manager: Arc<RelationshipManager>,
285 persisters: Arc<PersisterRouter>,
286 search_index: Arc<SearchIndex>,
287 peer_clients: Arc<DashMap<Arc<str>, Arc<MykoClient>>>,
288 event_sink: Option<flume::Sender<MEvent>>,
289 history_replay: Option<Arc<dyn crate::server::HistoryReplayProvider>>,
290 ) -> Self {
291 Self {
292 host_id,
293 registry,
294 handler_registry,
295 relationship_manager,
296 persisters,
297 search_index,
298 peer_clients,
299 peer_clients_tick: Cell::new(0).with_name("peer_clients_tick"),
300 event_sink,
301 query_cache: Arc::new(DashMap::with_hasher(ahash::RandomState::new())),
302 view_cache: Arc::new(DashMap::with_hasher(ahash::RandomState::new())),
303 report_cache: Arc::new(DashMap::with_hasher(ahash::RandomState::new())),
304 compute_gates: Arc::new(DashMap::with_hasher(ahash::RandomState::new())),
305 ingest_buffers: Arc::new(DashMap::with_hasher(ahash::RandomState::new())),
306 history_replay,
307 }
308 }
309
310 fn cache_key<T: CacheKey>(
311 &self,
312 kind: &str,
313 id: &str,
314 params: &T,
315 request: &RequestContext,
316 ) -> String {
317 let payload_hash = params.cache_key_hash();
318 format!("{}:{kind}:{id}:{payload_hash:016x}", request.host_id)
319 }
320
321 pub fn search_index(&self) -> &Arc<SearchIndex> {
323 &self.search_index
324 }
325
326 pub fn history_replay(&self) -> Option<&Arc<dyn crate::server::HistoryReplayProvider>> {
328 self.history_replay.as_ref()
329 }
330
331 pub fn register_peer_client<S: AsRef<str>>(&self, peer_id: S, client: Arc<MykoClient>) {
333 self.peer_clients
334 .insert(Arc::<str>::from(peer_id.as_ref()), client);
335 let next = self.peer_clients_tick.get().saturating_add(1);
336 self.peer_clients_tick.set(next);
337 }
338
339 pub fn unregister_peer_client(&self, peer_id: &str) {
341 if self.peer_clients.remove(peer_id).is_some() {
342 let next = self.peer_clients_tick.get().saturating_add(1);
343 self.peer_clients_tick.set(next);
344 }
345 }
346
347 pub fn peer_client(&self, peer_id: &str) -> Option<Arc<MykoClient>> {
349 self.peer_clients
350 .get(peer_id)
351 .map(|entry| entry.value().clone())
352 }
353
354 pub fn peer_connection_status(&self, peer_id: &str) -> Option<ConnectionStatus> {
356 self.peer_client(peer_id)
357 .map(|client| client.get_connection_status_sync())
358 }
359
360 pub fn peer_clients_tick(&self) -> Cell<u64, CellImmutable> {
362 self.peer_clients_tick.clone().lock()
363 }
364
365 pub fn peer_client_count(&self) -> usize {
367 self.peer_clients.len()
368 }
369
370 pub fn persist_health(&self) -> Arc<PersistHealth> {
372 self.persisters.default_health()
373 }
374
375 pub fn query_cache_len(&self) -> usize {
377 self.query_cache.len()
378 }
379
380 pub fn view_cache_len(&self) -> usize {
382 self.view_cache.len()
383 }
384
385 pub fn report_cache_len(&self) -> usize {
387 self.report_cache.len()
388 }
389
390 pub fn report_cache_live_count(&self) -> usize {
392 self.report_cache
393 .iter()
394 .filter(|entry| entry.value().is_alive())
395 .count()
396 }
397
398 pub fn query_cache_live_count(&self) -> usize {
400 self.query_cache
401 .iter()
402 .filter(|entry| entry.value().weak.upgrade().is_some())
403 .count()
404 }
405
406 pub fn view_cache_live_count(&self) -> usize {
408 self.view_cache
409 .iter()
410 .filter(|entry| entry.value().weak.upgrade().is_some())
411 .count()
412 }
413
414 pub fn sweep_dead_cache_entries(&self) {
421 self.query_cache
422 .retain(|_, entry| entry.weak.upgrade().is_some());
423 self.view_cache
424 .retain(|_, entry| entry.weak.upgrade().is_some());
425 self.report_cache.retain(|_, entry| entry.is_alive());
426 crate::query::sweep_all_belongs_to_source_indexes();
427 }
428
429 pub fn parse_item(
439 &self,
440 entity_type: &str,
441 json: serde_json::Value,
442 ) -> Option<Arc<dyn AnyItem>> {
443 let parse = self.handler_registry.get_item_parser(entity_type)?;
444 parse(json).ok()
445 }
446
447 pub fn set<T>(&self, entity: &T) -> Result<(), PersistError>
455 where
456 T: Eventable + 'static,
457 {
458 self.set_with_origin(entity, Origin::Local)
459 }
460
461 #[deprecated(note = "EventOptions is internal plumbing; use `set` instead")]
467 pub fn set_with_options<T>(
468 &self,
469 entity: &T,
470 options: Option<EventOptions>,
471 ) -> Result<(), PersistError>
472 where
473 T: Eventable + 'static,
474 {
475 self.set_with_origin(entity, Origin::from_options(&options.unwrap_or_default()))
476 }
477
478 pub(crate) fn set_with_origin<T>(&self, entity: &T, origin: Origin) -> Result<(), PersistError>
481 where
482 T: Eventable + 'static,
483 {
484 let item: Arc<dyn AnyItem> = Arc::new(entity.clone());
485 self.reduce_one(&item, MEventType::SET);
486 self.apply_effects(std::slice::from_ref(&item), MEventType::SET, origin)
487 }
488
489 pub fn del<T>(&self, entity: &T) -> Result<(), PersistError>
493 where
494 T: Eventable + Clone + 'static,
495 {
496 self.del_with_origin(entity, Origin::Local)
497 }
498
499 #[deprecated(note = "EventOptions is internal plumbing; use `del` instead")]
503 pub fn del_with_options<T>(
504 &self,
505 entity: &T,
506 options: Option<EventOptions>,
507 ) -> Result<(), PersistError>
508 where
509 T: Eventable + Clone + 'static,
510 {
511 self.del_with_origin(entity, Origin::from_options(&options.unwrap_or_default()))
512 }
513
514 pub(crate) fn del_with_origin<T>(&self, entity: &T, origin: Origin) -> Result<(), PersistError>
515 where
516 T: Eventable + Clone + 'static,
517 {
518 let item: Arc<dyn AnyItem> = Arc::new(entity.clone());
519 self.reduce_one(&item, MEventType::DEL);
520 self.apply_effects(std::slice::from_ref(&item), MEventType::DEL, origin)
521 }
522
523 pub fn batch_set<T>(&self, entities: &[T]) -> Result<(), PersistError>
527 where
528 T: Eventable + Clone + 'static,
529 {
530 self.batch_set_with_origin(entities, Origin::Local)
531 }
532
533 #[deprecated(note = "EventOptions is internal plumbing; use `batch_set` instead")]
537 pub fn batch_set_with_options<T>(
538 &self,
539 entities: &[T],
540 options: Option<EventOptions>,
541 ) -> Result<(), PersistError>
542 where
543 T: Eventable + Clone + 'static,
544 {
545 self.batch_set_with_origin(entities, Origin::from_options(&options.unwrap_or_default()))
546 }
547
548 pub(crate) fn batch_set_with_origin<T>(
550 &self,
551 entities: &[T],
552 origin: Origin,
553 ) -> Result<(), PersistError>
554 where
555 T: Eventable + Clone + 'static,
556 {
557 if entities.is_empty() {
558 return Ok(());
559 }
560 let items: Vec<Arc<dyn AnyItem>> = entities
561 .iter()
562 .map(|e| Arc::new(e.clone()) as Arc<dyn AnyItem>)
563 .collect();
564 self.emit_grouped(&items, MEventType::SET, origin)
565 }
566
567 pub fn batch_del<T>(&self, entities: &[T]) -> Result<(), PersistError>
571 where
572 T: Eventable + Clone + 'static,
573 {
574 self.batch_del_with_origin(entities, Origin::Local)
575 }
576
577 #[deprecated(note = "EventOptions is internal plumbing; use `batch_del` instead")]
581 pub fn batch_del_with_options<T>(
582 &self,
583 entities: &[T],
584 options: Option<EventOptions>,
585 ) -> Result<(), PersistError>
586 where
587 T: Eventable + Clone + 'static,
588 {
589 self.batch_del_with_origin(entities, Origin::from_options(&options.unwrap_or_default()))
590 }
591
592 pub(crate) fn batch_del_with_origin<T>(
594 &self,
595 entities: &[T],
596 origin: Origin,
597 ) -> Result<(), PersistError>
598 where
599 T: Eventable + Clone + 'static,
600 {
601 if entities.is_empty() {
602 return Ok(());
603 }
604 let items: Vec<Arc<dyn AnyItem>> = entities
605 .iter()
606 .map(|e| Arc::new(e.clone()) as Arc<dyn AnyItem>)
607 .collect();
608 self.emit_grouped(&items, MEventType::DEL, origin)
609 }
610
611 pub fn set_dyn(&self, item: Arc<dyn AnyItem>) -> Result<(), PersistError> {
619 self.set_dyn_with_origin(item, Origin::Local)
620 }
621
622 #[deprecated(note = "EventOptions is internal plumbing; use `set_dyn` instead")]
626 pub fn set_dyn_with_options(
627 &self,
628 item: Arc<dyn AnyItem>,
629 options: Option<EventOptions>,
630 ) -> Result<(), PersistError> {
631 self.set_dyn_with_origin(item, Origin::from_options(&options.unwrap_or_default()))
632 }
633
634 pub(crate) fn set_dyn_with_origin(
635 &self,
636 item: Arc<dyn AnyItem>,
637 origin: Origin,
638 ) -> Result<(), PersistError> {
639 self.reduce_one(&item, MEventType::SET);
640 self.apply_effects(std::slice::from_ref(&item), MEventType::SET, origin)
641 }
642
643 pub fn batch_set_dyn(&self, items: &[Arc<dyn AnyItem>]) -> Result<(), PersistError> {
645 self.batch_set_dyn_with_origin(items, Origin::Local)
646 }
647
648 #[deprecated(note = "EventOptions is internal plumbing; use `batch_set_dyn` instead")]
652 pub fn batch_set_dyn_with_options(
653 &self,
654 items: &[Arc<dyn AnyItem>],
655 options: Option<EventOptions>,
656 ) -> Result<(), PersistError> {
657 self.batch_set_dyn_with_origin(items, Origin::from_options(&options.unwrap_or_default()))
658 }
659
660 pub(crate) fn batch_set_dyn_with_origin(
661 &self,
662 items: &[Arc<dyn AnyItem>],
663 origin: Origin,
664 ) -> Result<(), PersistError> {
665 self.emit_grouped(items, MEventType::SET, origin)
666 }
667
668 pub fn del_dyn(&self, item: Arc<dyn AnyItem>) -> Result<(), PersistError> {
672 self.del_dyn_with_origin(item, Origin::Local)
673 }
674
675 #[deprecated(note = "EventOptions is internal plumbing; use `del_dyn` instead")]
679 pub fn del_dyn_with_options(
680 &self,
681 item: Arc<dyn AnyItem>,
682 options: Option<EventOptions>,
683 ) -> Result<(), PersistError> {
684 self.del_dyn_with_origin(item, Origin::from_options(&options.unwrap_or_default()))
685 }
686
687 pub(crate) fn del_dyn_with_origin(
688 &self,
689 item: Arc<dyn AnyItem>,
690 origin: Origin,
691 ) -> Result<(), PersistError> {
692 self.reduce_one(&item, MEventType::DEL);
693 self.apply_effects(std::slice::from_ref(&item), MEventType::DEL, origin)
694 }
695
696 pub fn batch_del_dyn(&self, items: &[Arc<dyn AnyItem>]) -> Result<(), PersistError> {
698 self.batch_del_dyn_with_origin(items, Origin::Local)
699 }
700
701 #[deprecated(note = "EventOptions is internal plumbing; use `batch_del_dyn` instead")]
705 pub fn batch_del_dyn_with_options(
706 &self,
707 items: &[Arc<dyn AnyItem>],
708 options: Option<EventOptions>,
709 ) -> Result<(), PersistError> {
710 self.batch_del_dyn_with_origin(items, Origin::from_options(&options.unwrap_or_default()))
711 }
712
713 pub(crate) fn batch_del_dyn_with_origin(
714 &self,
715 items: &[Arc<dyn AnyItem>],
716 origin: Origin,
717 ) -> Result<(), PersistError> {
718 self.emit_grouped(items, MEventType::DEL, origin)
719 }
720
721 pub fn del_by_id(&self, entity_type: &str, id: &str) -> Result<(), PersistError> {
728 self.del_by_id_with_origin(entity_type, id, Origin::Local)
729 }
730
731 #[deprecated(note = "EventOptions is internal plumbing; use `del_by_id` instead")]
735 pub fn del_by_id_with_options(
736 &self,
737 entity_type: &str,
738 id: &str,
739 options: Option<EventOptions>,
740 ) -> Result<(), PersistError> {
741 self.del_by_id_with_origin(
742 entity_type,
743 id,
744 Origin::from_options(&options.unwrap_or_default()),
745 )
746 }
747
748 pub(crate) fn del_by_id_with_origin(
749 &self,
750 entity_type: &str,
751 id: &str,
752 origin: Origin,
753 ) -> Result<(), PersistError> {
754 let id_arc: Arc<str> = id.into();
755
756 let existing = self
757 .registry
758 .get(entity_type)
759 .and_then(|store| store.get(&id_arc).get());
760
761 crate::server::entity_set_stats::record_del(entity_type);
762
763 self.registry.get_or_create(entity_type).remove(&id_arc);
765
766 self.search_index.remove_entity(entity_type, id);
768
769 if origin.should_produce() {
771 if let Some(item) = existing {
772 self.produce_del_dyn(&item)?;
773 } else {
774 tracing::warn!(
775 "del_by_id could not persist DEL without full entity: {}:{}",
776 entity_type,
777 id
778 );
779 }
780 }
781
782 tracing::trace!("Published DEL {}:{}", entity_type, id);
783 Ok(())
784 }
785
786 pub fn apply_event(&self, event: MEvent) -> Result<bool, PersistError> {
790 Ok(self.apply_event_batch(vec![event])? == 1)
791 }
792
793 pub fn apply_event_batch(&self, events: Vec<MEvent>) -> Result<usize, PersistError> {
798 if events.is_empty() {
799 return Ok(0);
800 }
801
802 let mut accepted = 0usize;
803 let mut immediate_events = Vec::new();
804 let mut buffered_by_type: HashMap<Arc<str>, (u64, Vec<MEvent>)> = HashMap::new();
805
806 for event in events {
807 match self
808 .handler_registry
809 .get_item_buffer_policy(&event.item_type)
810 {
811 IngestBufferPolicy::None => immediate_events.push(event),
812 IngestBufferPolicy::TimeWindow { window_ms } => {
813 let entity_type: Arc<str> = event.item_type.clone().into();
814 buffered_by_type
815 .entry(entity_type)
816 .or_insert_with(|| (window_ms, Vec::new()))
817 .1
818 .push(event);
819 }
820 }
821 }
822
823 if !immediate_events.is_empty() {
824 accepted += self.apply_event_batch_immediate(immediate_events)?;
825 }
826
827 for (entity_type, (window_ms, buffered_events)) in buffered_by_type {
828 accepted += buffered_events.len();
829 self.enqueue_buffered_events(entity_type, window_ms, buffered_events);
830 }
831
832 Ok(accepted)
833 }
834
835 fn apply_event_batch_immediate(&self, events: Vec<MEvent>) -> Result<usize, PersistError> {
836 if events.is_empty() {
837 return Ok(0);
838 }
839 let input_len = events.len();
840
841 let mut set_items: Vec<Arc<dyn AnyItem>> = Vec::new();
842 let mut del_items: Vec<Arc<dyn AnyItem>> = Vec::new();
843
844 for event in events {
845 let change = event.change_type;
846 let item_type = event.item_type;
847 let item_value = event.item;
848 let Some(item) = self.parse_item(&item_type, item_value) else {
849 tracing::warn!("Unknown entity type or parse error for ingest: {item_type}");
850 continue;
851 };
852 match change {
853 MEventType::SET => set_items.push(item),
854 MEventType::DEL => del_items.push(item),
855 }
856 }
857
858 let applied = set_items.len() + del_items.len();
859 if applied == 0 {
860 return Ok(0);
861 }
862
863 tracing::trace!(
864 target: "myko::server::context",
865 "apply_event_batch parsed: input_events={} sets={} dels={}",
866 input_len,
867 set_items.len(),
868 del_items.len()
869 );
870
871 let emit = || -> Result<(), PersistError> {
876 self.emit_grouped(&set_items, MEventType::SET, Origin::Local)?;
877 self.emit_grouped(&del_items, MEventType::DEL, Origin::Local)?;
878 Ok(())
879 };
880 emit()?;
881
882 Ok(applied)
883 }
884
885 fn ingest_buffer_for(&self, entity_type: Arc<str>) -> Arc<BufferedIngestType> {
886 self.ingest_buffers
887 .entry(entity_type)
888 .or_insert_with(|| Arc::new(BufferedIngestType::new()))
889 .clone()
890 }
891
892 fn enqueue_buffered_events(&self, entity_type: Arc<str>, window_ms: u64, events: Vec<MEvent>) {
893 let buffer = self.ingest_buffer_for(entity_type.clone());
894 let should_schedule = {
895 let Ok(mut state) = buffer.state.lock() else {
896 tracing::error!(
897 "Could not acquire ingest buffer lock for entity_type={}",
898 entity_type
899 );
900 if let Err(e) = self.apply_event_batch_immediate(events) {
901 tracing::error!("Failed to apply buffered events for {}: {}", entity_type, e);
902 }
903 return;
904 };
905
906 state.events.extend(events);
907 if state.flush_scheduled {
908 false
909 } else {
910 state.flush_scheduled = true;
911 true
912 }
913 };
914
915 if !should_schedule {
916 return;
917 }
918
919 let ctx = self.clone();
920 thread::spawn(move || {
921 thread::sleep(Duration::from_millis(window_ms));
922 ctx.flush_buffered_events_for_type(&entity_type);
923 });
924 }
925
926 fn flush_buffered_events_for_type(&self, entity_type: &Arc<str>) -> usize {
927 let Some(buffer) = self
928 .ingest_buffers
929 .get(entity_type.as_ref())
930 .map(|entry| entry.clone())
931 else {
932 return 0;
933 };
934
935 let events = {
936 let Ok(mut state) = buffer.state.lock() else {
937 tracing::error!(
938 "Could not acquire ingest buffer lock for flush entity_type={}",
939 entity_type
940 );
941 return 0;
942 };
943
944 state.flush_scheduled = false;
945 if state.events.is_empty() {
946 return 0;
947 }
948
949 std::mem::take(&mut state.events)
950 };
951
952 tracing::trace!(
953 target: "myko::server::context",
954 "flush_buffered_events entity_type={} count={}",
955 entity_type,
956 events.len()
957 );
958
959 match self.apply_event_batch_immediate(events) {
960 Ok(count) => count,
961 Err(e) => {
962 tracing::error!("Failed to flush buffered events for {}: {}", entity_type, e);
963 0
964 }
965 }
966 }
967
968 #[cfg(test)]
969 fn flush_all_buffered_events(&self) -> usize {
970 let entity_types: Vec<Arc<str>> = self
971 .ingest_buffers
972 .iter()
973 .map(|entry| entry.key().clone())
974 .collect();
975
976 entity_types
977 .into_iter()
978 .map(|entity_type| self.flush_buffered_events_for_type(&entity_type))
979 .sum()
980 }
981
982 fn reduce_one(&self, item: &Arc<dyn AnyItem>, change: MEventType) {
991 let entity_type = item.entity_type();
992 let _span = tracing::trace_span!("myko.reduce", ty = entity_type, op = ?change).entered();
1000 match change {
1001 MEventType::SET => {
1002 crate::server::entity_set_stats::record_set(entity_type);
1003 self.registry
1004 .get_or_create(entity_type)
1005 .insert(item.id(), item.clone());
1006 }
1007 MEventType::DEL => {
1008 crate::server::entity_set_stats::record_del(entity_type);
1009 self.registry.get_or_create(entity_type).remove(&item.id());
1010 }
1011 }
1012 }
1013
1014 fn emit_grouped(
1023 &self,
1024 items: &[Arc<dyn AnyItem>],
1025 change: MEventType,
1026 origin: Origin,
1027 ) -> Result<(), PersistError> {
1028 if items.is_empty() {
1029 return Ok(());
1030 }
1031
1032 let mut by_type: std::collections::BTreeMap<&'static str, Vec<Arc<dyn AnyItem>>> =
1033 std::collections::BTreeMap::new();
1034 for item in items {
1035 by_type
1036 .entry(item.entity_type())
1037 .or_default()
1038 .push(item.clone());
1039 }
1040
1041 hyphae::batch(|| {
1075 for (entity_type, group) in &by_type {
1076 let store = self.registry.get_or_create(entity_type);
1077 match change {
1078 MEventType::SET => {
1079 let mut entries: Vec<(Arc<str>, Arc<dyn AnyItem>)> =
1080 Vec::with_capacity(group.len());
1081 for item in group {
1082 crate::server::entity_set_stats::record_set(entity_type);
1083 entries.push((item.id(), item.clone()));
1084 }
1085 store.insert_many(entries);
1086 }
1087 MEventType::DEL => {
1088 let mut ids: Vec<Arc<str>> = Vec::with_capacity(group.len());
1089 for item in group {
1090 crate::server::entity_set_stats::record_del(entity_type);
1091 ids.push(item.id());
1092 }
1093 store.remove_many(ids);
1094 }
1095 }
1096 }
1097 });
1098
1099 for group in by_type.values() {
1101 self.apply_effects(group, change, origin)?;
1102 }
1103 Ok(())
1104 }
1105
1106 fn apply_effects(
1115 &self,
1116 items: &[Arc<dyn AnyItem>],
1117 change: MEventType,
1118 origin: Origin,
1119 ) -> Result<(), PersistError> {
1120 let _span = tracing::trace_span!(
1125 "myko.apply_effects",
1126 ty = items.first().map(|i| i.entity_type()).unwrap_or("empty"),
1127 op = ?change,
1128 )
1129 .entered();
1130
1131 match change {
1133 MEventType::SET => {
1134 for item in items {
1135 self.search_index.index_item(item);
1136 }
1137 }
1138 MEventType::DEL => {
1139 for item in items {
1140 self.search_index
1141 .remove_entity(item.entity_type(), &item.id());
1142 }
1143 }
1144 }
1145
1146 if origin.should_cascade(change) {
1148 match change {
1149 MEventType::SET => {
1150 for item in items {
1151 self.relationship_manager.forward_set(item.clone(), self)?;
1152 }
1153 }
1154 MEventType::DEL => self.relationship_manager.forward_del_batch(items, self)?,
1155 }
1156 }
1157
1158 if origin.should_produce() {
1160 match change {
1161 MEventType::SET => {
1162 for item in items {
1163 self.produce_set_dyn(item)?;
1164 }
1165 }
1166 MEventType::DEL => {
1167 for item in items {
1168 self.produce_del_dyn(item)?;
1169 }
1170 }
1171 }
1172 }
1173
1174 Ok(())
1175 }
1176
1177 fn produce_del_dyn(&self, item: &Arc<dyn AnyItem>) -> Result<(), PersistError> {
1182 if let Some(persister) = self.persisters.resolve(item.entity_type()) {
1183 let event = MEvent::del_from_any(item, &self.host_id.to_string());
1184 persister.persist(event)?;
1185 }
1186 if let Some(sink) = &self.event_sink {
1187 let event = MEvent::del_from_any(item, &self.host_id.to_string());
1188 let _ = sink.send(event);
1189 }
1190 Ok(())
1191 }
1192
1193 fn produce_set_dyn(&self, item: &Arc<dyn AnyItem>) -> Result<(), PersistError> {
1194 if let Some(persister) = self.persisters.resolve(item.entity_type()) {
1195 let event = MEvent::set_from_value(
1196 item.entity_type(),
1197 item.to_value(),
1198 &self.host_id.to_string(),
1199 );
1200 persister.persist(event)?;
1201 }
1202 if let Some(sink) = &self.event_sink {
1203 let event = MEvent::set_from_value(
1204 item.entity_type(),
1205 item.to_value(),
1206 &self.host_id.to_string(),
1207 );
1208 let _ = sink.send(event);
1209 }
1210 Ok(())
1211 }
1212
1213 pub fn query_map<Q>(
1222 &self,
1223 query: Q,
1224 request: Arc<RequestContext>,
1225 ) -> CellMap<<Q::Item as WithTypedId>::Id, Arc<Q::Item>, CellImmutable>
1226 where
1227 Q: QueryParams + 'static,
1228 Q::Item: Eventable
1229 + WithId
1230 + WithTypedId
1231 + DeserializeOwned
1232 + Clone
1233 + std::fmt::Debug
1234 + Send
1235 + Sync
1236 + 'static,
1237 {
1238 let key = self.cache_key("query", Q::query_id_static().as_ref(), &query, &request);
1239 let untyped = self.query_map_untyped(query, request);
1241 if let Some(entry) = self.query_cache.get(&key)
1242 && let Some(typed) = entry.value().get_or_create_typed(|source| {
1243 typed_map_from_any_item_with_typed_id(source, "CellServerCtx::query_map")
1244 })
1245 {
1246 return typed;
1247 }
1248 self.query_cache
1250 .insert(key.clone(), MapCacheEntry::new(&untyped));
1251 let entry = self.query_cache.get(&key).expect("just re-inserted");
1252 entry
1253 .value()
1254 .get_or_create_typed(|source| {
1255 typed_map_from_any_item_with_typed_id(source, "CellServerCtx::query_map")
1256 })
1257 .expect("typed projection from freshly inserted entry")
1258 }
1259
1260 pub fn query_map_by_str<Q>(
1264 &self,
1265 query: Q,
1266 request: Arc<RequestContext>,
1267 ) -> CellMap<Arc<str>, Arc<Q::Item>, CellImmutable>
1268 where
1269 Q: QueryParams + 'static,
1270 Q::Item:
1271 Eventable + WithId + DeserializeOwned + Clone + std::fmt::Debug + Send + Sync + 'static,
1272 {
1273 let key = self.cache_key("query", Q::query_id_static().as_ref(), &query, &request);
1274 let untyped = self.query_map_untyped(query, request);
1275 if let Some(entry) = self.query_cache.get(&key)
1276 && let Some(typed) = entry.value().get_or_create_typed(|source| {
1277 typed_map_arc_from_any_item(source, "CellServerCtx::query_map_by_str")
1278 })
1279 {
1280 return typed;
1281 }
1282 self.query_cache
1284 .insert(key.clone(), MapCacheEntry::new(&untyped));
1285 let entry = self.query_cache.get(&key).expect("just re-inserted");
1286 entry
1287 .value()
1288 .get_or_create_typed(|source| {
1289 typed_map_arc_from_any_item(source, "CellServerCtx::query_map_by_str")
1290 })
1291 .expect("typed projection from freshly inserted entry")
1292 }
1293
1294 pub fn query_map_untyped<Q>(&self, query: Q, request: Arc<RequestContext>) -> FilteredCellMap
1313 where
1314 Q: QueryFactory + QueryHandler + QueryParams + Clone + Send + Sync + 'static,
1315 Q::Item: DeserializeOwned + Clone + std::fmt::Debug + Send + Sync + 'static,
1316 {
1317 let key = self.cache_key("query", Q::query_id_static().as_ref(), &query, &request);
1318
1319 if let Some(cell) = self.try_get_cached_query(&key) {
1321 return cell;
1322 }
1323
1324 let gate = self
1325 .compute_gates
1326 .entry(key.clone())
1327 .or_insert_with(|| Arc::new(std::sync::Mutex::new(())))
1328 .clone();
1329 let _lock = gate.lock().unwrap();
1330
1331 if let Some(cell) = self.try_get_cached_query(&key) {
1333 return cell;
1334 }
1335
1336 let query_req = QueryRequest::with_tx(query, request.tx.clone());
1337 let any_query: Arc<dyn crate::query::AnyQuery> = Arc::new(query_req);
1338
1339 let built = Q::cell_factory(
1340 any_query,
1341 self.registry.clone(),
1342 request,
1343 Some(Arc::new(self.clone())),
1344 )
1345 .expect("query cell factory should not fail for typed query");
1346 self.query_cache
1347 .insert(key.clone(), MapCacheEntry::new(&built));
1348 self.compute_gates.remove(&key);
1357 built
1358 }
1359
1360 fn try_get_cached_query(&self, key: &str) -> Option<FilteredCellMap> {
1361 let existing = self.query_cache.get(key)?;
1362 if let Some(shared) = existing.value().get() {
1363 return Some(shared);
1364 }
1365 drop(existing);
1366 self.query_cache.remove(key);
1367 None
1368 }
1369
1370 pub fn view_map_untyped<V>(&self, view: V, request: Arc<RequestContext>) -> FilteredViewCellMap
1372 where
1373 V: ViewFactory + Clone + Send + Sync + 'static,
1374 V::Item: DeserializeOwned + Clone + std::fmt::Debug + Send + Sync + 'static,
1375 {
1376 let key = self.cache_key("view", V::view_id_static().as_ref(), &view, &request);
1377
1378 if let Some(cell) = self.try_get_cached_view(&key) {
1380 return cell;
1381 }
1382
1383 let gate = self
1384 .compute_gates
1385 .entry(key.clone())
1386 .or_insert_with(|| Arc::new(std::sync::Mutex::new(())))
1387 .clone();
1388 let _lock = gate.lock().unwrap();
1389
1390 if let Some(cell) = self.try_get_cached_view(&key) {
1392 return cell;
1393 }
1394
1395 let view_req = crate::view::ViewRequest::with_tx(view, request.tx.clone());
1396 let any_view: Arc<dyn crate::view::AnyView> = Arc::new(view_req);
1397
1398 let built = V::cell_factory(
1399 any_view,
1400 self.registry.clone(),
1401 request,
1402 Arc::new(self.clone()),
1403 )
1404 .expect("view cell factory should not fail for typed view");
1405 self.view_cache
1406 .insert(key.clone(), MapCacheEntry::new(&built));
1407 self.compute_gates.remove(&key);
1411 built
1412 }
1413
1414 fn try_get_cached_view(&self, key: &str) -> Option<FilteredViewCellMap> {
1415 let existing = self.view_cache.get(key)?;
1416 if let Some(shared) = existing.value().get() {
1417 return Some(shared);
1418 }
1419 drop(existing);
1420 self.view_cache.remove(key);
1421 None
1422 }
1423
1424 pub fn view_map<V>(&self, view: V, request: Arc<RequestContext>) -> FilteredViewCellMap
1426 where
1427 V: ViewFactory + Clone + Send + Sync + 'static,
1428 V::Item: DeserializeOwned + Clone + std::fmt::Debug + Send + Sync + 'static,
1429 {
1430 self.view_map_untyped(view, request)
1431 }
1432
1433 pub fn view<V>(&self, view: V, request: Arc<RequestContext>) -> TypedViewCellMap<V::Item>
1435 where
1436 V: ViewFactory + Clone + Send + Sync + 'static,
1437 V::Item: DeserializeOwned + Clone + std::fmt::Debug + Send + Sync + 'static,
1438 {
1439 let key = self.cache_key("view", V::view_id_static().as_ref(), &view, &request);
1440 let _untyped = self.view_map_untyped(view, request);
1441 if let Some(entry) = self.view_cache.get(&key)
1442 && let Some(typed) = entry.value().get_or_create_typed(|source| {
1443 typed_map_arc_from_any_item(source, "CellServerCtx::view")
1444 })
1445 {
1446 return typed;
1447 }
1448 unreachable!("view_map_untyped just populated the cache")
1449 }
1450
1451 pub fn entity_snapshot<T>(&self, id: &<T as WithTypedId>::Id) -> Option<Arc<T>>
1453 where
1454 T: Eventable + WithTypedId + Send + Sync + 'static,
1455 <T as WithTypedId>::Id: hyphae::IdFor<T, MapKey = Arc<str>>,
1456 {
1457 let store = self.registry.get_or_create(T::entity_name_static());
1458 let map_key = id.map_key();
1459 let item = store.get_value(&map_key)?;
1460 Some(downcast_any_item_arc::<T>(
1461 &item,
1462 "CellServerCtx::entity_snapshot",
1463 ))
1464 }
1465
1466 pub fn entity_snapshots<T>(&self) -> Vec<Arc<T>>
1468 where
1469 T: Eventable + WithTypedId + Send + Sync + 'static,
1470 <T as WithTypedId>::Id: hyphae::IdFor<T, MapKey = Arc<str>>,
1471 {
1472 let store = self.registry.get_or_create(T::entity_name_static());
1473 store
1474 .snapshot()
1475 .into_iter()
1476 .map(|(_, item)| downcast_any_item_arc::<T>(&item, "CellServerCtx::entity_snapshots"))
1477 .collect()
1478 }
1479
1480 pub fn entity_snapshots_by_id<T>(
1482 &self,
1483 ids: impl IntoIterator<Item = <T as WithTypedId>::Id>,
1484 ) -> Vec<Arc<T>>
1485 where
1486 T: Eventable + WithTypedId + Send + Sync + 'static,
1487 <T as WithTypedId>::Id: hyphae::IdFor<T, MapKey = Arc<str>>,
1488 {
1489 ids.into_iter()
1490 .filter_map(|id| self.entity_snapshot::<T>(&id))
1491 .collect()
1492 }
1493
1494 pub fn query_snapshot<Q>(&self, query: Q, request: Arc<RequestContext>) -> Vec<Arc<Q::Item>>
1500 where
1501 Q: QueryHandler + QueryParams + Clone + Send + Sync + 'static,
1502 Q::Item: DeserializeOwned + Clone + std::fmt::Debug + Send + Sync + 'static,
1503 {
1504 let query_item_type = Q::query_item_type_static();
1505 let store = self.registry.get_or_create(&query_item_type);
1506
1507 let query_context = Arc::new(QueryContext {
1508 req: request.clone(),
1509 });
1510 let query = Arc::new(query);
1511
1512 store
1513 .snapshot()
1514 .into_iter()
1515 .filter_map(|(_, item)| {
1516 let typed_item =
1517 downcast_any_item_arc::<Q::Item>(&item, "CellServerCtx::query_snapshot");
1518 let ctx = QueryTestCtx {
1519 item: typed_item.clone(),
1520 query: query.clone(),
1521 query_context: query_context.clone(),
1522 };
1523 if Q::test_entity(ctx) {
1524 Some(typed_item)
1525 } else {
1526 None
1527 }
1528 })
1529 .collect()
1530 }
1531
1532 pub fn report<R>(
1533 &self,
1534 report: R,
1535 request: Arc<RequestContext>,
1536 ) -> Cell<Arc<R::Output>, CellImmutable>
1537 where
1538 R: ReportHandler + ReportId + CacheKey + Clone + serde::Serialize + 'static,
1539 {
1540 let key = self.cache_key("report", report.report_id().as_ref(), &report, &request);
1541 let report_id = report.report_id();
1542
1543 if let Some(cell) = self.try_get_cached_report::<R>(&key) {
1545 crate::server::report_cache_stats::record_hit(&report_id);
1546 tracing::trace!(
1547 target: "myko::server::context::report_cache",
1548 "report_cache HIT report_id={} key={}",
1549 report_id,
1550 key,
1551 );
1552 return cell;
1553 }
1554
1555 let gate = self
1558 .compute_gates
1559 .entry(key.clone())
1560 .or_insert_with(|| Arc::new(std::sync::Mutex::new(())))
1561 .clone();
1562 let _lock = gate.lock().unwrap();
1563
1564 if let Some(cell) = self.try_get_cached_report::<R>(&key) {
1566 crate::server::report_cache_stats::record_hit_after_gate(&report_id);
1567 tracing::trace!(
1568 target: "myko::server::context::report_cache",
1569 "report_cache HIT_AFTER_GATE report_id={} key={}",
1570 report_id,
1571 key,
1572 );
1573 return cell;
1574 }
1575
1576 if tracing::enabled!(target: "myko::server::context::report_cache", tracing::Level::TRACE) {
1580 let payload = serde_json::to_string(&report)
1581 .unwrap_or_else(|e| format!("<serialize error: {e}>"));
1582 tracing::trace!(
1583 target: "myko::server::context::report_cache",
1584 "report_cache MISS_COMPUTE report_id={} key={} payload={}",
1585 report_id,
1586 key,
1587 payload,
1588 );
1589 }
1590
1591 let _span = tracing::trace_span!("myko.report", report = report_id.as_ref()).entered();
1595 crate::server::dispatch_metrics::record_report(report_id.as_ref(), request.origin());
1596 let nested_ctx = ReportContext::new(request, Arc::new(self.clone()));
1597 let built = report.compute(nested_ctx).materialize();
1601 #[cfg(feature = "profiling")]
1608 let built = built.with_name(report_id.as_ref());
1609 self.report_cache
1610 .insert(key.clone(), Arc::new(ReportCacheEntry::new(&built)));
1611 self.compute_gates.remove(&key);
1615
1616 crate::server::report_cache_stats::record_miss(&report_id);
1617
1618 built
1619 }
1620
1621 fn try_get_cached_report<R>(&self, key: &str) -> Option<Cell<Arc<R::Output>, CellImmutable>>
1623 where
1624 R: ReportHandler + 'static,
1625 {
1626 let existing = self.report_cache.get(key)?;
1627 if let Some(entry) = existing
1628 .value()
1629 .as_any()
1630 .downcast_ref::<ReportCacheEntry<Arc<R::Output>>>()
1631 && let Some(shared) = entry.get()
1632 {
1633 return Some(shared);
1634 }
1635 drop(existing);
1637 self.report_cache.remove(key);
1638 None
1639 }
1640
1641 pub fn new_server_transaction(&self) -> Arc<RequestContext> {
1642 Arc::new(RequestContext {
1643 tx: Arc::<str>::from(Uuid::new_v4().to_string()),
1644 client_id: None,
1645 lineage: vec![],
1646 host_id: self.host_id,
1647 created_at: chrono::Utc::now().to_string(),
1648 windback: None,
1649 })
1650 }
1651}
1652
1653impl std::fmt::Debug for CellServerCtx {
1654 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1655 f.debug_struct("CellServerCtx").finish()
1656 }
1657}
1658
1659#[cfg(test)]
1660mod tests {
1661 use std::sync::Arc;
1662
1663 use serde::{Deserialize, Serialize};
1664 use serde_json::json;
1665 use uuid::Uuid;
1666
1667 use super::CellServerCtx;
1668 use crate::{
1669 common::with_id::WithId,
1670 core::item::{
1671 AnyItem, Eventable, IngestBufferPolicy, IngestBufferRegistration, ItemRegistration,
1672 },
1673 hyphae::Gettable,
1674 search::SearchIndex,
1675 server::{HandlerRegistry, RelationshipManager, persister::PersisterRouter},
1676 store::StoreRegistry,
1677 test_util::scheduler_test_serial,
1678 wire::{MEvent, MEventType},
1679 };
1680
1681 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1682 struct BufferedTestItem {
1683 id: Arc<str>,
1684 value: i32,
1685 }
1686
1687 impl WithId for BufferedTestItem {
1688 fn id(&self) -> Arc<str> {
1689 self.id.clone()
1690 }
1691 }
1692
1693 impl AnyItem for BufferedTestItem {
1694 fn as_any(&self) -> &dyn std::any::Any {
1695 self
1696 }
1697
1698 fn entity_type(&self) -> &'static str {
1699 "BufferedTestItem"
1700 }
1701
1702 fn equals(&self, other: &dyn AnyItem) -> bool {
1703 other
1704 .as_any()
1705 .downcast_ref::<Self>()
1706 .map(|typed| self == typed)
1707 .unwrap_or(false)
1708 }
1709 }
1710
1711 impl Eventable for BufferedTestItem {
1712 const ENTITY_NAME_STATIC: &'static str = "BufferedTestItem";
1713 }
1714
1715 inventory::submit! {
1716 ItemRegistration {
1717 entity_type: "BufferedTestItem",
1718 crate_name: env!("CARGO_PKG_NAME"),
1719 parse: BufferedTestItem::parse,
1720 parse_bytes: BufferedTestItem::parse_bytes,
1721 serialize_json: |any| {
1722 let typed = any.as_any().downcast_ref::<BufferedTestItem>().unwrap();
1723 ::serde_json::value::to_raw_value(typed)
1724 },
1725 }
1726 }
1727
1728 inventory::submit! {
1729 IngestBufferRegistration {
1730 entity_type: "BufferedTestItem",
1731 policy: IngestBufferPolicy::TimeWindow { window_ms: 60_000 },
1732 }
1733 }
1734
1735 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1736 struct ImmediateTestItem {
1737 id: Arc<str>,
1738 value: i32,
1739 }
1740
1741 impl WithId for ImmediateTestItem {
1742 fn id(&self) -> Arc<str> {
1743 self.id.clone()
1744 }
1745 }
1746
1747 impl AnyItem for ImmediateTestItem {
1748 fn as_any(&self) -> &dyn std::any::Any {
1749 self
1750 }
1751
1752 fn entity_type(&self) -> &'static str {
1753 "ImmediateTestItem"
1754 }
1755
1756 fn equals(&self, other: &dyn AnyItem) -> bool {
1757 other
1758 .as_any()
1759 .downcast_ref::<Self>()
1760 .map(|typed| self == typed)
1761 .unwrap_or(false)
1762 }
1763 }
1764
1765 impl Eventable for ImmediateTestItem {
1766 const ENTITY_NAME_STATIC: &'static str = "ImmediateTestItem";
1767 }
1768
1769 inventory::submit! {
1770 ItemRegistration {
1771 entity_type: "ImmediateTestItem",
1772 crate_name: env!("CARGO_PKG_NAME"),
1773 parse: ImmediateTestItem::parse,
1774 parse_bytes: ImmediateTestItem::parse_bytes,
1775 serialize_json: |any| {
1776 let typed = any.as_any().downcast_ref::<ImmediateTestItem>().unwrap();
1777 ::serde_json::value::to_raw_value(typed)
1778 },
1779 }
1780 }
1781
1782 fn make_ctx() -> CellServerCtx {
1783 CellServerCtx::new(
1784 Uuid::new_v4(),
1785 Arc::new(StoreRegistry::new()),
1786 Arc::new(HandlerRegistry::new()),
1787 Arc::new(RelationshipManager::new()),
1788 Arc::new(PersisterRouter::default()),
1789 Arc::new(SearchIndex::new()),
1790 Arc::new(dashmap::DashMap::new()),
1791 None,
1792 None,
1793 )
1794 }
1795
1796 #[test]
1797 fn apply_event_batch_keeps_default_entities_immediate() {
1798 let _serial = scheduler_test_serial();
1799 let ctx = make_ctx();
1800 let applied = ctx
1801 .apply_event_batch(vec![MEvent {
1802 item: json!({
1803 "id": "immediate-1",
1804 "value": 7,
1805 }),
1806 change_type: MEventType::SET,
1807 item_type: "ImmediateTestItem".to_string(),
1808 created_at: "2026-03-12T00:00:00Z".to_string(),
1809 tx: "tx-immediate".to_string(),
1810 source_id: Some("test".to_string()),
1811 }])
1812 .expect("apply_event_batch should succeed");
1813
1814 assert_eq!(applied, 1);
1815 let store = ctx.registry.get_or_create("ImmediateTestItem");
1816 assert!(store.get(&Arc::<str>::from("immediate-1")).get().is_some());
1817 }
1818
1819 #[test]
1820 fn apply_event_batch_buffers_opted_in_entities() {
1821 let _serial = scheduler_test_serial();
1822 let ctx = make_ctx();
1823 let applied = ctx
1824 .apply_event_batch(vec![MEvent {
1825 item: json!({
1826 "id": "buffered-1",
1827 "value": 42,
1828 }),
1829 change_type: MEventType::SET,
1830 item_type: "BufferedTestItem".to_string(),
1831 created_at: "2026-03-12T00:00:00Z".to_string(),
1832 tx: "tx-buffered".to_string(),
1833 source_id: Some("test".to_string()),
1834 }])
1835 .expect("apply_event_batch should succeed");
1836
1837 assert_eq!(applied, 1);
1838 let store = ctx.registry.get_or_create("BufferedTestItem");
1839 assert!(store.get(&Arc::<str>::from("buffered-1")).get().is_none());
1840
1841 let flushed = ctx.flush_all_buffered_events();
1842 assert_eq!(flushed, 1);
1843 assert!(store.get(&Arc::<str>::from("buffered-1")).get().is_some());
1844 }
1845
1846 #[test]
1847 fn apply_event_batch_delivers_both_diffs_for_mixed_set_and_del_same_type() {
1848 let _serial = scheduler_test_serial();
1857 let ctx = make_ctx();
1858
1859 ctx.apply_event_batch(vec![MEvent {
1860 item: json!({ "id": "old-1", "value": 1 }),
1861 change_type: MEventType::SET,
1862 item_type: "ImmediateTestItem".to_string(),
1863 created_at: "2026-03-12T00:00:00Z".to_string(),
1864 tx: "tx-seed".to_string(),
1865 source_id: Some("test".to_string()),
1866 }])
1867 .expect("seed apply_event_batch should succeed");
1868
1869 let store = ctx.registry.get_or_create("ImmediateTestItem");
1870 let diffs_seen = Arc::new(std::sync::Mutex::new(Vec::new()));
1871 let diffs_seen_for_closure = diffs_seen.clone();
1872 let _guard = store.subscribe_diffs(move |diff| {
1873 diffs_seen_for_closure
1874 .lock()
1875 .unwrap()
1876 .push(format!("{diff:?}"));
1877 });
1878 diffs_seen.lock().unwrap().clear();
1881
1882 let applied = ctx
1883 .apply_event_batch(vec![
1884 MEvent {
1885 item: json!({ "id": "new-1", "value": 2 }),
1886 change_type: MEventType::SET,
1887 item_type: "ImmediateTestItem".to_string(),
1888 created_at: "2026-03-12T00:00:01Z".to_string(),
1889 tx: "tx-mixed".to_string(),
1890 source_id: Some("test".to_string()),
1891 },
1892 MEvent {
1893 item: json!({ "id": "old-1", "value": 1 }),
1894 change_type: MEventType::DEL,
1895 item_type: "ImmediateTestItem".to_string(),
1896 created_at: "2026-03-12T00:00:01Z".to_string(),
1897 tx: "tx-mixed".to_string(),
1898 source_id: Some("test".to_string()),
1899 },
1900 ])
1901 .expect("mixed apply_event_batch should succeed");
1902
1903 assert_eq!(applied, 2);
1904 assert!(store.get(&Arc::<str>::from("new-1")).get().is_some());
1905 assert!(store.get(&Arc::<str>::from("old-1")).get().is_none());
1906
1907 let seen = diffs_seen.lock().unwrap();
1908 assert_eq!(
1909 seen.len(),
1910 2,
1911 "both the SET and DEL diffs must reach subscribers, not just the last one: {:?}",
1912 *seen
1913 );
1914 }
1915
1916 #[test]
1917 fn compute_gates_does_not_leak_after_cache_populates() {
1918 use crate::{entities::server::GetPeerServers, request::RequestContext};
1924
1925 let _serial = scheduler_test_serial();
1926 let ctx = make_ctx();
1927 let request = Arc::new(RequestContext::internal(
1928 Arc::<str>::from(Uuid::new_v4().to_string()),
1929 ctx.host_id,
1930 "test",
1931 ));
1932
1933 let _ = ctx.query_map_untyped(GetPeerServers {}, request.clone());
1936 assert!(
1937 ctx.compute_gates.is_empty(),
1938 "compute_gates must be empty once the query cache is populated, got {:?}",
1939 ctx.compute_gates
1940 );
1941
1942 let _ = ctx.query_map_untyped(GetPeerServers {}, request);
1945 assert!(ctx.compute_gates.is_empty());
1946 }
1947
1948 #[test]
1949 fn compute_gates_does_not_leak_after_report_cache_populates() {
1950 use crate::{
1955 entities::client::{ClientId, ClientStatus},
1956 request::RequestContext,
1957 };
1958
1959 let _serial = scheduler_test_serial();
1960 let ctx = make_ctx();
1961 let request = Arc::new(RequestContext::internal(
1962 Arc::<str>::from(Uuid::new_v4().to_string()),
1963 ctx.host_id,
1964 "test",
1965 ));
1966
1967 let report = ClientStatus {
1968 client_id: ClientId::from(Arc::<str>::from("test-client")),
1969 };
1970 let _ = ctx.report(report.clone(), request.clone());
1971 assert!(
1972 ctx.compute_gates.is_empty(),
1973 "compute_gates must be empty once the report cache is populated, got {:?}",
1974 ctx.compute_gates
1975 );
1976
1977 let _ = ctx.report(report, request);
1978 assert!(ctx.compute_gates.is_empty());
1979 }
1980}