Skip to main content

zerodds_dcps/
instance_tracker.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright 2026 ZeroDDS Contributors
3//! `InstanceTracker` — the central bookkeeping for keyed topic
4//! instances, on both the writer and the reader side.
5//!
6//! Spec reference: OMG DDS-DCPS 1.4 §2.2.2.4.2.5 (`register_instance`),
7//! §2.2.2.4.2.7 (`unregister_instance`), §2.2.2.4.2.10 (`dispose`),
8//! §2.2.2.4.2.13 (`get_key_value`), §2.2.2.4.2.14 (`lookup_instance`),
9//! §2.2.2.5.1 (`InstanceStateKind`).
10//!
11//! # Data model
12//!
13//! We index instances by the **16-byte KeyHash** (XTypes 1.3 §7.6.8).
14//! Per instance we keep:
15//! * the assigned [`InstanceHandle`],
16//! * the lifecycle state ([`InstanceStateKind`]),
17//! * the generation counters (`disposed`, `no_writers`),
18//! * the number of still-registered writers (reader side),
19//! * the last observed sample timestamp.
20//!
21//! KeyHash → handle is a 1:1 map for the lifetime of the tracker
22//! (handles are not recycled, even once the instance is disposed +
23//! purged; the spec explicitly allows recycling, but we avoid it for
24//! test stability).
25
26extern crate alloc;
27
28use alloc::collections::BTreeMap;
29use alloc::sync::Arc;
30
31#[cfg(feature = "std")]
32use std::sync::Mutex;
33
34use zerodds_cdr::KEY_HASH_LEN;
35
36use crate::instance_handle::{InstanceHandle, InstanceHandleAllocator};
37use crate::sample_info::InstanceStateKind;
38use crate::time::Time;
39
40/// 16-byte KeyHash used as index.
41pub type KeyHash = [u8; KEY_HASH_LEN];
42
43/// Per-instance bookkeeping.
44#[derive(Debug, Clone)]
45pub struct InstanceState {
46    /// Local handle (stable across the tracker lifecycle).
47    pub handle: InstanceHandle,
48    /// Current lifecycle state.
49    pub kind: InstanceStateKind,
50    /// `NOT_ALIVE_DISPOSED → ALIVE` transitions since the first sample.
51    pub disposed_generation_count: i32,
52    /// `NOT_ALIVE_NO_WRITERS → ALIVE` transitions since the first sample.
53    pub no_writers_generation_count: i32,
54    /// Number of writers that currently consider this instance
55    /// registered. Reader side: each incoming sample increments the
56    /// counter (if it is a new writer); each unregister marker
57    /// decrements it. Writer side: 0 or 1 (a single writer's view).
58    pub writer_count: u32,
59    /// Wall-clock time of the last processed sample (or None).
60    pub last_sample_timestamp: Option<Time>,
61    /// Reader side: source_timestamp of this instance's last sample
62    /// delivered to the user API. Spec §2.2.3.12 TIME_BASED_FILTER: a
63    /// new sample is filtered out when
64    /// `t - last_delivered_ts < minimum_separation`.
65    pub last_delivered_ts: Option<Time>,
66    /// Reader side: wall-clock time at which the instance transitioned
67    /// into `NOT_ALIVE_DISPOSED`. Spec §2.2.3.22
68    /// `autopurge_disposed_samples_delay`: samples are purged after the
69    /// delay elapses.
70    pub disposed_at: Option<Time>,
71    /// Reader side: wall-clock time at which the instance transitioned
72    /// into `NOT_ALIVE_NO_WRITERS`. Spec §2.2.3.22
73    /// `autopurge_no_writer_samples_delay`.
74    pub no_writers_at: Option<Time>,
75    /// Reader side OWNERSHIP=EXCLUSIVE (Spec §2.2.3.10): current owner
76    /// writer as (GuidLike, Strength). On a tie, the lexicographically
77    /// higher `GuidLike` wins. On liveliness loss: reset explicitly via
78    /// `clear_owner`.
79    pub current_owner: Option<([u8; 16], i32)>,
80    /// Stored key holder (bytes), so `get_key_value` can replay the key
81    /// without re-decoding. This is the **PLAIN_CDR2-BE** stream, i.e.
82    /// exactly the input to `compute_key_hash`.
83    pub key_holder: alloc::vec::Vec<u8>,
84    /// View state on the reader side. Unused on the writer side.
85    pub reader_view_new: bool,
86    /// Number of reader samples so far in this instance (a sample_rank
87    /// helper value for the next read).
88    pub samples_in_cache: u32,
89}
90
91impl InstanceState {
92    fn fresh(handle: InstanceHandle, key_holder: alloc::vec::Vec<u8>) -> Self {
93        Self {
94            handle,
95            kind: InstanceStateKind::Alive,
96            disposed_generation_count: 0,
97            no_writers_generation_count: 0,
98            writer_count: 0,
99            last_sample_timestamp: None,
100            last_delivered_ts: None,
101            disposed_at: None,
102            no_writers_at: None,
103            current_owner: None,
104            key_holder,
105            reader_view_new: true,
106            samples_in_cache: 0,
107        }
108    }
109}
110
111/// Thread-safe tracker — instantiated in both the DataWriter and the
112/// DataReader.
113#[derive(Debug)]
114pub struct InstanceTracker {
115    inner: Arc<Mutex<TrackerInner>>,
116    allocator: Arc<InstanceHandleAllocator>,
117}
118
119#[derive(Debug, Default)]
120struct TrackerInner {
121    by_keyhash: BTreeMap<KeyHash, InstanceState>,
122    handle_to_keyhash: BTreeMap<InstanceHandle, KeyHash>,
123}
124
125impl Default for InstanceTracker {
126    fn default() -> Self {
127        Self::new()
128    }
129}
130
131impl InstanceTracker {
132    /// New tracker with its own [`InstanceHandleAllocator`].
133    #[must_use]
134    pub fn new() -> Self {
135        Self {
136            inner: Arc::new(Mutex::new(TrackerInner::default())),
137            allocator: Arc::new(InstanceHandleAllocator::new()),
138        }
139    }
140
141    /// Tracker with a shared allocator (e.g. when a writer and reader in
142    /// the same participant draw their handles from the same pool).
143    #[must_use]
144    pub fn with_allocator(allocator: Arc<InstanceHandleAllocator>) -> Self {
145        Self {
146            inner: Arc::new(Mutex::new(TrackerInner::default())),
147            allocator,
148        }
149    }
150
151    /// Registers the instance if it is not yet known; otherwise just
152    /// reactivates it (Spec §2.2.2.4.2.5).
153    ///
154    /// Always returns the (stable) [`InstanceHandle`].
155    pub fn register(
156        &self,
157        keyhash: KeyHash,
158        key_holder: alloc::vec::Vec<u8>,
159        timestamp: Option<Time>,
160    ) -> InstanceHandle {
161        let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
162        let entry = g.by_keyhash.entry(keyhash).or_insert_with(|| {
163            let h = self.allocator.allocate();
164            InstanceState::fresh(h, key_holder.clone())
165        });
166        // Reactivation from NOT_ALIVE → ALIVE produces a generation
167        // bump (Spec §2.2.2.5.1.7/8).
168        match entry.kind {
169            InstanceStateKind::NotAliveDisposed => {
170                entry.disposed_generation_count = entry.disposed_generation_count.saturating_add(1);
171                entry.kind = InstanceStateKind::Alive;
172            }
173            InstanceStateKind::NotAliveNoWriters => {
174                entry.no_writers_generation_count =
175                    entry.no_writers_generation_count.saturating_add(1);
176                entry.kind = InstanceStateKind::Alive;
177            }
178            InstanceStateKind::Alive => {}
179        }
180        entry.writer_count = entry.writer_count.saturating_add(1);
181        if let Some(ts) = timestamp {
182            entry.last_sample_timestamp = Some(ts);
183        }
184        let handle = entry.handle;
185        g.handle_to_keyhash.insert(handle, keyhash);
186        handle
187    }
188
189    /// Spec §2.2.3.12 TIME_BASED_FILTER — decides whether a sample with
190    /// `sample_ts` may be delivered to the user API. `false` if
191    /// `t - last_delivered_ts < min_separation` (drop); `true` otherwise
192    /// (deliver).
193    ///
194    /// For an unknown instance or the first sample (no
195    /// `last_delivered_ts`) it is always `true`.
196    #[must_use]
197    pub fn should_deliver_under_time_based_filter(
198        &self,
199        keyhash: &KeyHash,
200        sample_ts: Time,
201        min_separation_nanos: u128,
202    ) -> bool {
203        if min_separation_nanos == 0 {
204            return true;
205        }
206        let g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
207        let Some(s) = g.by_keyhash.get(keyhash) else {
208            return true;
209        };
210        let Some(last) = s.last_delivered_ts else {
211            return true;
212        };
213        // Nanosecond difference; for a sample in the past (clock skew),
214        // never filter — the caller has a reorder problem, not a filter
215        // problem.
216        let last_nanos = u128::from(u64::try_from(last.sec).unwrap_or(0)) * 1_000_000_000
217            + u128::from(last.nanosec);
218        let sample_nanos = u128::from(u64::try_from(sample_ts.sec).unwrap_or(0)) * 1_000_000_000
219            + u128::from(sample_ts.nanosec);
220        if sample_nanos < last_nanos {
221            return true;
222        }
223        sample_nanos - last_nanos >= min_separation_nanos
224    }
225
226    /// Spec §2.2.3.18 DESTINATION_ORDER — decides whether a sample with
227    /// `source_ts` may be delivered to the user API.
228    /// `BY_RECEPTION_TIMESTAMP`: always `true`. `BY_SOURCE_TIMESTAMP`:
229    /// only if `source_ts` is strictly greater than this instance's
230    /// `last_delivered_ts` (the tie-break on equal timestamps via the
231    /// writer GUID happens in the typed path).
232    #[must_use]
233    pub fn should_deliver_under_destination_order(
234        &self,
235        keyhash: &KeyHash,
236        source_ts: Time,
237        by_source_timestamp: bool,
238    ) -> bool {
239        if !by_source_timestamp {
240            return true;
241        }
242        let g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
243        let Some(s) = g.by_keyhash.get(keyhash) else {
244            return true;
245        };
246        let Some(last) = s.last_delivered_ts else {
247            return true;
248        };
249        let last_nanos = u128::from(u64::try_from(last.sec).unwrap_or(0)) * 1_000_000_000
250            + u128::from(last.nanosec);
251        let src_nanos = u128::from(u64::try_from(source_ts.sec).unwrap_or(0)) * 1_000_000_000
252            + u128::from(source_ts.nanosec);
253        src_nanos > last_nanos
254    }
255
256    /// Marks that a sample with `sample_ts` was delivered to the user
257    /// API (Spec §2.2.3.12 — for the next filter decision).
258    pub fn record_delivery(&self, keyhash: &KeyHash, sample_ts: Time) {
259        let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
260        if let Some(s) = g.by_keyhash.get_mut(keyhash) {
261            s.last_delivered_ts = Some(sample_ts);
262        }
263    }
264
265    /// Lookup without mutation (Spec §2.2.2.4.2.14 `lookup_instance`).
266    #[must_use]
267    pub fn lookup(&self, keyhash: &KeyHash) -> Option<InstanceHandle> {
268        let g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
269        g.by_keyhash.get(keyhash).map(|s| s.handle)
270    }
271
272    /// Returns a copy of the state snapshot for a handle.
273    #[must_use]
274    pub fn get_by_handle(&self, handle: InstanceHandle) -> Option<InstanceState> {
275        let g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
276        let kh = g.handle_to_keyhash.get(&handle)?;
277        g.by_keyhash.get(kh).cloned()
278    }
279
280    /// Returns a copy of the state snapshot for a KeyHash.
281    #[must_use]
282    pub fn get_by_keyhash(&self, keyhash: &KeyHash) -> Option<InstanceState> {
283        let g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
284        g.by_keyhash.get(keyhash).cloned()
285    }
286
287    /// Returns the key-holder byte stream for a handle (Spec
288    /// §2.2.2.4.2.13 `get_key_value`).
289    #[must_use]
290    pub fn get_key_holder(&self, handle: InstanceHandle) -> Option<alloc::vec::Vec<u8>> {
291        let g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
292        let kh = g.handle_to_keyhash.get(&handle)?;
293        g.by_keyhash.get(kh).map(|s| s.key_holder.clone())
294    }
295
296    /// Marks the instance as `NOT_ALIVE_DISPOSED` (Spec §2.2.2.4.2.10
297    /// `dispose`).
298    ///
299    /// Returns `false` if the instance is not known.
300    pub fn dispose(&self, handle: InstanceHandle, timestamp: Option<Time>) -> bool {
301        let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
302        let Some(kh) = g.handle_to_keyhash.get(&handle).copied() else {
303            return false;
304        };
305        if let Some(s) = g.by_keyhash.get_mut(&kh) {
306            s.kind = InstanceStateKind::NotAliveDisposed;
307            if let Some(ts) = timestamp {
308                s.last_sample_timestamp = Some(ts);
309                s.disposed_at = Some(ts);
310            }
311            return true;
312        }
313        false
314    }
315
316    /// Decrements the writer counter; if it drops to 0, the instance
317    /// transitions to `NOT_ALIVE_NO_WRITERS` (Spec §2.2.2.4.2.7).
318    pub fn unregister(&self, handle: InstanceHandle, timestamp: Option<Time>) -> bool {
319        let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
320        let Some(kh) = g.handle_to_keyhash.get(&handle).copied() else {
321            return false;
322        };
323        if let Some(s) = g.by_keyhash.get_mut(&kh) {
324            s.writer_count = s.writer_count.saturating_sub(1);
325            if s.writer_count == 0 && !matches!(s.kind, InstanceStateKind::NotAliveDisposed) {
326                s.kind = InstanceStateKind::NotAliveNoWriters;
327                if let Some(ts) = timestamp {
328                    s.no_writers_at = Some(ts);
329                }
330            }
331            if let Some(ts) = timestamp {
332                s.last_sample_timestamp = Some(ts);
333            }
334            return true;
335        }
336        false
337    }
338
339    /// Spec §2.2.3.10 OWNERSHIP=EXCLUSIVE strength selection.
340    /// Returns `true` if a sample from the writer with
341    /// `(writer_guid, writer_strength)` should be accepted for the
342    /// instance `keyhash`.
343    ///
344    /// Algorithm (Spec §2.2.3.10):
345    /// - No current owner → accept + set as owner.
346    /// - Strength > current → accept + replace owner.
347    /// - Strength == current and guid > current_guid → accept
348    ///   (spec tie-break via the lexicographically higher guid).
349    /// - Strength < current → reject.
350    /// - Strength == current and guid < current → reject.
351    /// - Strength == current and guid == current → accept
352    ///   (same writer).
353    pub fn should_accept_sample_under_exclusive_ownership(
354        &self,
355        keyhash: &KeyHash,
356        writer_guid: [u8; 16],
357        writer_strength: i32,
358    ) -> bool {
359        let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
360        let Some(s) = g.by_keyhash.get_mut(keyhash) else {
361            return true; // unknown instance → accept
362        };
363        match s.current_owner {
364            None => {
365                s.current_owner = Some((writer_guid, writer_strength));
366                true
367            }
368            Some((cur_guid, cur_str)) => {
369                if writer_strength > cur_str
370                    || (writer_strength == cur_str && writer_guid > cur_guid)
371                {
372                    s.current_owner = Some((writer_guid, writer_strength));
373                    true
374                } else {
375                    writer_strength == cur_str && writer_guid == cur_guid
376                }
377            }
378        }
379    }
380
381    /// Spec §2.2.3.23 — on liveliness loss of a writer: clear the owner
382    /// for all instances whose owner was this writer. The next sample
383    /// triggers failover selection.
384    pub fn clear_owner_for_writer(&self, writer_guid: [u8; 16]) -> usize {
385        let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
386        let mut cleared = 0;
387        for s in g.by_keyhash.values_mut() {
388            if let Some((g_, _)) = s.current_owner {
389                if g_ == writer_guid {
390                    s.current_owner = None;
391                    cleared += 1;
392                }
393            }
394        }
395        cleared
396    }
397
398    /// Like [`Self::clear_owner_for_writer`], but matches on the first
399    /// 12 bytes of the GUID (GuidPrefix). Allows failover when only the
400    /// participant identity (e.g. via SPDP lease expiry) is known.
401    pub fn clear_owner_for_writer_prefix(&self, prefix: [u8; 12]) -> usize {
402        let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
403        let mut cleared = 0;
404        for s in g.by_keyhash.values_mut() {
405            if let Some((g_, _)) = s.current_owner {
406                if g_[..12] == prefix {
407                    s.current_owner = None;
408                    cleared += 1;
409                }
410            }
411        }
412        cleared
413    }
414
415    /// Spec §2.2.3.22 READER_DATA_LIFECYCLE — purges instances whose
416    /// disposed/no-writer marker is older than the respective delay.
417    /// `now` is the caller-side wall-clock; the delays are in
418    /// nanoseconds. Returns the number of removed instances.
419    ///
420    /// Lazy purge: called from the read path (or a background tick). The
421    /// spec leaves the strategy open — we delete the affected instance
422    /// entirely, so that subsequent read/take no longer see it.
423    pub fn autopurge(
424        &self,
425        now: Time,
426        autopurge_disposed_delay_nanos: u128,
427        autopurge_nowriter_delay_nanos: u128,
428    ) -> usize {
429        let now_nanos = u128::from(u64::try_from(now.sec).unwrap_or(0)) * 1_000_000_000
430            + u128::from(now.nanosec);
431        let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
432        let mut to_purge: alloc::vec::Vec<KeyHash> = alloc::vec::Vec::new();
433        for (kh, s) in g.by_keyhash.iter() {
434            let purge = match s.kind {
435                InstanceStateKind::NotAliveDisposed
436                    if autopurge_disposed_delay_nanos != u128::MAX =>
437                {
438                    s.disposed_at.is_some_and(|t| {
439                        let t_nanos = u128::from(u64::try_from(t.sec).unwrap_or(0)) * 1_000_000_000
440                            + u128::from(t.nanosec);
441                        now_nanos.saturating_sub(t_nanos) >= autopurge_disposed_delay_nanos
442                    })
443                }
444                InstanceStateKind::NotAliveNoWriters
445                    if autopurge_nowriter_delay_nanos != u128::MAX =>
446                {
447                    s.no_writers_at.is_some_and(|t| {
448                        let t_nanos = u128::from(u64::try_from(t.sec).unwrap_or(0)) * 1_000_000_000
449                            + u128::from(t.nanosec);
450                        now_nanos.saturating_sub(t_nanos) >= autopurge_nowriter_delay_nanos
451                    })
452                }
453                _ => false,
454            };
455            if purge {
456                to_purge.push(*kh);
457            }
458        }
459        let count = to_purge.len();
460        for kh in to_purge {
461            if let Some(s) = g.by_keyhash.remove(&kh) {
462                g.handle_to_keyhash.remove(&s.handle);
463            }
464        }
465        count
466    }
467
468    /// Reader-side hook: marks that a sample has arrived for this
469    /// instance. Returns `(handle, was_new)`, where `was_new == true`
470    /// means this instance was previously unknown or the reader view was
471    /// freshly reset (`view_state = NEW`).
472    pub fn observe_sample(
473        &self,
474        keyhash: KeyHash,
475        key_holder: alloc::vec::Vec<u8>,
476        timestamp: Option<Time>,
477    ) -> (InstanceHandle, bool) {
478        let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
479        let mut was_new = false;
480        let entry = g.by_keyhash.entry(keyhash).or_insert_with(|| {
481            was_new = true;
482            let h = self.allocator.allocate();
483            InstanceState::fresh(h, key_holder.clone())
484        });
485        // Reactivate the reader view
486        if matches!(
487            entry.kind,
488            InstanceStateKind::NotAliveDisposed | InstanceStateKind::NotAliveNoWriters
489        ) {
490            // Generation bumps as on the writer side — this matches the
491            // spec, because readers derive it from the sample stream.
492            match entry.kind {
493                InstanceStateKind::NotAliveDisposed => {
494                    entry.disposed_generation_count =
495                        entry.disposed_generation_count.saturating_add(1);
496                }
497                InstanceStateKind::NotAliveNoWriters => {
498                    entry.no_writers_generation_count =
499                        entry.no_writers_generation_count.saturating_add(1);
500                }
501                InstanceStateKind::Alive => {}
502            }
503            entry.kind = InstanceStateKind::Alive;
504            entry.reader_view_new = true;
505        }
506        if let Some(ts) = timestamp {
507            entry.last_sample_timestamp = Some(ts);
508        }
509        entry.samples_in_cache = entry.samples_in_cache.saturating_add(1);
510        let handle = entry.handle;
511        g.handle_to_keyhash.insert(handle, keyhash);
512        (handle, was_new)
513    }
514
515    /// Reader side: after the first `read`/`take` of an instance, its
516    /// view state is set to `NOT_NEW`.
517    pub fn mark_view_seen(&self, handle: InstanceHandle) {
518        let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
519        if let Some(kh) = g.handle_to_keyhash.get(&handle).copied() {
520            if let Some(s) = g.by_keyhash.get_mut(&kh) {
521                s.reader_view_new = false;
522            }
523        }
524    }
525
526    /// Reader side: after `take`, `samples_in_cache` is reduced by `n`.
527    pub fn drain_samples(&self, handle: InstanceHandle, n: u32) {
528        let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
529        if let Some(kh) = g.handle_to_keyhash.get(&handle).copied() {
530            if let Some(s) = g.by_keyhash.get_mut(&kh) {
531                s.samples_in_cache = s.samples_in_cache.saturating_sub(n);
532            }
533        }
534    }
535
536    /// Lists all instance handles in stable order (BTreeMap order by
537    /// KeyHash). Spec §2.2.2.5.3.28 `read_next_instance`.
538    #[must_use]
539    pub fn ordered_handles(&self) -> alloc::vec::Vec<InstanceHandle> {
540        let g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
541        g.by_keyhash.values().map(|s| s.handle).collect()
542    }
543
544    /// Returns the **first** handle whose sort order lies strictly after
545    /// `previous_handle` (or the very first one if
546    /// `previous == HANDLE_NIL`). Spec §2.2.2.5.3.28.
547    #[must_use]
548    pub fn next_handle_after(&self, previous: InstanceHandle) -> Option<InstanceHandle> {
549        let g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
550        if previous.is_nil() {
551            return g.by_keyhash.values().next().map(|s| s.handle);
552        }
553        let prev_kh = g.handle_to_keyhash.get(&previous).copied()?;
554        let range: (core::ops::Bound<KeyHash>, core::ops::Bound<KeyHash>) = (
555            core::ops::Bound::Excluded(prev_kh),
556            core::ops::Bound::Unbounded,
557        );
558        g.by_keyhash.range(range).next().map(|(_, s)| s.handle)
559    }
560
561    /// Number of tracked instances.
562    #[must_use]
563    pub fn len(&self) -> usize {
564        let g = self.inner.lock().unwrap_or_else(|e| e.into_inner());
565        g.by_keyhash.len()
566    }
567
568    /// `true` if no instances are tracked.
569    #[must_use]
570    pub fn is_empty(&self) -> bool {
571        self.len() == 0
572    }
573}
574
575impl Clone for InstanceTracker {
576    fn clone(&self) -> Self {
577        Self {
578            inner: Arc::clone(&self.inner),
579            allocator: Arc::clone(&self.allocator),
580        }
581    }
582}
583
584#[cfg(test)]
585#[allow(clippy::expect_used, clippy::unwrap_used)]
586mod tests {
587    use super::*;
588
589    fn kh(byte: u8) -> KeyHash {
590        let mut k = [0u8; KEY_HASH_LEN];
591        k[0] = byte;
592        k
593    }
594
595    #[test]
596    fn register_assigns_stable_handle() {
597        let t = InstanceTracker::new();
598        let h1 = t.register(kh(1), alloc::vec![1], None);
599        let h2 = t.register(kh(1), alloc::vec![1], None);
600        assert_eq!(h1, h2);
601        assert!(!h1.is_nil());
602    }
603
604    #[test]
605    fn lookup_returns_handle_for_known_key() {
606        let t = InstanceTracker::new();
607        let h = t.register(kh(2), alloc::vec![2], None);
608        assert_eq!(t.lookup(&kh(2)), Some(h));
609        assert_eq!(t.lookup(&kh(99)), None);
610    }
611
612    #[test]
613    fn dispose_transitions_to_disposed() {
614        let t = InstanceTracker::new();
615        let h = t.register(kh(3), alloc::vec![3], None);
616        assert_eq!(t.get_by_handle(h).unwrap().kind, InstanceStateKind::Alive);
617        assert!(t.dispose(h, None));
618        assert_eq!(
619            t.get_by_handle(h).unwrap().kind,
620            InstanceStateKind::NotAliveDisposed
621        );
622    }
623
624    #[test]
625    fn unregister_decrements_writer_count() {
626        let t = InstanceTracker::new();
627        let h = t.register(kh(4), alloc::vec![4], None);
628        // Two writers have registered.
629        let _ = t.register(kh(4), alloc::vec![4], None);
630        assert_eq!(t.get_by_handle(h).unwrap().writer_count, 2);
631        assert!(t.unregister(h, None));
632        assert_eq!(t.get_by_handle(h).unwrap().kind, InstanceStateKind::Alive);
633        assert!(t.unregister(h, None));
634        // now 0 writers → no_writers
635        assert_eq!(
636            t.get_by_handle(h).unwrap().kind,
637            InstanceStateKind::NotAliveNoWriters
638        );
639    }
640
641    #[test]
642    fn re_register_after_dispose_bumps_disposed_generation() {
643        let t = InstanceTracker::new();
644        let h = t.register(kh(5), alloc::vec![5], None);
645        t.dispose(h, None);
646        let _ = t.register(kh(5), alloc::vec![5], None);
647        let s = t.get_by_handle(h).unwrap();
648        assert_eq!(s.kind, InstanceStateKind::Alive);
649        assert_eq!(s.disposed_generation_count, 1);
650    }
651
652    #[test]
653    fn observe_sample_creates_new_instance_on_first_call() {
654        let t = InstanceTracker::new();
655        let (h, was_new) = t.observe_sample(kh(6), alloc::vec![6], None);
656        assert!(was_new);
657        assert!(t.get_by_handle(h).unwrap().reader_view_new);
658        let (h2, was_new2) = t.observe_sample(kh(6), alloc::vec![6], None);
659        assert_eq!(h, h2);
660        assert!(!was_new2);
661    }
662
663    #[test]
664    fn ordered_handles_iterates_in_keyhash_order() {
665        let t = InstanceTracker::new();
666        let h_b = t.register(kh(2), alloc::vec![2], None);
667        let h_a = t.register(kh(1), alloc::vec![1], None);
668        let h_c = t.register(kh(3), alloc::vec![3], None);
669        assert_eq!(t.ordered_handles(), alloc::vec![h_a, h_b, h_c]);
670    }
671
672    #[test]
673    fn next_handle_after_walks_in_order() {
674        let t = InstanceTracker::new();
675        let h_a = t.register(kh(1), alloc::vec![1], None);
676        let h_b = t.register(kh(2), alloc::vec![2], None);
677        let h_c = t.register(kh(3), alloc::vec![3], None);
678        assert_eq!(t.next_handle_after(crate::HANDLE_NIL), Some(h_a));
679        assert_eq!(t.next_handle_after(h_a), Some(h_b));
680        assert_eq!(t.next_handle_after(h_b), Some(h_c));
681        assert_eq!(t.next_handle_after(h_c), None);
682    }
683
684    #[test]
685    fn get_key_holder_returns_stored_bytes() {
686        let t = InstanceTracker::new();
687        let h = t.register(kh(7), alloc::vec![1, 2, 3], None);
688        assert_eq!(t.get_key_holder(h), Some(alloc::vec![1u8, 2, 3]));
689    }
690
691    #[test]
692    fn mark_view_seen_clears_new_flag() {
693        let t = InstanceTracker::new();
694        let (h, _) = t.observe_sample(kh(8), alloc::vec![8], None);
695        assert!(t.get_by_handle(h).unwrap().reader_view_new);
696        t.mark_view_seen(h);
697        assert!(!t.get_by_handle(h).unwrap().reader_view_new);
698    }
699
700    #[test]
701    fn observe_after_dispose_bumps_disposed_generation() {
702        let t = InstanceTracker::new();
703        let (h, _) = t.observe_sample(kh(9), alloc::vec![9], None);
704        t.dispose(h, None);
705        let (_, _) = t.observe_sample(kh(9), alloc::vec![9], None);
706        assert_eq!(t.get_by_handle(h).unwrap().disposed_generation_count, 1);
707    }
708
709    #[test]
710    fn drain_samples_decrements_count() {
711        let t = InstanceTracker::new();
712        let (h, _) = t.observe_sample(kh(10), alloc::vec![10], None);
713        let (_, _) = t.observe_sample(kh(10), alloc::vec![10], None);
714        assert_eq!(t.get_by_handle(h).unwrap().samples_in_cache, 2);
715        t.drain_samples(h, 2);
716        assert_eq!(t.get_by_handle(h).unwrap().samples_in_cache, 0);
717    }
718
719    // ---- §2.2.3.12 TIME_BASED_FILTER Reader-Drop ----
720
721    #[test]
722    fn time_based_filter_first_sample_passes() {
723        let t = InstanceTracker::new();
724        // New instance — no last_delivered_ts, so deliver.
725        let _ = t.observe_sample(kh(20), alloc::vec![20], Some(Time::new(1, 0)));
726        let pass = t.should_deliver_under_time_based_filter(
727            &kh(20),
728            Time::new(1, 0),
729            100_000_000, // 100ms
730        );
731        assert!(pass);
732    }
733
734    #[test]
735    fn time_based_filter_too_close_drops() {
736        let t = InstanceTracker::new();
737        let _ = t.observe_sample(kh(20), alloc::vec![20], None);
738        t.record_delivery(&kh(20), Time::new(1, 0));
739        // 50ms later — drop at min_separation=100ms.
740        let pass = t.should_deliver_under_time_based_filter(
741            &kh(20),
742            Time::new(1, 50_000_000),
743            100_000_000,
744        );
745        assert!(!pass, "50ms < 100ms separation -> drop");
746    }
747
748    #[test]
749    fn time_based_filter_far_enough_passes() {
750        let t = InstanceTracker::new();
751        let _ = t.observe_sample(kh(20), alloc::vec![20], None);
752        t.record_delivery(&kh(20), Time::new(1, 0));
753        // 150ms later — pass.
754        let pass = t.should_deliver_under_time_based_filter(
755            &kh(20),
756            Time::new(1, 150_000_000),
757            100_000_000,
758        );
759        assert!(pass, "150ms > 100ms separation -> deliver");
760    }
761
762    #[test]
763    fn time_based_filter_zero_separation_always_passes() {
764        let t = InstanceTracker::new();
765        let _ = t.observe_sample(kh(20), alloc::vec![20], None);
766        t.record_delivery(&kh(20), Time::new(1, 0));
767        let pass = t.should_deliver_under_time_based_filter(&kh(20), Time::new(1, 0), 0);
768        assert!(pass, "min_separation=0 -> no filter");
769    }
770
771    #[test]
772    fn time_based_filter_per_instance_isolation() {
773        // The filter is per-instance; instance A's last_delivered does
774        // not affect instance B.
775        let t = InstanceTracker::new();
776        let _ = t.observe_sample(kh(1), alloc::vec![1], None);
777        let _ = t.observe_sample(kh(2), alloc::vec![2], None);
778        t.record_delivery(&kh(1), Time::new(5, 0));
779        // Instance 2 has no delivery yet → pass.
780        let pass =
781            t.should_deliver_under_time_based_filter(&kh(2), Time::new(5, 10_000_000), 100_000_000);
782        assert!(pass);
783    }
784
785    #[test]
786    fn time_based_filter_unknown_instance_passes() {
787        let t = InstanceTracker::new();
788        let pass = t.should_deliver_under_time_based_filter(&kh(99), Time::new(1, 0), 100_000_000);
789        assert!(pass, "unknown instance -> pass");
790    }
791
792    // ---- §2.2.3.22 READER_DATA_LIFECYCLE autopurge ----
793
794    #[test]
795    fn autopurge_disposed_after_delay() {
796        let t = InstanceTracker::new();
797        let h = t.register(kh(30), alloc::vec![30], None);
798        // Dispose with ts=t1.
799        t.dispose(h, Some(Time::new(10, 0)));
800        // Now=t1+5s, delay=3s → purge.
801        let purged = t.autopurge(Time::new(15, 0), 3_000_000_000, u128::MAX);
802        assert_eq!(purged, 1);
803        // After purge: lookup returns None.
804        assert!(t.lookup(&kh(30)).is_none());
805    }
806
807    #[test]
808    fn autopurge_disposed_before_delay_keeps_instance() {
809        let t = InstanceTracker::new();
810        let h = t.register(kh(31), alloc::vec![31], None);
811        t.dispose(h, Some(Time::new(10, 0)));
812        // Now=t1+1s, delay=5s → keep.
813        let purged = t.autopurge(Time::new(11, 0), 5_000_000_000, u128::MAX);
814        assert_eq!(purged, 0);
815        assert!(t.lookup(&kh(31)).is_some());
816    }
817
818    #[test]
819    fn autopurge_no_writers_after_delay() {
820        let t = InstanceTracker::new();
821        let h = t.register(kh(32), alloc::vec![32], None);
822        // unregister sets writer_count=0 → NotAliveNoWriters with timestamp.
823        t.unregister(h, Some(Time::new(20, 0)));
824        let purged = t.autopurge(Time::new(25, 0), u128::MAX, 3_000_000_000);
825        assert_eq!(purged, 1);
826        assert!(t.lookup(&kh(32)).is_none());
827    }
828
829    #[test]
830    fn autopurge_alive_instance_never_purged() {
831        let t = InstanceTracker::new();
832        let _h = t.register(kh(33), alloc::vec![33], None);
833        // ALIVE → purge ignored.
834        let purged = t.autopurge(Time::new(1000, 0), 0, 0);
835        assert_eq!(purged, 0);
836        assert!(t.lookup(&kh(33)).is_some());
837    }
838
839    #[test]
840    fn autopurge_infinity_delay_never_purges() {
841        // u128::MAX as delay = INFINITE = never purge (spec default).
842        let t = InstanceTracker::new();
843        let h = t.register(kh(34), alloc::vec![34], None);
844        t.dispose(h, Some(Time::new(10, 0)));
845        let purged = t.autopurge(Time::new(99999, 0), u128::MAX, u128::MAX);
846        assert_eq!(purged, 0);
847    }
848
849    // ---- §2.2.3.10 OWNERSHIP=EXCLUSIVE Strength-Selection ----
850
851    fn guid(byte: u8) -> [u8; 16] {
852        [byte; 16]
853    }
854
855    #[test]
856    fn exclusive_first_writer_wins() {
857        let t = InstanceTracker::new();
858        let _ = t.register(kh(40), alloc::vec![40], None);
859        // No owner → first writer accepted, becomes owner.
860        assert!(t.should_accept_sample_under_exclusive_ownership(&kh(40), guid(1), 10));
861        let s = t.get_by_keyhash(&kh(40)).unwrap();
862        assert_eq!(s.current_owner, Some((guid(1), 10)));
863    }
864
865    #[test]
866    fn exclusive_higher_strength_wins() {
867        let t = InstanceTracker::new();
868        let _ = t.register(kh(41), alloc::vec![41], None);
869        assert!(t.should_accept_sample_under_exclusive_ownership(&kh(41), guid(1), 10));
870        // Stronger writer arrives → replaces owner.
871        assert!(t.should_accept_sample_under_exclusive_ownership(&kh(41), guid(2), 20));
872        let s = t.get_by_keyhash(&kh(41)).unwrap();
873        assert_eq!(s.current_owner, Some((guid(2), 20)));
874    }
875
876    #[test]
877    fn exclusive_lower_strength_rejected() {
878        let t = InstanceTracker::new();
879        let _ = t.register(kh(42), alloc::vec![42], None);
880        assert!(t.should_accept_sample_under_exclusive_ownership(&kh(42), guid(2), 20));
881        // Weaker writer rejected.
882        assert!(!t.should_accept_sample_under_exclusive_ownership(&kh(42), guid(1), 5));
883        let s = t.get_by_keyhash(&kh(42)).unwrap();
884        assert_eq!(s.current_owner, Some((guid(2), 20)));
885    }
886
887    #[test]
888    fn exclusive_tie_break_by_higher_guid() {
889        let t = InstanceTracker::new();
890        let _ = t.register(kh(43), alloc::vec![43], None);
891        assert!(t.should_accept_sample_under_exclusive_ownership(&kh(43), guid(1), 10));
892        // Same strength, higher guid → wins.
893        assert!(t.should_accept_sample_under_exclusive_ownership(&kh(43), guid(2), 10));
894    }
895
896    #[test]
897    fn exclusive_tie_break_lower_guid_rejected() {
898        let t = InstanceTracker::new();
899        let _ = t.register(kh(44), alloc::vec![44], None);
900        assert!(t.should_accept_sample_under_exclusive_ownership(&kh(44), guid(2), 10));
901        // Same strength, lower guid → reject.
902        assert!(!t.should_accept_sample_under_exclusive_ownership(&kh(44), guid(1), 10));
903    }
904
905    #[test]
906    fn exclusive_same_writer_always_accepted() {
907        let t = InstanceTracker::new();
908        let _ = t.register(kh(45), alloc::vec![45], None);
909        assert!(t.should_accept_sample_under_exclusive_ownership(&kh(45), guid(7), 10));
910        // Same writer, same strength → always accepted.
911        assert!(t.should_accept_sample_under_exclusive_ownership(&kh(45), guid(7), 10));
912    }
913
914    // ---- §2.2.3.23 Liveliness-driven Failover ----
915
916    #[test]
917    fn clear_owner_for_writer_resets_owner() {
918        let t = InstanceTracker::new();
919        let _ = t.register(kh(50), alloc::vec![50], None);
920        let _ = t.register(kh(51), alloc::vec![51], None);
921        assert!(t.should_accept_sample_under_exclusive_ownership(&kh(50), guid(9), 100));
922        assert!(t.should_accept_sample_under_exclusive_ownership(&kh(51), guid(9), 100));
923        // Liveliness lost on writer guid(9) → clear from both instances.
924        let cleared = t.clear_owner_for_writer(guid(9));
925        assert_eq!(cleared, 2);
926        let s50 = t.get_by_keyhash(&kh(50)).unwrap();
927        let s51 = t.get_by_keyhash(&kh(51)).unwrap();
928        assert!(s50.current_owner.is_none());
929        assert!(s51.current_owner.is_none());
930    }
931
932    #[test]
933    fn failover_after_clear_accepts_weaker_writer() {
934        let t = InstanceTracker::new();
935        let _ = t.register(kh(52), alloc::vec![52], None);
936        // Strong writer becomes owner.
937        assert!(t.should_accept_sample_under_exclusive_ownership(&kh(52), guid(9), 100));
938        // Weaker writer normally rejected.
939        assert!(!t.should_accept_sample_under_exclusive_ownership(&kh(52), guid(1), 10));
940        // Clear (liveliness loss) → weaker writer can take over.
941        t.clear_owner_for_writer(guid(9));
942        assert!(t.should_accept_sample_under_exclusive_ownership(&kh(52), guid(1), 10));
943    }
944
945    #[test]
946    fn clear_owner_for_writer_prefix_matches_first_12_bytes() {
947        let t = InstanceTracker::new();
948        let _ = t.register(kh(60), alloc::vec![60], None);
949        // GUID = prefix [1;12] + entityId [9,9,9,9]
950        let mut full_a = [9u8; 16];
951        full_a[..12].fill(1);
952        let mut full_b = [9u8; 16];
953        full_b[..12].fill(2);
954        // owner with prefix=1
955        assert!(t.should_accept_sample_under_exclusive_ownership(&kh(60), full_a, 50));
956        // Clear by mismatching prefix → no clear.
957        assert_eq!(t.clear_owner_for_writer_prefix([2u8; 12]), 0);
958        let s = t.get_by_keyhash(&kh(60)).unwrap();
959        assert!(s.current_owner.is_some());
960        // Clear by matching prefix → cleared.
961        assert_eq!(t.clear_owner_for_writer_prefix([1u8; 12]), 1);
962        let s2 = t.get_by_keyhash(&kh(60)).unwrap();
963        assert!(s2.current_owner.is_none());
964        // Now full_b can take ownership (prefix=2, weaker).
965        let _ = full_b;
966    }
967
968    #[test]
969    fn clear_owner_for_writer_prefix_multi_instance() {
970        let t = InstanceTracker::new();
971        let _ = t.register(kh(70), alloc::vec![70], None);
972        let _ = t.register(kh(71), alloc::vec![71], None);
973        let _ = t.register(kh(72), alloc::vec![72], None);
974        let mut g = [0u8; 16];
975        g[..12].fill(7);
976        // Owner = same prefix on 3 instances.
977        for k in [kh(70), kh(71), kh(72)] {
978            assert!(t.should_accept_sample_under_exclusive_ownership(&k, g, 1));
979        }
980        // One participant disappears (lease lost) → all 3 cleared.
981        let cleared = t.clear_owner_for_writer_prefix([7u8; 12]);
982        assert_eq!(cleared, 3);
983    }
984}