Skip to main content

orbit_core/fleet/
cursor.rs

1use std::marker::PhantomData;
2
3use crate::OrbitTyped;
4use crate::fleet::{Fleet, NodeId};
5use crate::ring::cursor::{RingCursor, RingFrameSource, RingLoss, RingPoll, RingRead, poll_ring};
6use crate::ring::{Frame, RingTopology};
7
8struct FleetRingSource<'a, T: OrbitTyped> {
9    fleet: &'a Fleet,
10    _t: PhantomData<T>
11}
12
13struct FleetLaneSource<'a, T: OrbitTyped> {
14    fleet: &'a Fleet,
15    node_id: NodeId,
16    _t: PhantomData<T>
17}
18
19impl<'a, T: OrbitTyped> FleetLaneSource<'a, T> {
20    fn new(
21        fleet: &'a Fleet,
22        node_id: NodeId
23    ) -> Self {
24        Self { fleet, node_id, _t: PhantomData }
25    }
26}
27
28impl<T: OrbitTyped> RingFrameSource for FleetLaneSource<'_, T> {
29    fn kind(&self) -> u8 {
30        T::KIND
31    }
32
33    fn head(&self) -> u64 {
34        self.fleet.lane_head::<T>(self.node_id)
35    }
36
37    fn capacity(&self) -> usize {
38        self.fleet.ring_capacity::<T>()
39    }
40
41    fn read_at(
42        &self,
43        counter: u64
44    ) -> Option<Frame> {
45        self.fleet.read_lane_at::<T>(self.node_id, counter)
46    }
47
48    fn read_state_at(
49        &self,
50        counter: u64
51    ) -> RingRead {
52        self.fleet.read_lane_state_at::<T>(self.node_id, counter)
53    }
54}
55
56/// Caller-owned positions for every node lane of one ring type.
57#[derive(Clone, Debug, PartialEq, Eq)]
58pub struct FleetLaneCursor {
59    lanes: Vec<RingCursor>,
60    initial_counter: u64
61}
62
63impl FleetLaneCursor {
64    pub const fn from_counter(initial_counter: u64) -> Self {
65        Self { lanes: Vec::new(), initial_counter }
66    }
67
68    pub fn minimum_next_counter(&self) -> u64 {
69        self.lanes.iter().map(|cursor| cursor.next_counter()).min().unwrap_or(self.initial_counter)
70    }
71}
72
73/// Combined result of walking every node lane once.
74#[derive(Clone, Debug, Default, PartialEq, Eq)]
75pub struct FleetLanePoll {
76    pub frames: Vec<Frame>,
77    pub loss: RingLoss
78}
79
80impl<'a, T: OrbitTyped> FleetRingSource<'a, T> {
81    fn new(fleet: &'a Fleet) -> Self {
82        Self { fleet, _t: PhantomData }
83    }
84}
85
86impl<T: OrbitTyped> RingFrameSource for FleetRingSource<'_, T> {
87    fn kind(&self) -> u8 {
88        T::KIND
89    }
90
91    fn head(&self) -> u64 {
92        self.fleet.head::<T>()
93    }
94
95    fn capacity(&self) -> usize {
96        self.fleet.ring_capacity::<T>()
97    }
98
99    fn read_at(
100        &self,
101        counter: u64
102    ) -> Option<Frame> {
103        self.fleet.read_at::<T>(counter)
104    }
105
106    fn read_state_at(
107        &self,
108        counter: u64
109    ) -> RingRead {
110        self.fleet.read_state_at::<T>(counter)
111    }
112}
113
114impl Fleet {
115    /// Cursor that starts after every counter currently claimed for `T`.
116    /// Useful for subscribers that only want future writes.
117    pub fn cursor_at_head<T: OrbitTyped>(&self) -> RingCursor {
118        RingCursor::from_counter(self.head::<T>())
119    }
120
121    /// Cursor that starts at counter 0 for `T`.
122    pub const fn cursor_from_start<T: OrbitTyped>(&self) -> RingCursor {
123        let _ = self;
124        let _ = PhantomData::<T>;
125        RingCursor::from_start()
126    }
127
128    /// Walk `cursor` toward the current claim head for `T`, stopping at
129    /// an in-flight counter and reporting definitive losses as
130    /// [`crate::ring::cursor::RingLoss`].
131    pub fn poll_ring<T: OrbitTyped>(
132        &self,
133        cursor: &mut RingCursor
134    ) -> RingPoll {
135        poll_ring(&FleetRingSource::<T>::new(self), cursor)
136    }
137
138    /// Create one caller-owned cursor per physical node lane, starting at each
139    /// lane's current head.
140    pub fn lane_cursor_at_head<T: OrbitTyped>(&self) -> FleetLaneCursor {
141        self.assert_per_node::<T>();
142        let lanes = (0..self.fleet_capacity())
143            .map(|node| RingCursor::from_counter(self.lane_head::<T>(NodeId::new(node))))
144            .collect();
145        FleetLaneCursor { lanes, initial_counter: 0 }
146    }
147
148    /// Create one caller-owned cursor per physical node lane, starting at
149    /// counter zero.
150    pub fn lane_cursor_from_start<T: OrbitTyped>(&self) -> FleetLaneCursor {
151        self.assert_per_node::<T>();
152        FleetLaneCursor {
153            lanes: vec![RingCursor::from_start(); usize::from(self.fleet_capacity())],
154            initial_counter: 0
155        }
156    }
157
158    /// Poll every physical node lane and combine the retained frames and loss
159    /// counters into one result. Frames retain their writer node in `id`;
160    /// callers that need semantic ordering across lanes must provide it.
161    pub fn poll_lanes<T: OrbitTyped>(
162        &self,
163        cursor: &mut FleetLaneCursor
164    ) -> FleetLanePoll {
165        self.assert_per_node::<T>();
166        if cursor.lanes.is_empty() {
167            cursor.lanes = vec![
168                RingCursor::from_counter(cursor.initial_counter);
169                usize::from(self.fleet_capacity())
170            ];
171        }
172        assert_eq!(
173            cursor.lanes.len(),
174            usize::from(self.fleet_capacity()),
175            "lane cursor belongs to a different fleet capacity"
176        );
177
178        let mut combined = FleetLanePoll::default();
179        for (node, lane_cursor) in cursor.lanes.iter_mut().enumerate() {
180            let source = FleetLaneSource::<T>::new(self, NodeId::new(node as u16));
181            let poll = poll_ring(&source, lane_cursor);
182            combined.frames.extend(poll.frames);
183            combined.loss.overwritten =
184                combined.loss.overwritten.saturating_add(poll.loss.overwritten);
185            combined.loss.unavailable =
186                combined.loss.unavailable.saturating_add(poll.loss.unavailable);
187        }
188        combined
189    }
190
191    fn assert_per_node<T: OrbitTyped>(&self) {
192        assert_eq!(
193            T::RING_SPEC.topology,
194            RingTopology::PerNode,
195            "lane cursor requires a per-node ring"
196        );
197    }
198}