1use std::sync::atomic::{AtomicU64, Ordering};
50use std::sync::{Arc, Mutex, RwLock};
51
52use bytes::Bytes;
53
54use crate::id::NetId64;
55use crate::{NodeId, OrbitTyped};
56
57pub mod cursor;
58#[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
59mod readiness;
60#[cfg(unix)]
61pub mod shm;
62
63#[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
64pub use readiness::{ParkedRingEventFd, RingEventFd};
65
66#[derive(Clone, Copy, Debug, PartialEq, Eq)]
68#[repr(u8)]
69pub enum RingTopology {
70 Shared = 0,
72 PerNode = 1,
76 SharedOrdered = 2
80}
81
82#[derive(Clone, Copy, Debug, PartialEq, Eq)]
87pub struct RingSpec {
88 pub capacity: usize,
90 pub payload_capacity: usize,
92 pub topology: RingTopology
94}
95
96impl RingSpec {
97 pub const fn new(
98 capacity: usize,
99 payload_capacity: usize
100 ) -> Self {
101 Self { capacity, payload_capacity, topology: RingTopology::Shared }
102 }
103
104 pub const fn per_node(
106 capacity: usize,
107 payload_capacity: usize
108 ) -> Self {
109 Self { capacity, payload_capacity, topology: RingTopology::PerNode }
110 }
111
112 pub const fn shared_ordered(
114 capacity: usize,
115 payload_capacity: usize
116 ) -> Self {
117 Self { capacity, payload_capacity, topology: RingTopology::SharedOrdered }
118 }
119
120 pub(crate) fn assert_valid(self) {
121 assert!(self.capacity > 0, "ring capacity must be > 0");
122 assert!(self.capacity.is_power_of_two(), "ring capacity must be a power of two");
123 assert!(
124 self.payload_capacity <= u32::MAX as usize,
125 "ring payload capacity must fit in u32"
126 );
127 }
128}
129
130struct RingLane {
131 write_pos: AtomicU64,
132 write_lock: Mutex<()>,
133 slots: Vec<RwLock<Option<Frame>>>
134}
135
136impl RingLane {
137 fn new(capacity: usize) -> Self {
138 let mut slots = Vec::with_capacity(capacity);
139 for _ in 0..capacity {
140 slots.push(RwLock::new(None));
141 }
142 Self { write_pos: AtomicU64::new(0), write_lock: Mutex::new(()), slots }
143 }
144}
145
146#[derive(Clone, Debug, PartialEq, Eq)]
148pub struct Frame {
149 pub id: NetId64,
150 pub kind: u8,
151 pub ver: u64,
152 pub payload: Bytes
153}
154
155pub struct Ring {
159 kind: u8,
162 capacity: usize,
164 payload_capacity: usize,
166 topology: RingTopology,
167 version_counter: AtomicU64,
169 lanes: Vec<RingLane>
170}
171
172impl Ring {
173 pub fn new<T: OrbitTyped>() -> Self {
176 Self::new_for_fleet::<T>(1)
177 }
178
179 pub fn new_for_fleet<T: OrbitTyped>(fleet_capacity: u16) -> Self {
182 assert!(fleet_capacity > 0, "ring fleet capacity must be > 0");
183 let spec = T::RING_SPEC;
184 spec.assert_valid();
185 let capacity = spec.capacity;
186 let lane_count = match spec.topology {
187 RingTopology::Shared | RingTopology::SharedOrdered => 1,
188 RingTopology::PerNode => usize::from(fleet_capacity)
189 };
190 let mut lanes = Vec::with_capacity(lane_count);
191 for _ in 0..lane_count {
192 lanes.push(RingLane::new(capacity));
193 }
194 Self {
195 kind: T::KIND,
196 capacity,
197 payload_capacity: spec.payload_capacity,
198 topology: spec.topology,
199 version_counter: AtomicU64::new(0),
200 lanes
201 }
202 }
203
204 pub fn kind(&self) -> u8 {
206 self.kind
207 }
208
209 pub fn capacity(&self) -> usize {
211 self.capacity
212 }
213
214 pub fn payload_capacity(&self) -> usize {
216 self.payload_capacity
217 }
218
219 pub fn spec(&self) -> RingSpec {
220 RingSpec {
221 capacity: self.capacity,
222 payload_capacity: self.payload_capacity,
223 topology: self.topology
224 }
225 }
226
227 pub fn head(&self) -> u64 {
229 self.lanes[0].write_pos.load(Ordering::Acquire)
230 }
231
232 pub fn lane_count(&self) -> usize {
234 self.lanes.len()
235 }
236
237 pub fn lane_head(
239 &self,
240 node_id: NodeId
241 ) -> u64 {
242 self.lane(node_id).write_pos.load(Ordering::Acquire)
243 }
244
245 pub fn next_version(&self) -> u64 {
251 self.version_counter
252 .fetch_add(1, Ordering::AcqRel)
253 .checked_add(1)
254 .expect("ring semantic version exhausted")
255 }
256
257 pub fn current_version(&self) -> u64 {
259 self.version_counter.load(Ordering::Acquire)
260 }
261
262 pub fn write(
270 &self,
271 node_id: NodeId,
272 frame_kind: u8,
273 ver: u64,
274 payload: Bytes
275 ) -> NetId64 {
276 assert!(
277 payload.len() <= self.payload_capacity,
278 "payload {} > ring payload capacity {}",
279 payload.len(),
280 self.payload_capacity
281 );
282 let lane = self.lane(node_id);
283 match self.topology {
284 RingTopology::Shared => {
285 let counter = lane.write_pos.fetch_add(1, Ordering::AcqRel);
286 self.write_frame(lane, node_id, counter, frame_kind, ver, payload)
287 }
288 RingTopology::PerNode | RingTopology::SharedOrdered => {
289 let _write = lane.write_lock.lock().unwrap_or_else(|error| error.into_inner());
290 let counter = lane.write_pos.load(Ordering::Relaxed);
291 let id = self.write_frame(lane, node_id, counter, frame_kind, ver, payload);
292 lane.write_pos.store(counter.wrapping_add(1), Ordering::Release);
293 id
294 }
295 }
296 }
297
298 pub fn write_batch(
305 &self,
306 node_id: NodeId,
307 frame_kind: u8,
308 ver: u64,
309 payloads: Vec<Bytes>
310 ) -> Vec<NetId64> {
311 assert!(
312 payloads.len() <= self.capacity,
313 "batch {} > ring capacity {}",
314 payloads.len(),
315 self.capacity
316 );
317 for payload in &payloads {
318 assert!(
319 payload.len() <= self.payload_capacity,
320 "payload {} > ring payload capacity {}",
321 payload.len(),
322 self.payload_capacity
323 );
324 }
325 if payloads.is_empty() {
326 return Vec::new();
327 }
328
329 let lane = self.lane(node_id);
330 match self.topology {
331 RingTopology::Shared => {
332 let start = lane.write_pos.fetch_add(payloads.len() as u64, Ordering::AcqRel);
333 payloads
334 .into_iter()
335 .enumerate()
336 .map(|(offset, payload)| {
337 self.write_frame(
338 lane,
339 node_id,
340 start.wrapping_add(offset as u64),
341 frame_kind,
342 ver,
343 payload
344 )
345 })
346 .collect()
347 }
348 RingTopology::PerNode | RingTopology::SharedOrdered => {
349 let _write = lane.write_lock.lock().unwrap_or_else(|error| error.into_inner());
350 let start = lane.write_pos.load(Ordering::Relaxed);
351 let ids = payloads
352 .into_iter()
353 .enumerate()
354 .map(|(offset, payload)| {
355 self.write_frame(
356 lane,
357 node_id,
358 start.wrapping_add(offset as u64),
359 frame_kind,
360 ver,
361 payload
362 )
363 })
364 .collect::<Vec<_>>();
365 lane.write_pos.store(start.wrapping_add(ids.len() as u64), Ordering::Release);
366 ids
367 }
368 }
369 }
370
371 pub fn read(
379 &self,
380 id: NetId64
381 ) -> Option<Frame> {
382 if id.kind() != self.kind {
383 return None;
384 }
385 let lane = self.lane_for_frame(id)?;
386 let slot_idx = (id.counter() as usize) % self.capacity;
387 let guard = lane.slots[slot_idx].read().expect("ring slot poisoned");
388 match &*guard {
389 Some(f) if f.id == id => Some(f.clone()),
390 _ => None
391 }
392 }
393
394 pub fn read_head(&self) -> Option<Frame> {
398 let head = self.head();
399 if head == 0 {
400 return None;
401 }
402 let slot_idx = ((head - 1) as usize) % self.capacity;
403 self.lanes[0].slots[slot_idx].read().expect("ring slot poisoned").clone()
404 }
405
406 pub fn read_at(
413 &self,
414 counter: u64
415 ) -> Option<Frame> {
416 let slot_idx = (counter as usize) % self.capacity;
417 self.lanes[0].slots[slot_idx].read().expect("ring slot poisoned").clone()
418 }
419
420 pub(crate) fn read_state_at(
421 &self,
422 counter: u64
423 ) -> cursor::RingRead {
424 match self.read_at(counter) {
425 Some(frame) if frame.id.counter() == counter => cursor::RingRead::Ready(frame),
426 Some(frame) if frame.id.counter() > counter => cursor::RingRead::Unavailable,
427 Some(_) | None => cursor::RingRead::Pending
428 }
429 }
430
431 pub(crate) fn read_lane_at(
432 &self,
433 node_id: NodeId,
434 counter: u64
435 ) -> Option<Frame> {
436 let lane = self.lane(node_id);
437 let slot_idx = (counter as usize) % self.capacity;
438 lane.slots[slot_idx].read().expect("ring slot poisoned").clone()
439 }
440
441 pub(crate) fn read_lane_state_at(
442 &self,
443 node_id: NodeId,
444 counter: u64
445 ) -> cursor::RingRead {
446 match self.read_lane_at(node_id, counter) {
447 Some(frame) if frame.id.counter() == counter => cursor::RingRead::Ready(frame),
448 Some(frame) if frame.id.counter() > counter => cursor::RingRead::Unavailable,
449 Some(_) | None if self.topology != RingTopology::Shared => {
450 cursor::RingRead::Unavailable
451 }
452 Some(_) | None => cursor::RingRead::Pending
453 }
454 }
455
456 pub fn reset(&self) {
461 for lane in &self.lanes {
462 for slot in &lane.slots {
463 *slot.write().expect("ring slot poisoned") = None;
464 }
465 lane.write_pos.store(0, Ordering::Release);
466 }
467 self.version_counter.store(0, Ordering::Release);
468 }
469
470 fn lane(
471 &self,
472 node_id: NodeId
473 ) -> &RingLane {
474 let index = match self.topology {
475 RingTopology::Shared | RingTopology::SharedOrdered => 0,
476 RingTopology::PerNode => usize::from(node_id.get())
477 };
478 self.lanes.get(index).unwrap_or_else(|| {
479 panic!("node {} is outside ring lane count {}", node_id.get(), self.lanes.len())
480 })
481 }
482
483 fn lane_for_frame(
484 &self,
485 id: NetId64
486 ) -> Option<&RingLane> {
487 let index = match self.topology {
488 RingTopology::Shared | RingTopology::SharedOrdered => 0,
489 RingTopology::PerNode => usize::from(id.node())
490 };
491 self.lanes.get(index)
492 }
493
494 fn write_frame(
495 &self,
496 lane: &RingLane,
497 node_id: NodeId,
498 counter: u64,
499 frame_kind: u8,
500 ver: u64,
501 payload: Bytes
502 ) -> NetId64 {
503 let id = NetId64::make(self.kind, node_id.get(), counter);
504 let slot_idx = (counter as usize) % self.capacity;
505 let frame = Frame { id, kind: frame_kind, ver, payload };
506 let mut guard = lane.slots[slot_idx].write().expect("ring slot poisoned");
507 *guard = Some(frame);
508 id
509 }
510}
511
512impl cursor::RingFrameSource for Ring {
513 fn kind(&self) -> u8 {
514 Ring::kind(self)
515 }
516
517 fn head(&self) -> u64 {
518 Ring::head(self)
519 }
520
521 fn capacity(&self) -> usize {
522 Ring::capacity(self)
523 }
524
525 fn read_at(
526 &self,
527 counter: u64
528 ) -> Option<Frame> {
529 Ring::read_at(self, counter)
530 }
531
532 fn read_state_at(
533 &self,
534 counter: u64
535 ) -> cursor::RingRead {
536 Ring::read_state_at(self, counter)
537 }
538}
539
540#[cfg(unix)]
541impl cursor::RingFrameSource for shm::ShmRing {
542 fn kind(&self) -> u8 {
543 shm::ShmRing::kind(self)
544 }
545
546 fn head(&self) -> u64 {
547 shm::ShmRing::head(self)
548 }
549
550 fn capacity(&self) -> usize {
551 shm::ShmRing::capacity(self)
552 }
553
554 fn read_at(
555 &self,
556 counter: u64
557 ) -> Option<Frame> {
558 shm::ShmRing::read_at(self, counter)
559 }
560
561 fn read_state_at(
562 &self,
563 counter: u64
564 ) -> cursor::RingRead {
565 shm::ShmRing::read_state_at(self, counter)
566 }
567}
568
569impl std::fmt::Debug for Ring {
570 fn fmt(
571 &self,
572 f: &mut std::fmt::Formatter<'_>
573 ) -> std::fmt::Result {
574 f.debug_struct("Ring")
575 .field("kind", &self.kind)
576 .field("capacity", &self.capacity)
577 .field("payload_capacity", &self.payload_capacity)
578 .field("topology", &self.topology)
579 .field("lane_count", &self.lanes.len())
580 .field("head", &self.head())
581 .finish()
582 }
583}
584
585pub(crate) struct RingRegistry {
588 fleet_capacity: u16,
589 rings: dashmap::DashMap<u8, Arc<Ring>>
590}
591
592impl RingRegistry {
593 pub fn new(fleet_capacity: u16) -> Self {
594 Self { fleet_capacity, rings: dashmap::DashMap::new() }
595 }
596
597 pub fn get_or_create<T: OrbitTyped>(&self) -> Arc<Ring> {
599 let ring = self
600 .rings
601 .entry(T::KIND)
602 .or_insert_with(|| Arc::new(Ring::new_for_fleet::<T>(self.fleet_capacity)))
603 .clone();
604 assert_eq!(
605 ring.spec(),
606 T::RING_SPEC,
607 "OrbitTyped KIND {} was reused with a different ring spec",
608 T::KIND
609 );
610 ring
611 }
612
613 pub fn lookup(
615 &self,
616 kind: u8
617 ) -> Option<Arc<Ring>> {
618 self.rings.get(&kind).map(|e| e.clone())
619 }
620}