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 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 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 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#[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#[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 id_counters: DashMap<u8, Arc<AtomicU64>>,
172 backing: RingBacking,
176 #[cfg(unix)]
180 #[allow(dead_code)]
181 membership: Option<crate::shm::FleetMembership>
182}
183
184enum RingBacking {
188 InMemory(RingRegistry),
192 #[cfg(unix)]
196 Shm(ShmRingRegistry)
197}
198
199impl Fleet {
200 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 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 #[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 #[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 #[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 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 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 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 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 #[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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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}