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//! <!-- vale off -->
22//! ```
23//! use saluki_core::runtime::state::{DataspaceUpdate, Identifier, IdentifierFilter, DataspaceRegistry};
24//!
25//! # #[tokio::main]
26//! # async fn main() {
27//! let registry = DataspaceRegistry::new();
28//! let id = Identifier::named("my_value");
29//!
30//! // Subscribe before asserting:
31//! let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
32//!
33//! // Assert a value:
34//! registry.assert(42u32, id.clone());
35//!
36//! // Receive the assertion:
37//! let value = sub.recv().await;
38//! assert_eq!(value, Some(DataspaceUpdate::Asserted(id, 42)));
39//! # }
40//! ```
41//! <!-- vale on -->
42
43use 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/// An update received by a subscription, indicating that a value was asserted, retracted, or sent as a transient
62/// message.
63#[derive(Clone, Debug, PartialEq, Eq)]
64pub enum DataspaceUpdate<T> {
65    /// A value was asserted (made available), along with the identifier it's associated with.
66    Asserted(Identifier, T),
67
68    /// The value associated with the given identifier was retracted (withdrawn).
69    Retracted(Identifier),
70
71    /// A transient message was sent for the given identifier.
72    ///
73    /// Unlike assertions, messages are delivered only to subscribers present at the time the message is sent; they are
74    /// never stored or replayed to future subscribers.
75    Message(Identifier, T),
76}
77
78/// Internal key combining type and identifier.
79#[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
94/// A type-erased broadcast channel that supports sending retraction notifications without knowing `T`.
95trait AnyChannel: Send + Sync {
96    /// Sends a retraction for the given identifier.
97    fn send_retraction(&self, id: &Identifier);
98
99    /// Returns a downcasted reference to `self` that can be fallibly upcasted back to the original type.
100    fn as_any(&self) -> &dyn Any;
101}
102
103/// Concrete implementation of [`AnyChannel`] that wraps a `broadcast::Sender<DataspaceUpdate<T>>`.
104struct 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
118/// A type-erased filtered subscription channel.
119struct FilteredChannel {
120    type_id: TypeId,
121    filter: IdentifierFilter,
122    sender: Box<dyn AnyChannel>,
123}
124
125/// A stored assertion value along with the process that owns it.
126struct StoredValue {
127    value: Box<dyn Any + Send + Sync>,
128    owner: Id,
129}
130
131/// Internal registry state protected by a mutex.
132struct RegistryState {
133    /// Broadcast senders for exact-match subscriptions, keyed by (type, identifier).
134    channels: HashMap<StorageKey, Box<dyn AnyChannel>>,
135
136    /// Filtered subscription channels that are evaluated on every assert/retract.
137    filtered_channels: Vec<FilteredChannel>,
138
139    /// Current assertion values for each (type, identifier) pair, stored type-erased.
140    ///
141    /// Entries are removed on retraction. Used to replay current state to new subscribers.
142    current_values: HashMap<StorageKey, StoredValue>,
143
144    /// Tracks which storage keys each process has asserted, for automatic retraction on process exit.
145    process_assertions: HashMap<Id, HashSet<StorageKey>>,
146
147    /// Default capacity for new broadcast channels.
148    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    /// Gets or creates a broadcast sender for the given key, returning a new receiver.
163    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            // `StorageKey` includes `TypeId::of::<T>()`, so a channel stored under a
176            // given key is always a `TypedChannel<T>` for the same `T`.
177            .unwrap_or_else(|| unreachable!("type mismatch in dataspace registry"));
178
179        typed.tx.subscribe()
180    }
181
182    /// Creates a new filtered subscription channel, returning a new receiver.
183    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    /// Sends the given typed update on all channels (exact + filtered) matching the given key.
199    ///
200    /// Used to dispatch assertions and messages, both of which carry a concrete value `T`. The update is cloned once
201    /// per matching channel, as each channel is a separate broadcast sender.
202    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    /// Sends a retraction notification on all channels (exact + filtered) matching the given key.
222    ///
223    /// This is kept separate from [`notify_typed`](Self::notify_typed) because retractions carry no value: they are
224    /// dispatched via the type-erased [`AnyChannel::send_retraction`] path. That path is required by
225    /// [`retract_all_for_process`](DataspaceRegistry::retract_all_for_process), which iterates storage keys that carry
226    /// only a `TypeId`, with no concrete `T` in scope to downcast against.
227    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
240/// Shared inner state of the registry.
241struct DataspaceRegistryInner {
242    state: Mutex<RegistryState>,
243}
244
245/// A dataspace registry for async coordination between processes.
246///
247/// The registry stores broadcast channels indexed by type and [`Identifier`]. Processes can subscribe to receive
248/// assertion and retraction updates for a given type and identifier filter, and other processes can assert or retract
249/// values that are delivered to all matching subscribers.
250///
251/// # Thread Safety
252///
253/// `DataspaceRegistry` is `Clone` and can be safely shared across threads and tasks. All operations are thread-safe.
254#[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    /// Creates a new empty registry with the default channel capacity.
267    pub fn new() -> Self {
268        Self::with_channel_capacity(DEFAULT_CHANNEL_CAPACITY)
269    }
270
271    /// Returns the dataspace registry for the current supervision tree, if one exists.
272    pub fn try_current() -> Option<Self> {
273        CURRENT_DATASPACE.try_with(|ds| ds.clone()).ok()
274    }
275
276    /// Runs the given closure with this dataspace set as the current dataspace.
277    ///
278    /// This can be used to override the dataspace that gets returned in calls to
279    /// [`DataspaceRegistry::try_current`][Self::try_current], which can be difficult to achieve otherwise without fully
280    /// executing code under supervision.
281    #[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    /// Creates a new empty registry with the given channel capacity for broadcast channels.
290    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    /// Asserts a value with the given identifier, notifying all matching subscribers.
299    ///
300    /// The assertion is automatically associated with the current process, and will be automatically retracted when
301    /// that process exists. Only the owning process may update an existing assertion for a given type/identifier
302    /// combination.
303    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 an assertion already exists for this key, only the owning process may update it.
313        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        // Store the current value for future subscribers, along with the owning process.
325        state.current_values.insert(
326            key.clone(),
327            StoredValue {
328                value: Box::new(value.clone()),
329                owner: caller,
330            },
331        );
332
333        // Track this assertion against the owning process.
334        state.process_assertions.entry(caller).or_default().insert(key.clone());
335
336        // Notify all matching subscribers, exact and filtered.
337        let update = DataspaceUpdate::Asserted(id, value);
338        state.notify_typed::<T>(&key, &update);
339    }
340
341    /// Retracts the value of the given type and identifier, notifying all matching subscribers.
342    ///
343    /// Only the process that originally asserted the value may retract it.
344    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        // Check that the assertion exists and is owned by the calling process.
354        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        // Remove the stored value and clean up process tracking.
368        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    /// Sends a transient message with the given identifier to all matching subscribers.
381    ///
382    /// Unlike [`assert`](Self::assert), only the _current_ matching subscribers are notified: messages are never stored
383    /// or replayed. If not matching subscribers exist, the message is dropped.
384    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        // Messages are transient: notify only the subscribers matching right now, without storing anything.
393        let update = DataspaceUpdate::Message(id, value);
394        state.notify_typed::<T>(&key, &update);
395    }
396
397    /// Retracts all assertions made by the given process.
398    ///
399    /// This is called automatically when a process exits (via [`ProcessFuture`](crate::runtime::process::ProcessFuture)
400    /// drop) to ensure that no stale assertions remain in the registry after the owning process is gone.
401    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    /// This is a synchronous point-in-time read that locks the registry, collects matching values,
415    /// and returns immediately without creating any subscription or channel. Use this when you need
416    /// a snapshot of what is currently asserted rather than ongoing notifications of future changes.
417    ///
418    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    /// Subscribes to updates matching the given filter.
433    ///
434    /// Returns a [`Subscription`] that can be used to asynchronously receive updates. Any updates to assertions that
435    /// match the filter at the time of subscribing will be immediately replayed into the subscription's pending queue.
436    /// Messages are never replayed, so only those sent after subscribing are delivered.
437    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                // Replay current value if one exists.
449                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                // Replay all matching current values.
463                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
482/// A subscription to updates for a specific type/identifier combination.
483pub 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    /// Receives the next assertion, retraction, or message update.
493    ///
494    /// Returns `Some(update)` when an update is available, or `None` if the channel has been closed (all senders
495    /// dropped). If updates were missed due to the subscriber falling behind, the missed updates are skipped and the
496    /// next available update is returned.
497    pub async fn recv(&mut self) -> Option<DataspaceUpdate<T>> {
498        // TODO: Switch to bounded MPSC channels for delivering updates.
499
500        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        // Assert before any subscriber exists -- the value is stored.
558        registry.assert(42u32, id.clone());
559
560        // A later subscriber should receive the current value.
561        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 the registry, which drops the Arc. Since we only have one reference, the sender is dropped.
608        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        // Create a registry with a tiny buffer so we can force lag.
617        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        // Assert more values than the channel can hold.
623        for i in 0..10 {
624            registry.assert(i as u32, id.clone());
625        }
626
627        // The subscriber should skip lagged messages and still receive a value.
628        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        // Retract without any subscribers -- should not panic.
680        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        // Drain assertion notifications.
695        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        // Check retraction notifications.
704        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        // Subscribe to all identifiers before any specific-identifier activity exists.
793        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        // Assert before subscribing.
808        registry.assert(42u32, id.clone());
809
810        // Subscribe after asserting -- should immediately get the current value.
811        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        // Assert then retract -- no current value.
822        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 the registry so the channel closes, proving no pending value is delivered.
828        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        // Assert on two identifiers before subscribing.
841        registry.assert(1u32, id1.clone());
842        registry.assert(2u32, id2.clone());
843
844        let mut sub = registry.subscribe::<u32>(IdentifierFilter::all());
845
846        // Should receive both current values (order is not guaranteed since HashMap iteration is unordered).
847        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        // Assert before subscribing.
870        registry.assert(1u32, id.clone());
871
872        let mut sub = registry.subscribe::<u32>(IdentifierFilter::exact(id.clone()));
873
874        // Assert again after subscribing.
875        registry.assert(2u32, id.clone());
876
877        // First recv should return the initial value, second should return the broadcast value.
878        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        // Should only receive the two metrics identifiers.
900        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 to close channel, proving no value was delivered.
927        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        // Assert before subscribing.
941        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        // Should replay only the two matching values.
948        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            // Verify it's the same registry by asserting in one and reading from the other.
980            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        // Assert two values as if from the given process.
999        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        // Subscribe to both after assertion to get current values replayed.
1005        let mut sub1 = registry.subscribe::<u32>(IdentifierFilter::exact(id1.clone()));
1006        let mut sub2 = registry.subscribe::<u32>(IdentifierFilter::exact(id2.clone()));
1007
1008        // Drain the initial replayed values.
1009        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        // Retract all assertions for the process.
1018        registry.retract_all_for_process(process_id);
1019
1020        // Both subscribers should receive retraction notifications.
1021        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        // Retract only process A's assertions.
1045        registry.retract_all_for_process(pid_a);
1046
1047        // Process B's value should still be available to new subscribers.
1048        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        // Process A's value should not be available.
1054        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        // Receive the assertion.
1073        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        // Retract all for the process.
1078        registry.retract_all_for_process(process_id);
1079
1080        // Should receive the retraction on the wildcard subscription.
1081        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        // Manually retract (must be from the owning process).
1098        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(process_id, || {
1099            registry.retract::<u32>(id.clone());
1100        });
1101
1102        // Drain the assertion and retraction.
1103        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        // Now retract_all_for_process should be a no-op (no duplicate retraction).
1112        registry.retract_all_for_process(process_id);
1113
1114        // Drop the registry to close the channel, proving no further messages are pending.
1115        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        // Should not panic.
1124        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        // Spawn a subscriber future, asserting that it is quiesced.
1136        let mut recv_fut = test_spawn(sub.recv());
1137        assert!(!recv_fut.is_woken());
1138        assert_pending!(recv_fut.poll());
1139
1140        // Create an instrumented future that asserts a value, then completes.
1141        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        // Poll the worker to completion — the assertion happens during this poll.
1148        assert_ready!(worker.poll());
1149
1150        // The subscriber should now be woken and have the assertion update.
1151        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        // Set up a new subscriber call for getting the retraction.
1157        let mut recv_fut = test_spawn(sub.recv());
1158        assert_pending!(recv_fut.poll());
1159        assert!(!recv_fut.is_woken());
1160
1161        // Drop the InstrumentedProcess — this triggers automatic retraction.
1162        drop(worker);
1163
1164        // The drop should have woken the subscriber.
1165        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        // Spawn a subscriber future, asserting that it is quiesced.
1178        let mut recv_fut = test_spawn(sub.recv());
1179        assert!(!recv_fut.is_woken());
1180        assert_pending!(recv_fut.poll());
1181
1182        // Create an instrumented future that asserts a value, then pends forever.
1183        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        // Drive the worker to assert the value.
1191        assert_pending!(worker.poll());
1192
1193        // The subscriber should now be woken and have the assertion update.
1194        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        // Set up a new subscriber call for getting the retraction.
1200        let mut recv_fut = test_spawn(sub.recv());
1201        assert_pending!(recv_fut.poll());
1202        assert!(!recv_fut.is_woken());
1203
1204        // Drop the worker (simulates abort/cancel) — triggers automatic retraction.
1205        drop(worker);
1206
1207        // The drop should have woken the subscriber with a retraction.
1208        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        // Spawn futures for both subscribers, asserting that they are quiesced.
1223        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        // Create an instrumented future that asserts values of two different types.
1231        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        // Drive the worker to assert the values.
1239        assert_pending!(worker.poll());
1240
1241        // Both subscribers should now be woken and have the assertion updates.
1242        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        // Set up new subscriber calls for getting the retractions.
1254        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 the worker — both types should be retracted.
1262        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        // Assert as process A.
1279        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1280            registry.assert(42u32, id.clone());
1281        });
1282
1283        // Attempt to retract as process B -- should be silently ignored.
1284        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_b, || {
1285            registry.retract::<u32>(id.clone());
1286        });
1287
1288        // Value should still be present for new subscribers.
1289        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        // Retract as process A -- should succeed.
1295        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1296            registry.retract::<u32>(id.clone());
1297        });
1298
1299        // Subscriber should receive the retraction.
1300        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        // Assert as process A.
1331        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1332            registry.assert(42u32, id.clone());
1333        });
1334
1335        // Attempt to overwrite as process B -- should be silently ignored.
1336        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_b, || {
1337            registry.assert(99u32, id.clone());
1338        });
1339
1340        // Value should still be the original.
1341        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        // Update as process A -- should succeed.
1347        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1348            registry.assert(100u32, id.clone());
1349        });
1350
1351        // Subscriber should receive the update.
1352        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        // Send a message before anyone subscribes -- it should be dropped, not stored.
1392        registry.send(7u32, id.clone());
1393
1394        // The message must not have been stored as a current value.
1395        assert!(registry.current_values::<u32>(IdentifierFilter::all()).is_empty());
1396
1397        // A subscriber that turns up afterward receives nothing; dropping the registry closes the channel.
1398        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        // A message to a non-matching identifier must not be delivered; a matching one must be.
1426        registry.send(1u32, non_matching);
1427        registry.send(2u32, matching.clone());
1428
1429        // The first (and only) update received is the matching message, proving the non-matching one was filtered out.
1430        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        // A single subscription observes all three kinds of update, in send order, over the same channel.
1446        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        // Sending with no subscribers must not panic and must not store anything.
1464        registry.send(7u32, id.clone());
1465        assert!(registry.current_values::<u32>(IdentifierFilter::all()).is_empty());
1466
1467        // A later subscriber receives nothing (no replay); dropping the registry closes the channel.
1468        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        // A pre-existing subscriber, to observe live updates in order.
1481        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        // The stored current value is still the assertion, unaffected by the message.
1487        assert_eq!(
1488            registry.current_values::<u32>(IdentifierFilter::exact(id.clone())),
1489            vec![1]
1490        );
1491
1492        // A new subscriber replays only the assertion (messages are never replayed).
1493        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        // The pre-existing subscriber sees the assertion, then the message.
1499        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        // One exact and one wildcard subscriber, to exercise both dispatch paths.
1513        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        // Process A asserts a value.
1535        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_a, || {
1536            registry.assert(1u32, id.clone());
1537        });
1538
1539        // Process B sends a message for the same type/identifier -- messages bypass the ownership check and do not
1540        // panic, even in debug builds.
1541        crate::runtime::process::CURRENT_PROCESS_ID.sync_scope(pid_b, || {
1542            registry.send(2u32, id.clone());
1543        });
1544
1545        // The subscriber receives A's assertion, then B's message.
1546        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        // The assertion's stored value is untouched by the message.
1555        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        // `current_values` is a synchronous, point-in-time read: it collects the values currently asserted for the
1561        // given type that match the filter, without creating any subscription. Exercise its documented semantics
1562        // directly (type separation, each filter kind, and retraction being reflected in the next snapshot).
1563        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        // A value of a *different* type under a matching identifier must not leak into the `u32` snapshot.
1572        registry.assert("not a u32".to_string(), id_alpha.clone());
1573
1574        // Prefix filter: only the two `svc.`-prefixed `u32` values (order is unspecified, so sort before comparing).
1575        let mut prefixed = registry.current_values::<u32>(IdentifierFilter::prefix("svc."));
1576        prefixed.sort_unstable();
1577        assert_eq!(prefixed, vec![1, 2]);
1578
1579        // All filter: every `u32` value regardless of identifier.
1580        let mut all = registry.current_values::<u32>(IdentifierFilter::all());
1581        all.sort_unstable();
1582        assert_eq!(all, vec![1, 2, 3]);
1583
1584        // Exact filter: just the single matching value.
1585        assert_eq!(
1586            registry.current_values::<u32>(IdentifierFilter::exact(id_beta.clone())),
1587            vec![2]
1588        );
1589
1590        // The `String` value is visible only under its own type, confirming type-keyed separation.
1591        assert_eq!(
1592            registry.current_values::<String>(IdentifierFilter::all()),
1593            vec!["not a u32".to_string()]
1594        );
1595
1596        // A retraction is reflected in the *next* snapshot, since each call is an independent point-in-time read.
1597        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}