pub struct ShmRing { /* private fields */ }Expand description
Cross-process ring buffer backed by a POSIX SHM segment.
Use ShmRing::open_or_create to attach to (or create) the
segment by name. Multiple processes calling this with the same name and
matching RingSpec share the underlying memory; the first process to
call it does the one-time header initialization.
Implementations§
Source§impl ShmRing
impl ShmRing
Sourcepub fn open_or_create(
fleet_name: &str,
kind: u8,
spec: RingSpec,
) -> Result<Self>
pub fn open_or_create( fleet_name: &str, kind: u8, spec: RingSpec, ) -> Result<Self>
Open or create a SHM-backed ring under fleet_name for type
kind kind with spec. The first process to call
this initializes the header; later attachers reuse it.
Sourcepub fn open_or_create_for_fleet(
fleet_name: &str,
kind: u8,
spec: RingSpec,
fleet_capacity: u16,
) -> Result<Self>
pub fn open_or_create_for_fleet( fleet_name: &str, kind: u8, spec: RingSpec, fleet_capacity: u16, ) -> Result<Self>
Open or create a SHM-backed ring using fleet_capacity physical
writer lanes when spec is RingTopology::PerNode.
Sourcepub fn created(&self) -> bool
pub fn created(&self) -> bool
True when this handle was the one that created the SHM segment (vs attaching to an already-existing one).
Sourcepub fn unlink(&self) -> Result<()>
pub fn unlink(&self) -> Result<()>
Remove the SHM segment name. Existing mappings stay valid;
new opens will fail until open_or_create recreates it.
Call from the owner process at fleet shutdown.
Sourcepub fn lane_count(&self) -> usize
pub fn lane_count(&self) -> usize
Number of physical writer lanes in this segment.
Sourcepub fn payload_capacity(&self) -> usize
pub fn payload_capacity(&self) -> usize
Maximum inline payload bytes for this ring lane.
pub fn spec(&self) -> RingSpec
Sourcepub fn next_version(&self) -> u64
pub fn next_version(&self) -> u64
Allocate one non-zero semantic version shared by all writer lanes and processes attached to this ring.
Sourcepub fn current_version(&self) -> u64
pub fn current_version(&self) -> u64
Last semantic version allocated by any process attached to this ring.
Sourcepub fn write(
&self,
node_id: NodeId,
frame_kind: u8,
ver: u64,
payload: Bytes,
) -> Result<NetId64>
pub fn write( &self, node_id: NodeId, frame_kind: u8, ver: u64, payload: Bytes, ) -> Result<NetId64>
Append a frame. Atomically reserves the next counter, mints
the NetId64, writes the slot. Returns the minted id.
Sourcepub fn write_batch(
&self,
node_id: NodeId,
frame_kind: u8,
ver: u64,
payloads: Vec<Bytes>,
) -> Result<Vec<NetId64>>
pub fn write_batch( &self, node_id: NodeId, frame_kind: u8, ver: u64, payloads: Vec<Bytes>, ) -> Result<Vec<NetId64>>
Append one contiguous batch to a lane and return consecutive ids.
Per-node and shared-ordered lane heads advance only after the complete batch has committed. The batch may not exceed the ring capacity.
Sourcepub fn read(&self, id: NetId64) -> Option<Frame>
pub fn read(&self, id: NetId64) -> Option<Frame>
Read the slot whose counter matches id.counter(). Returns
None if the slot has been overwritten, was never written,
or a torn read could not be reconciled across two retries.
Sourcepub fn read_head(&self) -> Option<Frame>
pub fn read_head(&self) -> Option<Frame>
Read the most recent frame (head - 1). Returns None if no
write has happened yet, or a torn read could not resolve.
Sourcepub fn read_at(&self, counter: u64) -> Option<Frame>
pub fn read_at(&self, counter: u64) -> Option<Frame>
Read whatever frame currently occupies the slot at
counter % capacity, regardless of which counter is stored
there. Used by walking readers that need slot-by-slot access
without knowing the writer’s NetId64 ahead of time.