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