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: Arc<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: &str, fleet_capacity: u16) -> Result<Self> {
175 Self::join_as(name, fleet_capacity, NodeId::ZERO)
176 }
177
178 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 #[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 #[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 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 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 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 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 #[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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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}