Skip to main content

orbit_core/ring/
mod.rs

1//! Ring buffers — orbit-core's runtime substrate.
2//!
3//! > *"orbit runtime yani ring"* — the place where the fleet's
4//! > shared state actually lives at the lowest level. Higher-level
5//! > shapes (cache mutations, metrics snapshots, event streams, etc.)
6//! > reduce to *one or more rings*.
7//!
8//! ## Shape
9//!
10//! One [`Ring`] per [`OrbitTyped`] kind. A ring is one or more fixed-size
11//! circular lanes of [`Frame`]s. Shared rings use one lock-free multi-writer
12//! claim sequence. Shared-ordered rings serialize writers and expose only
13//! committed counters. Per-node rings give every fleet member a disjoint lane
14//! whose head advances only after a frame is committed. When a lane head
15//! exceeds its capacity, the oldest slot in that lane is overwritten.
16//!
17//! The frame layout mirrors the `nwd1` seed (see VISION §13):
18//!
19//! ```text
20//! ┌──────────┬──────┬──────┬─────────────┐
21//! │ id (8)   │ kind │ ver  │ payload (N) │
22//! │ NetId64  │  u8  │ u64  │   bytes     │
23//! └──────────┴──────┴──────┴─────────────┘
24//! ```
25//!
26//! Two `kind` bytes coexist on the wire and they mean different
27//! things (intentional, two-axis encoding):
28//!
29//! - `frame.id.kind()` — *which Rust type* (the data shape).
30//! - `frame.kind`      — *which message class* (state / event /
31//!   command / ack / invalidate / …). V0 leaves this open at `0`;
32//!   concrete classes appear when subscriber semantics arrive.
33//!
34//! ## V0 backing
35//!
36//! `RwLock<Option<Frame>>` per slot — simple, correct, slow. The SHM
37//! backing uses per-slot sequence numbers to prevent torn reads. Per-node
38//! lanes serialize concurrent publishers inside one process, commit the
39//! slot, then release-store their lane head.
40//!
41//! ## Who mints
42//!
43//! Writes go through [`Ring::write`], which mints the [`NetId64`]. The
44//! COUNTER part is the position inside either the shared sequence or the
45//! writer node's lane. NetId64s are therefore minted
46//! **server-side, by the writer process**. Browsers / external
47//! clients receive ids; they do not generate them.
48
49use 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/// Writer ownership for one [`OrbitTyped`] ring.
67#[derive(Clone, Copy, Debug, PartialEq, Eq)]
68#[repr(u8)]
69pub enum RingTopology {
70    /// Every fleet member publishes into one shared multi-writer sequence.
71    Shared = 0,
72    /// Every fleet member owns an independent single-process writer lane.
73    /// The embedder must not run two active processes with the same node id;
74    /// local concurrent tasks are serialized by the ring handle.
75    PerNode = 1,
76    /// Every fleet member publishes into one globally ordered sequence.
77    /// Writers are serialized by a process-recoverable OS lock associated
78    /// with the SHM name, and the head advances only after the slot commits.
79    SharedOrdered = 2
80}
81
82/// Physical policy for one [`OrbitTyped`] ring.
83///
84/// The policy is part of the wire contract: every process using the
85/// same `OrbitTyped::KIND` must declare the same values.
86#[derive(Clone, Copy, Debug, PartialEq, Eq)]
87pub struct RingSpec {
88    /// Number of slots retained by the ring, per lane.
89    pub capacity: usize,
90    /// Maximum payload bytes stored inline in each slot.
91    pub payload_capacity: usize,
92    /// How writers own and publish physical lanes.
93    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    /// Declare one independent writer lane per fleet node.
105    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    /// Declare one crash-recoverable, globally ordered writer lane.
113    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/// One record in a ring — the on-wire shape (mirrors `nwd1::Frame`).
147#[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
155/// A fixed-capacity, fleet-wide append-only log keyed on KIND byte.
156///
157/// V0 is single-process; V1 is SHM-backed. The API is the same.
158pub struct Ring {
159    /// The KIND this ring carries — equals `T::KIND` for the
160    /// `OrbitTyped` value-shape it's storing.
161    kind: u8,
162    /// Number of slots; constant for the ring's lifetime.
163    capacity: usize,
164    /// Maximum inline payload bytes for this ring lane.
165    payload_capacity: usize,
166    topology: RingTopology,
167    /// Ring-wide semantic version allocator shared by every writer lane.
168    version_counter: AtomicU64,
169    lanes: Vec<RingLane>
170}
171
172impl Ring {
173    /// Create the process-local ring declared by `T::RING_SPEC` for a
174    /// single-member fleet.
175    pub fn new<T: OrbitTyped>() -> Self {
176        Self::new_for_fleet::<T>(1)
177    }
178
179    /// Create the process-local ring declared by `T::RING_SPEC` with the
180    /// physical lane count required by `fleet_capacity`.
181    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    /// The KIND byte this ring carries (equals `T::KIND`).
205    pub fn kind(&self) -> u8 {
206        self.kind
207    }
208
209    /// Total slot count — fixed at construction.
210    pub fn capacity(&self) -> usize {
211        self.capacity
212    }
213
214    /// Maximum inline payload bytes for this ring lane.
215    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    /// Head of the sole shared lane, or lane zero for a per-node ring.
228    pub fn head(&self) -> u64 {
229        self.lanes[0].write_pos.load(Ordering::Acquire)
230    }
231
232    /// Number of physical writer lanes in this ring.
233    pub fn lane_count(&self) -> usize {
234        self.lanes.len()
235    }
236
237    /// Current head for `node_id`'s logical lane.
238    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    /// Allocate one non-zero semantic version shared by every writer lane.
246    ///
247    /// This counter is independent of physical ring positions. Semantic
248    /// layers can use it when per-node lanes need one deterministic
249    /// last-write-wins order.
250    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    /// Last semantic version allocated for this ring.
258    pub fn current_version(&self) -> u64 {
259        self.version_counter.load(Ordering::Acquire)
260    }
261
262    /// Append a frame. Atomically reserves the next counter, mints
263    /// the [`NetId64`], and writes the frame into the corresponding
264    /// slot. Returns the minted id.
265    ///
266    /// `frame_kind` is the message class byte (V0: pass `0`).
267    /// `ver` is the version / tick at write time (V0: caller's
268    /// choice).
269    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    /// Append a contiguous batch to one lane and return its consecutive ids.
299    ///
300    /// Per-node and shared-ordered lanes expose the new head only after the
301    /// whole batch commits. An empty batch is a no-op. A batch larger than the
302    /// ring is rejected because its first frames could not remain addressable
303    /// when the method returns.
304    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    /// Read the slot that the given [`NetId64`]'s counter points at.
372    ///
373    /// Returns:
374    /// - `Some(frame)` if the slot's stored id matches the queried id
375    ///   exactly (the slot has not been overwritten by a later writer).
376    /// - `None` if the slot is empty, has wrapped past, or holds a
377    ///   different id than the one asked for.
378    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    /// Read the most recent frame, regardless of who wrote it.
395    /// Useful for "what's the current state?" — ignores
396    /// counter-by-counter walking.
397    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    /// Read whatever frame currently occupies the slot at
407    /// `counter % capacity`, regardless of which counter is
408    /// stored in it. Used by walking readers that need slot-by-slot
409    /// access without knowing the writer's `NetId64` ahead of time.
410    ///
411    /// Returns `None` if the slot is empty.
412    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    /// Clear all slots and reset every lane head to zero.
457    ///
458    /// Intended for owner-controlled boot-time cleanup. Do not call
459    /// while other threads are publishing to this ring.
460    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
585/// Type-keyed registry of rings. A `Fleet` holds one of these and
586/// hands out `Arc<Ring>` per `OrbitTyped` kind on demand.
587pub(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    /// Get-or-create the ring declared by `T`.
598    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    /// Look up a ring by KIND byte (e.g. when only the id is known).
614    pub fn lookup(
615        &self,
616        kind: u8
617    ) -> Option<Arc<Ring>> {
618        self.rings.get(&kind).map(|e| e.clone())
619    }
620}