Skip to main content

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}