Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
26 commits
Select commit Hold shift + click to select a range
e700444
refactor(mpsc): simplify unbounded channel coordination
tisonkun Sep 9, 2026
960d382
ci: compile benchmarks without running them
tisonkun Sep 9, 2026
8733748
ci: compile benchmarks in the check job
tisonkun Sep 9, 2026
09f8ecb
test(mpsc): keep buffer reclamation checks inline
tisonkun Sep 9, 2026
9a7d91a
bench(watch): remove future erasure from adapters
tisonkun Sep 9, 2026
ef6646c
bench(waitgroup): remove incomplete worker round trips
tisonkun Sep 9, 2026
e64eca7
bench(oneshot): use the shared task waker
tisonkun Sep 9, 2026
82b230e
bench(singleflight): focus on complete caller lifecycles
tisonkun Sep 9, 2026
f1e1394
bench(once_map): poll cached computations once
tisonkun Sep 9, 2026
14accd5
bench(once_map): consolidate entry removal measurements
tisonkun Sep 9, 2026
26ba407
bench(pool): use one lazy construction configuration
tisonkun Sep 9, 2026
6019df6
test: share waker and deadlock support across integration suites
tisonkun Sep 9, 2026
469bbe8
bench(watch): keep one shared comparison suite
tisonkun Sep 9, 2026
f23ab1b
bench(waitgroup): remove duplicated lifecycle cases
tisonkun Sep 9, 2026
99a86be
bench(broadcast): reuse ecosystem concurrency coverage
tisonkun Sep 9, 2026
32d1820
bench(mpsc): retain representative workload dimensions
tisonkun Sep 9, 2026
21af06e
test: establish pending waiters before releasing conditions
tisonkun Sep 9, 2026
3bebe7b
test(mpsc): verify every producer message exactly once
tisonkun Sep 9, 2026
69b4e2b
test(once_map): retain public lifecycle and collision coverage
tisonkun Sep 9, 2026
bc83844
test(wakerset): rely on public deferred-drop coverage
tisonkun Sep 9, 2026
f661e6f
test(broadcast): remove subsumed delivery cases
tisonkun Sep 9, 2026
85b391a
test(mutex): consolidate mapped guard ownership coverage
tisonkun Sep 9, 2026
f79c080
test(semaphore): remove subsumed permit cases
tisonkun Sep 9, 2026
7a429bc
test(singleflight): keep the stronger basic computation case
tisonkun Sep 9, 2026
c5b907c
test(semaphore): keep node reclamation checks next to state
tisonkun Sep 9, 2026
f5fc479
docs(xtask): correct no-capture execution help
tisonkun Sep 9, 2026
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
22 changes: 3 additions & 19 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,8 @@ jobs:
- run: cargo x lint
- name: Check feature matrix
run: cargo x check
- name: Compile benchmarks
run: cargo x bench --no-run

test:
name: Run tests
Expand Down Expand Up @@ -112,37 +114,19 @@ jobs:
- name: Run Miri tests
run: cargo x miri

benchmark:
name: Run benchmarks
runs-on: ubuntu-24.04
timeout-minutes: 30
env:
RUSTUP_TOOLCHAIN: stable
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- name: Install stable toolchain
run: >-
rustup toolchain install stable
--profile minimal
--no-self-update
- uses: swatinem/rust-cache@6323deb102c322ba6fcbdcafc7e3dddab59af2b6 # v2.9.2
- run: cargo x bench

required:
name: Required
runs-on: ubuntu-24.04
if: ${{ always() }}
needs:
- benchmark
- check
- miri
- test
steps:
- name: Guardian
run: |
if [[ ! ( \
"${{ needs.benchmark.result }}" == "success" \
&& "${{ needs.check.result }}" == "success" \
"${{ needs.check.result }}" == "success" \
&& "${{ needs.miri.result }}" == "success" \
&& "${{ needs.test.result }}" == "success" \
) ]]; then
Expand Down
22 changes: 15 additions & 7 deletions asyncband/src/internal/semaphore.rs
Original file line number Diff line number Diff line change
Expand Up @@ -256,11 +256,6 @@ impl Semaphore {
wakers.wake_all();
}
}

