orbit_core/ring/cursor.rs
1//! Cursor and walking primitives for Orbit rings.
2//!
3//! This module owns the generic "walk counters from cursor toward head"
4//! semantics shared by event streams, cache compaction, metrics, and
5//! future ring-backed substrates. Semantic layers decode frames after
6//! this layer has handled wraparound, in-flight slots, loss, and cursor
7//! advance.
8
9use crate::ring::Frame;
10
11/// Result of reading one logical counter from a ring source.
12#[derive(Clone, Debug, PartialEq, Eq)]
13pub enum RingRead {
14 /// The expected counter is fully committed and safe to consume.
15 Ready(Frame),
16 /// The counter has been claimed but may still be in flight.
17 /// A cursor must stay on this counter and retry later.
18 Pending,
19 /// The counter can no longer be read from this slot.
20 Unavailable,
21}
22
23/// Read-only source that can be walked by [`RingCursor`].
24///
25/// Implementors expose only the ring facts needed by the generic walker:
26/// the ring kind, current head, fixed capacity, and counter-addressed
27/// frame reads.
28pub trait RingFrameSource {
29 /// The `OrbitTyped::KIND` carried by this ring.
30 fn kind(&self) -> u8;
31
32 /// Monotonic visible head. Depending on topology, this counts counters
33 /// reserved by writers or frames committed by writers.
34 fn head(&self) -> u64;
35
36 /// Fixed slot count for this ring.
37 fn capacity(&self) -> usize;
38
39 /// Read the frame currently occupying `counter % capacity`.
40 fn read_at(&self, counter: u64) -> Option<Frame>;
41
42 /// Classify the logical counter currently addressed by
43 /// `counter % capacity`.
44 ///
45 /// Sources that can distinguish an in-flight claim should override
46 /// this method. The compatibility default preserves the previous
47 /// `read_at` behavior for external implementations.
48 fn read_state_at(&self, counter: u64) -> RingRead {
49 self.read_at(counter)
50 .map_or(RingRead::Unavailable, RingRead::Ready)
51 }
52}
53
54/// Caller-owned position in a ring walk.
55///
56/// The cursor stores the next counter a reader should attempt. Different
57/// subscribers keep independent cursors over the same ring.
58#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
59pub struct RingCursor {
60 next_counter: u64,
61}
62
63impl RingCursor {
64 /// Start at counter 0 and replay whatever ring history is still
65 /// available.
66 pub const fn from_start() -> Self {
67 Self { next_counter: 0 }
68 }
69
70 /// Start from a known next counter.
71 pub const fn from_counter(next_counter: u64) -> Self {
72 Self { next_counter }
73 }
74
75 /// The next counter this cursor will read.
76 pub const fn next_counter(self) -> u64 {
77 self.next_counter
78 }
79
80 pub(crate) fn set_next_counter(&mut self, next_counter: u64) {
81 self.next_counter = next_counter;
82 }
83}
84
85/// Counters skipped while walking a ring.
86#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
87pub struct RingLoss {
88 /// Counters older than the current ring window.
89 pub overwritten: u64,
90 /// Counters inside the readable window whose slot was definitively
91 /// unavailable, wrapped, corrupt, or carried an unexpected frame id.
92 pub unavailable: u64,
93}
94
95impl RingLoss {
96 pub const fn total(self) -> u64 {
97 self.overwritten.saturating_add(self.unavailable)
98 }
99
100 pub const fn is_empty(self) -> bool {
101 self.overwritten == 0 && self.unavailable == 0
102 }
103}
104
105/// Result of walking a cursor toward a ring's visible head.
106#[derive(Clone, Debug, Default, PartialEq, Eq)]
107pub struct RingPoll {
108 pub frames: Vec<Frame>,
109 pub loss: RingLoss,
110 pub from_counter: u64,
111 pub to_counter: u64,
112}
113
114impl RingPoll {
115 pub fn is_empty(&self) -> bool {
116 self.frames.is_empty() && self.loss.is_empty()
117 }
118}
119
120/// Walk `cursor` toward the current visible head of `source`.
121///
122/// At most one ring window is inspected. If the cursor has fallen behind
123/// the oldest available counter, the skipped counters are recorded as
124/// overwritten and the walk resumes at the window floor. An in-flight
125/// counter stops the walk without advancing past it.
126pub fn poll_ring<S: RingFrameSource>(source: &S, cursor: &mut RingCursor) -> RingPoll {
127 let head = source.head();
128 let from_counter = cursor.next_counter();
129
130 if from_counter >= head {
131 cursor.set_next_counter(head);
132 return RingPoll {
133 from_counter,
134 to_counter: head,
135 ..RingPoll::default()
136 };
137 }
138
139 let capacity = source.capacity() as u64;
140 let oldest_available = head.saturating_sub(capacity);
141 let mut next = from_counter;
142 let mut loss = RingLoss::default();
143
144 if next < oldest_available {
145 loss.overwritten = oldest_available - next;
146 next = oldest_available;
147 }
148
149 let kind = source.kind();
150 let mut frames = Vec::new();
151 while next < head {
152 match source.read_state_at(next) {
153 RingRead::Ready(frame) => {
154 if frame.id.kind() != kind || frame.id.counter() != next {
155 loss.unavailable = loss.unavailable.saturating_add(1);
156 } else {
157 frames.push(frame);
158 }
159 next = next.saturating_add(1);
160 }
161 RingRead::Pending => break,
162 RingRead::Unavailable => {
163 loss.unavailable = loss.unavailable.saturating_add(1);
164 next = next.saturating_add(1);
165 }
166 }
167 }
168
169 cursor.set_next_counter(next);
170 RingPoll {
171 frames,
172 loss,
173 from_counter,
174 to_counter: next,
175 }
176}