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
15 changes: 11 additions & 4 deletions .github/workflows/test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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 --features tokio

test:
name: ${{ matrix.name }}
Expand Down
4 changes: 2 additions & 2 deletions Cargo.lock

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

14 changes: 14 additions & 0 deletions crates/executor/src/scheduler/tokio.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,19 @@ impl Shared {
owner: &Arc<Owner>,
future: FutureType,
) -> Result<TaskHandle<FutureType::Output>, SpawnError>
where
FutureType: Future + Send + 'static,
FutureType::Output: Send + 'static,
{
self.spawn_with_registration_hook(owner, future, || {})
}

fn spawn_with_registration_hook<FutureType>(
self: &Arc<Self>,
owner: &Arc<Owner>,
future: FutureType,
after_registration: impl FnOnce(),
) -> Result<TaskHandle<FutureType::Output>, SpawnError>
where
FutureType: Future + Send + 'static,
FutureType::Output: Send + 'static,
Expand Down Expand Up @@ -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 {
Expand Down
125 changes: 125 additions & 0 deletions crates/executor/src/scheduler/tokio/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,11 @@

use super::*;
use std::future::ready;
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()
Expand Down Expand Up @@ -191,6 +195,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();
Expand Down Expand Up @@ -359,6 +418,72 @@ 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();
server.set_nonblocking(true).unwrap();
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 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]));
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();
Expand Down
4 changes: 4 additions & 0 deletions crates/executor/src/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,10 @@ fn search(shared: &Shared, local: &LocalWorker, iteration: usize) -> Steal<Runna
return Steal::Success(runnable);
}
}
search_result(retry)
}

fn search_result(retry: bool) -> Steal<Runnable> {
if retry { Steal::Retry } else { Steal::Empty }
}

Expand Down
10 changes: 8 additions & 2 deletions crates/executor/src/worker/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -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!(
Expand Down
Loading