#[cfg(test)]
pub fn num_waiter_nodes(&self) -> usize {
self.waiters.lock().occupied_len()
}
}

#[derive(Debug)]
Expand Down Expand Up @@ -448,6 +443,19 @@ mod tests {
}
}

#[test]
fn fulfilled_reduce_permits_debt_reclaims_its_waiter_node() {
let semaphore = Semaphore::new(0);

for _ in 0..3 {
semaphore.reduce_permits(1);
assert_eq!(semaphore.waiters.lock().occupied_len(), 1);

semaphore.release(1);
assert_eq!(semaphore.waiters.lock().occupied_len(), 0);
}
}

#[test]
fn release_drains_more_than_one_wake_batch() {
const WAITER_COUNT: usize = WAKE_BATCH_SIZE + 3;
Expand All @@ -462,14 +470,14 @@ mod tests {
for acquire in &mut acquires {
assert!(acquire.poll_once(&waker).is_pending());
}
assert_eq!(semaphore.num_waiter_nodes(), WAITER_COUNT);
assert_eq!(semaphore.waiters.lock().occupied_len(), WAITER_COUNT);

semaphore.release(WAITER_COUNT);
assert_eq!(counter.0.load(Ordering::Relaxed), WAITER_COUNT);

for acquire in &mut acquires {
assert!(acquire.poll_once(&waker).is_ready());
}
assert_eq!(semaphore.num_waiter_nodes(), 0);
assert_eq!(semaphore.waiters.lock().occupied_len(), 0);
}
}
52 changes: 1 addition & 51 deletions asyncband/src/internal/wakerset.rs
Original file line number Diff line number Diff line change
Expand Up @@ -111,65 +111,15 @@ impl WakerSet {
pub fn unregister(&mut self, token: &mut Option<WakerToken>) -> Option<Waker> {
token.take().map(|token| self.wakers.remove(token.0))
}

#[cfg(test)]
fn registered_len(&self) -> usize {
self.wakers.len()
}
}

#[cfg(test)]
mod tests {
use std::mem::size_of;
use std::sync::Arc;
use std::sync::atomic::AtomicBool;
use std::sync::atomic::AtomicUsize;
use std::sync::atomic::Ordering;
use std::task::Wake;

use super::*;

struct DropWake {
dropped: Arc<AtomicBool>,
wake_count: AtomicUsize,
}

impl Wake for DropWake {
fn wake(self: Arc<Self>) {
self.wake_count.fetch_add(1, Ordering::Relaxed);
}
}

impl Drop for DropWake {
fn drop(&mut self) {
self.dropped.store(true, Ordering::Relaxed);
}
}
use super::WakerToken;

#[test]
fn waker_token_preserves_the_option_niche() {
assert_eq!(size_of::<WakerToken>(), size_of::<usize>());
assert_eq!(size_of::<WakerToken>(), size_of::<Option<WakerToken>>());
}

#[test]
fn unregister_returns_the_waker_for_deferred_drop() {
let dropped = Arc::new(AtomicBool::new(false));
let waker = Waker::from(Arc::new(DropWake {
dropped: dropped.clone(),
wake_count: AtomicUsize::new(0),
}));
let mut wakers = WakerSet::new();
let mut token = None;

drop(wakers.register(&mut token, &waker));
drop(waker);

let removed = wakers.unregister(&mut token);
assert_eq!(wakers.registered_len(), 0);
assert!(!dropped.load(Ordering::Relaxed));

drop(removed);
assert!(dropped.load(Ordering::Relaxed));
}
}
69 changes: 55 additions & 14 deletions asyncband/src/mpsc/unbounded/buffer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,15 +25,13 @@ pub const SEGMENT_BYTES: usize = 32 * 1024;
pub struct Buffer<T> {
writable: VecDeque<T>,
sealed: VecDeque<VecDeque<T>>,
spare: VecDeque<T>,
}

