saluki_core/runtime/state/
dataspace.rs

1//! A type-erased, async-aware dataspace registry for inter-process coordination.
2//!
3//! The [`DataspaceRegistry`] allows processes to assert and retract typed values by identifier, send transient
4//! messages, and subscribe to receive notifications of all three. Multiple subscribers can observe the same updates.
5//!
6//! - **Assertion**: a value of type `T` becomes available, associated with a given identifier.
7//! - **Retraction**: the value of type `T` associated with a given identifier is withdrawn.
8//! - **Message**: a transient value of type `T` sent for a given identifier.
9//!
10//! Subscribers can listen for updates matching a specific identifier, a prefix, or all identifiers for a given type.
11//!
12//! Assertions are *persistent*: they are stored, tied to the asserting process's lifecycle (automatically retracted
13//! when that process exits), and replayed to subscribers that appear later. Messages are *transient*: they are
14//! delivered only to the subscribers present at the time of sending, and are never stored or replayed.
15//!
16//! This enables decoupled coordination where processes don't need to know about each other, only the identifier and
17//! type of the values they're exchanging.
18//!
19//! # Example
20//!
21//! ```
22//! use saluki_core::runtime::state::{DataspaceUpdate, Identifier, IdentifierFilter, DataspaceRegistry};
23//!
24//! # #[tokio::main]
25//! # async fn main() {
26//! let registry = DataspaceRegistry::new();
27//! let id = Identifier::named("my_value");
28//!
29//! // Subscribe before asserting:
30//! let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
31//!
32//! // Assert a value:
33//! registry.assert(42u32, id.clone());
34//!
35//! // Receive the assertion:
36//! let value = sub.recv().await;
37//! assert_eq!(value, Some(DataspaceUpdate::Asserted(id, 42)));
38//! # }
39//! ```
40
41use 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/// An update received by a subscription, indicating that a value was asserted, retracted, or sent as a transient
60/// message.
61#[derive(Clone, Debug, PartialEq, Eq)]
62pub enum DataspaceUpdate<T> {
63    /// A value was asserted (made available), along with the identifier it's associated with.
64    Asserted(Identifier, T),
65
66    /// The value associated with the given identifier was retracted (withdrawn).
67    Retracted(Identifier),
68
69    /// A transient message was sent for the given identifier.
70    ///
71    /// Unlike assertions, messages are delivered only to subscribers present at the time the message is sent; they are
72    /// never stored or replayed to future subscribers.
73    Message(Identifier, T),
74}
75
76/// Internal key combining type and identifier.
77#[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
92/// A type-erased broadcast channel that supports sending retraction notifications without knowing `T`.
93trait AnyChannel: Send + Sync {
94    /// Sends a retraction for the given identifier.
95    fn send_retraction(&self, id: &Identifier);
96
97    /// Returns a downcasted reference to `self` that can be fallibly upcasted back to the original type.
98    fn as_any(&self) -> &dyn Any;
99}
100
101/// Concrete implementation of [`AnyChannel`] that wraps a `broadcast::Sender<DataspaceUpdate<T>>`.
102struct 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
116/// A type-erased filtered subscription channel.
117struct FilteredChannel {
118    type_id: TypeId,
119    filter: IdentifierFilter,
120    sender: Box<dyn AnyChannel>,
121}
122
123/// A stored assertion value along with the process that owns it.
124struct StoredValue {
125    value: Box<dyn Any + Send + Sync>,
126    owner: Id,
127}
128
129/// Internal registry state protected by a mutex.
130struct RegistryState {
131    /// Broadcast senders for exact-match subscriptions, keyed by (type, identifier).
132    channels: HashMap<StorageKey, Box<dyn AnyChannel>>,
133
134    /// Filtered subscription channels that are evaluated on every assert/retract.
135    filtered_channels: Vec<FilteredChannel>,
136
137    /// Current assertion values for each (type, identifier) pair, stored type-erased.
138    ///
139    /// Entries are removed on retraction. Used to replay current state to new subscribers.
140    current_values: HashMap<StorageKey, StoredValue>,
141
142    /// Tracks which storage keys each process has asserted, for automatic retraction on process exit.
143    process_assertions: HashMap<Id, HashSet<StorageKey>>,
144
145    /// Default capacity for new broadcast channels.
146    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    /// Gets or creates a broadcast sender for the given key, returning a new receiver.
161    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            // `StorageKey` includes `TypeId::of::<T>()`, so a channel stored under a
174            // given key is always a `TypedChannel<T>` for the same `T`.
175            .unwrap_or_else(|| unreachable!("type mismatch in dataspace registry"));
176
177        typed.tx.subscribe()
178    }
179
180    /// Creates a new filtered subscription channel, returning a new receiver.
181    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    /// Sends the given typed update on all channels (exact + filtered) matching the given key.
197    ///
198    /// Used to dispatch assertions and messages, both of which carry a concrete value `T`. The update is cloned once
199    /// per matching channel, as each channel is a separate broadcast sender.
200    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    /// Sends a retraction notification on all channels (exact + filtered) matching the given key.
220    ///
221    /// This is kept separate from [`notify_typed`](Self::notify_typed) because retractions carry no value: they are
222    /// dispatched via the type-erased [`AnyChannel::send_retraction`] path. That path is required by
223    /// [`retract_all_for_process`](DataspaceRegistry::retract_all_for_process), which iterates storage keys that carry
224    /// only a `TypeId`, with no concrete `T` in scope to downcast against.
225    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
238/// Shared inner state of the registry.
239struct DataspaceRegistryInner {
240    state: Mutex<RegistryState>,
241}
242
243/// A dataspace registry for async coordination between processes.
244///
245/// The registry stores broadcast channels indexed by type and [`Identifier`]. Processes can subscribe to receive
246/// assertion and retraction updates for a given type and identifier filter, and other processes can assert or retract
247/// values that are delivered to all matching subscribers.
248///
249/// # Thread Safety
250///
251/// `DataspaceRegistry` is `Clone` and can be safely shared across threads and tasks. All operations are thread-safe.
252#[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    /// Creates a new empty registry with the default channel capacity.
265    pub fn new() -> Self {
266        Self::with_channel_capacity(DEFAULT_CHANNEL_CAPACITY)
267    }
268
269    /// Returns the dataspace registry for the current supervision tree, if one exists.
270    pub fn try_current() -> Option<Self> {
271        CURRENT_DATASPACE.try_with(|ds| ds.clone()).ok()
272    }
273
274    /// Runs the given closure with this dataspace set as the current dataspace.
275    ///
276    /// This can be used to override the dataspace that gets returned in calls to
277    /// [`DataspaceRegistry::try_current`][Self::try_current], which can be difficult to achieve otherwise without fully
278    /// executing code under supervision.
279    #[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    /// Creates a new empty registry with the given channel capacity for broadcast channels.
288    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    /// Asserts a value with the given identifier, notifying all matching subscribers.
297    ///
298    /// The assertion is automatically associated with the current process, and will be automatically retracted when
299    /// that process exists. Only the owning process may update an existing assertion for a given type/identifier
300    /// combination.
301    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 an assertion already exists for this key, only the owning process may update it.
311        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        // Store the current value for future subscribers, along with the owning process.
323        state.current_values.insert(
324            key.clone(),
325            StoredValue {
326                value: Box::new(value.clone()),
327                owner: caller,
328            },
329        );
330
331        // Track this assertion against the owning process.
332        state.process_assertions.entry(caller).or_default().insert(key.clone());
333
334        // Notify all matching subscribers, exact and filtered.
335        let update = DataspaceUpdate::Asserted(id, value);
336        state.notify_typed::<T>(&key, &update);
337    }
338
339    /// Retracts the value of the given type and identifier, notifying all matching subscribers.
340    ///
341    /// Only the process that originally asserted the value may retract it.
342    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        // Check that the assertion exists and is owned by the calling process.
352        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        // Remove the stored value and clean up process tracking.
366        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    /// Sends a transient message with the given identifier to all matching subscribers.
379    ///
380    /// Unlike [`assert`](Self::assert), only the _current_ matching subscribers are notified: messages are never stored
381    /// or replayed. If not matching subscribers exist, the message is dropped.
382    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        // Messages are transient: notify only the subscribers matching right now, without storing anything.
391        let update = DataspaceUpdate::Message(id, value);
392        state.notify_typed::<T>(&key, &update);
393    }
394
395    /// Retracts all assertions made by the given process.
396    ///
397    /// This is called automatically when a process exits (via [`ProcessFuture`](crate::runtime::process::ProcessFuture)
398    /// drop) to ensure that no stale assertions remain in the registry after the owning process is gone.
399    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    /// This is a synchronous point-in-time read that locks the registry, collects matching values,
413    /// and returns immediately without creating any subscription or channel. Use this when you need
414    /// a snapshot of what is currently asserted rather than ongoing notifications of future changes.
415    ///
416    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    /// Subscribes to updates matching the given filter.
431    ///
432    /// Returns a [`Subscription`] that can be used to asynchronously receive updates. Any updates to assertions that
433    /// match the filter at the time of subscribing will be immediately replayed into the subscription's pending queue.
434    /// Messages are never replayed, so only those sent after subscribing are delivered.
435    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                // Replay current value if one exists.
447                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                // Replay all matching current values.
461                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
480/// A subscription to updates for a specific type/identifier combination.
481pub 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    /// Receives the next assertion, retraction, or message update.
491    ///
492    /// Returns `Some(update)` when an update is available, or `None` if the channel has been closed (all senders
493    /// dropped). If updates were missed due to the subscriber falling behind, the missed updates are skipped and the
494    /// next available update is returned.
495    pub async fn recv(&mut self) -> Option<DataspaceUpdate<T>> {
496        // TODO: Switch to bounded MPSC channels for delivering updates.
497
498        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        // Assert before any subscriber exists -- the value is stored.
556        registry.assert(42u32, id.clone());
557
558        // A later subscriber should receive the current value.
559        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 the registry, which drops the Arc. Since we only have one reference, the sender is dropped.
606        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        // Create a registry with a tiny buffer so we can force lag.
615        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        // Assert more values than the channel can hold.
621        for i in 0..10 {
622            registry.assert(i as u32, id.clone());
623        }
624
625        // The subscriber should skip lagged messages and still receive a value.
626        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        // Retract without any subscribers -- should not panic.
678        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        // Drain assertion notifications.
693        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        // Check retraction notifications.
702        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        // Subscribe to all identifiers before any specific-identifier activity exists.
791        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        // Assert before subscribing.
806        registry.assert(42u32, id.clone());
807
808        // Subscribe after asserting -- should immediately get the current value.
809        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        // Assert then retract -- no current value.
820        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 the registry so the channel closes, proving no pending value is delivered.
826        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        // Assert on two identifiers before subscribing.
839        registry.assert(1u32, id1.clone());
840        registry.assert(2u32, id2.clone());
841
842        let mut sub = registry.subscribe::<u32>(IdentifierFilter::all());
843
844        // Should receive both current values (order is not guaranteed since HashMap iteration is unordered).
845        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        // Assert before subscribing.
868        registry.assert(1u32, id.clone());
869
870        let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
871
872        // Assert again after subscribing.
873        registry.assert(2u32, id.clone());
874
875        // First recv should return the initial value, second should return the broadcast value.
876        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        // Should only receive the two metrics identifiers.
898        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 to close channel, proving no value was delivered.
925        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        // Assert before subscribing.
939        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        // Should replay only the two matching values.
946        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            // Verify it's the same registry by asserting in one and reading from the other.
978            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        // Assert two values as if from the given process.
997        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        // Subscribe to both after assertion to get current values replayed.
1003        let mut sub1 = registry.subscribe::<u32>(IdentifierFilter::exact(id1.clone()));
1004        let mut sub2 = registry.subscribe::<u32>(IdentifierFilter::exact(id2.clone()));
1005
1006        // Drain the initial replayed values.
1007        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        // Retract all assertions for the process.
1016        registry.retract_all_for_process(process_id);
1017
1018        // Both subscribers should receive retraction notifications.
1019        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        // Retract only process A's assertions.
1043        registry.retract_all_for_process(pid_a);
1044
1045        // Process B's value should still be available to new subscribers.
1046        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        // Process A's value should not be available.
1052        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        // Receive the assertion.
1071        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        // Retract all for the process.
1076        registry.retract_all_for_process(process_id);
1077
1078        // Should receive the retraction on the wildcard subscription.
1079        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        // Manually retract (must be from the owning process).
1096        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(process_id, || {
1097            registry.retract::<u32>(id.clone());
1098        });
1099
1100        // Drain the assertion and retraction.
1101        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        // Now retract_all_for_process should be a no-op (no duplicate retraction).
1110        registry.retract_all_for_process(process_id);
1111
1112        // Drop the registry to close the channel, proving no further messages are pending.
1113        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        // Should not panic.
1122        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        // Spawn a subscriber future, asserting that it is quiesced.
1134        let mut recv_fut = test_spawn(sub.recv());
1135        assert!(!recv_fut.is_woken());
1136        assert_pending!(recv_fut.poll());
1137
1138        // Create an instrumented future that asserts a value, then completes.
1139        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        // Poll the worker to completion — the assertion happens during this poll.
1146        assert_ready!(worker.poll());
1147
1148        // The subscriber should now be woken and have the assertion update.
1149        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        // Set up a new subscriber call for getting the retraction.
1155        let mut recv_fut = test_spawn(sub.recv());
1156        assert_pending!(recv_fut.poll());
1157        assert!(!recv_fut.is_woken());
1158
1159        // Drop the InstrumentedProcess — this triggers automatic retraction.
1160        drop(worker);
1161
1162        // The drop should have woken the subscriber.
1163        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        // Spawn a subscriber future, asserting that it is quiesced.
1176        let mut recv_fut = test_spawn(sub.recv());
1177        assert!(!recv_fut.is_woken());
1178        assert_pending!(recv_fut.poll());
1179
1180        // Create an instrumented future that asserts a value, then pends forever.
1181        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        // Drive the worker to assert the value.
1189        assert_pending!(worker.poll());
1190
1191        // The subscriber should now be woken and have the assertion update.
1192        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        // Set up a new subscriber call for getting the retraction.
1198        let mut recv_fut = test_spawn(sub.recv());
1199        assert_pending!(recv_fut.poll());
1200        assert!(!recv_fut.is_woken());
1201
1202        // Drop the worker (simulates abort/cancel) — triggers automatic retraction.
1203        drop(worker);
1204
1205        // The drop should have woken the subscriber with a retraction.
1206        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        // Spawn futures for both subscribers, asserting that they are quiesced.
1221        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        // Create an instrumented future that asserts values of two different types.
1229        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        // Drive the worker to assert the values.
1237        assert_pending!(worker.poll());
1238
1239        // Both subscribers should now be woken and have the assertion updates.
1240        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        // Set up new subscriber calls for getting the retractions.
1252        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 the worker — both types should be retracted.
1260        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        // Assert as process A.
1277        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1278            registry.assert(42u32, id.clone());
1279        });
1280
1281        // Attempt to retract as process B -- should be silently ignored.
1282        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_b, || {
1283            registry.retract::<u32>(id.clone());
1284        });
1285
1286        // Value should still be present for new subscribers.
1287        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        // Retract as process A -- should succeed.
1293        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1294            registry.retract::<u32>(id.clone());
1295        });
1296
1297        // Subscriber should receive the retraction.
1298        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        // Assert as process A.
1329        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1330            registry.assert(42u32, id.clone());
1331        });
1332
1333        // Attempt to overwrite as process B -- should be silently ignored.
1334        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_b, || {
1335            registry.assert(99u32, id.clone());
1336        });
1337
1338        // Value should still be the original.
1339        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        // Update as process A -- should succeed.
1345        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1346            registry.assert(100u32, id.clone());
1347        });
1348
1349        // Subscriber should receive the update.
1350        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        // Send a message before anyone subscribes -- it should be dropped, not stored.
1390        registry.send(7u32, id.clone());
1391
1392        // The message must not have been stored as a current value.
1393        assert!(registry.current_values::<u32>(IdentifierFilter::all()).is_empty());
1394
1395        // A subscriber that turns up afterward receives nothing; dropping the registry closes the channel.
1396        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        // A message to a non-matching identifier must not be delivered; a matching one must be.
1424        registry.send(1u32, non_matching);
1425        registry.send(2u32, matching.clone());
1426
1427        // The first (and only) update received is the matching message, proving the non-matching one was filtered out.
1428        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        // A single subscription observes all three kinds of update, in send order, over the same channel.
1444        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        // Sending with no subscribers must not panic and must not store anything.
1462        registry.send(7u32, id.clone());
1463        assert!(registry.current_values::<u32>(IdentifierFilter::all()).is_empty());
1464
1465        // A later subscriber receives nothing (no replay); dropping the registry closes the channel.
1466        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        // A pre-existing subscriber, to observe live updates in order.
1479        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        // The stored current value is still the assertion, unaffected by the message.
1485        assert_eq!(
1486            registry.current_values::<u32>(IdentifierFilter::exact(id.clone())),
1487            vec![1]
1488        );
1489
1490        // A new subscriber replays only the assertion (messages are never replayed).
1491        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        // The pre-existing subscriber sees the assertion, then the message.
1497        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        // One exact and one wildcard subscriber, to exercise both dispatch paths.
1511        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        // Process A asserts a value.
1533        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1534            registry.assert(1u32, id.clone());
1535        });
1536
1537        // Process B sends a message for the same type/identifier -- messages bypass the ownership check and do not
1538        // panic, even in debug builds.
1539        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_b, || {
1540            registry.send(2u32, id.clone());
1541        });
1542
1543        // The subscriber receives A's assertion, then B's message.
1544        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        // The assertion's stored value is untouched by the message.
1553        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        // `current_values` is a synchronous, point-in-time read: it collects the values currently asserted for the
1559        // given type that match the filter, without creating any subscription. Exercise its documented semantics
1560        // directly (type separation, each filter kind, and retraction being reflected in the next snapshot).
1561        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        // A value of a *different* type under a matching identifier must not leak into the `u32` snapshot.
1570        registry.assert("not a u32".to_string(), id_alpha.clone());
1571
1572        // Prefix filter: only the two `svc.`-prefixed `u32` values (order is unspecified, so sort before comparing).
1573        let mut prefixed = registry.current_values::<u32>(IdentifierFilter::prefix("svc."));
1574        prefixed.sort_unstable();
1575        assert_eq!(prefixed, vec![1, 2]);
1576
1577        // All filter: every `u32` value regardless of identifier.
1578        let mut all = registry.current_values::<u32>(IdentifierFilter::all());
1579        all.sort_unstable();
1580        assert_eq!(all, vec![1, 2, 3]);
1581
1582        // Exact filter: just the single matching value.
1583        assert_eq!(
1584            registry.current_values::<u32>(IdentifierFilter::exact(id_beta.clone())),
1585            vec![2]
1586        );
1587
1588        // The `String` value is visible only under its own type, confirming type-keyed separation.
1589        assert_eq!(
1590            registry.current_values::<String>(IdentifierFilter::all()),
1591            vec!["not a u32".to_string()]
1592        );
1593
1594        // A retraction is reflected in the *next* snapshot, since each call is an independent point-in-time read.
1595        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}