1use std::{
42 any::{Any, TypeId},
43 collections::{HashMap, HashSet, VecDeque},
44 hash::Hash,
45 sync::{Arc, Mutex},
46};
47
48use tokio::{sync::broadcast, task_local};
49
50use super::{Identifier, IdentifierFilter};
51use crate::runtime::process::Id;
52
53const DEFAULT_CHANNEL_CAPACITY: usize = 16;
54
55task_local! {
56 pub(crate) static CURRENT_DATASPACE: DataspaceRegistry;
57}
58
59#[derive(Clone, Debug, PartialEq, Eq)]
62pub enum DataspaceUpdate<T> {
63 Asserted(Identifier, T),
65
66 Retracted(Identifier),
68
69 Message(Identifier, T),
74}
75
76#[derive(Clone, PartialEq, Eq, Hash)]
78struct StorageKey {
79 type_id: TypeId,
80 identifier: Identifier,
81}
82
83impl StorageKey {
84 fn new<T: 'static>(identifier: Identifier) -> Self {
85 Self {
86 type_id: TypeId::of::<T>(),
87 identifier,
88 }
89 }
90}
91
92trait AnyChannel: Send + Sync {
94 fn send_retraction(&self, id: &Identifier);
96
97 fn as_any(&self) -> &dyn Any;
99}
100
101struct TypedChannel<T: Clone + Send + Sync + 'static> {
103 tx: broadcast::Sender<DataspaceUpdate<T>>,
104}
105
106impl<T: Clone + Send + Sync + 'static> AnyChannel for TypedChannel<T> {
107 fn send_retraction(&self, id: &Identifier) {
108 let _ = self.tx.send(DataspaceUpdate::Retracted(id.clone()));
109 }
110
111 fn as_any(&self) -> &dyn Any {
112 self
113 }
114}
115
116struct FilteredChannel {
118 type_id: TypeId,
119 filter: IdentifierFilter,
120 sender: Box<dyn AnyChannel>,
121}
122
123struct StoredValue {
125 value: Box<dyn Any + Send + Sync>,
126 owner: Id,
127}
128
129struct RegistryState {
131 channels: HashMap<StorageKey, Box<dyn AnyChannel>>,
133
134 filtered_channels: Vec<FilteredChannel>,
136
137 current_values: HashMap<StorageKey, StoredValue>,
141
142 process_assertions: HashMap<Id, HashSet<StorageKey>>,
144
145 channel_capacity: usize,
147}
148
149impl RegistryState {
150 fn new(channel_capacity: usize) -> Self {
151 Self {
152 channels: HashMap::new(),
153 filtered_channels: Vec::new(),
154 current_values: HashMap::new(),
155 process_assertions: HashMap::new(),
156 channel_capacity,
157 }
158 }
159
160 fn get_or_create_exact_sender<T>(&mut self, key: StorageKey) -> broadcast::Receiver<DataspaceUpdate<T>>
162 where
163 T: Clone + Send + Sync + 'static,
164 {
165 let channel = self.channels.entry(key).or_insert_with(|| {
166 let (tx, _) = broadcast::channel::<DataspaceUpdate<T>>(self.channel_capacity);
167 Box::new(TypedChannel { tx })
168 });
169
170 let typed = channel
171 .as_any()
172 .downcast_ref::<TypedChannel<T>>()
173 .unwrap_or_else(|| unreachable!("type mismatch in dataspace registry"));
176
177 typed.tx.subscribe()
178 }
179
180 fn create_filtered_sender<T>(&mut self, filter: IdentifierFilter) -> broadcast::Receiver<DataspaceUpdate<T>>
182 where
183 T: Clone + Send + Sync + 'static,
184 {
185 let (tx, rx) = broadcast::channel::<DataspaceUpdate<T>>(self.channel_capacity);
186
187 self.filtered_channels.push(FilteredChannel {
188 type_id: TypeId::of::<T>(),
189 filter,
190 sender: Box::new(TypedChannel { tx }),
191 });
192
193 rx
194 }
195
196 fn notify_typed<T>(&self, key: &StorageKey, update: &DataspaceUpdate<T>)
201 where
202 T: Clone + Send + Sync + 'static,
203 {
204 if let Some(ch) = self.channels.get(key) {
205 if let Some(typed) = ch.as_any().downcast_ref::<TypedChannel<T>>() {
206 let _ = typed.tx.send(update.clone());
207 }
208 }
209
210 for channel in &self.filtered_channels {
211 if channel.type_id == key.type_id && channel.filter.matches(&key.identifier) {
212 if let Some(typed) = channel.sender.as_any().downcast_ref::<TypedChannel<T>>() {
213 let _ = typed.tx.send(update.clone());
214 }
215 }
216 }
217 }
218
219 fn notify_retraction(&self, key: &StorageKey) {
226 if let Some(ch) = self.channels.get(key) {
227 ch.send_retraction(&key.identifier);
228 }
229
230 for filtered in &self.filtered_channels {
231 if filtered.type_id == key.type_id && filtered.filter.matches(&key.identifier) {
232 filtered.sender.send_retraction(&key.identifier);
233 }
234 }
235 }
236}
237
238struct DataspaceRegistryInner {
240 state: Mutex<RegistryState>,
241}
242
243#[derive(Clone)]
253pub struct DataspaceRegistry {
254 inner: Arc<DataspaceRegistryInner>,
255}
256
257impl Default for DataspaceRegistry {
258 fn default() -> Self {
259 Self::new()
260 }
261}
262
263impl DataspaceRegistry {
264 pub fn new() -> Self {
266 Self::with_channel_capacity(DEFAULT_CHANNEL_CAPACITY)
267 }
268
269 pub fn try_current() -> Option<Self> {
271 CURRENT_DATASPACE.try_with(|ds| ds.clone()).ok()
272 }
273
274 #[cfg(any(test, feature = "test-util"))]
280 pub fn with_current<F, R>(&self, f: F) -> R
281 where
282 F: FnOnce() -> R,
283 {
284 CURRENT_DATASPACE.sync_scope(self.clone(), f)
285 }
286
287 pub fn with_channel_capacity(capacity: usize) -> Self {
289 Self {
290 inner: Arc::new(DataspaceRegistryInner {
291 state: Mutex::new(RegistryState::new(capacity)),
292 }),
293 }
294 }
295
296 pub fn assert<T>(&self, value: T, id: impl Into<Identifier>)
302 where
303 T: Clone + Send + Sync + 'static,
304 {
305 let id = id.into();
306 let key = StorageKey::new::<T>(id.clone());
307 let caller = Id::current();
308 let mut state = self.inner.state.lock().unwrap();
309
310 if let Some(existing) = state.current_values.get(&key) {
312 debug_assert_eq!(
313 existing.owner, caller,
314 "process {caller:?} attempted to update assertion owned by {:?}",
315 existing.owner
316 );
317 if existing.owner != caller {
318 return;
319 }
320 }
321
322 state.current_values.insert(
324 key.clone(),
325 StoredValue {
326 value: Box::new(value.clone()),
327 owner: caller,
328 },
329 );
330
331 state.process_assertions.entry(caller).or_default().insert(key.clone());
333
334 let update = DataspaceUpdate::Asserted(id, value);
336 state.notify_typed::<T>(&key, &update);
337 }
338
339 pub fn retract<T>(&self, id: impl Into<Identifier>)
343 where
344 T: Clone + Send + Sync + 'static,
345 {
346 let id = id.into();
347 let key = StorageKey::new::<T>(id.clone());
348 let caller = Id::current();
349 let mut state = self.inner.state.lock().unwrap();
350
351 let Some(stored) = state.current_values.get(&key) else {
353 return;
354 };
355
356 debug_assert_eq!(
357 stored.owner, caller,
358 "process {caller:?} attempted to retract assertion owned by {:?}",
359 stored.owner
360 );
361 if stored.owner != caller {
362 return;
363 }
364
365 state.current_values.remove(&key);
367
368 if let Some(keys) = state.process_assertions.get_mut(&caller) {
369 keys.remove(&key);
370 if keys.is_empty() {
371 state.process_assertions.remove(&caller);
372 }
373 }
374
375 state.notify_retraction(&key);
376 }
377
378 pub fn send<T>(&self, value: T, id: impl Into<Identifier>)
383 where
384 T: Clone + Send + Sync + 'static,
385 {
386 let id = id.into();
387 let key = StorageKey::new::<T>(id.clone());
388 let state = self.inner.state.lock().unwrap();
389
390 let update = DataspaceUpdate::Message(id, value);
392 state.notify_typed::<T>(&key, &update);
393 }
394
395 pub(crate) fn retract_all_for_process(&self, process_id: Id) {
400 let mut state = self.inner.state.lock().unwrap();
401
402 let Some(keys) = state.process_assertions.remove(&process_id) else {
403 return;
404 };
405
406 for key in keys {
407 state.current_values.remove(&key);
408 state.notify_retraction(&key);
409 }
410 }
411
412 pub fn current_values<T>(&self, filter: IdentifierFilter) -> Vec<T>
417 where
418 T: Clone + Send + Sync + 'static,
419 {
420 let type_id = TypeId::of::<T>();
421 let state = self.inner.state.lock().unwrap();
422 state
423 .current_values
424 .iter()
425 .filter(|(key, _)| key.type_id == type_id && filter.matches(&key.identifier))
426 .filter_map(|(_, stored)| stored.value.downcast_ref::<T>().cloned())
427 .collect()
428 }
429
430 pub fn subscribe<T>(&self, filter: IdentifierFilter) -> Subscription<T>
436 where
437 T: Clone + Send + Sync + 'static,
438 {
439 let mut state = self.inner.state.lock().unwrap();
440
441 match filter {
442 IdentifierFilter::Exact(ref id) => {
443 let key = StorageKey::new::<T>(id.clone());
444 let rx = state.get_or_create_exact_sender::<T>(key.clone());
445
446 let pending: VecDeque<_> = state
448 .current_values
449 .get(&key)
450 .and_then(|stored| stored.value.downcast_ref::<T>())
451 .map(|value| DataspaceUpdate::Asserted(id.clone(), value.clone()))
452 .into_iter()
453 .collect();
454
455 Subscription { pending, rx }
456 }
457 filter @ (IdentifierFilter::All | IdentifierFilter::Prefix(_)) => {
458 let rx = state.create_filtered_sender::<T>(filter.clone());
459
460 let type_id = TypeId::of::<T>();
462 let pending: VecDeque<_> = state
463 .current_values
464 .iter()
465 .filter(|(key, _)| key.type_id == type_id && filter.matches(&key.identifier))
466 .filter_map(|(key, stored)| {
467 stored
468 .value
469 .downcast_ref::<T>()
470 .map(|value| DataspaceUpdate::Asserted(key.identifier.clone(), value.clone()))
471 })
472 .collect();
473
474 Subscription { pending, rx }
475 }
476 }
477 }
478}
479
480pub struct Subscription<T> {
482 pending: VecDeque<DataspaceUpdate<T>>,
483 rx: broadcast::Receiver<DataspaceUpdate<T>>,
484}
485
486impl<T> Subscription<T>
487where
488 T: Clone + Send + Sync + 'static,
489{
490 pub async fn recv(&mut self) -> Option<DataspaceUpdate<T>> {
496 if let Some(value) = self.pending.pop_front() {
499 return Some(value);
500 }
501
502 loop {
503 match self.rx.recv().await {
504 Ok(value) => return Some(value),
505 Err(broadcast::error::RecvError::Lagged(_)) => continue,
506 Err(broadcast::error::RecvError::Closed) => return None,
507 }
508 }
509 }
510}
511
512#[cfg(test)]
513mod tests {
514 use std::future::pending;
515
516 use tokio_test::{assert_pending, assert_ready, assert_ready_eq, task::spawn as test_spawn};
517
518 use super::*;
519 use crate::runtime::process::Process;
520
521 #[test]
522 fn subscribe_then_assert() {
523 let registry = DataspaceRegistry::new();
524 let id = Identifier::numeric(1);
525
526 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
527 registry.assert(42u32, id.clone());
528
529 let mut recv = test_spawn(sub.recv());
530 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id, 42)));
531 }
532
533 #[test]
534 fn multiple_subscribers_receive_same_assertion() {
535 let registry = DataspaceRegistry::new();
536 let id = Identifier::numeric(1);
537
538 let mut sub1 = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
539 let mut sub2 = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
540
541 registry.assert(42u32, id.clone());
542
543 let mut recv1 = test_spawn(sub1.recv());
544 assert_ready_eq!(recv1.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 42)));
545
546 let mut recv2 = test_spawn(sub2.recv());
547 assert_ready_eq!(recv2.poll(), Some(DataspaceUpdate::Asserted(id, 42)));
548 }
549
550 #[test]
551 fn assert_without_subscribers_stores_value() {
552 let registry = DataspaceRegistry::new();
553 let id = Identifier::numeric(1);
554
555 registry.assert(42u32, id.clone());
557
558 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
560 let mut recv = test_spawn(sub.recv());
561 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id, 42)));
562 }
563
564 #[test]
565 fn different_types_same_identifier() {
566 let registry = DataspaceRegistry::new();
567 let id = Identifier::numeric(1);
568
569 let mut sub_u32 = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
570 let mut sub_string = registry.subscribe::<String>(IdentifierFilter::exact(id.clone()));
571
572 registry.assert(42u32, id.clone());
573 registry.assert("hello".to_string(), id.clone());
574
575 let mut recv_u32 = test_spawn(sub_u32.recv());
576 assert_ready_eq!(recv_u32.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 42)));
577
578 let mut recv_string = test_spawn(sub_string.recv());
579 assert_ready_eq!(
580 recv_string.poll(),
581 Some(DataspaceUpdate::Asserted(id, "hello".to_string()))
582 );
583 }
584
585 #[test]
586 fn process_identifier() {
587 let registry = DataspaceRegistry::new();
588 let pid = crate::runtime::process::Id::new();
589 let id = Identifier::from(pid);
590
591 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
592 registry.assert(42u32, id.clone());
593
594 let mut recv = test_spawn(sub.recv());
595 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id, 42)));
596 }
597
598 #[test]
599 fn channel_closed_returns_none() {
600 let registry = DataspaceRegistry::new();
601 let id = Identifier::numeric(1);
602
603 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id));
604
605 drop(registry);
607
608 let mut recv = test_spawn(sub.recv());
609 assert_ready_eq!(recv.poll(), None);
610 }
611
612 #[test]
613 fn lagged_subscriber_recovers() {
614 let registry = DataspaceRegistry::with_channel_capacity(2);
616 let id = Identifier::numeric(1);
617
618 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
619
620 for i in 0..10 {
622 registry.assert(i as u32, id.clone());
623 }
624
625 let mut recv = test_spawn(sub.recv());
627 let value = assert_ready!(recv.poll());
628 assert!(value.is_some());
629 }
630
631 #[test]
632 fn multiple_values_received_in_order() {
633 let registry = DataspaceRegistry::new();
634 let id = Identifier::numeric(1);
635
636 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
637
638 registry.assert(1u32, id.clone());
639 registry.assert(2u32, id.clone());
640 registry.assert(3u32, id.clone());
641
642 let mut recv = test_spawn(sub.recv());
643 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 1)));
644 drop(recv);
645
646 let mut recv = test_spawn(sub.recv());
647 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 2)));
648 drop(recv);
649
650 let mut recv = test_spawn(sub.recv());
651 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id, 3)));
652 }
653
654 #[test]
655 fn assert_then_retract() {
656 let registry = DataspaceRegistry::new();
657 let id = Identifier::numeric(1);
658
659 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
660
661 registry.assert(42u32, id.clone());
662 registry.retract::<u32>(id.clone());
663
664 let mut recv = test_spawn(sub.recv());
665 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 42)));
666 drop(recv);
667
668 let mut recv = test_spawn(sub.recv());
669 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Retracted(id)));
670 }
671
672 #[test]
673 fn retract_without_subscribers_succeeds() {
674 let registry = DataspaceRegistry::new();
675 let id = Identifier::numeric(1);
676
677 registry.retract::<u32>(id);
679 }
680
681 #[test]
682 fn multiple_subscribers_receive_retraction() {
683 let registry = DataspaceRegistry::new();
684 let id = Identifier::numeric(1);
685
686 let mut sub1 = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
687 let mut sub2 = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
688
689 registry.assert(42u32, id.clone());
690 registry.retract::<u32>(id.clone());
691
692 let mut recv1 = test_spawn(sub1.recv());
694 assert_ready_eq!(recv1.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 42)));
695 drop(recv1);
696
697 let mut recv2 = test_spawn(sub2.recv());
698 assert_ready_eq!(recv2.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 42)));
699 drop(recv2);
700
701 let mut recv1 = test_spawn(sub1.recv());
703 assert_ready_eq!(recv1.poll(), Some(DataspaceUpdate::Retracted(id.clone())));
704
705 let mut recv2 = test_spawn(sub2.recv());
706 assert_ready_eq!(recv2.poll(), Some(DataspaceUpdate::Retracted(id)));
707 }
708
709 #[test]
710 fn assert_retract_assert_sequence() {
711 let registry = DataspaceRegistry::new();
712 let id = Identifier::numeric(1);
713
714 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
715
716 registry.assert(1u32, id.clone());
717 registry.retract::<u32>(id.clone());
718 registry.assert(2u32, id.clone());
719
720 let mut recv = test_spawn(sub.recv());
721 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 1)));
722 drop(recv);
723
724 let mut recv = test_spawn(sub.recv());
725 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Retracted(id.clone())));
726 drop(recv);
727
728 let mut recv = test_spawn(sub.recv());
729 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id, 2)));
730 }
731
732 #[test]
733 fn all_filter_receives_from_multiple_identifiers() {
734 let registry = DataspaceRegistry::new();
735 let id1 = Identifier::numeric(1);
736 let id2 = Identifier::numeric(2);
737
738 let mut sub = registry.subscribe::<u32>(IdentifierFilter::all());
739
740 registry.assert(1u32, id1.clone());
741 registry.assert(2u32, id2.clone());
742
743 let mut recv = test_spawn(sub.recv());
744 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id1, 1)));
745 drop(recv);
746
747 let mut recv = test_spawn(sub.recv());
748 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id2, 2)));
749 }
750
751 #[test]
752 fn all_filter_receives_retraction() {
753 let registry = DataspaceRegistry::new();
754 let id = Identifier::numeric(1);
755
756 let mut sub = registry.subscribe::<u32>(IdentifierFilter::all());
757
758 registry.assert(42u32, id.clone());
759 registry.retract::<u32>(id.clone());
760
761 let mut recv = test_spawn(sub.recv());
762 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 42)));
763 drop(recv);
764
765 let mut recv = test_spawn(sub.recv());
766 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Retracted(id)));
767 }
768
769 #[test]
770 fn all_filter_and_exact_both_receive() {
771 let registry = DataspaceRegistry::new();
772 let id = Identifier::numeric(1);
773
774 let mut specific = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
775 let mut all = registry.subscribe::<u32>(IdentifierFilter::all());
776
777 registry.assert(42u32, id.clone());
778
779 let mut recv_specific = test_spawn(specific.recv());
780 assert_ready_eq!(recv_specific.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 42)));
781
782 let mut recv_all = test_spawn(all.recv());
783 assert_ready_eq!(recv_all.poll(), Some(DataspaceUpdate::Asserted(id, 42)));
784 }
785
786 #[test]
787 fn all_filter_created_before_exact_channels() {
788 let registry = DataspaceRegistry::new();
789
790 let mut all = registry.subscribe::<u32>(IdentifierFilter::all());
792
793 let id = Identifier::numeric(1);
794 registry.assert(42u32, id.clone());
795
796 let mut recv = test_spawn(all.recv());
797 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id, 42)));
798 }
799
800 #[test]
801 fn subscribe_receives_current_value() {
802 let registry = DataspaceRegistry::new();
803 let id = Identifier::numeric(1);
804
805 registry.assert(42u32, id.clone());
807
808 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
810 let mut recv = test_spawn(sub.recv());
811 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id, 42)));
812 }
813
814 #[test]
815 fn subscribe_after_retract_gets_nothing_pending() {
816 let registry = DataspaceRegistry::new();
817 let id = Identifier::numeric(1);
818
819 registry.assert(42u32, id.clone());
821 registry.retract::<u32>(id.clone());
822
823 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id));
824
825 drop(registry);
827
828 let mut recv = test_spawn(sub.recv());
829 assert_ready_eq!(recv.poll(), None);
830 }
831
832 #[test]
833 fn all_filter_receives_current_values() {
834 let registry = DataspaceRegistry::new();
835 let id1 = Identifier::numeric(1);
836 let id2 = Identifier::numeric(2);
837
838 registry.assert(1u32, id1.clone());
840 registry.assert(2u32, id2.clone());
841
842 let mut sub = registry.subscribe::<u32>(IdentifierFilter::all());
843
844 let mut recv = test_spawn(sub.recv());
846 let v1 = assert_ready!(recv.poll());
847 drop(recv);
848
849 let mut recv = test_spawn(sub.recv());
850 let v2 = assert_ready!(recv.poll());
851
852 let mut received = [v1.unwrap(), v2.unwrap()];
853 received.sort_by_key(|update| match update {
854 DataspaceUpdate::Asserted(_, v) | DataspaceUpdate::Message(_, v) => *v,
855 DataspaceUpdate::Retracted(_) => 0,
856 });
857
858 assert_eq!(received[0], DataspaceUpdate::Asserted(id1, 1));
859 assert_eq!(received[1], DataspaceUpdate::Asserted(id2, 2));
860 }
861
862 #[test]
863 fn subscribe_receives_current_then_future() {
864 let registry = DataspaceRegistry::new();
865 let id = Identifier::numeric(1);
866
867 registry.assert(1u32, id.clone());
869
870 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
871
872 registry.assert(2u32, id.clone());
874
875 let mut recv = test_spawn(sub.recv());
877 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 1)));
878 drop(recv);
879
880 let mut recv = test_spawn(sub.recv());
881 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id, 2)));
882 }
883
884 #[test]
885 fn prefix_filter_matches_named_identifiers() {
886 let registry = DataspaceRegistry::new();
887 let id1 = Identifier::named("metrics.cpu");
888 let id2 = Identifier::named("metrics.mem");
889 let id3 = Identifier::named("logs.error");
890
891 let mut sub = registry.subscribe::<u32>(IdentifierFilter::prefix("metrics."));
892
893 registry.assert(1u32, id1.clone());
894 registry.assert(2u32, id2.clone());
895 registry.assert(3u32, id3);
896
897 let mut recv = test_spawn(sub.recv());
899 let v1 = assert_ready!(recv.poll());
900 drop(recv);
901
902 let mut recv = test_spawn(sub.recv());
903 let v2 = assert_ready!(recv.poll());
904
905 let mut received = [v1.unwrap(), v2.unwrap()];
906 received.sort_by_key(|update| match update {
907 DataspaceUpdate::Asserted(_, v) | DataspaceUpdate::Message(_, v) => *v,
908 DataspaceUpdate::Retracted(_) => 0,
909 });
910
911 assert_eq!(received[0], DataspaceUpdate::Asserted(id1, 1));
912 assert_eq!(received[1], DataspaceUpdate::Asserted(id2, 2));
913 }
914
915 #[test]
916 fn prefix_filter_does_not_match_numeric() {
917 let registry = DataspaceRegistry::new();
918 let id = Identifier::numeric(42);
919
920 let mut sub = registry.subscribe::<u32>(IdentifierFilter::prefix("any"));
921
922 registry.assert(1u32, id);
923
924 drop(registry);
926
927 let mut recv = test_spawn(sub.recv());
928 assert_ready_eq!(recv.poll(), None);
929 }
930
931 #[test]
932 fn prefix_filter_replays_matching_current_values() {
933 let registry = DataspaceRegistry::new();
934 let id1 = Identifier::named("svc.alpha");
935 let id2 = Identifier::named("svc.beta");
936 let id3 = Identifier::named("other.gamma");
937
938 registry.assert(1u32, id1.clone());
940 registry.assert(2u32, id2.clone());
941 registry.assert(3u32, id3);
942
943 let mut sub = registry.subscribe::<u32>(IdentifierFilter::prefix("svc."));
944
945 let mut recv = test_spawn(sub.recv());
947 let v1 = assert_ready!(recv.poll());
948 drop(recv);
949
950 let mut recv = test_spawn(sub.recv());
951 let v2 = assert_ready!(recv.poll());
952
953 let mut received = [v1.unwrap(), v2.unwrap()];
954 received.sort_by_key(|update| match update {
955 DataspaceUpdate::Asserted(_, v) | DataspaceUpdate::Message(_, v) => *v,
956 DataspaceUpdate::Retracted(_) => 0,
957 });
958
959 assert_eq!(received[0], DataspaceUpdate::Asserted(id1, 1));
960 assert_eq!(received[1], DataspaceUpdate::Asserted(id2, 2));
961 }
962
963 #[test]
964 fn try_current_returns_none_outside_context() {
965 assert!(DataspaceRegistry::try_current().is_none());
966 }
967
968 #[test]
969 fn current_and_try_current_work_within_scope() {
970 let registry = DataspaceRegistry::new();
971 let registry_clone = registry.clone();
972
973 let mut scope_fut = test_spawn(CURRENT_DATASPACE.scope(registry, async {
974 let current = DataspaceRegistry::try_current();
975 assert!(current.is_some());
976
977 let current = current.unwrap();
979 let id = Identifier::named("test");
980 current.assert(42u32, id.clone());
981
982 let mut sub = registry_clone.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
983 let value = sub.recv().await;
984 assert_eq!(value, Some(DataspaceUpdate::Asserted(id, 42)));
985 }));
986 assert_ready!(scope_fut.poll());
987 }
988
989 #[test]
990 fn retract_all_for_process_retracts_all_owned_assertions() {
991 let registry = DataspaceRegistry::new();
992 let process_id = Id::new();
993 let id1 = Identifier::named("val1");
994 let id2 = Identifier::named("val2");
995
996 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(process_id, || {
998 registry.assert(1u32, id1.clone());
999 registry.assert(2u32, id2.clone());
1000 });
1001
1002 let mut sub1 = registry.subscribe::<u32>(IdentifierFilter::exact(id1.clone()));
1004 let mut sub2 = registry.subscribe::<u32>(IdentifierFilter::exact(id2.clone()));
1005
1006 let mut recv = test_spawn(sub1.recv());
1008 let _ = assert_ready!(recv.poll());
1009 drop(recv);
1010
1011 let mut recv = test_spawn(sub2.recv());
1012 let _ = assert_ready!(recv.poll());
1013 drop(recv);
1014
1015 registry.retract_all_for_process(process_id);
1017
1018 let mut recv = test_spawn(sub1.recv());
1020 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Retracted(id1)));
1021 drop(recv);
1022
1023 let mut recv = test_spawn(sub2.recv());
1024 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Retracted(id2)));
1025 }
1026
1027 #[test]
1028 fn retract_all_for_process_does_not_affect_other_processes() {
1029 let registry = DataspaceRegistry::new();
1030 let pid_a = Id::new();
1031 let pid_b = Id::new();
1032 let id_a = Identifier::named("a");
1033 let id_b = Identifier::named("b");
1034
1035 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1036 registry.assert(1u32, id_a.clone());
1037 });
1038 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_b, || {
1039 registry.assert(2u32, id_b.clone());
1040 });
1041
1042 registry.retract_all_for_process(pid_a);
1044
1045 let mut sub_b = registry.subscribe::<u32>(IdentifierFilter::exact(id_b.clone()));
1047 let mut recv = test_spawn(sub_b.recv());
1048 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id_b, 2)));
1049 drop(recv);
1050
1051 let mut sub_a = registry.subscribe::<u32>(IdentifierFilter::exact(id_a));
1053 drop(registry);
1054 let mut recv = test_spawn(sub_a.recv());
1055 assert_ready_eq!(recv.poll(), None);
1056 }
1057
1058 #[test]
1059 fn retract_all_for_process_notifies_filtered_subscribers() {
1060 let registry = DataspaceRegistry::new();
1061 let process_id = Id::new();
1062 let id = Identifier::named("val");
1063
1064 let mut all_sub = registry.subscribe::<u32>(IdentifierFilter::all());
1065
1066 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(process_id, || {
1067 registry.assert(42u32, id.clone());
1068 });
1069
1070 let mut recv = test_spawn(all_sub.recv());
1072 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 42)));
1073 drop(recv);
1074
1075 registry.retract_all_for_process(process_id);
1077
1078 let mut recv = test_spawn(all_sub.recv());
1080 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Retracted(id)));
1081 }
1082
1083 #[test]
1084 fn manual_retract_cleans_up_process_tracking() {
1085 let registry = DataspaceRegistry::new();
1086 let process_id = Id::new();
1087 let id = Identifier::named("val");
1088
1089 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
1090
1091 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(process_id, || {
1092 registry.assert(42u32, id.clone());
1093 });
1094
1095 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(process_id, || {
1097 registry.retract::<u32>(id.clone());
1098 });
1099
1100 let mut recv = test_spawn(sub.recv());
1102 let _ = assert_ready!(recv.poll());
1103 drop(recv);
1104
1105 let mut recv = test_spawn(sub.recv());
1106 let _ = assert_ready!(recv.poll());
1107 drop(recv);
1108
1109 registry.retract_all_for_process(process_id);
1111
1112 drop(registry);
1114 let mut recv = test_spawn(sub.recv());
1115 assert_ready_eq!(recv.poll(), None);
1116 }
1117
1118 #[test]
1119 fn retract_all_for_unknown_process_is_noop() {
1120 let registry = DataspaceRegistry::new();
1121 registry.retract_all_for_process(Id::new());
1123 }
1124
1125 #[test]
1126 fn instrumented_process_retracts_on_normal_completion() {
1127 let registry = DataspaceRegistry::new();
1128 let process = Process::supervisor_with_dataspace("test_worker", None, Some(registry.clone())).unwrap();
1129 let id = "from_worker";
1130
1131 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id));
1132
1133 let mut recv_fut = test_spawn(sub.recv());
1135 assert!(!recv_fut.is_woken());
1136 assert_pending!(recv_fut.poll());
1137
1138 let mut worker = test_spawn(process.into_process_future(async {
1140 DataspaceRegistry::try_current()
1141 .expect("dataspace registry should be available")
1142 .assert(42u32, "from_worker");
1143 }));
1144
1145 assert_ready!(worker.poll());
1147
1148 assert!(recv_fut.is_woken());
1150 assert_ready_eq!(recv_fut.poll(), Some(DataspaceUpdate::Asserted(id.into(), 42)));
1151
1152 drop(recv_fut);
1153
1154 let mut recv_fut = test_spawn(sub.recv());
1156 assert_pending!(recv_fut.poll());
1157 assert!(!recv_fut.is_woken());
1158
1159 drop(worker);
1161
1162 assert!(recv_fut.is_woken());
1164 assert_ready_eq!(recv_fut.poll(), Some(DataspaceUpdate::Retracted(id.into())));
1165 }
1166
1167 #[test]
1168 fn instrumented_process_retracts_on_drop_before_completion() {
1169 let registry = DataspaceRegistry::new();
1170 let process = Process::supervisor_with_dataspace("test_worker", None, Some(registry.clone())).unwrap();
1171 let id = "from_worker";
1172
1173 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id));
1174
1175 let mut recv_fut = test_spawn(sub.recv());
1177 assert!(!recv_fut.is_woken());
1178 assert_pending!(recv_fut.poll());
1179
1180 let mut worker = test_spawn(process.into_process_future(async {
1182 DataspaceRegistry::try_current()
1183 .expect("dataspace registry should be available")
1184 .assert(42u32, "from_worker");
1185 pending::<()>().await;
1186 }));
1187
1188 assert_pending!(worker.poll());
1190
1191 assert!(recv_fut.is_woken());
1193 assert_ready_eq!(recv_fut.poll(), Some(DataspaceUpdate::Asserted(id.into(), 42)));
1194
1195 drop(recv_fut);
1196
1197 let mut recv_fut = test_spawn(sub.recv());
1199 assert_pending!(recv_fut.poll());
1200 assert!(!recv_fut.is_woken());
1201
1202 drop(worker);
1204
1205 assert!(recv_fut.is_woken());
1207 assert_ready_eq!(recv_fut.poll(), Some(DataspaceUpdate::Retracted(id.into())));
1208 }
1209
1210 #[test]
1211 fn instrumented_process_retracts_multiple_types_on_drop() {
1212 let registry = DataspaceRegistry::new();
1213 let process = Process::supervisor_with_dataspace("test_worker", None, Some(registry.clone())).unwrap();
1214 let id_num = "number";
1215 let id_str = "text";
1216
1217 let mut sub_u32 = registry.subscribe::<u32>(IdentifierFilter::exact(id_num));
1218 let mut sub_str = registry.subscribe::<String>(IdentifierFilter::exact(id_str));
1219
1220 let mut recv_u32_fut = test_spawn(sub_u32.recv());
1222 let mut recv_str_fut = test_spawn(sub_str.recv());
1223 assert!(!recv_u32_fut.is_woken());
1224 assert!(!recv_str_fut.is_woken());
1225 assert_pending!(recv_u32_fut.poll());
1226 assert_pending!(recv_str_fut.poll());
1227
1228 let mut worker = test_spawn(process.into_process_future(async {
1230 let ds = DataspaceRegistry::try_current().expect("dataspace registry should be available");
1231 ds.assert(42u32, id_num);
1232 ds.assert("hello".to_string(), id_str);
1233 pending::<()>().await;
1234 }));
1235
1236 assert_pending!(worker.poll());
1238
1239 assert!(recv_u32_fut.is_woken());
1241 assert!(recv_str_fut.is_woken());
1242 assert_ready_eq!(recv_u32_fut.poll(), Some(DataspaceUpdate::Asserted(id_num.into(), 42)));
1243 assert_ready_eq!(
1244 recv_str_fut.poll(),
1245 Some(DataspaceUpdate::Asserted(id_str.into(), "hello".to_string()))
1246 );
1247
1248 drop(recv_u32_fut);
1249 drop(recv_str_fut);
1250
1251 let mut recv_u32 = test_spawn(sub_u32.recv());
1253 let mut recv_str = test_spawn(sub_str.recv());
1254 assert_pending!(recv_u32.poll());
1255 assert_pending!(recv_str.poll());
1256 assert!(!recv_u32.is_woken());
1257 assert!(!recv_str.is_woken());
1258
1259 drop(worker);
1261
1262 assert!(recv_u32.is_woken());
1263 assert!(recv_str.is_woken());
1264 assert_ready_eq!(recv_u32.poll(), Some(DataspaceUpdate::Retracted(id_num.into())));
1265 assert_ready_eq!(recv_str.poll(), Some(DataspaceUpdate::Retracted(id_str.into())));
1266 }
1267
1268 #[test]
1269 #[cfg(not(debug_assertions))]
1270 fn retract_from_non_owner_is_ignored() {
1271 let registry = DataspaceRegistry::new();
1272 let pid_a = Id::new();
1273 let pid_b = Id::new();
1274 let id = Identifier::named("owned_by_a");
1275
1276 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1278 registry.assert(42u32, id.clone());
1279 });
1280
1281 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_b, || {
1283 registry.retract::<u32>(id.clone());
1284 });
1285
1286 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
1288 let mut recv = test_spawn(sub.recv());
1289 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 42)));
1290 drop(recv);
1291
1292 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1294 registry.retract::<u32>(id.clone());
1295 });
1296
1297 let mut recv = test_spawn(sub.recv());
1299 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Retracted(id)));
1300 }
1301
1302 #[test]
1303 #[cfg(debug_assertions)]
1304 #[should_panic(expected = "attempted to retract assertion owned by")]
1305 fn retract_from_non_owner_panics_in_debug() {
1306 let registry = DataspaceRegistry::new();
1307 let pid_a = Id::new();
1308 let pid_b = Id::new();
1309 let id = Identifier::named("owned_by_a");
1310
1311 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1312 registry.assert(42u32, id.clone());
1313 });
1314
1315 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_b, || {
1316 registry.retract::<u32>(id.clone());
1317 });
1318 }
1319
1320 #[test]
1321 #[cfg(not(debug_assertions))]
1322 fn reassert_from_non_owner_is_ignored() {
1323 let registry = DataspaceRegistry::new();
1324 let pid_a = Id::new();
1325 let pid_b = Id::new();
1326 let id = Identifier::named("owned_by_a");
1327
1328 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1330 registry.assert(42u32, id.clone());
1331 });
1332
1333 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_b, || {
1335 registry.assert(99u32, id.clone());
1336 });
1337
1338 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
1340 let mut recv = test_spawn(sub.recv());
1341 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 42)));
1342 drop(recv);
1343
1344 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1346 registry.assert(100u32, id.clone());
1347 });
1348
1349 let mut recv = test_spawn(sub.recv());
1351 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id, 100)));
1352 }
1353
1354 #[test]
1355 #[cfg(debug_assertions)]
1356 #[should_panic(expected = "attempted to update assertion owned by")]
1357 fn reassert_from_non_owner_panics_in_debug() {
1358 let registry = DataspaceRegistry::new();
1359 let pid_a = Id::new();
1360 let pid_b = Id::new();
1361 let id = Identifier::named("owned_by_a");
1362
1363 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1364 registry.assert(42u32, id.clone());
1365 });
1366
1367 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_b, || {
1368 registry.assert(99u32, id.clone());
1369 });
1370 }
1371
1372 #[test]
1373 fn message_delivered_to_current_subscriber() {
1374 let registry = DataspaceRegistry::new();
1375 let id = Identifier::numeric(1);
1376
1377 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
1378 registry.send(7u32, id.clone());
1379
1380 let mut recv = test_spawn(sub.recv());
1381 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Message(id, 7)));
1382 }
1383
1384 #[test]
1385 fn message_not_replayed_to_late_subscriber() {
1386 let registry = DataspaceRegistry::new();
1387 let id = Identifier::numeric(1);
1388
1389 registry.send(7u32, id.clone());
1391
1392 assert!(registry.current_values::<u32>(IdentifierFilter::all()).is_empty());
1394
1395 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id));
1397 drop(registry);
1398
1399 let mut recv = test_spawn(sub.recv());
1400 assert_ready_eq!(recv.poll(), None);
1401 }
1402
1403 #[test]
1404 fn message_delivered_to_all_filter() {
1405 let registry = DataspaceRegistry::new();
1406 let id = Identifier::numeric(1);
1407
1408 let mut sub = registry.subscribe::<u32>(IdentifierFilter::all());
1409 registry.send(7u32, id.clone());
1410
1411 let mut recv = test_spawn(sub.recv());
1412 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Message(id, 7)));
1413 }
1414
1415 #[test]
1416 fn message_delivered_to_prefix_filter() {
1417 let registry = DataspaceRegistry::new();
1418 let matching = Identifier::named("svc.alpha");
1419 let non_matching = Identifier::named("other.beta");
1420
1421 let mut sub = registry.subscribe::<u32>(IdentifierFilter::prefix("svc."));
1422
1423 registry.send(1u32, non_matching);
1425 registry.send(2u32, matching.clone());
1426
1427 let mut recv = test_spawn(sub.recv());
1429 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Message(matching, 2)));
1430 }
1431
1432 #[test]
1433 fn subscription_interleaves_assert_message_retract_in_order() {
1434 let registry = DataspaceRegistry::new();
1435 let id = Identifier::numeric(1);
1436
1437 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
1438
1439 registry.assert(1u32, id.clone());
1440 registry.send(2u32, id.clone());
1441 registry.retract::<u32>(id.clone());
1442
1443 let mut recv = test_spawn(sub.recv());
1445 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 1)));
1446 drop(recv);
1447
1448 let mut recv = test_spawn(sub.recv());
1449 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Message(id.clone(), 2)));
1450 drop(recv);
1451
1452 let mut recv = test_spawn(sub.recv());
1453 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Retracted(id)));
1454 }
1455
1456 #[test]
1457 fn message_without_subscribers_is_noop() {
1458 let registry = DataspaceRegistry::new();
1459 let id = Identifier::numeric(1);
1460
1461 registry.send(7u32, id.clone());
1463 assert!(registry.current_values::<u32>(IdentifierFilter::all()).is_empty());
1464
1465 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id));
1467 drop(registry);
1468
1469 let mut recv = test_spawn(sub.recv());
1470 assert_ready_eq!(recv.poll(), None);
1471 }
1472
1473 #[test]
1474 fn message_does_not_clobber_current_assertion() {
1475 let registry = DataspaceRegistry::new();
1476 let id = Identifier::numeric(1);
1477
1478 let mut existing = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
1480
1481 registry.assert(1u32, id.clone());
1482 registry.send(2u32, id.clone());
1483
1484 assert_eq!(
1486 registry.current_values::<u32>(IdentifierFilter::exact(id.clone())),
1487 vec![1]
1488 );
1489
1490 let mut late = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
1492 let mut recv = test_spawn(late.recv());
1493 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 1)));
1494 drop(recv);
1495
1496 let mut recv = test_spawn(existing.recv());
1498 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 1)));
1499 drop(recv);
1500
1501 let mut recv = test_spawn(existing.recv());
1502 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Message(id, 2)));
1503 }
1504
1505 #[test]
1506 fn message_delivered_to_multiple_subscribers() {
1507 let registry = DataspaceRegistry::new();
1508 let id = Identifier::numeric(1);
1509
1510 let mut sub_exact = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
1512 let mut sub_all = registry.subscribe::<u32>(IdentifierFilter::all());
1513
1514 registry.send(7u32, id.clone());
1515
1516 let mut recv = test_spawn(sub_exact.recv());
1517 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Message(id.clone(), 7)));
1518
1519 let mut recv = test_spawn(sub_all.recv());
1520 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Message(id, 7)));
1521 }
1522
1523 #[test]
1524 fn send_from_non_owner_process_succeeds() {
1525 let registry = DataspaceRegistry::new();
1526 let pid_a = Id::new();
1527 let pid_b = Id::new();
1528 let id = Identifier::named("owned_by_a");
1529
1530 let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
1531
1532 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1534 registry.assert(1u32, id.clone());
1535 });
1536
1537 crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_b, || {
1540 registry.send(2u32, id.clone());
1541 });
1542
1543 let mut recv = test_spawn(sub.recv());
1545 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Asserted(id.clone(), 1)));
1546 drop(recv);
1547
1548 let mut recv = test_spawn(sub.recv());
1549 assert_ready_eq!(recv.poll(), Some(DataspaceUpdate::Message(id.clone(), 2)));
1550 drop(recv);
1551
1552 assert_eq!(registry.current_values::<u32>(IdentifierFilter::exact(id)), vec![1]);
1554 }
1555
1556 #[test]
1557 fn current_values_returns_point_in_time_snapshot() {
1558 let registry = DataspaceRegistry::new();
1562 let id_alpha = Identifier::named("svc.alpha");
1563 let id_beta = Identifier::named("svc.beta");
1564 let id_other = Identifier::named("other.gamma");
1565
1566 registry.assert(1u32, id_alpha.clone());
1567 registry.assert(2u32, id_beta.clone());
1568 registry.assert(3u32, id_other.clone());
1569 registry.assert("not a u32".to_string(), id_alpha.clone());
1571
1572 let mut prefixed = registry.current_values::<u32>(IdentifierFilter::prefix("svc."));
1574 prefixed.sort_unstable();
1575 assert_eq!(prefixed, vec![1, 2]);
1576
1577 let mut all = registry.current_values::<u32>(IdentifierFilter::all());
1579 all.sort_unstable();
1580 assert_eq!(all, vec![1, 2, 3]);
1581
1582 assert_eq!(
1584 registry.current_values::<u32>(IdentifierFilter::exact(id_beta.clone())),
1585 vec![2]
1586 );
1587
1588 assert_eq!(
1590 registry.current_values::<String>(IdentifierFilter::all()),
1591 vec!["not a u32".to_string()]
1592 );
1593
1594 registry.retract::<u32>(id_alpha);
1596 let mut after_retract = registry.current_values::<u32>(IdentifierFilter::all());
1597 after_retract.sort_unstable();
1598 assert_eq!(after_retract, vec![2, 3]);
1599 }
1600}