From 9af8b3938203a9c491a458980b053ab6844fb876 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Mon, 5 Oct 2026 22:39:13 +1300 Subject: [PATCH 1/4] Require source-region coverage across supported systems. --- .github/workflows/test.yml | 15 ++++-- Cargo.lock | 4 +- crates/executor/src/scheduler/tokio.rs | 14 +++++ crates/executor/src/scheduler/tokio/tests.rs | 57 ++++++++++++++++++++ crates/executor/src/worker.rs | 4 ++ crates/executor/src/worker/tests.rs | 10 +++- 6 files changed, 96 insertions(+), 8 deletions(-) diff --git a/.github/workflows/test.yml b/.github/workflows/test.yml index 3c9eac0..bfc4986 100644 --- a/.github/workflows/test.yml +++ b/.github/workflows/test.yml @@ -12,9 +12,14 @@ permissions: jobs: coverage: - name: Formatting, Clippy, and coverage - runs-on: ubuntu-latest + name: Coverage (${{ matrix.runner }}) + runs-on: ${{ matrix.runner }} timeout-minutes: 20 + strategy: + fail-fast: false + matrix: + # Native selectors and positioned file I/O are target-specific. + runner: [ubuntu-24.04, macos-15, windows-latest] steps: - uses: actions/checkout@v7 - uses: actions-rust-lang/setup-rust-toolchain@v2 @@ -26,9 +31,11 @@ jobs: - name: Install coverage tool run: cargo install cargo-llvm-cov --locked - run: cargo fmt --all -- --check + if: matrix.runner == 'ubuntu-24.04' - run: cargo clippy --workspace --all-targets --locked -- -D warnings - - name: Run tests and require complete line coverage - run: cargo bake --locked test:coverage --all-targets true + if: matrix.runner == 'ubuntu-24.04' + - name: Run tests and require complete source-region coverage + run: cargo bake --locked test:coverage --all-targets true --all-features true test: name: ${{ matrix.name }} diff --git a/Cargo.lock b/Cargo.lock index 5c6a5cf..0f4513f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -154,9 +154,9 @@ dependencies = [ [[package]] name = "bake-test-rust" -version = "0.2.3" +version = "0.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1a9c5b715c668d525d144fedae81c64c1871de6f4d01bb9ce76376ce85240ef8" +checksum = "e9d5ceec9cc08a3bb421dd3755ad9e073831934793aa88d8f55095ad98c1dc45" dependencies = [ "bake", "serde_json", diff --git a/crates/executor/src/scheduler/tokio.rs b/crates/executor/src/scheduler/tokio.rs index 7c4fda6..7f1d128 100644 --- a/crates/executor/src/scheduler/tokio.rs +++ b/crates/executor/src/scheduler/tokio.rs @@ -65,6 +65,19 @@ impl Shared { owner: &Arc, future: FutureType, ) -> Result, SpawnError> + where + FutureType: Future + Send + 'static, + FutureType::Output: Send + 'static, + { + self.spawn_with_registration_hook(owner, future, || {}) + } + + fn spawn_with_registration_hook( + self: &Arc, + owner: &Arc, + future: FutureType, + after_registration: impl FnOnce(), + ) -> Result, SpawnError> where FutureType: Future + Send + 'static, FutureType::Output: Send + 'static, @@ -96,6 +109,7 @@ impl Shared { ); owner.remaining.fetch_add(1, Ordering::Release); drop(registry); + after_registration(); // Publish ownership before spawn, but do not hold the registry lock: // a shut-down runtime can destroy the submitted future immediately. let inner = self.runtime.spawn(TrackedFuture { diff --git a/crates/executor/src/scheduler/tokio/tests.rs b/crates/executor/src/scheduler/tokio/tests.rs index 25cbe5a..13cf5d9 100644 --- a/crates/executor/src/scheduler/tokio/tests.rs +++ b/crates/executor/src/scheduler/tokio/tests.rs @@ -3,7 +3,9 @@ use super::*; use std::future::ready; +use std::sync::mpsc; use std::task::{Context, Poll, Waker}; +use std::thread; fn runtime() -> ::tokio::runtime::Runtime { ::tokio::runtime::Builder::new_current_thread() @@ -191,6 +193,61 @@ fn task_identifier_exhaustion_is_reported() { )); } +#[test] +fn shutdown_cancels_a_task_registered_before_tokio_spawn() { + let runtime = runtime(); + let scheduler = Scheduler::new(runtime.handle().clone()); + let shared = Arc::clone(&scheduler.handle.shared); + let owner = Arc::new(Owner::new()); + let (registered, wait_for_registration) = mpsc::sync_channel(0); + let (resume, wait_to_resume) = mpsc::sync_channel(0); + + let spawning = { + let shared = Arc::clone(&shared); + let owner = Arc::clone(&owner); + thread::spawn(move || { + shared.spawn_with_registration_hook(&owner, std::future::pending::<()>(), || { + registered.send(()).unwrap(); + wait_to_resume.recv().unwrap(); + }) + }) + }; + + wait_for_registration.recv().unwrap(); + { + let registry = shared.registry.lock().unwrap(); + let registration = registry.tasks.values().next().unwrap(); + assert!(registration.abort.is_none()); + assert!(!owner.is_empty()); + } + + shared.close(None, true); + + assert!(shared.closed.load(Ordering::Acquire)); + assert!( + shared + .registry + .lock() + .unwrap() + .tasks + .values() + .next() + .unwrap() + .state + .cancelled + .load(Ordering::Acquire) + ); + + resume.send(()).unwrap(); + let task = match spawning.join().unwrap() { + Ok(task) => task, + Err(error) => panic!("task registration failed: {error}"), + }; + + assert!(matches!(runtime.block_on(task), Err(TaskError::Cancelled))); + assert!(owner.is_empty()); +} + #[test] fn socket_registration_preserves_configuration_and_conversion_errors() { let runtime = runtime(); diff --git a/crates/executor/src/worker.rs b/crates/executor/src/worker.rs index 21be8cc..4d654cf 100644 --- a/crates/executor/src/worker.rs +++ b/crates/executor/src/worker.rs @@ -133,6 +133,10 @@ fn search(shared: &Shared, local: &LocalWorker, iteration: usize) -> Steal Steal { if retry { Steal::Retry } else { Steal::Empty } } diff --git a/crates/executor/src/worker/tests.rs b/crates/executor/src/worker/tests.rs index cd00d3e..a6a7a3b 100644 --- a/crates/executor/src/worker/tests.rs +++ b/crates/executor/src/worker/tests.rs @@ -2,8 +2,8 @@ // Copyright, 2026, by Samuel Williams. use super::{ - IdleSearch, SearchResult, WorkerState, classify_search, finish_idle_search, take_idle_runnable, - take_stolen, + IdleSearch, SearchResult, WorkerState, classify_search, finish_idle_search, search_result, + take_idle_runnable, take_stolen, }; use crate::owner::Owner; use crate::task::{Runnable, TaskState, UNASSIGNED_WORKER}; @@ -45,6 +45,12 @@ fn records_retry_separately_from_empty_and_success() { assert!(retry); } +#[test] +fn search_reports_retry_only_when_a_queue_raced() { + assert!(matches!(search_result(true), Steal::Retry)); + assert!(matches!(search_result(false), Steal::Empty)); +} + #[test] fn worker_yields_after_a_concurrent_steal_retry() { assert!(matches!( From ff4396b597011b3ff080ffe7d7184ad10a163417 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Mon, 5 Oct 2026 22:48:18 +1300 Subject: [PATCH 2/4] Cover native and Tokio source regions on each platform. --- .github/workflows/test.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/test.yml b/.github/workflows/test.yml index bfc4986..5ad6fa9 100644 --- a/.github/workflows/test.yml +++ b/.github/workflows/test.yml @@ -35,7 +35,7 @@ jobs: - run: cargo clippy --workspace --all-targets --locked -- -D warnings if: matrix.runner == 'ubuntu-24.04' - name: Run tests and require complete source-region coverage - run: cargo bake --locked test:coverage --all-targets true --all-features true + run: cargo bake --locked test:coverage --all-targets true --features tokio test: name: ${{ matrix.name }} From a702c970d257d025d1f49cd76bd2e520007fefcf Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Mon, 5 Oct 2026 22:53:20 +1300 Subject: [PATCH 3/4] Cover Tokio socket write readiness. --- crates/executor/src/scheduler/tokio/tests.rs | 52 ++++++++++++++++++++ 1 file changed, 52 insertions(+) diff --git a/crates/executor/src/scheduler/tokio/tests.rs b/crates/executor/src/scheduler/tokio/tests.rs index 13cf5d9..328666b 100644 --- a/crates/executor/src/scheduler/tokio/tests.rs +++ b/crates/executor/src/scheduler/tokio/tests.rs @@ -3,6 +3,7 @@ use super::*; use std::future::ready; +use std::io::Read; use std::sync::mpsc; use std::task::{Context, Poll, Waker}; use std::thread; @@ -416,6 +417,57 @@ fn write_retries_interrupted_and_would_block_operations() { }); } +#[test] +fn socket_write_waits_until_the_send_buffer_is_writable() { + let runtime = runtime(); + let scheduler = Scheduler::new(runtime.handle().clone()); + let handle = scheduler.handle(); + let (client, mut server) = socket_pair(); + let socket = handle.register_socket(client).unwrap(); + let fill = vec![0; 64 * 1024]; + + runtime.block_on(async { + socket.inner.writable().await.unwrap(); + let mut written = 0; + loop { + match socket.inner.try_write(&fill) { + Ok(0) => panic!("socket write made no progress"), + Ok(count) => written += count, + Err(error) if error.kind() == io::ErrorKind::WouldBlock => break, + Err(error) => panic!("failed to fill socket send buffer: {error}"), + } + } + assert!(written > 0); + }); + + let (start_reader, wait_for_start) = mpsc::sync_channel(0); + let (finish_reader, wait_for_finish) = mpsc::sync_channel(0); + let reader = thread::spawn(move || { + wait_for_start.recv().unwrap(); + let mut buffer = [0; 64 * 1024]; + let count = server.read(&mut buffer).unwrap(); + wait_for_finish.recv().unwrap(); + count + }); + + let mut write = Box::pin(handle.io_write(&socket, vec![42])); + runtime.block_on(async { + std::future::poll_fn(|context| { + assert!(write.as_mut().poll(context).is_pending()); + start_reader.send(()).unwrap(); + Poll::Ready(()) + }) + .await; + + let (result, buffer) = write.await; + assert_eq!(result.unwrap(), 1); + assert_eq!(buffer, [42]); + }); + + finish_reader.send(()).unwrap(); + assert!(reader.join().unwrap() > 0); +} + #[test] fn write_returns_readiness_errors_and_non_retryable_errors() { let runtime = runtime(); From 0599854280811e69529ebe79b80e04c899f89ec9 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Mon, 5 Oct 2026 23:26:17 +1300 Subject: [PATCH 4/4] Drain the peer while testing socket write readiness. --- crates/executor/src/scheduler/tokio/tests.rs | 22 +++++++++++++++++--- 1 file changed, 19 insertions(+), 3 deletions(-) diff --git a/crates/executor/src/scheduler/tokio/tests.rs b/crates/executor/src/scheduler/tokio/tests.rs index 328666b..1cd4a8a 100644 --- a/crates/executor/src/scheduler/tokio/tests.rs +++ b/crates/executor/src/scheduler/tokio/tests.rs @@ -7,6 +7,7 @@ use std::io::Read; use std::sync::mpsc; use std::task::{Context, Poll, Waker}; use std::thread; +use std::time::Duration; fn runtime() -> ::tokio::runtime::Runtime { ::tokio::runtime::Builder::new_current_thread() @@ -423,6 +424,7 @@ fn socket_write_waits_until_the_send_buffer_is_writable() { let scheduler = Scheduler::new(runtime.handle().clone()); let handle = scheduler.handle(); let (client, mut server) = socket_pair(); + server.set_nonblocking(true).unwrap(); let socket = handle.register_socket(client).unwrap(); let fill = vec![0; 64 * 1024]; @@ -445,9 +447,23 @@ fn socket_write_waits_until_the_send_buffer_is_writable() { let reader = thread::spawn(move || { wait_for_start.recv().unwrap(); let mut buffer = [0; 64 * 1024]; - let count = server.read(&mut buffer).unwrap(); - wait_for_finish.recv().unwrap(); - count + let mut total = 0; + loop { + match wait_for_finish.try_recv() { + Ok(()) | Err(mpsc::TryRecvError::Disconnected) => break, + Err(mpsc::TryRecvError::Empty) => {} + } + + match server.read(&mut buffer) { + Ok(0) => break, + Ok(count) => total += count, + Err(error) if error.kind() == io::ErrorKind::WouldBlock => { + thread::sleep(Duration::from_millis(1)); + } + Err(error) => panic!("failed to drain socket receive buffer: {error}"), + } + } + total }); let mut write = Box::pin(handle.io_write(&socket, vec![42]));