Skip to main content

orbit_core/fleet/
mod.rs

1//! `Fleet` — the per-process handle for one node address in a fleet's
2//! shared physical geometry.
3
4use std::sync::Arc;
5use std::sync::atomic::{AtomicU64, Ordering};
6
7use bytes::Bytes;
8use dashmap::DashMap;
9
10use crate::OrbitTyped;
11use crate::error::{Error, Result};
12use crate::id::NetId64;
13#[cfg(any(target_os = "linux", target_os = "freebsd"))]
14use crate::ring::RingEventFd;
15#[cfg(unix)]
16use crate::ring::shm::{ShmRing, ShmRingRegistry};
17use crate::ring::{Frame, Ring, RingRegistry, RingTopology};
18
19mod cursor;
20pub use cursor::{FleetLaneCursor, FleetLanePoll};
21
22/// A node's physical writer slot inside the fleet.
23///
24/// Orbit validates the address against fleet capacity but does not allocate
25/// it. The embedding runtime must ensure that simultaneously active writers
26/// receive distinct ids. Read-only inspectors can examine SHM without joining
27/// as a writer.
28#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
29#[repr(transparent)]
30pub struct NodeId(pub u16);
31
32impl NodeId {
33    pub const ZERO: Self = Self(0);
34
35    pub const fn new(value: u16) -> Self {
36        Self(value)
37    }
38
39    pub const fn get(self) -> u16 {
40        self.0
41    }
42}
43
44impl std::fmt::Display for NodeId {
45    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
46        write!(f, "node:{}", self.0)
47    }
48}
49
50/// Per-process handle into the fleet. Cheap to clone — the inner
51/// state is `Arc`-shared.
52#[derive(Clone)]
53pub struct Fleet {
54    inner: Arc<FleetInner>,
55}
56
57struct FleetInner {
58    name: &'static str,
59    fleet_capacity: u16,
60    node_id: NodeId,
61    /// Per-KIND counter for `next_id` calls that don't go through a
62    /// ring (i.e. when the caller wants a fleet-unique id without
63    /// allocating a ring slot). V0: process-local atomic. V1: still
64    /// here, parallel to the ring's own write-position.
65    id_counters: DashMap<u8, Arc<AtomicU64>>,
66    /// Per-KIND ring buffers — orbit's runtime substrate. Either
67    /// in-process for unit-test / single-process use, or POSIX SHM
68    /// for real cross-process visibility (V1, master+worker fleet).
69    backing: RingBacking,
70}
71
72/// Backing storage for the fleet's ring buffers — chosen at
73/// `Fleet::join` / `Fleet::join_shm` time and frozen for the
74/// fleet's lifetime.
75enum RingBacking {
76    /// Process-local DashMap of `Ring` instances. No cross-process
77    /// visibility — peers running other processes do not see this
78    /// fleet's writes. Useful for unit tests and embedded scenarios.
79    InMemory(RingRegistry),
80    /// POSIX-SHM-backed `ShmRing` instances. Multiple processes
81    /// joining the same fleet name share the same kernel-level
82    /// memory; writes from one are visible to all immediately.
83    #[cfg(unix)]
84    Shm(ShmRingRegistry),
85}
86
87impl Fleet {
88    /// Join (or create) a fleet under `name` with `fleet_capacity` physical
89    /// node lanes. In-memory backings remain process-local.
90    pub fn join(name: &'static str, fleet_capacity: u16) -> Result<Self> {
91        Self::join_as(name, fleet_capacity, NodeId::ZERO)
92    }
93
94    /// Join (or create) a process-local fleet with an explicit node id.
95    pub fn join_as(name: &'static str, fleet_capacity: u16, node_id: NodeId) -> Result<Self> {
96        if fleet_capacity == 0 {
97            return Err(Error::EmptyFleet);
98        }
99        if node_id.get() >= fleet_capacity {
100            return Err(Error::NodeOutsideFleet {
101                node_id: node_id.get(),
102                fleet_capacity,
103            });
104        }
105        Ok(Self {
106            inner: Arc::new(FleetInner {
107                name,
108                fleet_capacity,
109                node_id,
110                id_counters: DashMap::new(),
111                backing: RingBacking::InMemory(RingRegistry::new(fleet_capacity)),
112            }),
113        })
114    }
115
116    /// Join (or create) a fleet whose ring storage is backed by
117    /// POSIX shared memory. Multiple processes calling this with
118    /// the same `name` share the same kernel-level
119    /// segments — the fleet sees each other's writes.
120    ///
121    /// Each `OrbitTyped` kind gets its own SHM segment whose layout is
122    /// declared by `OrbitTyped::RING_SPEC`.
123    ///
124    /// Cross-process naming: segments are `/orbit-{name}-{kind}-{uid}`.
125    /// macOS limits POSIX SHM names to 31 chars (PSHMNAMLEN); a
126    /// short fleet name is required there.
127    #[cfg(unix)]
128    pub fn join_shm(name: &'static str, fleet_capacity: u16) -> Result<Self> {
129        Self::join_shm_as(name, fleet_capacity, NodeId::ZERO)
130    }
131
132    /// Join (or create) a SHM-backed fleet with an explicit node id.
133    ///
134    /// Per-node rings require one active process membership per node id.
135    /// Orbit validates the id range but does not own process lifecycle and
136    /// therefore cannot prevent duplicate live memberships.
137    #[cfg(unix)]
138    pub fn join_shm_as(name: &'static str, fleet_capacity: u16, node_id: NodeId) -> Result<Self> {
139        if fleet_capacity == 0 {
140            return Err(Error::EmptyFleet);
141        }
142        if node_id.get() >= fleet_capacity {
143            return Err(Error::NodeOutsideFleet {
144                node_id: node_id.get(),
145                fleet_capacity,
146            });
147        }
148        Ok(Self {
149            inner: Arc::new(FleetInner {
150                name,
151                fleet_capacity,
152                node_id,
153                id_counters: DashMap::new(),
154                backing: RingBacking::Shm(ShmRingRegistry::new(name, fleet_capacity)),
155            }),
156        })
157    }
158
159    pub fn name(&self) -> &'static str {
160        self.inner.name
161    }
162
163    /// Number of physical node lanes reserved for this fleet.
164    pub fn fleet_capacity(&self) -> u16 {
165        self.inner.fleet_capacity
166    }
167
168    pub fn node_id(&self) -> NodeId {
169        self.inner.node_id
170    }
171
172    /// Mint a fresh fleet-unique [`NetId64`] for type `T` *without*
173    /// publishing anything. Use this when the caller only needs the
174    /// id (e.g. minting an id to attach to data being persisted to
175    /// DB before going through the ring).
176    ///
177    /// For most use cases prefer [`Fleet::publish`] — it mints AND
178    /// stores in a single atomic step.
179    pub fn next_id<T: OrbitTyped>(&self) -> NetId64 {
180        let counter_arc = self
181            .inner
182            .id_counters
183            .entry(T::KIND)
184            .or_insert_with(|| Arc::new(AtomicU64::new(0)))
185            .clone();
186        let counter = counter_arc.fetch_add(1, Ordering::Relaxed);
187        NetId64::make(T::KIND, self.node_id().get(), counter)
188    }
189
190    /// True when this fleet's ring storage is backed by POSIX SHM
191    /// (visible across processes). False for in-memory fleets.
192    pub fn is_shm(&self) -> bool {
193        #[cfg(unix)]
194        {
195            matches!(self.inner.backing, RingBacking::Shm(_))
196        }
197        #[cfg(not(unix))]
198        {
199            false
200        }
201    }
202
203    /// Get-or-create the in-memory ring for type `T`. Only valid
204    /// for fleets created via [`Fleet::join`]; SHM-backed fleets
205    /// should use [`Fleet::shm_ring`] instead.
206    ///
207    /// # Panics
208    ///
209    /// Panics if called on a SHM-backed fleet.
210    pub fn ring<T: OrbitTyped>(&self) -> Arc<Ring> {
211        match &self.inner.backing {
212            RingBacking::InMemory(r) => r.get_or_create::<T>(),
213            #[cfg(unix)]
214            RingBacking::Shm(_) => {
215                panic!("Fleet::ring called on SHM-backed fleet — use Fleet::shm_ring instead");
216            }
217        }
218    }
219
220    /// Get-or-create the SHM ring for type `T`. Only valid on fleets
221    /// created via [`Fleet::join_shm`].
222    ///
223    /// # Errors
224    ///
225    /// Returns an `io::Error` if the SHM segment cannot be opened
226    /// (permissions, name too long, etc.).
227    ///
228    /// # Panics
229    ///
230    /// Panics if called on an in-memory fleet.
231    #[cfg(unix)]
232    pub fn shm_ring<T: OrbitTyped>(&self) -> std::io::Result<Arc<ShmRing>> {
233        match &self.inner.backing {
234            RingBacking::Shm(r) => r.get_or_create_for::<T>(),
235            RingBacking::InMemory(_) => {
236                panic!("Fleet::shm_ring called on in-memory fleet — use Fleet::ring instead");
237            }
238        }
239    }
240
241    /// Publish a payload to the ring for type `T`. Mints a [`NetId64`],
242    /// writes the [`Frame`] to the appropriate shared or node-owned lane,
243    /// and returns the id.
244    ///
245    /// # Panics
246    ///
247    /// Panics if the ring cannot be opened or the payload exceeds
248    /// `T::RING_SPEC.payload_capacity`. Ring failures are
249    /// operator-visible, not silently ignored.
250    pub fn publish<T: OrbitTyped>(&self, frame_kind: u8, ver: u64, payload: Bytes) -> NetId64 {
251        match &self.inner.backing {
252            RingBacking::InMemory(r) => {
253                let ring = r.get_or_create::<T>();
254                ring.write(self.node_id(), frame_kind, ver, payload)
255            }
256            #[cfg(unix)]
257            RingBacking::Shm(r) => {
258                let ring = r
259                    .get_or_create_for::<T>()
260                    .expect("SHM ring open failed — fleet unusable");
261                ring.write(self.node_id(), frame_kind, ver, payload)
262                    .expect("SHM ring write failed")
263            }
264        }
265    }
266
267    /// Publish a contiguous batch into one ring lane.
268    ///
269    /// The returned ids are ordered and consecutive. For per-node rings the
270    /// lane head becomes visible only after every frame in the batch has been
271    /// committed. Semantic layers can use this to publish a multi-slot blob,
272    /// then publish a separate descriptor that references the first id and
273    /// frame count.
274    ///
275    /// # Panics
276    ///
277    /// Panics if the ring cannot be opened, a payload exceeds the declared
278    /// slot capacity, or the batch itself is larger than the ring.
279    pub fn publish_batch<T: OrbitTyped>(
280        &self,
281        frame_kind: u8,
282        ver: u64,
283        payloads: Vec<Bytes>,
284    ) -> Vec<NetId64> {
285        match &self.inner.backing {
286            RingBacking::InMemory(r) => {
287                let ring = r.get_or_create::<T>();
288                ring.write_batch(self.node_id(), frame_kind, ver, payloads)
289            }
290            #[cfg(unix)]
291            RingBacking::Shm(r) => {
292                let ring = r
293                    .get_or_create_for::<T>()
294                    .expect("SHM ring open failed — fleet unusable");
295                ring.write_batch(self.node_id(), frame_kind, ver, payloads)
296                    .expect("SHM ring batch write failed")
297            }
298        }
299    }
300
301    /// Look up a previously-published frame by its id. Returns the
302    /// frame if its slot still holds the same id (i.e. the ring has
303    /// not wrapped past it).
304    pub fn read(&self, id: NetId64) -> Option<Frame> {
305        match &self.inner.backing {
306            RingBacking::InMemory(r) => r.lookup(id.kind())?.read(id),
307            #[cfg(unix)]
308            RingBacking::Shm(r) => r.lookup(id.kind())?.read(id),
309        }
310    }
311
312    /// Read the most recent frame for type `T`. Per-node rings read this
313    /// fleet handle's local lane; shared rings read their sole lane.
314    pub fn read_head<T: OrbitTyped>(&self) -> Option<Frame> {
315        if T::RING_SPEC.topology == RingTopology::PerNode {
316            let head = self.lane_head::<T>(self.node_id());
317            return (head > 0)
318                .then(|| self.read_lane_at::<T>(self.node_id(), head - 1))
319                .flatten();
320        }
321        match &self.inner.backing {
322            RingBacking::InMemory(r) => {
323                let ring = r.get_or_create::<T>();
324                ring.read_head()
325            }
326            #[cfg(unix)]
327            RingBacking::Shm(r) => {
328                let ring = r.get_or_create_for::<T>().ok()?;
329                ring.read_head()
330            }
331        }
332    }
333
334    /// Current head for type `T`'s ring. Per-node rings report this fleet
335    /// handle's local committed head; shared rings report their sole lane's
336    /// visible head.
337    /// Lazily
338    /// creates / attaches the ring on first access — important for
339    /// cross-process readers, where a child process may need to
340    /// attach to a SHM segment a peer already populated. Returns 0
341    /// when the ring is fresh / no counters have been claimed.
342    pub fn head<T: OrbitTyped>(&self) -> u64 {
343        if T::RING_SPEC.topology == RingTopology::PerNode {
344            return self.lane_head::<T>(self.node_id());
345        }
346        match &self.inner.backing {
347            RingBacking::InMemory(r) => r.get_or_create::<T>().head(),
348            #[cfg(unix)]
349            RingBacking::Shm(r) => r
350                .get_or_create_for::<T>()
351                .map(|ring| ring.head())
352                .unwrap_or(0),
353        }
354    }
355
356    /// Read whatever frame currently occupies `counter % capacity`.
357    /// Per-node rings read this fleet handle's local lane. Lazily attaches
358    /// the ring on first access (same rationale as [`Fleet::head`]).
359    /// Returns `None` if the slot is empty/torn or attach fails.
360    ///
361    /// Used by walking readers; for typed handle-based reads,
362    /// prefer [`Fleet::read`].
363    pub fn read_at<T: OrbitTyped>(&self, counter: u64) -> Option<Frame> {
364        if T::RING_SPEC.topology == RingTopology::PerNode {
365            return self.read_lane_at::<T>(self.node_id(), counter);
366        }
367        match &self.inner.backing {
368            RingBacking::InMemory(r) => r.get_or_create::<T>().read_at(counter),
369            #[cfg(unix)]
370            RingBacking::Shm(r) => r.get_or_create_for::<T>().ok()?.read_at(counter),
371        }
372    }
373
374    pub(crate) fn read_state_at<T: OrbitTyped>(
375        &self,
376        counter: u64,
377    ) -> crate::ring::cursor::RingRead {
378        if T::RING_SPEC.topology == RingTopology::PerNode {
379            return self.read_lane_state_at::<T>(self.node_id(), counter);
380        }
381        match &self.inner.backing {
382            RingBacking::InMemory(r) => r.get_or_create::<T>().read_state_at(counter),
383            #[cfg(unix)]
384            RingBacking::Shm(r) => r
385                .get_or_create_for::<T>()
386                .map(|ring| ring.read_state_at(counter))
387                .unwrap_or(crate::ring::cursor::RingRead::Unavailable),
388        }
389    }
390
391    /// Current head for one physical node lane.
392    ///
393    /// On a shared ring every node id addresses the sole shared lane.
394    pub fn lane_head<T: OrbitTyped>(&self, node_id: NodeId) -> u64 {
395        match &self.inner.backing {
396            RingBacking::InMemory(r) => r.get_or_create::<T>().lane_head(node_id),
397            #[cfg(unix)]
398            RingBacking::Shm(r) => r
399                .get_or_create_for::<T>()
400                .map(|ring| ring.lane_head(node_id))
401                .unwrap_or(0),
402        }
403    }
404
405    /// Read the frame currently occupying one node lane's counter slot.
406    pub fn read_lane_at<T: OrbitTyped>(&self, node_id: NodeId, counter: u64) -> Option<Frame> {
407        match &self.inner.backing {
408            RingBacking::InMemory(r) => r.get_or_create::<T>().read_lane_at(node_id, counter),
409            #[cfg(unix)]
410            RingBacking::Shm(r) => r
411                .get_or_create_for::<T>()
412                .ok()?
413                .read_lane_at(node_id, counter),
414        }
415    }
416
417    pub(crate) fn read_lane_state_at<T: OrbitTyped>(
418        &self,
419        node_id: NodeId,
420        counter: u64,
421    ) -> crate::ring::cursor::RingRead {
422        match &self.inner.backing {
423            RingBacking::InMemory(r) => r.get_or_create::<T>().read_lane_state_at(node_id, counter),
424            #[cfg(unix)]
425            RingBacking::Shm(r) => r
426                .get_or_create_for::<T>()
427                .map(|ring| ring.read_lane_state_at(node_id, counter))
428                .unwrap_or(crate::ring::cursor::RingRead::Unavailable),
429        }
430    }
431
432    /// Capacity of the ring for type `T`. Lazily attaches the ring
433    /// on first access. Falls back to `T::RING_SPEC.capacity` when
434    /// SHM attach fails.
435    pub fn ring_capacity<T: OrbitTyped>(&self) -> usize {
436        match &self.inner.backing {
437            RingBacking::InMemory(r) => r.get_or_create::<T>().capacity(),
438            #[cfg(unix)]
439            RingBacking::Shm(r) => r
440                .get_or_create_for::<T>()
441                .map(|ring| ring.capacity())
442                .unwrap_or(T::RING_SPEC.capacity),
443        }
444    }
445
446    /// Allocate one semantic version shared by every writer lane of `T`.
447    ///
448    /// This is separate from each lane's physical frame counter. It is useful
449    /// for semantic layers that retain per-node write scalability but require
450    /// a deterministic fleet-wide last-write-wins order.
451    pub fn next_ring_version<T: OrbitTyped>(&self) -> u64 {
452        match &self.inner.backing {
453            RingBacking::InMemory(r) => r.get_or_create::<T>().next_version(),
454            #[cfg(unix)]
455            RingBacking::Shm(r) => r
456                .get_or_create_for::<T>()
457                .expect("SHM ring open failed — fleet unusable")
458                .next_version(),
459        }
460    }
461
462    /// Return the last semantic version allocated for `T` without advancing it.
463    pub fn current_ring_version<T: OrbitTyped>(&self) -> u64 {
464        match &self.inner.backing {
465            RingBacking::InMemory(r) => r.get_or_create::<T>().current_version(),
466            #[cfg(unix)]
467            RingBacking::Shm(r) => r
468                .get_or_create_for::<T>()
469                .expect("SHM ring open failed — fleet unusable")
470                .current_version(),
471        }
472    }
473
474    /// Clear every lane for `T` and reset all heads to zero.
475    ///
476    /// This is an owner-side boot cleanup primitive. It is safe for
477    /// runtime state such as events and periodic metrics when the
478    /// embedding application calls it before peer processes begin
479    /// publishing. It is not a coordination protocol; callers must not
480    /// reset a ring while other fleet members are actively writing it.
481    pub fn reset_ring<T: OrbitTyped>(&self) -> std::io::Result<()> {
482        match &self.inner.backing {
483            RingBacking::InMemory(r) => {
484                r.get_or_create::<T>().reset();
485                Ok(())
486            }
487            #[cfg(unix)]
488            RingBacking::Shm(r) => {
489                r.get_or_create_for::<T>()?.reset();
490                Ok(())
491            }
492        }
493    }
494
495    #[cfg(any(target_os = "linux", target_os = "freebsd"))]
496    /// Create a process-local readiness fd for one notified SHM ring.
497    ///
498    /// The fd only signals that the ring generation changed. After draining
499    /// it, callers must poll the ring with their own cursor. Multiple writes
500    /// may coalesce into one readiness notification.
501    pub fn ring_event_fd<T: OrbitTyped>(&self) -> std::io::Result<RingEventFd> {
502        match &self.inner.backing {
503            RingBacking::Shm(rings) => RingEventFd::new(rings.get_or_create_for::<T>()?),
504            RingBacking::InMemory(_) => Err(std::io::Error::new(
505                std::io::ErrorKind::Unsupported,
506                "Orbit eventfd requires a shared-memory fleet",
507            )),
508        }
509    }
510
511    #[cfg(any(target_os = "linux", target_os = "freebsd"))]
512    /// Publish one frame and notify native waiters after it commits.
513    pub fn publish_notified<T: OrbitTyped>(
514        &self,
515        frame_kind: u8,
516        ver: u64,
517        payload: Bytes,
518    ) -> std::io::Result<NetId64> {
519        match &self.inner.backing {
520            RingBacking::Shm(rings) => {
521                let ring = rings.get_or_create_for::<T>()?;
522                let id = ring.write(self.node_id(), frame_kind, ver, payload)?;
523                RingEventFd::notify(&ring)?;
524                Ok(id)
525            }
526            // Process-local fleets do not need a kernel wake bridge.
527            RingBacking::InMemory(rings) => {
528                Ok(rings
529                    .get_or_create::<T>()
530                    .write(self.node_id(), frame_kind, ver, payload))
531            }
532        }
533    }
534
535    #[cfg(any(target_os = "linux", target_os = "freebsd"))]
536    /// Publish one contiguous batch and notify native waiters once after the
537    /// complete batch commits.
538    pub fn publish_batch_notified<T: OrbitTyped>(
539        &self,
540        frame_kind: u8,
541        ver: u64,
542        payloads: Vec<Bytes>,
543    ) -> std::io::Result<Vec<NetId64>> {
544        match &self.inner.backing {
545            RingBacking::Shm(rings) => {
546                let ring = rings.get_or_create_for::<T>()?;
547                let ids = ring.write_batch(self.node_id(), frame_kind, ver, payloads)?;
548                if !ids.is_empty() {
549                    RingEventFd::notify(&ring)?;
550                }
551                Ok(ids)
552            }
553            RingBacking::InMemory(rings) => Ok(rings.get_or_create::<T>().write_batch(
554                self.node_id(),
555                frame_kind,
556                ver,
557                payloads,
558            )),
559        }
560    }
561}
562
563impl std::fmt::Debug for Fleet {
564    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
565        f.debug_struct("Fleet")
566            .field("name", &self.inner.name)
567            .field("fleet_capacity", &self.inner.fleet_capacity)
568            .field("node_id", &self.inner.node_id)
569            .field("id_counters", &self.inner.id_counters.len())
570            .finish()
571    }
572}