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}
149
150enum RingBacking {
154 InMemory(RingRegistry),
158 #[cfg(unix)]
162 Shm(ShmRingRegistry),
163}
164
165impl Fleet {
166 pub fn join(name: &'static str, fleet_capacity: u16) -> Result<Self> {
169 Self::join_as(name, fleet_capacity, NodeId::ZERO)
170 }
171
172 pub fn join_as(name: &'static str, fleet_capacity: u16, node_id: NodeId) -> Result<Self> {
174 if fleet_capacity == 0 {
175 return Err(Error::EmptyFleet);
176 }
177 if node_id.get() >= fleet_capacity {
178 return Err(Error::NodeOutsideFleet {
179 node_id: node_id.get(),
180 fleet_capacity,
181 });
182 }
183 Ok(Self {
184 inner: Arc::new(FleetInner {
185 name,
186 fleet_capacity,
187 node_id,
188 id_counters: DashMap::new(),
189 backing: RingBacking::InMemory(RingRegistry::new(fleet_capacity)),
190 }),
191 })
192 }
193
194 #[cfg(unix)]
206 pub fn join_shm(name: &'static str, fleet_capacity: u16) -> Result<Self> {
207 Self::join_shm_as(name, fleet_capacity, NodeId::ZERO)
208 }
209
210 #[cfg(unix)]
216 pub fn join_shm_as(name: &'static str, fleet_capacity: u16, node_id: NodeId) -> Result<Self> {
217 if fleet_capacity == 0 {
218 return Err(Error::EmptyFleet);
219 }
220 if node_id.get() >= fleet_capacity {
221 return Err(Error::NodeOutsideFleet {
222 node_id: node_id.get(),
223 fleet_capacity,
224 });
225 }
226 Ok(Self {
227 inner: Arc::new(FleetInner {
228 name,
229 fleet_capacity,
230 node_id,
231 id_counters: DashMap::new(),
232 backing: RingBacking::Shm(ShmRingRegistry::new(name, fleet_capacity)),
233 }),
234 })
235 }
236
237 pub fn name(&self) -> &'static str {
238 self.inner.name
239 }
240
241 pub fn fleet_capacity(&self) -> u16 {
243 self.inner.fleet_capacity
244 }
245
246 pub fn node_id(&self) -> NodeId {
247 self.inner.node_id
248 }
249
250 pub fn next_id<T: OrbitTyped>(&self) -> NetId64 {
258 let counter_arc = self
259 .inner
260 .id_counters
261 .entry(T::KIND)
262 .or_insert_with(|| Arc::new(AtomicU64::new(0)))
263 .clone();
264 let counter = counter_arc.fetch_add(1, Ordering::Relaxed);
265 NetId64::make(T::KIND, self.node_id().get(), counter)
266 }
267
268 pub fn is_shm(&self) -> bool {
271 #[cfg(unix)]
272 {
273 matches!(self.inner.backing, RingBacking::Shm(_))
274 }
275 #[cfg(not(unix))]
276 {
277 false
278 }
279 }
280
281 pub fn ring<T: OrbitTyped>(&self) -> Arc<Ring> {
289 match &self.inner.backing {
290 RingBacking::InMemory(r) => r.get_or_create::<T>(),
291 #[cfg(unix)]
292 RingBacking::Shm(_) => {
293 panic!("Fleet::ring called on SHM-backed fleet — use Fleet::shm_ring instead");
294 }
295 }
296 }
297
298 #[cfg(unix)]
310 pub fn shm_ring<T: OrbitTyped>(&self) -> std::io::Result<Arc<ShmRing>> {
311 match &self.inner.backing {
312 RingBacking::Shm(r) => r.get_or_create_for::<T>(),
313 RingBacking::InMemory(_) => {
314 panic!("Fleet::shm_ring called on in-memory fleet — use Fleet::ring instead");
315 }
316 }
317 }
318
319 pub fn publish<T: OrbitTyped>(&self, frame_kind: u8, ver: u64, payload: Bytes) -> NetId64 {
329 match &self.inner.backing {
330 RingBacking::InMemory(r) => {
331 let ring = r.get_or_create::<T>();
332 ring.write(self.node_id(), frame_kind, ver, payload)
333 }
334 #[cfg(unix)]
335 RingBacking::Shm(r) => {
336 let ring = r
337 .get_or_create_for::<T>()
338 .expect("SHM ring open failed — fleet unusable");
339 ring.write(self.node_id(), frame_kind, ver, payload)
340 .expect("SHM ring write failed")
341 }
342 }
343 }
344
345 pub fn publish_batch<T: OrbitTyped>(
358 &self,
359 frame_kind: u8,
360 ver: u64,
361 payloads: Vec<Bytes>,
362 ) -> Vec<NetId64> {
363 match &self.inner.backing {
364 RingBacking::InMemory(r) => {
365 let ring = r.get_or_create::<T>();
366 ring.write_batch(self.node_id(), frame_kind, ver, payloads)
367 }
368 #[cfg(unix)]
369 RingBacking::Shm(r) => {
370 let ring = r
371 .get_or_create_for::<T>()
372 .expect("SHM ring open failed — fleet unusable");
373 ring.write_batch(self.node_id(), frame_kind, ver, payloads)
374 .expect("SHM ring batch write failed")
375 }
376 }
377 }
378
379 pub fn read(&self, id: NetId64) -> Option<Frame> {
383 match &self.inner.backing {
384 RingBacking::InMemory(r) => r.lookup(id.kind())?.read(id),
385 #[cfg(unix)]
386 RingBacking::Shm(r) => r.lookup(id.kind())?.read(id),
387 }
388 }
389
390 pub fn read_head<T: OrbitTyped>(&self) -> Option<Frame> {
393 if T::RING_SPEC.topology == RingTopology::PerNode {
394 let head = self.lane_head::<T>(self.node_id());
395 return (head > 0)
396 .then(|| self.read_lane_at::<T>(self.node_id(), head - 1))
397 .flatten();
398 }
399 match &self.inner.backing {
400 RingBacking::InMemory(r) => {
401 let ring = r.get_or_create::<T>();
402 ring.read_head()
403 }
404 #[cfg(unix)]
405 RingBacking::Shm(r) => {
406 let ring = r.get_or_create_for::<T>().ok()?;
407 ring.read_head()
408 }
409 }
410 }
411
412 pub fn head<T: OrbitTyped>(&self) -> u64 {
421 if T::RING_SPEC.topology == RingTopology::PerNode {
422 return self.lane_head::<T>(self.node_id());
423 }
424 match &self.inner.backing {
425 RingBacking::InMemory(r) => r.get_or_create::<T>().head(),
426 #[cfg(unix)]
427 RingBacking::Shm(r) => r
428 .get_or_create_for::<T>()
429 .map(|ring| ring.head())
430 .unwrap_or(0),
431 }
432 }
433
434 pub fn read_at<T: OrbitTyped>(&self, counter: u64) -> Option<Frame> {
442 if T::RING_SPEC.topology == RingTopology::PerNode {
443 return self.read_lane_at::<T>(self.node_id(), counter);
444 }
445 match &self.inner.backing {
446 RingBacking::InMemory(r) => r.get_or_create::<T>().read_at(counter),
447 #[cfg(unix)]
448 RingBacking::Shm(r) => r.get_or_create_for::<T>().ok()?.read_at(counter),
449 }
450 }
451
452 pub(crate) fn read_state_at<T: OrbitTyped>(
453 &self,
454 counter: u64,
455 ) -> crate::ring::cursor::RingRead {
456 if T::RING_SPEC.topology == RingTopology::PerNode {
457 return self.read_lane_state_at::<T>(self.node_id(), counter);
458 }
459 match &self.inner.backing {
460 RingBacking::InMemory(r) => r.get_or_create::<T>().read_state_at(counter),
461 #[cfg(unix)]
462 RingBacking::Shm(r) => r
463 .get_or_create_for::<T>()
464 .map(|ring| ring.read_state_at(counter))
465 .unwrap_or(crate::ring::cursor::RingRead::Unavailable),
466 }
467 }
468
469 pub fn lane_head<T: OrbitTyped>(&self, node_id: NodeId) -> u64 {
473 match &self.inner.backing {
474 RingBacking::InMemory(r) => r.get_or_create::<T>().lane_head(node_id),
475 #[cfg(unix)]
476 RingBacking::Shm(r) => r
477 .get_or_create_for::<T>()
478 .map(|ring| ring.lane_head(node_id))
479 .unwrap_or(0),
480 }
481 }
482
483 pub fn read_lane_at<T: OrbitTyped>(&self, node_id: NodeId, counter: u64) -> Option<Frame> {
485 match &self.inner.backing {
486 RingBacking::InMemory(r) => r.get_or_create::<T>().read_lane_at(node_id, counter),
487 #[cfg(unix)]
488 RingBacking::Shm(r) => r
489 .get_or_create_for::<T>()
490 .ok()?
491 .read_lane_at(node_id, counter),
492 }
493 }
494
495 pub(crate) fn read_lane_state_at<T: OrbitTyped>(
496 &self,
497 node_id: NodeId,
498 counter: u64,
499 ) -> crate::ring::cursor::RingRead {
500 match &self.inner.backing {
501 RingBacking::InMemory(r) => r.get_or_create::<T>().read_lane_state_at(node_id, counter),
502 #[cfg(unix)]
503 RingBacking::Shm(r) => r
504 .get_or_create_for::<T>()
505 .map(|ring| ring.read_lane_state_at(node_id, counter))
506 .unwrap_or(crate::ring::cursor::RingRead::Unavailable),
507 }
508 }
509
510 pub fn ring_capacity<T: OrbitTyped>(&self) -> usize {
514 match &self.inner.backing {
515 RingBacking::InMemory(r) => r.get_or_create::<T>().capacity(),
516 #[cfg(unix)]
517 RingBacking::Shm(r) => r
518 .get_or_create_for::<T>()
519 .map(|ring| ring.capacity())
520 .unwrap_or(T::RING_SPEC.capacity),
521 }
522 }
523
524 pub fn next_ring_version<T: OrbitTyped>(&self) -> u64 {
530 match &self.inner.backing {
531 RingBacking::InMemory(r) => r.get_or_create::<T>().next_version(),
532 #[cfg(unix)]
533 RingBacking::Shm(r) => r
534 .get_or_create_for::<T>()
535 .expect("SHM ring open failed — fleet unusable")
536 .next_version(),
537 }
538 }
539
540 pub fn current_ring_version<T: OrbitTyped>(&self) -> u64 {
542 match &self.inner.backing {
543 RingBacking::InMemory(r) => r.get_or_create::<T>().current_version(),
544 #[cfg(unix)]
545 RingBacking::Shm(r) => r
546 .get_or_create_for::<T>()
547 .expect("SHM ring open failed — fleet unusable")
548 .current_version(),
549 }
550 }
551
552 pub fn reset_ring<T: OrbitTyped>(&self) -> std::io::Result<()> {
560 match &self.inner.backing {
561 RingBacking::InMemory(r) => {
562 r.get_or_create::<T>().reset();
563 Ok(())
564 }
565 #[cfg(unix)]
566 RingBacking::Shm(r) => {
567 r.get_or_create_for::<T>()?.reset();
568 Ok(())
569 }
570 }
571 }
572
573 #[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
574 pub fn ring_event_fd<T: OrbitTyped>(&self) -> std::io::Result<RingEventFd> {
580 match &self.inner.backing {
581 RingBacking::Shm(rings) => RingEventFd::new(rings.get_or_create_for::<T>()?),
582 RingBacking::InMemory(_) => Err(std::io::Error::new(
583 std::io::ErrorKind::Unsupported,
584 "Orbit eventfd requires a shared-memory fleet",
585 )),
586 }
587 }
588
589 #[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
590 pub fn publish_notified<T: OrbitTyped>(
592 &self,
593 frame_kind: u8,
594 ver: u64,
595 payload: Bytes,
596 ) -> std::io::Result<NetId64> {
597 match &self.inner.backing {
598 RingBacking::Shm(rings) => {
599 let ring = rings.get_or_create_for::<T>()?;
600 let id = ring.write(self.node_id(), frame_kind, ver, payload)?;
601 RingEventFd::notify(&ring)?;
602 Ok(id)
603 }
604 RingBacking::InMemory(rings) => {
606 Ok(rings
607 .get_or_create::<T>()
608 .write(self.node_id(), frame_kind, ver, payload))
609 }
610 }
611 }
612
613 #[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
614 pub fn publish_batch_notified<T: OrbitTyped>(
617 &self,
618 frame_kind: u8,
619 ver: u64,
620 payloads: Vec<Bytes>,
621 ) -> std::io::Result<Vec<NetId64>> {
622 match &self.inner.backing {
623 RingBacking::Shm(rings) => {
624 let ring = rings.get_or_create_for::<T>()?;
625 let ids = ring.write_batch(self.node_id(), frame_kind, ver, payloads)?;
626 if !ids.is_empty() {
627 RingEventFd::notify(&ring)?;
628 }
629 Ok(ids)
630 }
631 RingBacking::InMemory(rings) => Ok(rings.get_or_create::<T>().write_batch(
632 self.node_id(),
633 frame_kind,
634 ver,
635 payloads,
636 )),
637 }
638 }
639}
640
641impl std::fmt::Debug for Fleet {
642 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
643 f.debug_struct("Fleet")
644 .field("name", &self.inner.name)
645 .field("fleet_capacity", &self.inner.fleet_capacity)
646 .field("node_id", &self.inner.node_id)
647 .field("id_counters", &self.inner.id_counters.len())
648 .finish()
649 }
650}