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