From e03ca80104706af7a113ad4dccdeaeb9da2f025c Mon Sep 17 00:00:00 2001 From: Razz4780 Date: Tue, 4 Aug 2026 00:31:43 +0200 Subject: [PATCH] loom test added --- .github/workflows/ci.yaml | 18 ++++++++++++++++++ Cargo.toml | 6 ++++++ src/queue.rs | 40 +++++++++++++++++++++++++++++++++++---- 3 files changed, 60 insertions(+), 4 deletions(-) diff --git a/.github/workflows/ci.yaml b/.github/workflows/ci.yaml index bbb5567..4cd4171 100644 --- a/.github/workflows/ci.yaml +++ b/.github/workflows/ci.yaml @@ -71,3 +71,21 @@ jobs: run: > cargo +nightly miri test --lib queue::test::concurrent_enqueues_work + + loom: + runs-on: ubuntu-24.04 + timeout-minutes: 30 + + steps: + - uses: actions/checkout@v4 + + - uses: dtolnay/rust-toolchain@stable + + - uses: Swatinem/rust-cache@v2 + + - name: Run Loom model + env: + RUSTFLAGS: --cfg loom + run: > + cargo test --release --lib + queue::test::loom_concurrent_enqueues_work diff --git a/Cargo.toml b/Cargo.toml index b9b094e..45d3931 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -18,3 +18,9 @@ static_assertions = "1" [dev-dependencies] tokio = { version = "1", features = ["full"] } + +[target.'cfg(loom)'.dependencies] +loom = "0.7" + +[lints.rust] +unexpected_cfgs = { level = "warn", check-cfg = ['cfg(loom)'] } diff --git a/src/queue.rs b/src/queue.rs index ee3f964..eb4dccf 100644 --- a/src/queue.rs +++ b/src/queue.rs @@ -2,14 +2,15 @@ use std::{ cell::UnsafeCell, mem::{ManuallyDrop, MaybeUninit}, ops::Not, - sync::{ - Arc, Weak, - atomic::{AtomicBool, AtomicPtr, Ordering}, - }, + sync::{Arc, Weak}, task::Waker, }; use futures::task::AtomicWaker; +#[cfg(loom)] +use loom::sync::atomic::{AtomicBool, AtomicPtr, Ordering}; +#[cfg(not(loom))] +use std::sync::atomic::{AtomicBool, AtomicPtr, Ordering}; /// Receiver handle for a [`Queue`]. /// @@ -358,6 +359,37 @@ mod test { use crate::queue::Receiver; + #[cfg(loom)] + #[test] + fn loom_concurrent_enqueues_work() { + loom::model(|| { + let mut receiver = Receiver::default(); + let queue = receiver.queue().clone(); + let first = queue.create(1); + let second = queue.create(2); + + let first_sender = loom::thread::spawn(move || first.enqueue()); + let second_sender = loom::thread::spawn(move || second.enqueue()); + + let mut values = Vec::new(); + while values.len() < 2 { + let Some(node) = receiver.dequeue() else { + loom::thread::yield_now(); + continue; + }; + values.push(*node.get().value()); + } + + first_sender.join().unwrap(); + second_sender.join().unwrap(); + + assert!(receiver.dequeue().is_none()); + + values.sort_unstable(); + assert_eq!(values, [1, 2]); + }); + } + #[tokio::test(flavor = "multi_thread", worker_threads = 5)] async fn snapshots_work() { const SENDERS: usize = 4;