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(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#[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(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 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 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#[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#[derive(Clone)]
131pub struct Fleet {
132 inner: Arc<FleetInner>,
133}
134
135struct FleetInner {
136 name: &'static str,
137 fleet_capacity: u16,
138 node_id: NodeId,
139 id_counters: DashMap<u8, Arc<AtomicU64>>,
144 backing: RingBacking,
148 #[cfg(unix)]
152 #[allow(dead_code)]
153 membership: Option<crate::shm::FleetMembership>,
154}
155
156enum RingBacking {
160 InMemory(RingRegistry),
164 #[cfg(unix)]
168 Shm(ShmRingRegistry),
169}
170
171impl Fleet {
172 pub fn join(name: &'static str, fleet_capacity: u16) -> Result<Self> {
175 Self::join_as(name, fleet_capacity, NodeId::ZERO)
176 }
177
178 pub fn join_as(name: &'static 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,
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 #[cfg(unix)]
214 pub fn join_shm(name: &'static str, fleet_capacity: u16) -> Result<Self> {
215 Self::join_shm_as(name, fleet_capacity, NodeId::ZERO)
216 }
217
218 #[cfg(unix)]
224 pub fn join_shm_as(name: &'static 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 membership = crate::shm::join_fleet_membership(name).map_err(Error::Io)?;
235 Ok(Self {
236 inner: Arc::new(FleetInner {
237 name,
238 fleet_capacity,
239 node_id,
240 id_counters: DashMap::new(),
241 backing: RingBacking::Shm(ShmRingRegistry::new(name, fleet_capacity)),
242 membership: Some(membership),
243 }),
244 })
245 }
246
247 pub fn name(&self) -> &'static str {
248 self.inner.name
249 }
250
251 pub fn fleet_capacity(&self) -> u16 {
253 self.inner.fleet_capacity
254 }
255
256 pub fn node_id(&self) -> NodeId {
257 self.inner.node_id
258 }
259
260 pub fn next_id<T: OrbitTyped>(&self) -> NetId64 {
268 let counter_arc = self
269 .inner
270 .id_counters
271 .entry(T::KIND)
272 .or_insert_with(|| Arc::new(AtomicU64::new(0)))
273 .clone();
274 let counter = counter_arc.fetch_add(1, Ordering::Relaxed);
275 NetId64::make(T::KIND, self.node_id().get(), counter)
276 }
277
278 pub fn is_shm(&self) -> bool {
281 #[cfg(unix)]
282 {
283 matches!(self.inner.backing, RingBacking::Shm(_))
284 }
285 #[cfg(not(unix))]
286 {
287 false
288 }
289 }
290
291 pub fn ring<T: OrbitTyped>(&self) -> Arc<Ring> {
299 match &self.inner.backing {
300 RingBacking::InMemory(r) => r.get_or_create::<T>(),
301 #[cfg(unix)]
302 RingBacking::Shm(_) => {
303 panic!("Fleet::ring called on SHM-backed fleet — use Fleet::shm_ring instead");
304 }
305 }
306 }
307
308 #[cfg(unix)]
320 pub fn shm_ring<T: OrbitTyped>(&self) -> std::io::Result<Arc<ShmRing>> {
321 match &self.inner.backing {
322 RingBacking::Shm(r) => r.get_or_create_for::<T>(),
323 RingBacking::InMemory(_) => {
324 panic!("Fleet::shm_ring called on in-memory fleet — use Fleet::ring instead");
325 }
326 }
327 }
328
329 pub fn publish<T: OrbitTyped>(&self, frame_kind: u8, ver: u64, payload: Bytes) -> NetId64 {
339 match &self.inner.backing {
340 RingBacking::InMemory(r) => {
341 let ring = r.get_or_create::<T>();
342 ring.write(self.node_id(), frame_kind, ver, payload)
343 }
344 #[cfg(unix)]
345 RingBacking::Shm(r) => {
346 let ring = r
347 .get_or_create_for::<T>()
348 .expect("SHM ring open failed — fleet unusable");
349 ring.write(self.node_id(), frame_kind, ver, payload)
350 .expect("SHM ring write failed")
351 }
352 }
353 }
354
355 pub fn publish_batch<T: OrbitTyped>(
368 &self,
369 frame_kind: u8,
370 ver: u64,
371 payloads: Vec<Bytes>,
372 ) -> Vec<NetId64> {
373 match &self.inner.backing {
374 RingBacking::InMemory(r) => {
375 let ring = r.get_or_create::<T>();
376 ring.write_batch(self.node_id(), frame_kind, ver, payloads)
377 }
378 #[cfg(unix)]
379 RingBacking::Shm(r) => {
380 let ring = r
381 .get_or_create_for::<T>()
382 .expect("SHM ring open failed — fleet unusable");
383 ring.write_batch(self.node_id(), frame_kind, ver, payloads)
384 .expect("SHM ring batch write failed")
385 }
386 }
387 }
388
389 pub fn read(&self, id: NetId64) -> Option<Frame> {
393 match &self.inner.backing {
394 RingBacking::InMemory(r) => r.lookup(id.kind())?.read(id),
395 #[cfg(unix)]
396 RingBacking::Shm(r) => r.lookup(id.kind())?.read(id),
397 }
398 }
399
400 pub fn read_head<T: OrbitTyped>(&self) -> Option<Frame> {
403 if T::RING_SPEC.topology == RingTopology::PerNode {
404 let head = self.lane_head::<T>(self.node_id());
405 return (head > 0)
406 .then(|| self.read_lane_at::<T>(self.node_id(), head - 1))
407 .flatten();
408 }
409 match &self.inner.backing {
410 RingBacking::InMemory(r) => {
411 let ring = r.get_or_create::<T>();
412 ring.read_head()
413 }
414 #[cfg(unix)]
415 RingBacking::Shm(r) => {
416 let ring = r.get_or_create_for::<T>().ok()?;
417 ring.read_head()
418 }
419 }
420 }
421
422 pub fn head<T: OrbitTyped>(&self) -> u64 {
431 if T::RING_SPEC.topology == RingTopology::PerNode {
432 return self.lane_head::<T>(self.node_id());
433 }
434 match &self.inner.backing {
435 RingBacking::InMemory(r) => r.get_or_create::<T>().head(),
436 #[cfg(unix)]
437 RingBacking::Shm(r) => r
438 .get_or_create_for::<T>()
439 .map(|ring| ring.head())
440 .unwrap_or(0),
441 }
442 }
443
444 pub fn read_at<T: OrbitTyped>(&self, counter: u64) -> Option<Frame> {
452 if T::RING_SPEC.topology == RingTopology::PerNode {
453 return self.read_lane_at::<T>(self.node_id(), counter);
454 }
455 match &self.inner.backing {
456 RingBacking::InMemory(r) => r.get_or_create::<T>().read_at(counter),
457 #[cfg(unix)]
458 RingBacking::Shm(r) => r.get_or_create_for::<T>().ok()?.read_at(counter),
459 }
460 }
461
462 pub(crate) fn read_state_at<T: OrbitTyped>(
463 &self,
464 counter: u64,
465 ) -> crate::ring::cursor::RingRead {
466 if T::RING_SPEC.topology == RingTopology::PerNode {
467 return self.read_lane_state_at::<T>(self.node_id(), counter);
468 }
469 match &self.inner.backing {
470 RingBacking::InMemory(r) => r.get_or_create::<T>().read_state_at(counter),
471 #[cfg(unix)]
472 RingBacking::Shm(r) => r
473 .get_or_create_for::<T>()
474 .map(|ring| ring.read_state_at(counter))
475 .unwrap_or(crate::ring::cursor::RingRead::Unavailable),
476 }
477 }
478
479 pub fn lane_head<T: OrbitTyped>(&self, node_id: NodeId) -> u64 {
483 match &self.inner.backing {
484 RingBacking::InMemory(r) => r.get_or_create::<T>().lane_head(node_id),
485 #[cfg(unix)]
486 RingBacking::Shm(r) => r
487 .get_or_create_for::<T>()
488 .map(|ring| ring.lane_head(node_id))
489 .unwrap_or(0),
490 }
491 }
492
493 pub fn read_lane_at<T: OrbitTyped>(&self, node_id: NodeId, counter: u64) -> Option<Frame> {
495 match &self.inner.backing {
496 RingBacking::InMemory(r) => r.get_or_create::<T>().read_lane_at(node_id, counter),
497 #[cfg(unix)]
498 RingBacking::Shm(r) => r
499 .get_or_create_for::<T>()
500 .ok()?
501 .read_lane_at(node_id, counter),
502 }
503 }
504
505 pub(crate) fn read_lane_state_at<T: OrbitTyped>(
506 &self,
507 node_id: NodeId,
508 counter: u64,
509 ) -> crate::ring::cursor::RingRead {
510 match &self.inner.backing {
511 RingBacking::InMemory(r) => r.get_or_create::<T>().read_lane_state_at(node_id, counter),
512 #[cfg(unix)]
513 RingBacking::Shm(r) => r
514 .get_or_create_for::<T>()
515 .map(|ring| ring.read_lane_state_at(node_id, counter))
516 .unwrap_or(crate::ring::cursor::RingRead::Unavailable),
517 }
518 }
519
520 pub fn ring_capacity<T: OrbitTyped>(&self) -> usize {
524 match &self.inner.backing {
525 RingBacking::InMemory(r) => r.get_or_create::<T>().capacity(),
526 #[cfg(unix)]
527 RingBacking::Shm(r) => r
528 .get_or_create_for::<T>()
529 .map(|ring| ring.capacity())
530 .unwrap_or(T::RING_SPEC.capacity),
531 }
532 }
533
534 pub fn next_ring_version<T: OrbitTyped>(&self) -> u64 {
540 match &self.inner.backing {
541 RingBacking::InMemory(r) => r.get_or_create::<T>().next_version(),
542 #[cfg(unix)]
543 RingBacking::Shm(r) => r
544 .get_or_create_for::<T>()
545 .expect("SHM ring open failed — fleet unusable")
546 .next_version(),
547 }
548 }
549
550 pub fn current_ring_version<T: OrbitTyped>(&self) -> u64 {
552 match &self.inner.backing {
553 RingBacking::InMemory(r) => r.get_or_create::<T>().current_version(),
554 #[cfg(unix)]
555 RingBacking::Shm(r) => r
556 .get_or_create_for::<T>()
557 .expect("SHM ring open failed — fleet unusable")
558 .current_version(),
559 }
560 }
561
562 pub fn reset_ring<T: OrbitTyped>(&self) -> std::io::Result<()> {
570 match &self.inner.backing {
571 RingBacking::InMemory(r) => {
572 r.get_or_create::<T>().reset();
573 Ok(())
574 }
575 #[cfg(unix)]
576 RingBacking::Shm(r) => {
577 r.get_or_create_for::<T>()?.reset();
578 Ok(())
579 }
580 }
581 }
582
583 #[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
584 pub fn ring_event_fd<T: OrbitTyped>(&self) -> std::io::Result<RingEventFd> {
590 match &self.inner.backing {
591 RingBacking::Shm(rings) => RingEventFd::new(rings.get_or_create_for::<T>()?),
592 RingBacking::InMemory(_) => Err(std::io::Error::new(
593 std::io::ErrorKind::Unsupported,
594 "Orbit eventfd requires a shared-memory fleet",
595 )),
596 }
597 }
598
599 #[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
600 pub fn publish_notified<T: OrbitTyped>(
602 &self,
603 frame_kind: u8,
604 ver: u64,
605 payload: Bytes,
606 ) -> std::io::Result<NetId64> {
607 match &self.inner.backing {
608 RingBacking::Shm(rings) => {
609 let ring = rings.get_or_create_for::<T>()?;
610 let id = ring.write(self.node_id(), frame_kind, ver, payload)?;
611 RingEventFd::notify(&ring)?;
612 Ok(id)
613 }
614 RingBacking::InMemory(rings) => {
616 Ok(rings
617 .get_or_create::<T>()
618 .write(self.node_id(), frame_kind, ver, payload))
619 }
620 }
621 }
622
623 #[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
624 pub fn publish_batch_notified<T: OrbitTyped>(
627 &self,
628 frame_kind: u8,
629 ver: u64,
630 payloads: Vec<Bytes>,
631 ) -> std::io::Result<Vec<NetId64>> {
632 match &self.inner.backing {
633 RingBacking::Shm(rings) => {
634 let ring = rings.get_or_create_for::<T>()?;
635 let ids = ring.write_batch(self.node_id(), frame_kind, ver, payloads)?;
636 if !ids.is_empty() {
637 RingEventFd::notify(&ring)?;
638 }
639 Ok(ids)
640 }
641 RingBacking::InMemory(rings) => Ok(rings.get_or_create::<T>().write_batch(
642 self.node_id(),
643 frame_kind,
644 ver,
645 payloads,
646 )),
647 }
648 }
649}
650
651impl std::fmt::Debug for Fleet {
652 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
653 f.debug_struct("Fleet")
654 .field("name", &self.inner.name)
655 .field("fleet_capacity", &self.inner.fleet_capacity)
656 .field("node_id", &self.inner.node_id)
657 .field("id_counters", &self.inner.id_counters.len())
658 .finish()
659 }
660}