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(
41        &self,
42        counter: u64
43    ) -> Option<Frame>;
44
45    /// Classify the logical counter currently addressed by
46    /// `counter % capacity`.
47    ///
48    /// Sources that can distinguish an in-flight claim should override
49    /// this method. The compatibility default preserves the previous
50    /// `read_at` behavior for external implementations.
51    fn read_state_at(
52        &self,
53        counter: u64
54    ) -> RingRead {
55        self.read_at(counter).map_or(RingRead::Unavailable, RingRead::Ready)
56    }
57}
58
59/// Caller-owned position in a ring walk.
60///
61/// The cursor stores the next counter a reader should attempt. Different
62/// subscribers keep independent cursors over the same ring.
63#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
64pub struct RingCursor {
65    next_counter: u64
66}
67
68impl RingCursor {
69    /// Start at counter 0 and replay whatever ring history is still
70    /// available.
71    pub const fn from_start() -> Self {
72        Self { next_counter: 0 }
73    }
74
75    /// Start from a known next counter.
76    pub const fn from_counter(next_counter: u64) -> Self {
77        Self { next_counter }
78    }
79
80    /// The next counter this cursor will read.
81    pub const fn next_counter(self) -> u64 {
82        self.next_counter
83    }
84
85    pub(crate) fn set_next_counter(
86        &mut self,
87        next_counter: u64
88    ) {
89        self.next_counter = next_counter;
90    }
91}
92
93/// Counters skipped while walking a ring.
94#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
95pub struct RingLoss {
96    /// Counters older than the current ring window.
97    pub overwritten: u64,
98    /// Counters inside the readable window whose slot was definitively
99    /// unavailable, wrapped, corrupt, or carried an unexpected frame id.
100    pub unavailable: u64
101}
102
103impl RingLoss {
104    pub const fn total(self) -> u64 {
105        self.overwritten.saturating_add(self.unavailable)
106    }
107
108    pub const fn is_empty(self) -> bool {
109        self.overwritten == 0 && self.unavailable == 0
110    }
111}
112
113/// Result of walking a cursor toward a ring's visible head.
114#[derive(Clone, Debug, Default, PartialEq, Eq)]
115pub struct RingPoll {
116    pub frames: Vec<Frame>,
117    pub loss: RingLoss,
118    pub from_counter: u64,
119    pub to_counter: u64
120}
121
122impl RingPoll {
123    pub fn is_empty(&self) -> bool {
124        self.frames.is_empty() && self.loss.is_empty()
125    }
126}
127
128/// Walk `cursor` toward the current visible head of `source`.
129///
130/// At most one ring window is inspected. If the cursor has fallen behind
131/// the oldest available counter, the skipped counters are recorded as
132/// overwritten and the walk resumes at the window floor. An in-flight
133/// counter stops the walk without advancing past it.
134pub fn poll_ring<S: RingFrameSource>(
135    source: &S,
136    cursor: &mut RingCursor
137) -> RingPoll {
138    let head = source.head();
139    let from_counter = cursor.next_counter();
140
141    if from_counter >= head {
142        cursor.set_next_counter(head);
143        return RingPoll { from_counter, to_counter: head, ..RingPoll::default() };
144    }
145
146    let capacity = source.capacity() as u64;
147    let oldest_available = head.saturating_sub(capacity);
148    let mut next = from_counter;
149    let mut loss = RingLoss::default();
150
151    if next < oldest_available {
152        loss.overwritten = oldest_available - next;
153        next = oldest_available;
154    }
155
156    let kind = source.kind();
157    let mut frames = Vec::new();
158    while next < head {
159        match source.read_state_at(next) {
160            RingRead::Ready(frame) => {
161                if frame.id.kind() != kind || frame.id.counter() != next {
162                    loss.unavailable = loss.unavailable.saturating_add(1);
163                } else {
164                    frames.push(frame);
165                }
166                next = next.saturating_add(1);
167            }
168            RingRead::Pending => break,
169            RingRead::Unavailable => {
170                loss.unavailable = loss.unavailable.saturating_add(1);
171                next = next.saturating_add(1);
172            }
173        }
174    }
175
176    cursor.set_next_counter(next);
177    RingPoll { frames, loss, from_counter, to_counter: next }
178}