1extern 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
40pub type KeyHash = [u8; KEY_HASH_LEN];
42
43#[derive(Debug, Clone)]
45pub struct InstanceState {
46 pub handle: InstanceHandle,
48 pub kind: InstanceStateKind,
50 pub disposed_generation_count: i32,
52 pub no_writers_generation_count: i32,
54 pub writer_count: u32,
59 pub last_sample_timestamp: Option<Time>,
61 pub last_delivered_ts: Option<Time>,
66 pub disposed_at: Option<Time>,
71 pub no_writers_at: Option<Time>,
75 pub current_owner: Option<([u8; 16], i32)>,
80 pub key_holder: alloc::vec::Vec<u8>,
84 pub reader_view_new: bool,
86 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#[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 #[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 #[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 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 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 #[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 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 #[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 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 #[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 #[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 #[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 #[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 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 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 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; };
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 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 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 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 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 if matches!(
487 entry.kind,
488 InstanceStateKind::NotAliveDisposed | InstanceStateKind::NotAliveNoWriters
489 ) {
490 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 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 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 #[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 #[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 #[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 #[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 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 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 #[test]
722 fn time_based_filter_first_sample_passes() {
723 let t = InstanceTracker::new();
724 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, );
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 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 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 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 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 #[test]
795 fn autopurge_disposed_after_delay() {
796 let t = InstanceTracker::new();
797 let h = t.register(kh(30), alloc::vec![30], None);
798 t.dispose(h, Some(Time::new(10, 0)));
800 let purged = t.autopurge(Time::new(15, 0), 3_000_000_000, u128::MAX);
802 assert_eq!(purged, 1);
803 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 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 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 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 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 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 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 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 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 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 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 assert!(t.should_accept_sample_under_exclusive_ownership(&kh(45), guid(7), 10));
912 }
913
914 #[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 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 assert!(t.should_accept_sample_under_exclusive_ownership(&kh(52), guid(9), 100));
938 assert!(!t.should_accept_sample_under_exclusive_ownership(&kh(52), guid(1), 10));
940 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 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 assert!(t.should_accept_sample_under_exclusive_ownership(&kh(60), full_a, 50));
956 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 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 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 for k in [kh(70), kh(71), kh(72)] {
978 assert!(t.should_accept_sample_under_exclusive_ownership(&k, g, 1));
979 }
980 let cleared = t.clear_owner_for_writer_prefix([7u8; 12]);
982 assert_eq!(cleared, 3);
983 }
984}