impl<T> Buffer<T> {
pub fn new() -> Self {
Self {
writable: VecDeque::new(),
sealed: VecDeque::new(),
spare: VecDeque::new(),
}
}

Expand All @@ -48,11 +46,7 @@ impl<T> Buffer<T> {

pub fn push(&mut self, value: T) {
if self.writable.len() == Self::segment_capacity() {
let next = if self.spare.capacity() == 0 {
VecDeque::with_capacity(Self::segment_capacity())
} else {
mem::take(&mut self.spare)
};
let next = VecDeque::with_capacity(Self::segment_capacity());
let sealed = mem::replace(&mut self.writable, next);
self.sealed.push_back(sealed);
}
Expand All @@ -62,22 +56,19 @@ impl<T> Buffer<T> {
pub fn refill(&mut self, batch: &mut VecDeque<T>) {
debug_assert!(batch.is_empty());
if let Some(sealed) = self.sealed.pop_front() {
// Keep one empty segment for the next producer rollover. Every other consumed
// segment is released, so retained payload storage does not track peak occupancy.
self.spare = mem::replace(batch, sealed);
*batch = sealed;
if self.sealed.is_empty() && self.sealed.capacity() * size_of::<VecDeque<T>>() > 1024 {
self.sealed = VecDeque::new();
}
} else if !self.writable.is_empty() {
self.spare = VecDeque::new();
mem::swap(batch, &mut self.writable);
}
}
}

pub fn pop_batch<T>(batch: &mut VecDeque<T>) -> T {
if batch.len() == 1 && batch.capacity().saturating_mul(size_of::<T>()) > SEGMENT_BYTES {
// Retire the allocation on the last value, outside the inbox lock. Keep this as a tail
// Retire the allocation on the last value, outside the shared lock. Keep this as a tail
// expression to avoid intermediate storage for large inline values.
mem::take(batch).pop_front()
} else {
Expand All @@ -87,5 +78,55 @@ pub fn pop_batch<T>(batch: &mut VecDeque<T>) -> T {
}

#[cfg(test)]
#[path = "buffer_tests.rs"]
mod tests;
mod tests {
use std::collections::VecDeque;

use super::Buffer;
use super::SEGMENT_BYTES;
use super::pop_batch;

fn allocated_bytes<T>(buffer: &Buffer<T>, batch: &VecDeque<T>) -> usize {
let slots = batch.capacity()
+ buffer.writable.capacity()
+ buffer.sealed.iter().map(VecDeque::capacity).sum::<usize>();
slots * size_of::<T>()
}

fn receive<T>(buffer: &mut Buffer<T>, batch: &mut VecDeque<T>) -> T {
if batch.is_empty() {
buffer.refill(batch);
}
pop_batch(batch)
}

#[test]
fn a_partial_drain_reclaims_segments_and_preserves_new_sends() {
let mut buffer = Buffer::new();
let mut batch = VecDeque::new();
for value in 0..1024usize {
buffer.push([value; 128]);
}
let peak = allocated_bytes(&buffer, &batch);
for value in 0..512 {
assert_eq!(receive(&mut buffer, &mut batch), [value; 128]);
}
assert!(allocated_bytes(&buffer, &batch) <= peak * 3 / 4);
// This value must stay behind both the current batch and the sealed segments.
buffer.push([1024; 128]);
for value in 512..=1024 {
assert_eq!(receive(&mut buffer, &mut batch), [value; 128]);
}
assert!(allocated_bytes(&buffer, &batch) <= 2 * SEGMENT_BYTES);
buffer.refill(&mut batch);
assert!(batch.is_empty());
}

#[test]
fn oversized_inline_values_release_the_allocation_on_the_last_receive() {
let mut buffer = Buffer::new();
let mut batch = VecDeque::new();
buffer.push([7u8; SEGMENT_BYTES + 1]);
assert_eq!(receive(&mut buffer, &mut batch), [7u8; SEGMENT_BYTES + 1]);
assert_eq!(allocated_bytes(&buffer, &batch), 0);
}
}
114 changes: 0 additions & 114 deletions asyncband/src/mpsc/unbounded/buffer_tests.rs

This file was deleted.

Loading