Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

3 changes: 3 additions & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,4 +1,7 @@
[workspace]
# Every member is edition 2021, which implies this; stating it silences the
# warning cargo prints on a virtual workspace that leaves it unset.
resolver = "2"
members = [
"ipc_lib",
"writer",
Expand Down
2 changes: 1 addition & 1 deletion ipc_lib/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "ipc_lib"
version = "0.3.0"
version = "0.3.1"
edition = "2021"

[dependencies]
Expand Down
16 changes: 11 additions & 5 deletions ipc_lib/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,7 @@ fn fnv1a_32(data: &[u8]) -> u32 {
fn topic_to_shm_name(topic: &str) -> String {
let id = topic.trim_start_matches('/');
let hash = fnv1a_32(id.as_bytes());
let name = format!("{}_{}_{}",IPC_SHMEM_PREFIX, id, hash);
let name = format!("{IPC_SHMEM_PREFIX}_{id}_{hash}");
name.replace('/', "_")
}

Expand Down Expand Up @@ -90,10 +90,10 @@ impl<T: Copy + Default> IPC<T> {
/// a unique shared memory path in `/dev/shm/`.
pub fn new(topic: &str) -> Result<Self, String> {
let shm_name = topic_to_shm_name(topic);
log::info!("Opening IPC topic '{}' -> shm '/{}'", topic, shm_name);
log::info!("Opening IPC topic '{topic}' -> shm '/{shm_name}'");

let cname =
CString::new(format!("/{}", shm_name)).map_err(|_| "Invalid shm name".to_string())?;
CString::new(format!("/{shm_name}")).map_err(|_| "Invalid shm name".to_string())?;

let mut create = false;
let mut fd = unsafe { shm_open(cname.as_ptr(), O_RDWR, 0o666) };
Expand Down Expand Up @@ -239,8 +239,14 @@ mod tests {
assert!(!name.contains('/')); // no slashes in final name

// Same topic with or without leading slash gives same name
assert_eq!(topic_to_shm_name("/my/topic"), topic_to_shm_name("my/topic"));
assert_eq!(topic_to_shm_name("//my/topic"), topic_to_shm_name("my/topic"));
assert_eq!(
topic_to_shm_name("/my/topic"),
topic_to_shm_name("my/topic")
);
assert_eq!(
topic_to_shm_name("//my/topic"),
topic_to_shm_name("my/topic")
);

// Different topics with same underscore pattern get different hashes
assert_ne!(topic_to_shm_name("my/topic"), topic_to_shm_name("my_topic"));
Expand Down
170 changes: 168 additions & 2 deletions ipc_lib/src/seqlock.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,10 +33,49 @@
//! `cargo check` for a target whose ABI does not produce it, and the
//! alignment guard on the storage type (tyndall's `static_assert`) fails the
//! build of any program that opens a record over an over-aligned type.
//!
//! # Memory ordering
//!
//! The barriers match tyndall's, which are the Linux seqlock's. Four of them,
//! and each one closes a specific hole on a weakly ordered machine (armv7 and
//! aarch64 — x86's store order makes two of them free):
//!
//! ```text
//! writer tyndall here
//! seq -> odd WRITE_ONCE fetch_add(AcqRel)
//! ---- store-store ---- smp_wmb() the AcqRel above
//! size, entry plain stores plain stores
//! ---- store-store ---- smp_wmb() the Release below
//! seq -> even WRITE_ONCE fetch_add(Release)
//!
//! reader
//! load seq1 (spin while odd) READ_ONCE load(Acquire)
//! ---- load-load ---- smp_rmb() the Acquire above
//! read entry plain load plain load
//! ---- load-load ---- smp_rmb() fence(Acquire)
//! load seq2, retry if != seq1 READ_ONCE load(Acquire)
//! ```
//!
//! The two that are easy to get wrong are the *first* writer barrier and the
//! *second* reader barrier, because neither is implied by the release/acquire
//! pair that carries the payload:
//!
//! * a plain `Release` on the odd increment orders everything *before* it,
//! which is the previous write, not the payload stores that follow. Without
//! `AcqRel` those stores may become visible first, and a reader that samples
//! `seq` either side of them sees the same even number over a half-written
//! entry and accepts it;
//! * an `Acquire` on the second `seq` load orders what comes *after* it, not
//! the entry read before it. Without [`fence`] the entry read may drift past
//! the validation and pick up bytes from the next write.
//!
//! Both are silent when they go wrong: a torn record, not a crash. Nothing
//! here is UB — the storage is `Copy` plain data — so the failure is wrong
//! numbers reaching a consumer.

use std::cell::UnsafeCell;
use std::mem::{align_of, offset_of, size_of};
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::atomic::{fence, AtomicU32, Ordering};

use log::debug;

Expand Down Expand Up @@ -126,18 +165,27 @@ impl<T: Copy + Default> SeqLock<T> {

pub fn write(&self, value: T) {
let () = Self::STORAGE_FITS_CACHELINE;
self.seq.fetch_add(1, Ordering::Release);
// `AcqRel`, not `Release`: the acquire half is tyndall's `smp_wmb()`
// between the odd sequence number and the payload (see the module
// doc). A plain `Release` would let the stores below become visible
// first, and a reader sampling `seq` either side of them would accept
// a half-written entry.
self.seq.fetch_add(1, Ordering::AcqRel);
unsafe {
*self.size.get() = size_of::<T>();
*self.entry.0.get() = value;
}
// `Release` is tyndall's second `smp_wmb()`: the payload stores above
// are visible before the even sequence number that publishes them.
self.seq.fetch_add(1, Ordering::Release);
}

pub fn read(&self, state: &mut SeqLockState, always_update_entry: bool) -> Result<T, i32> {
let () = Self::STORAGE_FITS_CACHELINE;
let mut entry = T::default();
loop {
// `Acquire` is tyndall's first `smp_rmb()`: the entry read below
// cannot be hoisted above this sample of the sequence number.
let seq1: u32 = self.seq.load(Ordering::Acquire);
if seq1 & 1 != 0 {
debug!("SeqLock is being written to, retrying...");
Expand All @@ -149,6 +197,11 @@ impl<T: Copy + Default> SeqLock<T> {
entry = unsafe { *self.entry.0.get() };
}

// tyndall's second `smp_rmb()`: the entry read above must not
// drift below this load, or the validation would be checking a
// sequence number the bytes did not come from. The `Acquire` on
// the load itself orders what follows it, not what precedes it.
fence(Ordering::Acquire);
let seq2 = self.seq.load(Ordering::Acquire);
if seq1 != seq2 {
debug!("SeqLock read failed, seq mismatch: {seq1} != {seq2}");
Expand Down Expand Up @@ -254,6 +307,119 @@ mod tests {

/// A record on the stack, zero-initialised as a fresh shared-memory
/// segment is: the write/read contract over the new layout.
/// A payload wide enough to tear: every word carries the same counter, so
/// any mix of two writes is visible as a word that disagrees with the
/// first one.
#[derive(Clone, Copy, Default, Debug)]
#[repr(C)]
struct Wide([u64; 32]);

impl Wide {
fn of(counter: u64) -> Self {
Self([counter; 32])
}

/// The counter, if every word agrees.
fn counter(&self) -> Option<u64> {
let first = self.0[0];
self.0.iter().all(|&word| word == first).then_some(first)
}
}

/// A record shared between the threads of the test, the way [`crate::IPC`]
/// shares one between processes.
///
/// `SeqLock` is deliberately neither `Send` nor `Sync` on its own — it is
/// a raw shared-memory record, and it is `IPC` that owns the mapping and
/// asserts the contract (one writer, any number of readers). The test
/// makes the same assertion for the same reason.
struct Shared(SeqLock<Wide>);

// SAFETY: exactly the contract `IPC` documents. One writer thread, three
// reader threads, and the seqlock's own atomics are the synchronisation.
unsafe impl Send for Shared {}
unsafe impl Sync for Shared {}

/// A reader never observes a mix of two writes.
///
/// This is a stress test, not a proof, and it is worth being precise about
/// what it does and does not establish. A missing barrier is a reordering
/// the hardware is *allowed* to perform, not one it must: the test can
/// only fail when a machine actually takes the liberty. Measured, with the
/// writer's `AcqRel` reverted to `Release`, it still passes on an Apple
/// M-series host, whose cores rarely reorder stores in practice, and it
/// would never fail on x86, whose store order hides the writer's half of
/// the bug outright.
///
/// So the evidence that the barriers are right is the code generation, not
/// this test: on aarch64 the writer's first increment must emit `ldaddal`
/// (acquire-release) rather than `ldaddl` (release only), and the reader's
/// [`fence`] must emit `dmb ishld`, which is exactly tyndall's `smp_rmb()`.
///
/// The test earns its two seconds by catching the gross regressions that
/// no amount of reasoning protects against — losing the retry loop,
/// publishing before the payload is stored, a future rewrite of `read`
/// that drops the second sequence check — and it may yet catch a real
/// reordering on the i.MX or a Jetson runner.
#[test]
fn a_reader_never_sees_a_half_written_entry() {
use std::sync::atomic::AtomicBool;
use std::sync::Arc;

// SAFETY: as in the round-trip test, every field is valid zeroed.
let lock: Arc<Shared> = Arc::new(Shared(unsafe { std::mem::zeroed() }));
let stop = Arc::new(AtomicBool::new(false));

let writer = {
let lock = Arc::clone(&lock);
let stop = Arc::clone(&stop);
std::thread::spawn(move || {
let mut counter: u64 = 1;
while !stop.load(Ordering::Relaxed) {
lock.0.write(Wide::of(counter));
counter = counter.wrapping_add(1);
}
counter
})
};

let readers: Vec<_> = (0..3)
.map(|_| {
let lock = Arc::clone(&lock);
let stop = Arc::clone(&stop);
std::thread::spawn(move || {
let mut state = SeqLockState::default();
let mut reads = 0u64;
while !stop.load(Ordering::Relaxed) {
if let Ok(entry) = lock.0.read(&mut state, true) {
assert!(
entry.counter().is_some(),
"torn entry: words disagree, first four are {:?}",
&entry.0[..4]
);
reads += 1;
}
}
reads
})
})
.collect();

std::thread::sleep(std::time::Duration::from_secs(2));
stop.store(true, Ordering::Relaxed);

let writes = writer.join().expect("writer panicked");
let reads: u64 = readers
.into_iter()
.map(|reader| reader.join().expect("reader saw a torn entry"))
.sum();
println!("{writes} writes, {reads} consistent reads across 3 readers");
// A run that raced nothing proves nothing; fail loudly rather than
// pass on an empty test.
assert!(writes > 1000, "the writer barely ran ({writes} writes)");
assert!(reads > 1000, "the readers barely ran ({reads} reads)");
}

#[test]
fn write_then_read_round_trips_through_the_aligned_entry() {
// SAFETY: every field is valid when zeroed (an atomic, a usize, and a
Expand Down
4 changes: 2 additions & 2 deletions reader/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,10 +13,10 @@ fn main() {
Ok(mut shm) => {
let mut data = Data::default();
let success = shm.get(&mut data);
println!("Read Succes {:?}: {:?}", success, data);
println!("Read Succes {success:?}: {data:?}");
}
Err(e) => {
eprintln!("Failed to read shared memory: {}", e);
eprintln!("Failed to read shared memory: {e}");
}
}
}
4 changes: 2 additions & 2 deletions writer/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,9 +13,9 @@ fn main() {

let data = Data {
counter: 3,
value: 3.14,
value: std::f32::consts::PI,
};

shm.set(data);
println!("Wrote: {:?}", data);
println!("Wrote: {data:?}");
}