1use 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#[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 pub fn attach_existing(name: impl Into<String>) -> std::io::Result<Self> {
43 let uid = unsafe { libc::geteuid() };
45 Self::attach_existing_for_uid(name, uid)
46 }
47
48 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 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 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#[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#[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 id_counters: DashMap<u8, Arc<AtomicU64>>,
153 backing: RingBacking,
157 #[cfg(unix)]
161 #[allow(dead_code)]
162 membership: Option<crate::shm::FleetMembership>
163}
164
165enum RingBacking {
169 InMemory(RingRegistry),
173 #[cfg(unix)]
177 Shm(ShmRingRegistry)
178}
179
180impl Fleet {
181 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 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 #[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 #[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 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 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 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 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 #[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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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}