orbit_core/fleet/
cursor.rs1use 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#[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#[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 pub fn cursor_at_head<T: OrbitTyped>(&self) -> RingCursor {
118 RingCursor::from_counter(self.head::<T>())
119 }
120
121 pub const fn cursor_from_start<T: OrbitTyped>(&self) -> RingCursor {
123 let _ = self;
124 let _ = PhantomData::<T>;
125 RingCursor::from_start()
126 }
127
128 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 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 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 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}