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