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(unix)]
14use crate::ring::shm::{ShmRing, ShmRingRegistry};
15use crate::ring::{Frame, Ring, RingRegistry, RingTopology};
16#[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
17use crate::ring::{ParkedRingEventFd, RingEventFd};
18
19mod cursor;
20pub use cursor::{FleetLaneCursor, FleetLanePoll};
21
22/// Read-only namespace handle for inspecting an existing SHM fleet.
23///
24/// Unlike [`Fleet`], this handle has no node id, creates no rings, and can
25/// never publish, reset, or unlink. Opening a ring maps only an object that
26/// already exists; dropping the observer or a ring view only unmaps local
27/// memory.
28#[cfg(unix)]
29#[derive(Clone, Debug, PartialEq, Eq)]
30pub struct FleetObserver {
31    name: String,
32    uid: u32
33}
34
35#[cfg(unix)]
36impl FleetObserver {
37    /// Create an observer for the effective user's existing fleet namespace.
38    ///
39    /// This constructor itself performs no SHM operation. Each [`Self::ring`]
40    /// call attaches to one exact existing kind and returns `NotFound` when it
41    /// is absent.
42    pub fn attach_existing(name: impl Into<String>) -> std::io::Result<Self> {
43        // SAFETY: `geteuid` has no error path.
44        let uid = unsafe { libc::geteuid() };
45        Self::attach_existing_for_uid(name, uid)
46    }
47
48    /// Address an explicit uid-scoped fleet namespace.
49    ///
50    /// POSIX permissions still decide whether the caller may read another
51    /// user's SHM objects.
52    pub fn attach_existing_for_uid(
53        name: impl Into<String>,
54        uid: u32
55    ) -> std::io::Result<Self> {
56        let name = name.into();
57        if name.is_empty() {
58            return Err(std::io::Error::new(
59                std::io::ErrorKind::InvalidInput,
60                "fleet name must not be empty"
61            ));
62        }
63        if name.contains('/') || name.contains('\0') {
64            return Err(std::io::Error::new(
65                std::io::ErrorKind::InvalidInput,
66                "fleet name must not contain '/' or a NUL byte"
67            ));
68        }
69        Ok(Self { name, uid })
70    }
71
72    pub fn name(&self) -> &str {
73        &self.name
74    }
75
76    pub fn uid(&self) -> u32 {
77        self.uid
78    }
79
80    /// Attach read-only to one exact existing ring kind.
81    pub fn ring(
82        &self,
83        kind: u8
84    ) -> std::io::Result<crate::ring::shm::ShmRingView> {
85        crate::ring::shm::ShmRingView::attach_existing_for_uid(&self.name, kind, self.uid)
86    }
87
88    /// Attach read-only and verify a linked [`OrbitTyped`] contract.
89    pub fn typed_ring<T: OrbitTyped>(&self) -> std::io::Result<crate::ring::shm::ShmRingView> {
90        let view = self.ring(T::KIND)?;
91        if view.metadata().spec != T::RING_SPEC {
92            return Err(std::io::Error::new(
93                std::io::ErrorKind::InvalidData,
94                format!(
95                    "OrbitTyped KIND {} declares {:?}; existing spec is {:?}",
96                    T::KIND,
97                    T::RING_SPEC,
98                    view.metadata().spec
99                )
100            ));
101        }
102        Ok(view)
103    }
104}
105
106/// A node's physical writer slot inside the fleet.
107///
108/// Orbit validates the address against fleet capacity but does not allocate
109/// it. The embedding runtime must ensure that simultaneously active writers
110/// receive distinct ids. Read-only inspectors can examine SHM without joining
111/// as a writer.
112#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
113#[repr(transparent)]
114pub struct NodeId(pub u16);
115
116impl NodeId {
117    pub const ZERO: Self = Self(0);
118
119    pub const fn new(value: u16) -> Self {
120        Self(value)
121    }
122
123    pub const fn get(self) -> u16 {
124        self.0
125    }
126}
127
128impl std::fmt::Display for NodeId {
129    fn fmt(
130        &self,
131        f: &mut std::fmt::Formatter<'_>
132    ) -> std::fmt::Result {
133        write!(f, "node:{}", self.0)
134    }
135}
136
137/// Per-process handle into the fleet. Cheap to clone — the inner
138/// state is `Arc`-shared.
139#[derive(Clone)]
140pub struct Fleet {
141    inner: Arc<FleetInner>
142}
143
144struct FleetInner {
145    name: Arc<str>,
146    fleet_capacity: u16,
147    node_id: NodeId,
148    /// Per-KIND counter for `next_id` calls that don't go through a
149    /// ring (i.e. when the caller wants a fleet-unique id without
150    /// allocating a ring slot). V0: process-local atomic. V1: still
151    /// here, parallel to the ring's own write-position.
152    id_counters: DashMap<u8, Arc<AtomicU64>>,
153    /// Per-KIND ring buffers — orbit's runtime substrate. Either
154    /// in-process for unit-test / single-process use, or POSIX SHM
155    /// for real cross-process visibility (V1, master+worker fleet).
156    backing: RingBacking,
157    /// Held for as long as this process is in a shared-memory fleet, so
158    /// a lifecycle tool can tell a live fleet from a stopped one and refuse
159    /// to remove what is in use. Read-only attachments hold none.
160    #[cfg(unix)]
161    #[allow(dead_code)]
162    membership: Option<crate::shm::FleetMembership>
163}
164
165/// Backing storage for the fleet's ring buffers — chosen at
166/// `Fleet::join` / `Fleet::join_shm` time and frozen for the
167/// fleet's lifetime.
168enum RingBacking {
169    /// Process-local DashMap of `Ring` instances. No cross-process
170    /// visibility — peers running other processes do not see this
171    /// fleet's writes. Useful for unit tests and embedded scenarios.
172    InMemory(RingRegistry),
173    /// POSIX-SHM-backed `ShmRing` instances. Multiple processes
174    /// joining the same fleet name share the same kernel-level
175    /// memory; writes from one are visible to all immediately.
176    #[cfg(unix)]
177    Shm(ShmRingRegistry)
178}
179
180impl Fleet {
181    /// Join (or create) a fleet under `name` with `fleet_capacity` physical
182    /// node lanes. In-memory backings remain process-local.
183    pub fn join(
184        name: &str,
185        fleet_capacity: u16
186    ) -> Result<Self> {
187        Self::join_as(name, fleet_capacity, NodeId::ZERO)
188    }
189
190    /// Join (or create) a process-local fleet with an explicit node id.
191    pub fn join_as(
192        name: &str,
193        fleet_capacity: u16,
194        node_id: NodeId
195    ) -> Result<Self> {
196        if fleet_capacity == 0 {
197            return Err(Error::EmptyFleet);
198        }
199        if node_id.get() >= fleet_capacity {
200            return Err(Error::NodeOutsideFleet { node_id: node_id.get(), fleet_capacity });
201        }
202        Ok(Self {
203            inner: Arc::new(FleetInner {
204                name: Arc::from(name),
205                fleet_capacity,
206                node_id,
207                id_counters: DashMap::new(),
208                backing: RingBacking::InMemory(RingRegistry::new(fleet_capacity)),
209                #[cfg(unix)]
210                membership: None
211            })
212        })
213    }
214
215    /// Join (or create) a fleet whose ring storage is backed by
216    /// POSIX shared memory. Multiple processes calling this with
217    /// the same `name` share the same kernel-level
218    /// segments — the fleet sees each other's writes.
219    ///
220    /// Each `OrbitTyped` kind gets its own SHM segment whose layout is
221    /// declared by `OrbitTyped::RING_SPEC`.
222    ///
223    /// Cross-process naming: segments are `/orbit-{name}-{kind}-{uid}`.
224    /// macOS limits POSIX SHM names to 31 chars (PSHMNAMLEN); a
225    /// short fleet name is required there.
226    #[cfg(unix)]
227    pub fn join_shm(
228        name: &str,
229        fleet_capacity: u16
230    ) -> Result<Self> {
231        Self::join_shm_as(name, fleet_capacity, NodeId::ZERO)
232    }
233
234    /// Join (or create) a SHM-backed fleet with an explicit node id.
235    ///
236    /// Per-node rings require one active process membership per node id.
237    /// Orbit validates the id range but does not own process lifecycle and
238    /// therefore cannot prevent duplicate live memberships.
239    #[cfg(unix)]
240    pub fn join_shm_as(
241        name: &str,
242        fleet_capacity: u16,
243        node_id: NodeId
244    ) -> Result<Self> {
245        if fleet_capacity == 0 {
246            return Err(Error::EmptyFleet);
247        }
248        if node_id.get() >= fleet_capacity {
249            return Err(Error::NodeOutsideFleet { node_id: node_id.get(), fleet_capacity });
250        }
251        let name: Arc<str> = Arc::from(name);
252        let membership = crate::shm::join_fleet_membership(&name).map_err(Error::Io)?;
253        Ok(Self {
254            inner: Arc::new(FleetInner {
255                name: Arc::clone(&name),
256                fleet_capacity,
257                node_id,
258                id_counters: DashMap::new(),
259                backing: RingBacking::Shm(ShmRingRegistry::new(name.as_ref(), fleet_capacity)),
260                membership: Some(membership)
261            })
262        })
263    }
264
265    pub fn name(&self) -> &str {
266        &self.inner.name
267    }
268
269    /// Number of physical node lanes reserved for this fleet.
270    pub fn fleet_capacity(&self) -> u16 {
271        self.inner.fleet_capacity
272    }
273
274    pub fn node_id(&self) -> NodeId {
275        self.inner.node_id
276    }
277
278    /// Mint a fresh fleet-unique [`NetId64`] for type `T` *without*
279    /// publishing anything. Use this when the caller only needs the
280    /// id (e.g. minting an id to attach to data being persisted to
281    /// DB before going through the ring).
282    ///
283    /// For most use cases prefer [`Fleet::publish`] — it mints AND
284    /// stores in a single atomic step.
285    pub fn next_id<T: OrbitTyped>(&self) -> NetId64 {
286        let counter_arc = self
287            .inner
288            .id_counters
289            .entry(T::KIND)
290            .or_insert_with(|| Arc::new(AtomicU64::new(0)))
291            .clone();
292        let counter = counter_arc.fetch_add(1, Ordering::Relaxed);
293        NetId64::make(T::KIND, self.node_id().get(), counter)
294    }
295
296    /// True when this fleet's ring storage is backed by POSIX SHM
297    /// (visible across processes). False for in-memory fleets.
298    pub fn is_shm(&self) -> bool {
299        #[cfg(unix)]
300        {
301            matches!(self.inner.backing, RingBacking::Shm(_))
302        }
303        #[cfg(not(unix))]
304        {
305            false
306        }
307    }
308
309    /// Get-or-create the in-memory ring for type `T`. Only valid
310    /// for fleets created via [`Fleet::join`]; SHM-backed fleets
311    /// should use [`Fleet::shm_ring`] instead.
312    ///
313    /// # Panics
314    ///
315    /// Panics if called on a SHM-backed fleet.
316    pub fn ring<T: OrbitTyped>(&self) -> Arc<Ring> {
317        match &self.inner.backing {
318            RingBacking::InMemory(r) => r.get_or_create::<T>(),
319            #[cfg(unix)]
320            RingBacking::Shm(_) => {
321                panic!("Fleet::ring called on SHM-backed fleet — use Fleet::shm_ring instead");
322            }
323        }
324    }
325
326    /// Get-or-create the SHM ring for type `T`. Only valid on fleets
327    /// created via [`Fleet::join_shm`].
328    ///
329    /// # Errors
330    ///
331    /// Returns an `io::Error` if the SHM segment cannot be opened
332    /// (permissions, name too long, etc.).
333    ///
334    /// # Panics
335    ///
336    /// Panics if called on an in-memory fleet.
337    #[cfg(unix)]
338    pub fn shm_ring<T: OrbitTyped>(&self) -> std::io::Result<Arc<ShmRing>> {
339        match &self.inner.backing {
340            RingBacking::Shm(r) => r.get_or_create_for::<T>(),
341            RingBacking::InMemory(_) => {
342                panic!("Fleet::shm_ring called on in-memory fleet — use Fleet::ring instead");
343            }
344        }
345    }
346
347    /// Publish a payload to the ring for type `T`. Mints a [`NetId64`],
348    /// writes the [`Frame`] to the appropriate shared or node-owned lane,
349    /// and returns the id.
350    ///
351    /// # Panics
352    ///
353    /// Panics if the ring cannot be opened or the payload exceeds
354    /// `T::RING_SPEC.payload_capacity`. Ring failures are
355    /// operator-visible, not silently ignored.
356    pub fn publish<T: OrbitTyped>(
357        &self,
358        frame_kind: u8,
359        ver: u64,
360        payload: Bytes
361    ) -> NetId64 {
362        match &self.inner.backing {
363            RingBacking::InMemory(r) => {
364                let ring = r.get_or_create::<T>();
365                ring.write(self.node_id(), frame_kind, ver, payload)
366            }
367            #[cfg(unix)]
368            RingBacking::Shm(r) => {
369                let ring =
370                    r.get_or_create_for::<T>().expect("SHM ring open failed — fleet unusable");
371                ring.write(self.node_id(), frame_kind, ver, payload).expect("SHM ring write failed")
372            }
373        }
374    }
375
376    /// Publish a contiguous batch into one ring lane.
377    ///
378    /// The returned ids are ordered and consecutive. For per-node rings the
379    /// lane head becomes visible only after every frame in the batch has been
380    /// committed. Semantic layers can use this to publish a multi-slot blob,
381    /// then publish a separate descriptor that references the first id and
382    /// frame count.
383    ///
384    /// # Panics
385    ///
386    /// Panics if the ring cannot be opened, a payload exceeds the declared
387    /// slot capacity, or the batch itself is larger than the ring.
388    pub fn publish_batch<T: OrbitTyped>(
389        &self,
390        frame_kind: u8,
391        ver: u64,
392        payloads: Vec<Bytes>
393    ) -> Vec<NetId64> {
394        match &self.inner.backing {
395            RingBacking::InMemory(r) => {
396                let ring = r.get_or_create::<T>();
397                ring.write_batch(self.node_id(), frame_kind, ver, payloads)
398            }
399            #[cfg(unix)]
400            RingBacking::Shm(r) => {
401                let ring =
402                    r.get_or_create_for::<T>().expect("SHM ring open failed — fleet unusable");
403                ring.write_batch(self.node_id(), frame_kind, ver, payloads)
404                    .expect("SHM ring batch write failed")
405            }
406        }
407    }
408
409    /// Look up a previously-published frame by its id. Returns the
410    /// frame if its slot still holds the same id (i.e. the ring has
411    /// not wrapped past it).
412    pub fn read(
413        &self,
414        id: NetId64
415    ) -> Option<Frame> {
416        match &self.inner.backing {
417            RingBacking::InMemory(r) => r.lookup(id.kind())?.read(id),
418            #[cfg(unix)]
419            RingBacking::Shm(r) => r.lookup(id.kind())?.read(id)
420        }
421    }
422
423    /// Read the most recent frame for type `T`. Per-node rings read this
424    /// fleet handle's local lane; shared rings read their sole lane.
425    pub fn read_head<T: OrbitTyped>(&self) -> Option<Frame> {
426        if T::RING_SPEC.topology == RingTopology::PerNode {
427            let head = self.lane_head::<T>(self.node_id());
428            return (head > 0).then(|| self.read_lane_at::<T>(self.node_id(), head - 1)).flatten();
429        }
430        match &self.inner.backing {
431            RingBacking::InMemory(r) => {
432                let ring = r.get_or_create::<T>();
433                ring.read_head()
434            }
435            #[cfg(unix)]
436            RingBacking::Shm(r) => {
437                let ring = r.get_or_create_for::<T>().ok()?;
438                ring.read_head()
439            }
440        }
441    }
442
443    /// Current head for type `T`'s ring. Per-node rings report this fleet
444    /// handle's local committed head; shared rings report their sole lane's
445    /// visible head.
446    /// Lazily
447    /// creates / attaches the ring on first access — important for
448    /// cross-process readers, where a child process may need to
449    /// attach to a SHM segment a peer already populated. Returns 0
450    /// when the ring is fresh / no counters have been claimed.
451    pub fn head<T: OrbitTyped>(&self) -> u64 {
452        if T::RING_SPEC.topology == RingTopology::PerNode {
453            return self.lane_head::<T>(self.node_id());
454        }
455        match &self.inner.backing {
456            RingBacking::InMemory(r) => r.get_or_create::<T>().head(),
457            #[cfg(unix)]
458            RingBacking::Shm(r) => r.get_or_create_for::<T>().map(|ring| ring.head()).unwrap_or(0)
459        }
460    }
461
462    /// Read whatever frame currently occupies `counter % capacity`.
463    /// Per-node rings read this fleet handle's local lane. Lazily attaches
464    /// the ring on first access (same rationale as [`Fleet::head`]).
465    /// Returns `None` if the slot is empty/torn or attach fails.
466    ///
467    /// Used by walking readers; for typed handle-based reads,
468    /// prefer [`Fleet::read`].
469    pub fn read_at<T: OrbitTyped>(
470        &self,
471        counter: u64
472    ) -> Option<Frame> {
473        if T::RING_SPEC.topology == RingTopology::PerNode {
474            return self.read_lane_at::<T>(self.node_id(), counter);
475        }
476        match &self.inner.backing {
477            RingBacking::InMemory(r) => r.get_or_create::<T>().read_at(counter),
478            #[cfg(unix)]
479            RingBacking::Shm(r) => r.get_or_create_for::<T>().ok()?.read_at(counter)
480        }
481    }
482
483    pub(crate) fn read_state_at<T: OrbitTyped>(
484        &self,
485        counter: u64
486    ) -> crate::ring::cursor::RingRead {
487        if T::RING_SPEC.topology == RingTopology::PerNode {
488            return self.read_lane_state_at::<T>(self.node_id(), counter);
489        }
490        match &self.inner.backing {
491            RingBacking::InMemory(r) => r.get_or_create::<T>().read_state_at(counter),
492            #[cfg(unix)]
493            RingBacking::Shm(r) => r
494                .get_or_create_for::<T>()
495                .map(|ring| ring.read_state_at(counter))
496                .unwrap_or(crate::ring::cursor::RingRead::Unavailable)
497        }
498    }
499
500    /// Current head for one physical node lane.
501    ///
502    /// On a shared ring every node id addresses the sole shared lane.
503    pub fn lane_head<T: OrbitTyped>(
504        &self,
505        node_id: NodeId
506    ) -> u64 {
507        match &self.inner.backing {
508            RingBacking::InMemory(r) => r.get_or_create::<T>().lane_head(node_id),
509            #[cfg(unix)]
510            RingBacking::Shm(r) => {
511                r.get_or_create_for::<T>().map(|ring| ring.lane_head(node_id)).unwrap_or(0)
512            }
513        }
514    }
515
516    /// Read the frame currently occupying one node lane's counter slot.
517    pub fn read_lane_at<T: OrbitTyped>(
518        &self,
519        node_id: NodeId,
520        counter: u64
521    ) -> Option<Frame> {
522        match &self.inner.backing {
523            RingBacking::InMemory(r) => r.get_or_create::<T>().read_lane_at(node_id, counter),
524            #[cfg(unix)]
525            RingBacking::Shm(r) => r.get_or_create_for::<T>().ok()?.read_lane_at(node_id, counter)
526        }
527    }
528
529    pub(crate) fn read_lane_state_at<T: OrbitTyped>(
530        &self,
531        node_id: NodeId,
532        counter: u64
533    ) -> crate::ring::cursor::RingRead {
534        match &self.inner.backing {
535            RingBacking::InMemory(r) => r.get_or_create::<T>().read_lane_state_at(node_id, counter),
536            #[cfg(unix)]
537            RingBacking::Shm(r) => r
538                .get_or_create_for::<T>()
539                .map(|ring| ring.read_lane_state_at(node_id, counter))
540                .unwrap_or(crate::ring::cursor::RingRead::Unavailable)
541        }
542    }
543
544    /// Capacity of the ring for type `T`. Lazily attaches the ring
545    /// on first access. Falls back to `T::RING_SPEC.capacity` when
546    /// SHM attach fails.
547    pub fn ring_capacity<T: OrbitTyped>(&self) -> usize {
548        match &self.inner.backing {
549            RingBacking::InMemory(r) => r.get_or_create::<T>().capacity(),
550            #[cfg(unix)]
551            RingBacking::Shm(r) => r
552                .get_or_create_for::<T>()
553                .map(|ring| ring.capacity())
554                .unwrap_or(T::RING_SPEC.capacity)
555        }
556    }
557
558    /// Allocate one semantic version shared by every writer lane of `T`.
559    ///
560    /// This is separate from each lane's physical frame counter. It is useful
561    /// for semantic layers that retain per-node write scalability but require
562    /// a deterministic fleet-wide last-write-wins order.
563    pub fn next_ring_version<T: OrbitTyped>(&self) -> u64 {
564        match &self.inner.backing {
565            RingBacking::InMemory(r) => r.get_or_create::<T>().next_version(),
566            #[cfg(unix)]
567            RingBacking::Shm(r) => r
568                .get_or_create_for::<T>()
569                .expect("SHM ring open failed — fleet unusable")
570                .next_version()
571        }
572    }
573
574    /// Return the last semantic version allocated for `T` without advancing it.
575    pub fn current_ring_version<T: OrbitTyped>(&self) -> u64 {
576        match &self.inner.backing {
577            RingBacking::InMemory(r) => r.get_or_create::<T>().current_version(),
578            #[cfg(unix)]
579            RingBacking::Shm(r) => r
580                .get_or_create_for::<T>()
581                .expect("SHM ring open failed — fleet unusable")
582                .current_version()
583        }
584    }
585
586    /// Clear every lane for `T` and reset all heads to zero.
587    ///
588    /// This is an owner-side boot cleanup primitive. It is safe for
589    /// runtime state such as events and periodic metrics when the
590    /// embedding application calls it before peer processes begin
591    /// publishing. It is not a coordination protocol; callers must not
592    /// reset a ring while other fleet members are actively writing it.
593    pub fn reset_ring<T: OrbitTyped>(&self) -> std::io::Result<()> {
594        match &self.inner.backing {
595            RingBacking::InMemory(r) => {
596                r.get_or_create::<T>().reset();
597                Ok(())
598            }
599            #[cfg(unix)]
600            RingBacking::Shm(r) => {
601                r.get_or_create_for::<T>()?.reset();
602                Ok(())
603            }
604        }
605    }
606
607    #[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
608    /// Create a process-local readiness fd for one notified SHM ring.
609    ///
610    /// The fd only signals that the ring generation changed. After draining
611    /// it, callers must poll the ring with their own cursor. Multiple writes
612    /// may coalesce into one readiness notification.
613    pub fn ring_event_fd<T: OrbitTyped>(&self) -> std::io::Result<RingEventFd> {
614        match &self.inner.backing {
615            RingBacking::Shm(rings) => RingEventFd::new(rings.get_or_create_for::<T>()?),
616            RingBacking::InMemory(_) => Err(std::io::Error::new(
617                std::io::ErrorKind::Unsupported,
618                "Orbit eventfd requires a shared-memory fleet"
619            ))
620        }
621    }
622
623    #[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
624    /// Publish one frame and notify native waiters after it commits.
625    pub fn publish_notified<T: OrbitTyped>(
626        &self,
627        frame_kind: u8,
628        ver: u64,
629        payload: Bytes
630    ) -> std::io::Result<NetId64> {
631        match &self.inner.backing {
632            RingBacking::Shm(rings) => {
633                let ring = rings.get_or_create_for::<T>()?;
634                let id = ring.write(self.node_id(), frame_kind, ver, payload)?;
635                RingEventFd::notify(&ring)?;
636                Ok(id)
637            }
638            // Process-local fleets do not need a kernel wake bridge.
639            RingBacking::InMemory(rings) => {
640                Ok(rings.get_or_create::<T>().write(self.node_id(), frame_kind, ver, payload))
641            }
642        }
643    }
644
645    #[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
646    /// Readiness for a ring published with
647    /// [`publish_notified_parked`](Self::publish_notified_parked). See
648    /// [`ParkedRingEventFd`]: every reader of such a ring uses this, not
649    /// [`ring_event_fd`](Self::ring_event_fd).
650    pub fn ring_event_fd_parked<T: OrbitTyped>(&self) -> std::io::Result<ParkedRingEventFd> {
651        match &self.inner.backing {
652            RingBacking::Shm(rings) => ParkedRingEventFd::new(rings.get_or_create_for::<T>()?),
653            RingBacking::InMemory(_) => Err(std::io::Error::new(
654                std::io::ErrorKind::Unsupported,
655                "Orbit parked readiness requires a shared-memory fleet"
656            ))
657        }
658    }
659
660    #[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
661    /// Publish one frame and wake the ring's readers only if one is parked.
662    ///
663    /// For rings read through [`ring_event_fd_parked`](Self::ring_event_fd_parked):
664    /// while every reader is busy draining, a publish is a write and an
665    /// atomic, with no syscall. Rings published with
666    /// [`publish_notified`](Self::publish_notified) are not affected.
667    pub fn publish_notified_parked<T: OrbitTyped>(
668        &self,
669        frame_kind: u8,
670        ver: u64,
671        payload: Bytes
672    ) -> std::io::Result<NetId64> {
673        match &self.inner.backing {
674            RingBacking::Shm(rings) => {
675                let ring = rings.get_or_create_for::<T>()?;
676                let id = ring.write(self.node_id(), frame_kind, ver, payload)?;
677                ParkedRingEventFd::notify(&ring)?;
678                Ok(id)
679            }
680            RingBacking::InMemory(rings) => {
681                Ok(rings.get_or_create::<T>().write(self.node_id(), frame_kind, ver, payload))
682            }
683        }
684    }
685
686    #[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
687    /// Publish one contiguous batch and wake the ring's readers only if one
688    /// is parked — [`publish_notified_parked`](Self::publish_notified_parked)
689    /// for a batch. A per-node lane commits the whole batch before a reader
690    /// can see any of it.
691    pub fn publish_batch_notified_parked<T: OrbitTyped>(
692        &self,
693        frame_kind: u8,
694        ver: u64,
695        payloads: Vec<Bytes>
696    ) -> std::io::Result<Vec<NetId64>> {
697        match &self.inner.backing {
698            RingBacking::Shm(rings) => {
699                let ring = rings.get_or_create_for::<T>()?;
700                let ids = ring.write_batch(self.node_id(), frame_kind, ver, payloads)?;
701                if !ids.is_empty() {
702                    ParkedRingEventFd::notify(&ring)?;
703                }
704                Ok(ids)
705            }
706            RingBacking::InMemory(rings) => Ok(rings.get_or_create::<T>().write_batch(
707                self.node_id(),
708                frame_kind,
709                ver,
710                payloads
711            ))
712        }
713    }
714
715    #[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
716    /// Publish one contiguous batch and notify native waiters once after the
717    /// complete batch commits.
718    pub fn publish_batch_notified<T: OrbitTyped>(
719        &self,
720        frame_kind: u8,
721        ver: u64,
722        payloads: Vec<Bytes>
723    ) -> std::io::Result<Vec<NetId64>> {
724        match &self.inner.backing {
725            RingBacking::Shm(rings) => {
726                let ring = rings.get_or_create_for::<T>()?;
727                let ids = ring.write_batch(self.node_id(), frame_kind, ver, payloads)?;
728                if !ids.is_empty() {
729                    RingEventFd::notify(&ring)?;
730                }
731                Ok(ids)
732            }
733            RingBacking::InMemory(rings) => Ok(rings.get_or_create::<T>().write_batch(
734                self.node_id(),
735                frame_kind,
736                ver,
737                payloads
738            ))
739        }
740    }
741}
742
743impl std::fmt::Debug for Fleet {
744    fn fmt(
745        &self,
746        f: &mut std::fmt::Formatter<'_>
747    ) -> std::fmt::Result {
748        f.debug_struct("Fleet")
749            .field("name", &self.inner.name)
750            .field("fleet_capacity", &self.inner.fleet_capacity)
751            .field("node_id", &self.inner.node_id)
752            .field("id_counters", &self.inner.id_counters.len())
753            .finish()
754    }
755}