Skip to content

Commit e448b73

Browse files
Require source-region coverage across supported systems. (#10)
* Require source-region coverage across supported systems. * Cover native and Tokio source regions on each platform. * Cover Tokio socket write readiness. * Drain the peer while testing socket write readiness.
1 parent 7db81b8 commit e448b73

6 files changed

Lines changed: 164 additions & 8 deletions

File tree

‎.github/workflows/test.yml‎

Lines changed: 11 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -12,9 +12,14 @@ permissions:
1212

1313
jobs:
1414
coverage:
15-
name: Formatting, Clippy, and coverage
16-
runs-on: ubuntu-latest
15+
name: Coverage (${{ matrix.runner }})
16+
runs-on: ${{ matrix.runner }}
1717
timeout-minutes: 20
18+
strategy:
19+
fail-fast: false
20+
matrix:
21+
# Native selectors and positioned file I/O are target-specific.
22+
runner: [ubuntu-24.04, macos-15, windows-latest]
1823
steps:
1924
- uses: actions/checkout@v7
2025
- uses: actions-rust-lang/setup-rust-toolchain@v2
@@ -26,9 +31,11 @@ jobs:
2631
- name: Install coverage tool
2732
run: cargo install cargo-llvm-cov --locked
2833
- run: cargo fmt --all -- --check
34+
if: matrix.runner == 'ubuntu-24.04'
2935
- run: cargo clippy --workspace --all-targets --locked -- -D warnings
30-
- name: Run tests and require complete line coverage
31-
run: cargo bake --locked test:coverage --all-targets true
36+
if: matrix.runner == 'ubuntu-24.04'
37+
- name: Run tests and require complete source-region coverage
38+
run: cargo bake --locked test:coverage --all-targets true --features tokio
3239

3340
test:
3441
name: ${{ matrix.name }}

‎Cargo.lock‎

Lines changed: 2 additions & 2 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎crates/executor/src/scheduler/tokio.rs‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -65,6 +65,19 @@ impl Shared {
6565
owner: &Arc<Owner>,
6666
future: FutureType,
6767
) -> Result<TaskHandle<FutureType::Output>, SpawnError>
68+
where
69+
FutureType: Future + Send + 'static,
70+
FutureType::Output: Send + 'static,
71+
{
72+
self.spawn_with_registration_hook(owner, future, || {})
73+
}
74+
75+
fn spawn_with_registration_hook<FutureType>(
76+
self: &Arc<Self>,
77+
owner: &Arc<Owner>,
78+
future: FutureType,
79+
after_registration: impl FnOnce(),
80+
) -> Result<TaskHandle<FutureType::Output>, SpawnError>
6881
where
6982
FutureType: Future + Send + 'static,
7083
FutureType::Output: Send + 'static,
@@ -96,6 +109,7 @@ impl Shared {
96109
);
97110
owner.remaining.fetch_add(1, Ordering::Release);
98111
drop(registry);
112+
after_registration();
99113
// Publish ownership before spawn, but do not hold the registry lock:
100114
// a shut-down runtime can destroy the submitted future immediately.
101115
let inner = self.runtime.spawn(TrackedFuture {

‎crates/executor/src/scheduler/tokio/tests.rs‎

Lines changed: 125 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,11 @@
33

44
use super::*;
55
use std::future::ready;
6+
use std::io::Read;
7+
use std::sync::mpsc;
68
use std::task::{Context, Poll, Waker};
9+
use std::thread;
10+
use std::time::Duration;
711

812
fn runtime() -> ::tokio::runtime::Runtime {
913
::tokio::runtime::Builder::new_current_thread()
@@ -191,6 +195,61 @@ fn task_identifier_exhaustion_is_reported() {
191195
));
192196
}
193197

198+
#[test]
199+
fn shutdown_cancels_a_task_registered_before_tokio_spawn() {
200+
let runtime = runtime();
201+
let scheduler = Scheduler::new(runtime.handle().clone());
202+
let shared = Arc::clone(&scheduler.handle.shared);
203+
let owner = Arc::new(Owner::new());
204+
let (registered, wait_for_registration) = mpsc::sync_channel(0);
205+
let (resume, wait_to_resume) = mpsc::sync_channel(0);
206+
207+
let spawning = {
208+
let shared = Arc::clone(&shared);
209+
let owner = Arc::clone(&owner);
210+
thread::spawn(move || {
211+
shared.spawn_with_registration_hook(&owner, std::future::pending::<()>(), || {
212+
registered.send(()).unwrap();
213+
wait_to_resume.recv().unwrap();
214+
})
215+
})
216+
};
217+
218+
wait_for_registration.recv().unwrap();
219+
{
220+
let registry = shared.registry.lock().unwrap();
221+
let registration = registry.tasks.values().next().unwrap();
222+
assert!(registration.abort.is_none());
223+
assert!(!owner.is_empty());
224+
}
225+
226+
shared.close(None, true);
227+
228+
assert!(shared.closed.load(Ordering::Acquire));
229+
assert!(
230+
shared
231+
.registry
232+
.lock()
233+
.unwrap()
234+
.tasks
235+
.values()
236+
.next()
237+
.unwrap()
238+
.state
239+
.cancelled
240+
.load(Ordering::Acquire)
241+
);
242+
243+
resume.send(()).unwrap();
244+
let task = match spawning.join().unwrap() {
245+
Ok(task) => task,
246+
Err(error) => panic!("task registration failed: {error}"),
247+
};
248+
249+
assert!(matches!(runtime.block_on(task), Err(TaskError::Cancelled)));
250+
assert!(owner.is_empty());
251+
}
252+
194253
#[test]
195254
fn socket_registration_preserves_configuration_and_conversion_errors() {
196255
let runtime = runtime();
@@ -359,6 +418,72 @@ fn write_retries_interrupted_and_would_block_operations() {
359418
});
360419
}
361420

421+
#[test]
422+
fn socket_write_waits_until_the_send_buffer_is_writable() {
423+
let runtime = runtime();
424+
let scheduler = Scheduler::new(runtime.handle().clone());
425+
let handle = scheduler.handle();
426+
let (client, mut server) = socket_pair();
427+
server.set_nonblocking(true).unwrap();
428+
let socket = handle.register_socket(client).unwrap();
429+
let fill = vec![0; 64 * 1024];
430+
431+
runtime.block_on(async {
432+
socket.inner.writable().await.unwrap();
433+
let mut written = 0;
434+
loop {
435+
match socket.inner.try_write(&fill) {
436+
Ok(0) => panic!("socket write made no progress"),
437+
Ok(count) => written += count,
438+
Err(error) if error.kind() == io::ErrorKind::WouldBlock => break,
439+
Err(error) => panic!("failed to fill socket send buffer: {error}"),
440+
}
441+
}
442+
assert!(written > 0);
443+
});
444+
445+
let (start_reader, wait_for_start) = mpsc::sync_channel(0);
446+
let (finish_reader, wait_for_finish) = mpsc::sync_channel(0);
447+
let reader = thread::spawn(move || {
448+
wait_for_start.recv().unwrap();
449+
let mut buffer = [0; 64 * 1024];
450+
let mut total = 0;
451+
loop {
452+
match wait_for_finish.try_recv() {
453+
Ok(()) | Err(mpsc::TryRecvError::Disconnected) => break,
454+
Err(mpsc::TryRecvError::Empty) => {}
455+
}
456+
457+
match server.read(&mut buffer) {
458+
Ok(0) => break,
459+
Ok(count) => total += count,
460+
Err(error) if error.kind() == io::ErrorKind::WouldBlock => {
461+
thread::sleep(Duration::from_millis(1));
462+
}
463+
Err(error) => panic!("failed to drain socket receive buffer: {error}"),
464+
}
465+
}
466+
total
467+
});
468+
469+
let mut write = Box::pin(handle.io_write(&socket, vec![42]));
470+
runtime.block_on(async {
471+
std::future::poll_fn(|context| {
472+
assert!(write.as_mut().poll(context).is_pending());
473+
start_reader.send(()).unwrap();
474+
Poll::Ready(())
475+
})
476+
.await;
477+
478+
let (result, buffer) = write.await;
479+
assert_eq!(result.unwrap(), 1);
480+
assert_eq!(buffer, [42]);
481+
});
482+
483+
finish_reader.send(()).unwrap();
484+
assert!(reader.join().unwrap() > 0);
485+
}
486+
362487
#[test]
363488
fn write_returns_readiness_errors_and_non_retryable_errors() {
364489
let runtime = runtime();

‎crates/executor/src/worker.rs‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -133,6 +133,10 @@ fn search(shared: &Shared, local: &LocalWorker, iteration: usize) -> Steal<Runna
133133
return Steal::Success(runnable);
134134
}
135135
}
136+
search_result(retry)
137+
}
138+
139+
fn search_result(retry: bool) -> Steal<Runnable> {
136140
if retry { Steal::Retry } else { Steal::Empty }
137141
}
138142

‎crates/executor/src/worker/tests.rs‎

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,8 +2,8 @@
22
// Copyright, 2026, by Samuel Williams.
33

44
use super::{
5-
IdleSearch, SearchResult, WorkerState, classify_search, finish_idle_search, take_idle_runnable,
6-
take_stolen,
5+
IdleSearch, SearchResult, WorkerState, classify_search, finish_idle_search, search_result,
6+
take_idle_runnable, take_stolen,
77
};
88
use crate::owner::Owner;
99
use crate::task::{Runnable, TaskState, UNASSIGNED_WORKER};
@@ -45,6 +45,12 @@ fn records_retry_separately_from_empty_and_success() {
4545
assert!(retry);
4646
}
4747

48+
#[test]
49+
fn search_reports_retry_only_when_a_queue_raced() {
50+
assert!(matches!(search_result(true), Steal::Retry));
51+
assert!(matches!(search_result(false), Steal::Empty));
52+
}
53+
4854
#[test]
4955
fn worker_yields_after_a_concurrent_steal_retry() {
5056
assert!(matches!(

0 commit comments

Comments
 (0)