diff --git a/.semgrep/rules/no-clock-read-in-drop.yaml b/.semgrep/rules/no-clock-read-in-drop.yaml
new file mode 100644
index 0000000000..eb2e38fb2c
--- /dev/null
+++ b/.semgrep/rules/no-clock-read-in-drop.yaml
@@ -0,0 +1,27 @@
+rules:
+ - id: rust-no-clock-read-in-drop
+ languages: [rust]
+ severity: ERROR
+ message: |
+ Do not read the clock from a `Drop` implementation.
+
+ Once anything in the process has paused the virtual clock, tokio routes
+ every read through the calling thread's runtime context. A `Drop` that
+ runs during thread-local teardown may find that context already
+ destroyed, and tokio's response is a panic inside a destructor -- which
+ aborts the process rather than failing the test.
+
+ Take the reading before the value is dropped and pass it in, or record
+ the instant when the value is created.
+ paths:
+ exclude:
+ - .codeql/tests/
+ - clock/src/
+ patterns:
+ - pattern-inside: |
+ fn drop(&mut self) {
+ ...
+ }
+ - pattern-either:
+ - pattern: clock::now()
+ - pattern: clock::system_now()
diff --git a/Cargo.lock b/Cargo.lock
index eb3f3da814..ffd3dfc0f4 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -1316,6 +1316,7 @@ dependencies = [
name = "dataplane-clock"
version = "0.25.2"
dependencies = [
+ "dataplane-concurrency",
"tokio",
]
diff --git a/Cargo.toml b/Cargo.toml
index f0783857a8..3ab7535f2f 100644
--- a/Cargo.toml
+++ b/Cargo.toml
@@ -256,7 +256,7 @@ overflow-checks = false
codegen-units = 1
rpath = true
-[profile.fuzz]
+[profile.checked]
inherits = "release"
opt-level = 2
debug-assertions = true
diff --git a/acl/tests/property_predicate.rs b/acl/tests/property_predicate.rs
index 3e245db8d4..007462100d 100644
--- a/acl/tests/property_predicate.rs
+++ b/acl/tests/property_predicate.rs
@@ -211,100 +211,102 @@ where
}
const MIN_ASSERTED_HITS: u64 = 20;
const MIN_ASSERTED_MISSES: u64 = 20;
-fn run_property(
- name_prefix: &str,
- install_dpdk: impl Fn(String, &FiveTupleRule) -> T + core::panic::RefUnwindSafe,
-) where
- A: KeyAddr,
- PrefixSpec: FieldHit + FieldMiss + IsUniversal,
- T: Lookup, Verdict>,
- RawRule: TypeGenerator,
-{
- let asserted_hits = AtomicU64::new(0);
- let asserted_misses = AtomicU64::new(0);
+macro_rules! run_property {
+ ($a:ty, $name_prefix:expr, $install_dpdk:expr) => {{
+ let asserted_hits = AtomicU64::new(0);
+ let asserted_misses = AtomicU64::new(0);
- bolero::check!()
- .with_type::<(RawRule, Box<[u8]>, Box<[u8]>)>()
- .for_each(|(raw, hit_bytes, miss_bytes)| {
- let rule = build_rule(raw);
- let dpdk = install_dpdk(unique_name(name_prefix), &rule);
- let reference = ReferenceTable::, Verdict>::new(vec![RefRule::new(
- rule.into_backend_fields::(),
- Verdict::Drop,
- )]);
-
- let hits = HitsGen { rule };
- let n_hits = sweep(&hits, hit_bytes, |k| {
- assert!(rule.accepts(k), "hits gen produced a rejected key: {k:?}");
- assert_eq!(reference.lookup(k), Some(&Verdict::Drop));
- assert_eq!(dpdk.lookup(k), Some(&Verdict::Drop));
- });
- asserted_hits.fetch_add(n_hits, Ordering::Relaxed);
+ bolero::check!()
+ .with_type::<(RawRule<$a>, Box<[u8]>, Box<[u8]>)>()
+ .for_each(|(raw, hit_bytes, miss_bytes)| {
+ let rule = build_rule(raw);
+ let dpdk = $install_dpdk(unique_name($name_prefix), &rule);
+ let reference = ReferenceTable::, Verdict>::new(vec![RefRule::new(
+ rule.into_backend_fields::(),
+ Verdict::Drop,
+ )]);
- if !rule.is_universal() {
- let misses = MissesGen { rule };
- let n_misses = sweep(&misses, miss_bytes, |k| {
- assert!(
- !rule.accepts(k),
- "misses gen produced an accepted key: {k:?}",
- );
- assert_eq!(reference.lookup(k), None);
- assert_eq!(dpdk.lookup(k), None);
+ let hits = HitsGen { rule };
+ let n_hits = sweep(&hits, hit_bytes, |k| {
+ assert!(rule.accepts(k), "hits gen produced a rejected key: {k:?}");
+ assert_eq!(reference.lookup(k), Some(&Verdict::Drop));
+ assert_eq!(dpdk.lookup(k), Some(&Verdict::Drop));
});
- asserted_misses.fetch_add(n_misses, Ordering::Relaxed);
- }
- });
+ asserted_hits.fetch_add(n_hits, Ordering::Relaxed);
- let h = asserted_hits.load(Ordering::Relaxed);
- let m = asserted_misses.load(Ordering::Relaxed);
- assert!(
- h >= MIN_ASSERTED_HITS,
- "asserted only {h} hits (< {MIN_ASSERTED_HITS}); generator may have gone inert",
- );
- assert!(
- m >= MIN_ASSERTED_MISSES,
- "asserted only {m} misses (< {MIN_ASSERTED_MISSES}); generator may have gone inert",
- );
+ if !rule.is_universal() {
+ let misses = MissesGen { rule };
+ let n_misses = sweep(&misses, miss_bytes, |k| {
+ assert!(
+ !rule.accepts(k),
+ "misses gen produced an accepted key: {k:?}",
+ );
+ assert_eq!(reference.lookup(k), None);
+ assert_eq!(dpdk.lookup(k), None);
+ });
+ asserted_misses.fetch_add(n_misses, Ordering::Relaxed);
+ }
+ });
+
+ let h = asserted_hits.load(Ordering::Relaxed);
+ let m = asserted_misses.load(Ordering::Relaxed);
+ assert!(
+ h >= MIN_ASSERTED_HITS,
+ "asserted only {h} hits (< {MIN_ASSERTED_HITS}); generator may have gone inert",
+ );
+ assert!(
+ m >= MIN_ASSERTED_MISSES,
+ "asserted only {m} misses (< {MIN_ASSERTED_MISSES}); generator may have gone inert",
+ );
+ }};
}
#[test]
#[dpdk::with_eal]
fn property_v4() {
- run_property::>("prop_v4", |name, rule| {
- install_table(
- &name,
- NonZero::new(2).expect("nonzero"),
- vec![
- RuleSpec::, Verdict>::new(
- Priority::new(1).expect("nonzero priority"),
- CategoryMask::new(1).expect("nonzero mask"),
- rule.into_backend_fields::(),
- Verdict::Drop,
- )
- .expect("RuleSpec"),
- ],
- )
- .expect("install_table")
- });
+ run_property!(
+ Ipv4Addr,
+ "prop_v4",
+ |name: String, rule: &FiveTupleRule| {
+ install_table(
+ &name,
+ NonZero::new(2).expect("nonzero"),
+ vec![
+ RuleSpec::, Verdict>::new(
+ Priority::new(1).expect("nonzero priority"),
+ CategoryMask::new(1).expect("nonzero mask"),
+ rule.into_backend_fields::(),
+ Verdict::Drop,
+ )
+ .expect("RuleSpec"),
+ ],
+ )
+ .expect("install_table")
+ }
+ );
}
#[test]
#[dpdk::with_eal]
fn property_v6() {
- run_property::>("prop_v6", |name, rule| {
- install_table(
- &name,
- NonZero::new(2).expect("nonzero"),
- vec![
- RuleSpec::, Verdict>::new(
- Priority::new(1).expect("nonzero priority"),
- CategoryMask::new(1).expect("nonzero mask"),
- rule.into_backend_fields::(),
- Verdict::Drop,
- )
- .expect("RuleSpec"),
- ],
- )
- .expect("install_table")
- });
+ run_property!(
+ Ipv6Addr,
+ "prop_v6",
+ |name: String, rule: &FiveTupleRule| {
+ install_table(
+ &name,
+ NonZero::new(2).expect("nonzero"),
+ vec![
+ RuleSpec::, Verdict>::new(
+ Priority::new(1).expect("nonzero priority"),
+ CategoryMask::new(1).expect("nonzero mask"),
+ rule.into_backend_fields::(),
+ Verdict::Drop,
+ )
+ .expect("RuleSpec"),
+ ],
+ )
+ .expect("install_table")
+ }
+ );
}
diff --git a/clock/Cargo.toml b/clock/Cargo.toml
index 17644aa53f..75776e69e1 100644
--- a/clock/Cargo.toml
+++ b/clock/Cargo.toml
@@ -11,3 +11,9 @@ virtual = ["dep:tokio"]
[dependencies]
tokio = { workspace = true, optional = true, features = ["test-util", "time"] }
+
+[dev-dependencies]
+concurrency = { workspace = true }
+
+[lints.rust]
+unexpected_cfgs = { level = "warn", check-cfg = ['cfg(wall_clock)'] }
diff --git a/clock/build.rs b/clock/build.rs
new file mode 100644
index 0000000000..4a4a7fc660
--- /dev/null
+++ b/clock/build.rs
@@ -0,0 +1,39 @@
+// SPDX-License-Identifier: Apache-2.0
+// Copyright Open Network Fabric Authors
+
+use std::process::Command;
+use std::{env, fs, path::PathBuf};
+
+fn main() {
+ println!("cargo::rerun-if-env-changed=RUSTC_BOOTSTRAP");
+ println!("cargo::rustc-check-cfg=cfg(has_spawn_hook)");
+
+ let out = PathBuf::from(env::var_os("OUT_DIR").expect("cargo sets OUT_DIR"));
+ let probe = out.join("spawn_hook_probe.rs");
+ if fs::write(
+ &probe,
+ "#![feature(thread_spawn_hook)]\n\
+ pub fn probe() { std::thread::add_spawn_hook(|_| || {}); }\n",
+ )
+ .is_err()
+ {
+ return;
+ }
+
+ let rustc = env::var_os("RUSTC").unwrap_or_else(|| "rustc".into());
+ let accepted = Command::new(rustc)
+ .args(["--crate-type=lib", "--emit=metadata", "-o"])
+ .arg(out.join("spawn_hook_probe.rmeta"))
+ .arg(&probe)
+ .status()
+ .is_ok_and(|status| status.success());
+
+ if accepted {
+ println!("cargo::rustc-cfg=has_spawn_hook");
+ } else {
+ println!(
+ "cargo::warning=thread_spawn_hook is unavailable, so a test that drives the clock \
+ cannot check the threads it spawns. Set RUSTC_BOOTSTRAP=1 (the dev shell does)."
+ );
+ }
+}
diff --git a/clock/src/lib.rs b/clock/src/lib.rs
index 9493f60e41..1bd29872d4 100644
--- a/clock/src/lib.rs
+++ b/clock/src/lib.rs
@@ -1,35 +1,83 @@
// SPDX-License-Identifier: Apache-2.0
// Copyright Open Network Fabric Authors
+#![cfg_attr(
+ all(has_spawn_hook, feature = "virtual", not(wall_clock)),
+ feature(thread_spawn_hook)
+)]
#![deny(clippy::all, clippy::pedantic)]
#![deny(rustdoc::all)]
#![deny(unsafe_code)]
pub use std::time::{Duration, Instant, SystemTime, SystemTimeError, TryFromFloatSecsError};
+#[cfg(feature = "virtual")]
+pub mod virtual_time;
+
#[must_use]
pub fn now() -> Instant {
- #[cfg(feature = "virtual")]
+ #[cfg(all(feature = "virtual", not(wall_clock)))]
{
- tokio::time::Instant::now().into_std()
+ checked_now().unwrap_or_else(|| virtual_time::refuse())
}
- #[cfg(not(feature = "virtual"))]
+ #[cfg(not(all(feature = "virtual", not(wall_clock))))]
{
Instant::now()
}
}
+#[must_use]
+pub fn checked_now() -> Option {
+ #[cfg(all(feature = "virtual", not(wall_clock)))]
+ {
+ if virtual_time::armed() && tokio::runtime::Handle::try_current().is_err() {
+ return None;
+ }
+ Some(tokio::time::Instant::now().into_std())
+ }
+ #[cfg(not(all(feature = "virtual", not(wall_clock))))]
+ {
+ Some(Instant::now())
+ }
+}
+
+#[must_use]
+pub const fn is_routed() -> bool {
+ cfg!(all(feature = "virtual", not(wall_clock)))
+}
+
+#[must_use]
+pub fn elapsed_since_first_reading() -> Option<(bool, Duration)> {
+ // nosemgrep: rust-no-direct-std-sync-import
+ static ORIGIN: std::sync::OnceLock = std::sync::OnceLock::new();
+ let reading = checked_now()?;
+ let origin = *ORIGIN.get_or_init(|| reading);
+ Some(if reading >= origin {
+ (false, reading.saturating_duration_since(origin))
+ } else {
+ (true, origin.saturating_duration_since(reading))
+ })
+}
+
#[must_use]
pub fn system_now() -> SystemTime {
SystemTime::now()
}
+#[cfg(test)]
+pub(crate) fn serially() -> concurrency::sync::MutexGuard<'static, ()> {
+ static SERIAL: concurrency::sync::Mutex<()> = concurrency::sync::Mutex::new(());
+ SERIAL.lock()
+}
+
#[cfg(test)]
mod tests {
+ use super::serially;
use super::{Duration, now, system_now};
#[test]
fn now_is_monotonic() {
+ let _serial = serially();
let first = now();
let second = now();
assert!(second >= first, "the monotonic clock went backwards");
@@ -37,6 +85,7 @@ mod tests {
#[test]
fn now_works_with_no_runtime() {
+ let _serial = serially();
let _ = now();
let _ = system_now();
}
diff --git a/clock/src/virtual_time.rs b/clock/src/virtual_time.rs
new file mode 100644
index 0000000000..8e79725c5e
--- /dev/null
+++ b/clock/src/virtual_time.rs
@@ -0,0 +1,274 @@
+// SPDX-License-Identifier: Apache-2.0
+// Copyright Open Network Fabric Authors
+
+use crate::Duration;
+use std::cell::Cell;
+// nosemgrep: rust-no-direct-std-sync-import
+use std::sync::atomic::{AtomicUsize, Ordering};
+
+static LIVE: AtomicUsize = AtomicUsize::new(0);
+
+const YIELDS: usize = 4;
+
+thread_local! {
+ static IN_WORLD: Cell = const { Cell::new(0) };
+}
+
+#[cfg(not(wall_clock))]
+#[inline]
+#[must_use]
+pub(crate) fn armed() -> bool {
+ LIVE.load(Ordering::Acquire) != 0 && IN_WORLD.with(Cell::get) != 0
+}
+
+thread_local! {
+ #[cfg(all(has_spawn_hook, not(wall_clock)))]
+ static HOOKED: Cell = const { Cell::new(false) };
+}
+
+#[cfg(all(has_spawn_hook, not(wall_clock)))]
+fn inherit_across_spawns() {
+ if !HOOKED.replace(true) {
+ std::thread::add_spawn_hook(|_parent| {
+ let handle = tokio::runtime::Handle::try_current().ok();
+ let in_world = IN_WORLD.with(Cell::get);
+ move || {
+ IN_WORLD.with(|depth| depth.set(in_world));
+ if let Some(handle) = handle {
+ std::mem::forget(Box::leak(Box::new(handle)).enter());
+ }
+ }
+ });
+ }
+}
+
+#[cfg(not(all(has_spawn_hook, not(wall_clock))))]
+fn inherit_across_spawns() {}
+
+#[cfg(not(wall_clock))]
+#[cold]
+#[inline(never)]
+pub(crate) fn refuse() -> ! {
+ panic!(
+ "clock::now() on a thread with no tokio runtime while the virtual clock is paused.\n\
+ \n\
+ This read would have answered from the wall clock, which is a different timeline from the \
+ paused one -- they disagree by however far the test has advanced -- so comparing it \
+ against a deadline taken on the other side is silently wrong, in either direction.\n\
+ \n\
+ A thread spawned with `std::thread` inherits its parent's clock automatically, so reaching \
+ this means this one did not come from there -- DPDK's EAL, or a C library calling \
+ `pthread_create`. Enter `clock::virtual_time::Paused::handle()` on the thread itself; a \
+ guard held by whoever created it does nothing, because tokio's context is thread-local."
+ );
+}
+
+#[derive(Debug)]
+pub struct Paused {
+ runtime: tokio::runtime::Runtime,
+}
+
+impl Paused {
+ #[must_use]
+ pub fn new() -> Self {
+ let runtime = tokio::runtime::Builder::new_current_thread()
+ .enable_time()
+ .start_paused(!cfg!(wall_clock))
+ .build()
+ .unwrap_or_else(|e| panic!("a current-thread runtime with timers does not build: {e}"));
+
+ if !cfg!(wall_clock) {
+ inherit_across_spawns();
+ IN_WORLD.with(|depth| depth.set(depth.get() + 1));
+ LIVE.fetch_add(1, Ordering::AcqRel);
+ }
+
+ Self { runtime }
+ }
+
+ pub fn block_on(&self, future: F) -> F::Output {
+ self.runtime.block_on(future)
+ }
+
+ #[must_use]
+ pub fn handle(&self) -> tokio::runtime::Handle {
+ self.runtime.handle().clone()
+ }
+}
+
+impl Default for Paused {
+ fn default() -> Self {
+ Self::new()
+ }
+}
+
+impl Drop for Paused {
+ fn drop(&mut self) {
+ if !cfg!(wall_clock) {
+ IN_WORLD.with(|depth| depth.set(depth.get().saturating_sub(1)));
+ LIVE.fetch_sub(1, Ordering::Release);
+ }
+ }
+}
+
+pub async fn advance(by: Duration) {
+ #[cfg(not(wall_clock))]
+ {
+ tokio::time::advance(by).await;
+ for _ in 0..YIELDS {
+ tokio::task::yield_now().await;
+ }
+ }
+ #[cfg(wall_clock)]
+ {
+ let _ = YIELDS;
+ tokio::time::sleep(by).await;
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::{Paused, advance};
+ use crate::serially;
+ use crate::{Duration, now};
+ use std::thread;
+
+ const LONG: Duration = if cfg!(wall_clock) {
+ Duration::from_millis(50)
+ } else {
+ Duration::from_hours(1)
+ };
+
+ const NEARLY_LONG: Duration = if cfg!(wall_clock) {
+ Duration::from_millis(20)
+ } else {
+ Duration::from_mins(2)
+ };
+
+ #[test]
+ fn the_clock_moves_when_a_test_says_so() {
+ let _serial = serially();
+ let clock = Paused::new();
+ clock.block_on(async {
+ let before = now();
+ advance(LONG).await;
+ assert!(
+ now().duration_since(before) >= LONG,
+ "the clock was advanced by {LONG:?} and did not follow"
+ );
+ });
+ }
+
+ #[test]
+ fn a_timer_fires_when_the_clock_passes_it() {
+ let _serial = serially();
+ let clock = Paused::new();
+ clock.block_on(async {
+ let deadline = now() + NEARLY_LONG;
+ let waiting = tokio::spawn(async move {
+ while now() < deadline {
+ tokio::task::yield_now().await;
+ }
+ });
+ advance(LONG).await;
+ waiting.await.expect("the waiter panicked");
+ });
+ }
+
+ #[test]
+ #[cfg(all(has_spawn_hook, not(wall_clock)))]
+ fn a_spawned_thread_inherits_the_clock_without_being_told() {
+ let _serial = serially();
+ let clock = Paused::new();
+ let (driver, worker) = clock.block_on(async {
+ advance(LONG).await;
+ let worker = thread::spawn(now)
+ .join()
+ .expect("an inherited read was refused");
+ (now(), worker)
+ });
+ assert_eq!(
+ driver,
+ worker,
+ "a spawned thread read {:?} away from the thread that advanced the clock",
+ driver.saturating_duration_since(worker)
+ );
+ }
+
+ #[test]
+ #[cfg(not(wall_clock))]
+ fn a_thread_in_the_world_with_no_clock_is_refused() {
+ let _serial = serially();
+ let clock = Paused::new();
+ clock.block_on(async { advance(LONG).await });
+
+ let refused = std::panic::catch_unwind(now);
+ let panic = refused.expect_err("a read from the wrong timeline was allowed");
+ let message = panic
+ .downcast_ref::<&'static str>()
+ .copied()
+ .or_else(|| panic.downcast_ref::().map(String::as_str))
+ .unwrap_or("");
+ assert!(
+ message.contains("no tokio runtime"),
+ "the refusal did not explain itself: {message}"
+ );
+ }
+
+ #[test]
+ #[cfg(not(wall_clock))]
+ fn a_thread_outside_the_tree_is_left_alone() {
+ let _serial = serially();
+ let (started, wait_for_start) = std::sync::mpsc::channel();
+ let (finish, wait_to_finish) = std::sync::mpsc::channel();
+
+ let driving = thread::spawn(move || {
+ let clock = Paused::new();
+ clock.block_on(async { advance(LONG).await });
+ started.send(()).expect("the test is waiting");
+ wait_to_finish.recv().expect("the test releases this");
+ });
+ wait_for_start.recv().expect("the driver starts");
+
+ let outsider = thread::spawn(now).join();
+ finish.send(()).expect("the driver is waiting");
+ driving.join().expect("the driver panicked");
+ outsider.expect("a thread outside the world was refused a clock read");
+ }
+
+ #[test]
+ fn a_thread_that_entered_reads_the_same_clock() {
+ let _serial = serially();
+ let clock = Paused::new();
+ let driver = clock.block_on(async {
+ advance(LONG).await;
+ now()
+ });
+
+ let handle = clock.handle();
+ let worker = thread::spawn(move || {
+ let _guard = handle.enter();
+ now()
+ })
+ .join()
+ .expect("an entered read was refused");
+
+ assert!(
+ worker >= driver,
+ "an entered worker read {:?} behind the thread that advanced the clock",
+ driver.saturating_duration_since(worker)
+ );
+ }
+
+ #[test]
+ fn dropping_the_clock_disarms_the_check() {
+ let _serial = serially();
+ {
+ let clock = Paused::new();
+ clock.block_on(async { advance(LONG).await });
+ }
+ thread::spawn(now)
+ .join()
+ .expect("an ordinary read was refused after the clock was dropped");
+ }
+}
diff --git a/concurrency/tests/quiescent_shuttle.rs b/concurrency/tests/quiescent_shuttle.rs
index ab1647514c..31d252e0b3 100644
--- a/concurrency/tests/quiescent_shuttle.rs
+++ b/concurrency/tests/quiescent_shuttle.rs
@@ -22,8 +22,6 @@
#![cfg(not(feature = "loom"))]
-use std::panic::RefUnwindSafe;
-
use bolero::TypeGenerator;
use dataplane_concurrency::sync::Arc;
use dataplane_concurrency::sync::atomic::{AtomicUsize, Ordering};
@@ -155,20 +153,20 @@ fn run_plan(plan: &Plan) {
const TEST_TIME: std::time::Duration = std::time::Duration::from_secs(10);
-fn fuzz_test(
- test: impl Fn(Arg) + RefUnwindSafe,
-) {
- bolero::check!()
- .with_type()
- .cloned()
- .with_test_time(TEST_TIME)
- .for_each(test);
+macro_rules! fuzz_test {
+ ($test:expr) => {{
+ bolero::check!()
+ .with_type()
+ .cloned()
+ .with_test_time(TEST_TIME)
+ .for_each($test);
+ }};
}
#[test]
#[cfg(feature = "shuttle")]
fn protocol_under_shuttle() {
- fuzz_test(|plan: Plan| {
+ fuzz_test!(|plan: Plan| {
let runner = shuttle::Runner::new(
shuttle::scheduler::RandomScheduler::new(1),
dataplane_concurrency::shuttle_config(),
@@ -180,7 +178,7 @@ fn protocol_under_shuttle() {
#[test]
#[cfg(feature = "shuttle")]
fn protocol_under_shuttle_pct() {
- fuzz_test(|plan: Plan| {
+ fuzz_test!(|plan: Plan| {
// PCT requires both threads to actually do atomic ops; if
// either side is effectively empty, shuttle's PCT scheduler
// panics with "test closure did not exercise any concurrency".
@@ -208,5 +206,5 @@ fn protocol_under_shuttle_pct() {
#[test]
#[cfg(not(feature = "shuttle"))]
fn protocol_under_std() {
- fuzz_test(|plan: Plan| run_plan(&plan));
+ fuzz_test!(|plan: Plan| run_plan(&plan));
}
diff --git a/concurrency/tests/scope_property.rs b/concurrency/tests/scope_property.rs
index 9e323e3808..3b9bba5ac6 100644
--- a/concurrency/tests/scope_property.rs
+++ b/concurrency/tests/scope_property.rs
@@ -28,8 +28,6 @@
#![cfg(not(feature = "loom"))]
-use std::panic::RefUnwindSafe;
-
use bolero::TypeGenerator;
use dataplane_concurrency::sync::Arc;
use dataplane_concurrency::sync::atomic::{AtomicUsize, Ordering};
@@ -98,26 +96,26 @@ fn run_plan(plan: &Plan) {
const TEST_TIME: std::time::Duration = std::time::Duration::from_secs(10);
-fn fuzz_test(
- test: impl Fn(Arg) + RefUnwindSafe,
-) {
- bolero::check!()
- .with_type()
- .cloned()
- .with_test_time(TEST_TIME)
- .for_each(test);
+macro_rules! fuzz_test {
+ ($test:expr) => {{
+ bolero::check!()
+ .with_type()
+ .cloned()
+ .with_test_time(TEST_TIME)
+ .for_each($test);
+ }};
}
#[test]
#[cfg(feature = "shuttle")]
fn scope_conservation_under_shuttle() {
- fuzz_test(|plan: Plan| shuttle::check_random(move || run_plan(&plan), 1));
+ fuzz_test!(|plan: Plan| shuttle::check_random(move || run_plan(&plan), 1));
}
#[test]
#[cfg(feature = "shuttle")]
fn scope_conservation_under_shuttle_pct() {
- fuzz_test(|plan: Plan| {
+ fuzz_test!(|plan: Plan| {
// PCT requires every thread to do at least one atomic op;
// skip degenerate shapes that wouldn't exercise concurrency.
let nontrivial = plan
@@ -136,5 +134,5 @@ fn scope_conservation_under_shuttle_pct() {
#[test]
#[cfg(not(feature = "shuttle"))]
fn scope_conservation_under_std() {
- fuzz_test(|plan: Plan| run_plan(&plan));
+ fuzz_test!(|plan: Plan| run_plan(&plan));
}
diff --git a/config/src/external/overlay/completeness.rs b/config/src/external/overlay/completeness.rs
index c74744d0b9..4837da6d9e 100644
--- a/config/src/external/overlay/completeness.rs
+++ b/config/src/external/overlay/completeness.rs
@@ -362,28 +362,27 @@ fn survey_nat(nat: &VpcExposeNat, seen: &mut Observed) {
const CASES: usize = 512;
-fn survey_drawn(seen: &RefCell) {
- let seen = std::panic::AssertUnwindSafe(seen);
- bolero::check!()
- .with_generator(Sequence::default())
- .with_iterations(CASES)
- .for_each(|ops| {
- let overlay = Sequence::fold(ops)
- .overlay()
- .unwrap_or_else(|e| panic!("{ops:?} does not assemble: {e}"));
- survey(&overlay, &mut seen.borrow_mut());
- });
-}
-
-fn census() -> Observed {
- let seen = RefCell::new(Observed::default());
- survey_drawn(&seen);
- seen.into_inner()
+macro_rules! survey_drawn {
+ ($seen:expr) => {{
+ let seen = $seen;
+ let seen = std::panic::AssertUnwindSafe(seen);
+ bolero::check!()
+ .with_generator(Sequence::default())
+ .with_iterations(CASES)
+ .for_each(|ops| {
+ let overlay = Sequence::fold(ops)
+ .overlay()
+ .unwrap_or_else(|e| panic!("{ops:?} does not assemble: {e}"));
+ survey(&overlay, &mut seen.borrow_mut());
+ });
+ }};
}
#[test]
fn every_surveyed_field_is_classified() {
- let seen = census();
+ let seen = RefCell::new(Observed::default());
+ survey_drawn!(&seen);
+ let seen = seen.into_inner();
let surveyed: BTreeSet<&str> = seen.0.keys().copied().collect();
let classified: BTreeSet<&str> = REACH.iter().map(|(field, _)| *field).collect();
@@ -417,7 +416,9 @@ fn every_surveyed_field_is_classified() {
#[test]
fn the_algebra_reaches_what_it_is_recorded_to_reach() {
- let seen = census();
+ let seen = RefCell::new(Observed::default());
+ survey_drawn!(&seen);
+ let seen = seen.into_inner();
for (field, reach) in REACH {
let Some(values) = seen.0.get(field) else {
diff --git a/dataplane/src/packet_processor/fuzz.rs b/dataplane/src/packet_processor/fuzz.rs
index f5090b877e..095c91c641 100644
--- a/dataplane/src/packet_processor/fuzz.rs
+++ b/dataplane/src/packet_processor/fuzz.rs
@@ -648,6 +648,24 @@ pub(crate) struct Pick {
#[cfg(test)]
pub(crate) type Poll = Vec;
+#[cfg(test)]
+pub(crate) fn settled(body: impl FnOnce()) {
+ static RUNTIME: concurrency::sync::LazyLock =
+ concurrency::sync::LazyLock::new(|| {
+ tokio::runtime::Builder::new_current_thread()
+ .enable_all()
+ .build()
+ .expect("build tokio runtime")
+ });
+
+ RUNTIME.block_on(async {
+ body();
+ for _ in 0..4 {
+ tokio::task::yield_now().await;
+ }
+ });
+}
+
#[cfg(test)]
pub(crate) fn run_schedule(
worker: &mut Worker,
@@ -686,6 +704,47 @@ pub(crate) fn run_schedule(
bursts
}
+#[cfg(test)]
+pub(crate) async fn run_schedule_over_time(
+ worker: &mut Worker,
+ loads: &mut [Box],
+ schedule: &[Poll],
+ waits: &[Duration],
+) -> Vec> {
+ let mut bursts = Vec::new();
+ for (nth, poll) in schedule.iter().enumerate() {
+ let mut burst = Vec::new();
+ let mut origin = Vec::new();
+ for pick in poll {
+ if loads.is_empty() {
+ break;
+ }
+ let which = usize::from(pick.load) % loads.len();
+ for _ in 0..pick.take {
+ let Some(packet) = loads[which].next() else {
+ break;
+ };
+ burst.push(packet);
+ origin.push(which);
+ }
+ }
+ if !burst.is_empty() {
+ for (answer, which) in worker.send_batch(burst).iter().zip(&origin) {
+ loads[*which].observe(answer);
+ }
+ bursts.push(origin);
+ }
+ if let Some(wait) = waits.get(nth) {
+ clock::virtual_time::advance(*wait).await;
+ }
+ }
+
+ for load in loads {
+ drive(worker, load.as_mut());
+ }
+ bursts
+}
+
#[cfg(test)]
pub(crate) mod derive {
use super::routed::{Blast, Conversation, Inbound};
@@ -1400,56 +1459,59 @@ mod shapes {
super::assert_within_budget("shapes::Batch", &Batch);
}
- #[tokio::test]
- #[dpdk::with_eal]
- async fn every_shape_leaves_the_pipeline_with_a_verdict() {
+ #[test]
+ fn every_shape_leaves_the_pipeline_with_a_verdict() {
static FORWARDED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
static DROPPED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
static BY_SHAPE: LazyLock<[AtomicU64; Shape::ALL.len()]> =
LazyLock::new(|| std::array::from_fn(|_| AtomicU64::new(0)));
+ let _eal = dpdk::test_support::start_eal();
+
bolero::check!()
.with_max_len(MAX_INPUT_LEN)
.with_generator(Batch)
.for_each(|(exposes, stacks)| {
- let Some(mut fabric) = Fabric::build(exposes) else {
- return;
- };
- let private = exposes.first().and_then(|e| {
- e.ips
- .first()
- .map(lpm::prefix::PrefixWithOptionalPorts::prefix)
- });
-
- for (shape, headers) in stacks {
- let mut headers = headers.clone();
- aim(&mut headers, private);
- let Some(mut packet) = wire(&headers) else {
- continue;
+ settled(|| {
+ let Some(mut fabric) = Fabric::build(exposes) else {
+ return;
};
- BY_SHAPE[*shape as usize].fetch_add(1, Ordering::Relaxed);
+ let private = exposes.first().and_then(|e| {
+ e.ips
+ .first()
+ .map(lpm::prefix::PrefixWithOptionalPorts::prefix)
+ });
- arrive(&mut packet, local());
- let out = fabric.send(packet);
+ for (shape, headers) in stacks {
+ let mut headers = headers.clone();
+ aim(&mut headers, private);
+ let Some(mut packet) = wire(&headers) else {
+ continue;
+ };
+ BY_SHAPE[*shape as usize].fetch_add(1, Ordering::Relaxed);
- match verdict(&out) {
- Verdict::Forwarded { dst_vpcd, .. } => {
- assert_eq!(
- dst_vpcd,
- Some(remote()),
- "forwarded without a destination VPC, on a {shape:?} stack: \
+ arrive(&mut packet, local());
+ let out = fabric.send(packet);
+
+ match verdict(&out) {
+ Verdict::Forwarded { dst_vpcd, .. } => {
+ assert_eq!(
+ dst_vpcd,
+ Some(remote()),
+ "forwarded without a destination VPC, on a {shape:?} stack: \
nothing chose where this packet goes"
- );
- FORWARDED.fetch_add(1, Ordering::Relaxed);
- }
- Verdict::Dropped(_) => {
- DROPPED.fetch_add(1, Ordering::Relaxed);
- }
- Verdict::Delivered { .. } => {
- unreachable!("the overlay slice has no egress stage")
+ );
+ FORWARDED.fetch_add(1, Ordering::Relaxed);
+ }
+ Verdict::Dropped(_) => {
+ DROPPED.fetch_add(1, Ordering::Relaxed);
+ }
+ Verdict::Delivered { .. } => {
+ unreachable!("the overlay slice has no egress stage")
+ }
}
}
- }
+ });
});
let forwarded = FORWARDED.load(Ordering::Relaxed);
@@ -1598,93 +1660,97 @@ mod round_trip {
super::assert_within_budget("round_trip::Batch", &Batch);
}
- #[tokio::test]
- #[dpdk::with_eal]
- async fn a_translated_flow_comes_back_to_where_it_started() {
+ #[test]
+ fn a_translated_flow_comes_back_to_where_it_started() {
static ROUND_TRIPPED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
static NOT_FORWARDED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
+ let _eal = dpdk::test_support::start_eal();
+
bolero::check!()
.with_max_len(MAX_INPUT_LEN)
.with_generator(Batch)
.for_each(|(exposes, flows)| {
- let Some(mut fabric) = Fabric::build(exposes) else {
- return;
- };
- let privates = private_addresses(exposes);
- if privates.is_empty() {
- return;
- }
-
- for flow in flows {
- let prefix = privates[usize::from(flow.prefix) % privates.len()];
- let src = match prefix.as_address() {
- IpAddr::V4(a) => {
- let mut o = a.octets();
- o[3] = o[3].wrapping_add(flow.host % 8);
- IpAddr::V4(Ipv4Addr::from(o))
- }
- IpAddr::V6(a) => {
- let mut o = a.octets();
- o[15] = o[15].wrapping_add(flow.host % 8);
- IpAddr::V6(Ipv6Addr::from(o))
- }
+ settled(|| {
+ let Some(mut fabric) = Fabric::build(exposes) else {
+ return;
};
- let dst = peer(src);
+ let privates = private_addresses(exposes);
+ if privates.is_empty() {
+ return;
+ }
- let Some(mut request) = udp(src, dst, flow.sport, flow.dport) else {
- continue;
- };
- arrive(&mut request, local());
- let out = fabric.send(request);
-
- let Verdict::Forwarded {
- src: public_src,
- dst: reached,
- ..
- } = verdict(&out)
- else {
- NOT_FORWARDED.fetch_add(1, Ordering::Relaxed);
- continue;
- };
- let (Some(public_src), Some(reached)) = (public_src, reached) else {
- continue;
- };
- let public_port = out
- .transport_src_port()
- .unwrap_or_else(|| unreachable!("a udp packet has a source port"))
- .get();
+ for flow in flows {
+ let prefix = privates[usize::from(flow.prefix) % privates.len()];
+ let src = match prefix.as_address() {
+ IpAddr::V4(a) => {
+ let mut o = a.octets();
+ o[3] = o[3].wrapping_add(flow.host % 8);
+ IpAddr::V4(Ipv4Addr::from(o))
+ }
+ IpAddr::V6(a) => {
+ let mut o = a.octets();
+ o[15] = o[15].wrapping_add(flow.host % 8);
+ IpAddr::V6(Ipv6Addr::from(o))
+ }
+ };
+ let dst = peer(src);
- let Some(mut reply) = udp(reached, public_src, flow.dport, public_port) else {
- continue;
- };
- arrive(&mut reply, remote());
- let back = fabric.send(reply);
+ let Some(mut request) = udp(src, dst, flow.sport, flow.dport) else {
+ continue;
+ };
+ arrive(&mut request, local());
+ let out = fabric.send(request);
+
+ let Verdict::Forwarded {
+ src: public_src,
+ dst: reached,
+ ..
+ } = verdict(&out)
+ else {
+ NOT_FORWARDED.fetch_add(1, Ordering::Relaxed);
+ continue;
+ };
+ let (Some(public_src), Some(reached)) = (public_src, reached) else {
+ continue;
+ };
+ let public_port = out
+ .transport_src_port()
+ .unwrap_or_else(|| unreachable!("a udp packet has a source port"))
+ .get();
- match verdict(&back) {
- Verdict::Forwarded { src: s, dst: d, .. } => {
- assert_eq!(
- d,
- Some(src),
- "the reply did not come back to the host that sent the request"
- );
- assert_eq!(s, Some(dst), "the reply's source was rewritten");
- assert_eq!(
- back.transport_dst_port().map(std::num::NonZero::get),
- Some(flow.sport),
- "the reply did not get the original source port back"
- );
- ROUND_TRIPPED.fetch_add(1, Ordering::Relaxed);
- }
- Verdict::Dropped(reason) => panic!(
- "the reply of a forwarded flow was dropped: {reason:?} \
+ let Some(mut reply) = udp(reached, public_src, flow.dport, public_port)
+ else {
+ continue;
+ };
+ arrive(&mut reply, remote());
+ let back = fabric.send(reply);
+
+ match verdict(&back) {
+ Verdict::Forwarded { src: s, dst: d, .. } => {
+ assert_eq!(
+ d,
+ Some(src),
+ "the reply did not come back to the host that sent the request"
+ );
+ assert_eq!(s, Some(dst), "the reply's source was rewritten");
+ assert_eq!(
+ back.transport_dst_port().map(std::num::NonZero::get),
+ Some(flow.sport),
+ "the reply did not get the original source port back"
+ );
+ ROUND_TRIPPED.fetch_add(1, Ordering::Relaxed);
+ }
+ Verdict::Dropped(reason) => panic!(
+ "the reply of a forwarded flow was dropped: {reason:?} \
(request {src} -> {dst} became {public_src}:{public_port})"
- ),
- Verdict::Delivered { .. } => {
- unreachable!("the overlay slice has no egress stage")
+ ),
+ Verdict::Delivered { .. } => {
+ unreachable!("the overlay slice has no egress stage")
+ }
}
}
- }
+ });
});
let round_tripped = ROUND_TRIPPED.load(Ordering::Relaxed);
@@ -1894,80 +1960,84 @@ mod acl {
super::assert_within_budget("acl::Batch", &Batch);
}
- #[tokio::test]
- #[dpdk::with_eal]
- async fn the_acl_verdict_follows_the_protocol_the_packet_carries() {
+ #[test]
+ fn the_acl_verdict_follows_the_protocol_the_packet_carries() {
static DENIED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
static PERMITTED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
static BEHIND_EXT: LazyLock = LazyLock::new(|| AtomicU64::new(0));
static PERMITTED_OUT: LazyLock = LazyLock::new(|| AtomicU64::new(0));
static DENIED_BY_ACL: LazyLock = LazyLock::new(|| AtomicU64::new(0));
+ let _eal = dpdk::test_support::start_eal();
+
bolero::check!()
.with_max_len(MAX_INPUT_LEN)
.with_generator(Batch)
.for_each(|(exposes, default_allow, rule_proto, packets)| {
- let default = if *default_allow {
- AclAction::Allow
- } else {
- AclAction::Deny
- };
- let rule = rule_proto.as_match();
- let Some(mut fabric) =
- Fabric::build_with_acl(exposes, Some(&peering_acl(default, rule)))
- else {
- return;
- };
- let Some(private) = exposes
- .iter()
- .flat_map(|e| e.ips.iter().map(PrefixWithOptionalPorts::prefix))
- .next()
- .map(|p: Prefix| p.as_address())
- else {
- return;
- };
- let v6 = private.is_ipv6();
- let dst = peer(private);
-
- for (spec, headers) in packets {
- let Some(mut packet) = wire(headers, *spec, private, dst) else {
- continue;
+ settled(|| {
+ let default = if *default_allow {
+ AclAction::Allow
+ } else {
+ AclAction::Deny
+ };
+ let rule = rule_proto.as_match();
+ let Some(mut fabric) =
+ Fabric::build_with_acl(exposes, Some(&peering_acl(default, rule)))
+ else {
+ return;
+ };
+ let Some(private) = exposes
+ .iter()
+ .flat_map(|e| e.ips.iter().map(PrefixWithOptionalPorts::prefix))
+ .next()
+ .map(|p: Prefix| p.as_address())
+ else {
+ return;
};
- arrive(&mut packet, local());
- let out = fabric.send(packet);
+ let v6 = private.is_ipv6();
+ let dst = peer(private);
- let permitted = rule_matches(rule, carried(spec.proto, v6)) != *default_allow;
- if spec.behind_extension {
- BEHIND_EXT.fetch_add(1, Ordering::Relaxed);
- }
+ for (spec, headers) in packets {
+ let Some(mut packet) = wire(headers, *spec, private, dst) else {
+ continue;
+ };
+ arrive(&mut packet, local());
+ let out = fabric.send(packet);
- let seen = verdict(&out);
- let acl_dropped = seen == Verdict::Dropped(DoneReason::AclDropped);
- let forwarded = matches!(seen, Verdict::Forwarded { .. });
- if permitted {
- assert!(
- !acl_dropped,
- "the acl dropped a {:?} packet it permits (rule={rule:?} \
- default={default:?} behind_extension={})",
- spec.proto, spec.behind_extension
- );
- PERMITTED.fetch_add(1, Ordering::Relaxed);
- if forwarded {
- PERMITTED_OUT.fetch_add(1, Ordering::Relaxed);
+ let permitted =
+ rule_matches(rule, carried(spec.proto, v6)) != *default_allow;
+ if spec.behind_extension {
+ BEHIND_EXT.fetch_add(1, Ordering::Relaxed);
}
- } else {
- assert!(
- !forwarded,
- "a {:?} packet the acl denies was forwarded (rule={rule:?} \
+
+ let seen = verdict(&out);
+ let acl_dropped = seen == Verdict::Dropped(DoneReason::AclDropped);
+ let forwarded = matches!(seen, Verdict::Forwarded { .. });
+ if permitted {
+ assert!(
+ !acl_dropped,
+ "the acl dropped a {:?} packet it permits (rule={rule:?} \
default={default:?} behind_extension={})",
- spec.proto, spec.behind_extension
- );
- DENIED.fetch_add(1, Ordering::Relaxed);
- if acl_dropped {
- DENIED_BY_ACL.fetch_add(1, Ordering::Relaxed);
+ spec.proto, spec.behind_extension
+ );
+ PERMITTED.fetch_add(1, Ordering::Relaxed);
+ if forwarded {
+ PERMITTED_OUT.fetch_add(1, Ordering::Relaxed);
+ }
+ } else {
+ assert!(
+ !forwarded,
+ "a {:?} packet the acl denies was forwarded (rule={rule:?} \
+ default={default:?} behind_extension={})",
+ spec.proto, spec.behind_extension
+ );
+ DENIED.fetch_add(1, Ordering::Relaxed);
+ if acl_dropped {
+ DENIED_BY_ACL.fetch_add(1, Ordering::Relaxed);
+ }
}
}
- }
+ });
});
let (permitted, permitted_out, denied, denied_by_acl, behind) = (
@@ -2227,71 +2297,74 @@ mod port_forward {
true
}
- #[tokio::test]
- #[dpdk::with_eal]
- async fn a_forwarded_port_reaches_the_host_behind_it() {
+ #[test]
+ fn a_forwarded_port_reaches_the_host_behind_it() {
static FORWARDED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
static ANSWERED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
static REFUSED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
+ let _eal = dpdk::test_support::start_eal();
+
bolero::check!()
.with_max_len(MAX_INPUT_LEN)
.with_generator(Reaches)
.for_each(|reaches| {
- let Some(mut fabric) = Fabric::routed(&[expose()], None) else {
- unreachable!("the port-forwarding fixture does not configure")
- };
-
- for reach in reaches {
- let external: IpAddr = format!("172.16.5.{}", reach.host)
- .parse()
- .unwrap_or_else(|_| unreachable!());
- let dport = if reach.past_the_range {
- EXTERNAL_PORT + PORTS + (reach.port % PORTS)
- } else {
- EXTERNAL_PORT + reach.port
+ settled(|| {
+ let Some(mut fabric) = Fabric::routed(&[expose()], None) else {
+ unreachable!("the port-forwarding fixture does not configure")
};
- let Some(inbound) = udp(outside(), external, reach.src_port, dport) else {
- continue;
- };
- let out = fabric.send(tunnelled_from(vni(REMOTE_VNI), &inbound));
+ for reach in reaches {
+ let external: IpAddr = format!("172.16.5.{}", reach.host)
+ .parse()
+ .unwrap_or_else(|_| unreachable!());
+ let dport = if reach.past_the_range {
+ EXTERNAL_PORT + PORTS + (reach.port % PORTS)
+ } else {
+ EXTERNAL_PORT + reach.port
+ };
- if reach.past_the_range {
- assert!(
- !matches!(verdict(&out), Verdict::Delivered { .. }),
- "a packet to {external}:{dport}, past the declared range, was \
+ let Some(inbound) = udp(outside(), external, reach.src_port, dport) else {
+ continue;
+ };
+ let out = fabric.send(tunnelled_from(vni(REMOTE_VNI), &inbound));
+
+ if reach.past_the_range {
+ assert!(
+ !matches!(verdict(&out), Verdict::Delivered { .. }),
+ "a packet to {external}:{dport}, past the declared range, was \
forwarded anyway"
- );
- REFUSED.fetch_add(1, Ordering::Relaxed);
- continue;
- }
+ );
+ REFUSED.fetch_add(1, Ordering::Relaxed);
+ continue;
+ }
- assert!(
- matches!(verdict(&out), Verdict::Delivered { .. }),
- "a packet to the declared {external}:{dport} was not forwarded: {:?}",
- verdict(&out)
- );
- let arrived = inside(&out).expect("a forwarded packet was not tunnelled");
- let expected_host: IpAddr = format!("10.0.5.{}", reach.host)
- .parse()
- .unwrap_or_else(|_| unreachable!());
- assert_eq!(
- arrived.ip_destination(),
- Some(expected_host),
- "{external}:{dport} reached the wrong host"
- );
- assert_eq!(
- arrived.transport_dst_port().map(std::num::NonZero::get),
- Some(INTERNAL_PORT + reach.port),
- "{external}:{dport} reached the right host on the wrong port"
- );
- FORWARDED.fetch_add(1, Ordering::Relaxed);
+ assert!(
+ matches!(verdict(&out), Verdict::Delivered { .. }),
+ "a packet to the declared {external}:{dport} was not forwarded: {:?}",
+ verdict(&out)
+ );
+ let arrived = inside(&out).expect("a forwarded packet was not tunnelled");
+ let expected_host: IpAddr = format!("10.0.5.{}", reach.host)
+ .parse()
+ .unwrap_or_else(|_| unreachable!());
+ assert_eq!(
+ arrived.ip_destination(),
+ Some(expected_host),
+ "{external}:{dport} reached the wrong host"
+ );
+ assert_eq!(
+ arrived.transport_dst_port().map(std::num::NonZero::get),
+ Some(INTERNAL_PORT + reach.port),
+ "{external}:{dport} reached the right host on the wrong port"
+ );
+ FORWARDED.fetch_add(1, Ordering::Relaxed);
- if answers(&mut fabric, expected_host, *reach, external, dport) {
- ANSWERED.fetch_add(1, Ordering::Relaxed);
+ if answers(&mut fabric, expected_host, *reach, external, dport) {
+ ANSWERED.fetch_add(1, Ordering::Relaxed);
+ }
}
- }
+ });
});
let (forwarded, answered, refused) = (
@@ -2380,71 +2453,74 @@ mod interleaved {
super::assert_within_budget("interleaved::Interleaving", &Interleaving);
}
- #[tokio::test]
- #[dpdk::with_eal]
- async fn interleaved_traffic_is_each_satisfied() {
+ #[test]
+ fn interleaved_traffic_is_each_satisfied() {
static CHECKED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
static ABANDONED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
static MIXED_LOADS: LazyLock = LazyLock::new(|| AtomicU64::new(0));
static MIXED_KINDS: LazyLock = LazyLock::new(|| AtomicU64::new(0));
+ let _eal = dpdk::test_support::start_eal();
+
bolero::check!()
.with_max_len(MAX_INPUT_LEN)
.with_generator(Interleaving)
.for_each(|(senders, schedule)| {
- let Some(mut fabric) = Fabric::routed(&exposes(), None) else {
- return;
- };
-
- let dst: IpAddr = "3.3.3.1".parse().unwrap_or_else(|_| unreachable!());
- let mut kinds = Vec::new();
- let mut loads: Vec> = Vec::new();
- for (i, sender) in senders.iter().enumerate() {
- let Ok(src) = format!("1.1.{i}.{}", sender.host).parse::() else {
- continue;
+ settled(|| {
+ let Some(mut fabric) = Fabric::routed(&exposes(), None) else {
+ return;
};
- kinds.push(sender.kind);
- loads.push(match sender.kind {
- Kind::Conversation => Box::new(Conversation::new(
- Path::fixture(),
- src,
- dst,
- sender.sport,
- sender.dport,
- )),
- Kind::Blast => Box::new(Blast::new(
- Path::fixture(),
- src,
- dst,
- sender.sport,
- sender.dport,
- sender.count,
- )) as Box,
- });
- }
- for burst in run_schedule(fabric.worker(), &mut loads, schedule) {
- let mut loads_in: Vec = burst.clone();
- loads_in.sort_unstable();
- loads_in.dedup();
- if loads_in.len() > 1 {
- MIXED_LOADS.fetch_add(1, Ordering::Relaxed);
+ let dst: IpAddr = "3.3.3.1".parse().unwrap_or_else(|_| unreachable!());
+ let mut kinds = Vec::new();
+ let mut loads: Vec> = Vec::new();
+ for (i, sender) in senders.iter().enumerate() {
+ let Ok(src) = format!("1.1.{i}.{}", sender.host).parse::() else {
+ continue;
+ };
+ kinds.push(sender.kind);
+ loads.push(match sender.kind {
+ Kind::Conversation => Box::new(Conversation::new(
+ Path::fixture(),
+ src,
+ dst,
+ sender.sport,
+ sender.dport,
+ )),
+ Kind::Blast => Box::new(Blast::new(
+ Path::fixture(),
+ src,
+ dst,
+ sender.sport,
+ sender.dport,
+ sender.count,
+ )) as Box,
+ });
}
- let mut kinds_in: Vec = burst.iter().map(|i| kinds[*i]).collect();
- kinds_in.sort_unstable_by_key(|k| format!("{k:?}"));
- kinds_in.dedup();
- if kinds_in.len() > 1 {
- MIXED_KINDS.fetch_add(1, Ordering::Relaxed);
+
+ for burst in run_schedule(fabric.worker(), &mut loads, schedule) {
+ let mut loads_in: Vec = burst.clone();
+ loads_in.sort_unstable();
+ loads_in.dedup();
+ if loads_in.len() > 1 {
+ MIXED_LOADS.fetch_add(1, Ordering::Relaxed);
+ }
+ let mut kinds_in: Vec = burst.iter().map(|i| kinds[*i]).collect();
+ kinds_in.sort_unstable_by_key(|k| format!("{k:?}"));
+ kinds_in.dedup();
+ if kinds_in.len() > 1 {
+ MIXED_KINDS.fetch_add(1, Ordering::Relaxed);
+ }
}
- }
- for load in &loads {
- if load.checked() {
- CHECKED.fetch_add(1, Ordering::Relaxed);
- } else {
- ABANDONED.fetch_add(1, Ordering::Relaxed);
+ for load in &loads {
+ if load.checked() {
+ CHECKED.fetch_add(1, Ordering::Relaxed);
+ } else {
+ ABANDONED.fetch_add(1, Ordering::Relaxed);
+ }
}
- }
+ });
});
let (checked, abandoned, mixed_loads, mixed_kinds) = (
@@ -2564,9 +2640,8 @@ mod offers {
super::assert_within_budget("offers::Offered", &Offered);
}
- #[tokio::test]
- #[dpdk::with_eal]
- async fn a_configuration_carries_everything_it_offers() {
+ #[test]
+ fn a_configuration_carries_everything_it_offers() {
static CHECKED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
static ABANDONED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
static DERIVED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
@@ -2574,43 +2649,47 @@ mod offers {
static INBOUND: LazyLock = LazyLock::new(|| AtomicU64::new(0));
static OUTBOUND: LazyLock = LazyLock::new(|| AtomicU64::new(0));
+ let _eal = dpdk::test_support::start_eal();
+
let overlay = overlay();
bolero::check!()
.with_max_len(MAX_INPUT_LEN)
.with_generator(Offered)
.for_each(|(vary, schedule)| {
- let mut fabric = Fabric::routed_over_validated(
- &overlay,
- topology(&[vni(LOCAL_VNI), vni(REMOTE_VNI)]),
- );
+ settled(|| {
+ let mut fabric = Fabric::routed_over_validated(
+ &overlay,
+ topology(&[vni(LOCAL_VNI), vni(REMOTE_VNI)]),
+ );
- let mut loads = loads_for(&overlay, vary);
- DERIVED.fetch_add(loads.len() as u64, Ordering::Relaxed);
- for load in &loads {
- if load.describe().starts_with("[inbound") {
- INBOUND.fetch_add(1, Ordering::Relaxed);
- } else {
- OUTBOUND.fetch_add(1, Ordering::Relaxed);
+ let mut loads = loads_for(&overlay, vary);
+ DERIVED.fetch_add(loads.len() as u64, Ordering::Relaxed);
+ for load in &loads {
+ if load.describe().starts_with("[inbound") {
+ INBOUND.fetch_add(1, Ordering::Relaxed);
+ } else {
+ OUTBOUND.fetch_add(1, Ordering::Relaxed);
+ }
}
- }
- for burst in run_schedule(fabric.worker(), &mut loads, schedule) {
- let mut seen = burst.clone();
- seen.sort_unstable();
- seen.dedup();
- if seen.len() > 1 {
- MIXED.fetch_add(1, Ordering::Relaxed);
+ for burst in run_schedule(fabric.worker(), &mut loads, schedule) {
+ let mut seen = burst.clone();
+ seen.sort_unstable();
+ seen.dedup();
+ if seen.len() > 1 {
+ MIXED.fetch_add(1, Ordering::Relaxed);
+ }
}
- }
- for load in &loads {
- if load.checked() {
- CHECKED.fetch_add(1, Ordering::Relaxed);
- } else {
- ABANDONED.fetch_add(1, Ordering::Relaxed);
+ for load in &loads {
+ if load.checked() {
+ CHECKED.fetch_add(1, Ordering::Relaxed);
+ } else {
+ ABANDONED.fetch_add(1, Ordering::Relaxed);
+ }
}
- }
+ });
});
let (checked, abandoned, derived, mixed) = (
@@ -2669,6 +2748,10 @@ mod generated {
static BY_FLOW: LazyLock = LazyLock::new(|| AtomicU64::new(0));
static EXCEPTING: LazyLock = LazyLock::new(|| AtomicU64::new(0));
+ static AGED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
+ static AGED_MILLIS: LazyLock = LazyLock::new(|| AtomicU64::new(0));
+ static SURVIVED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
+
fn report_and_assert_coverage() {
let (checked, derived, mixed) = (
CHECKED.load(Ordering::Relaxed),
@@ -2794,94 +2877,116 @@ mod generated {
}
}
- #[tokio::test]
- #[dpdk::with_eal]
- async fn a_generated_configuration_carries_its_own_traffic() {
+ #[test]
+ fn a_generated_configuration_carries_its_own_traffic() {
+ let _eal = dpdk::test_support::start_eal();
bolero::check!()
.with_max_len(MAX_INPUT_LEN)
.with_generator(Generated)
.for_each(|(ops, vary, schedule)| {
- let draft = Sequence::fold(ops);
- let overlay = draft
- .overlay()
- .unwrap_or_else(|e| panic!("{ops:?} does not assemble: {e}"));
- let validated = overlay
- .validate()
- .unwrap_or_else(|e| panic!("{ops:?} does not validate: {e}"));
+ settled(|| {
+ let draft = Sequence::fold(ops);
+ let overlay = draft
+ .overlay()
+ .unwrap_or_else(|e| panic!("{ops:?} does not assemble: {e}"));
+ let validated = overlay
+ .validate()
+ .unwrap_or_else(|e| panic!("{ops:?} does not validate: {e}"));
+
+ let vnis: Vec = validated
+ .vpc_table()
+ .values()
+ .map(config::external::overlay::vpc::ValidatedVpc::vni)
+ .collect();
+ if vnis.is_empty() {
+ return;
+ }
+ if validated.vpc_table().peerings().next().is_some() {
+ PEERED.fetch_add(1, Ordering::Relaxed);
+ }
+ if vnis.len() > 2 {
+ MULTI.fetch_add(1, Ordering::Relaxed);
+ }
- let vnis: Vec = validated
- .vpc_table()
- .values()
- .map(config::external::overlay::vpc::ValidatedVpc::vni)
- .collect();
- if vnis.is_empty() {
- return;
- }
- if validated.vpc_table().peerings().next().is_some() {
- PEERED.fetch_add(1, Ordering::Relaxed);
- }
- if vnis.len() > 2 {
- MULTI.fetch_add(1, Ordering::Relaxed);
- }
+ let mut fabric = Fabric::routed_over_validated(&validated, topology(&vnis));
- let mut fabric = Fabric::routed_over_validated(&validated, topology(&vnis));
-
- let (permitting, by_flow, excepting) = (Cell::new(0), Cell::new(0), Cell::new(0));
- let mut loads = loads_where(
- &validated,
- vary,
- &carried_counting(&draft, &permitting, &by_flow, &excepting),
- );
- PERMITTING.fetch_add(permitting.get(), Ordering::Relaxed);
- BY_FLOW.fetch_add(by_flow.get(), Ordering::Relaxed);
- EXCEPTING.fetch_add(excepting.get(), Ordering::Relaxed);
- DERIVED.fetch_add(loads.len() as u64, Ordering::Relaxed);
- for load in &loads {
- if load.describe().starts_with("[inbound") {
- INBOUND.fetch_add(1, Ordering::Relaxed);
+ let (permitting, by_flow, excepting) =
+ (Cell::new(0), Cell::new(0), Cell::new(0));
+ let mut loads = loads_where(
+ &validated,
+ vary,
+ &carried_counting(&draft, &permitting, &by_flow, &excepting),
+ );
+ PERMITTING.fetch_add(permitting.get(), Ordering::Relaxed);
+ BY_FLOW.fetch_add(by_flow.get(), Ordering::Relaxed);
+ EXCEPTING.fetch_add(excepting.get(), Ordering::Relaxed);
+ DERIVED.fetch_add(loads.len() as u64, Ordering::Relaxed);
+ for load in &loads {
+ if load.describe().starts_with("[inbound") {
+ INBOUND.fetch_add(1, Ordering::Relaxed);
+ }
}
- }
- for burst in run_schedule(fabric.worker(), &mut loads, schedule) {
- let mut seen = burst.clone();
- seen.sort_unstable();
- seen.dedup();
- if seen.len() > 1 {
- MIXED.fetch_add(1, Ordering::Relaxed);
+ for burst in run_schedule(fabric.worker(), &mut loads, schedule) {
+ let mut seen = burst.clone();
+ seen.sort_unstable();
+ seen.dedup();
+ if seen.len() > 1 {
+ MIXED.fetch_add(1, Ordering::Relaxed);
+ }
}
- }
- for load in &loads {
- assert!(
- load.checked(),
- "a load derived from the configuration did not complete: {}",
- load.describe()
- );
- CHECKED.fetch_add(1, Ordering::Relaxed);
- }
+ for load in &loads {
+ assert!(
+ load.checked(),
+ "a load derived from the configuration did not complete: {}",
+ load.describe()
+ );
+ CHECKED.fetch_add(1, Ordering::Relaxed);
+ }
+ });
});
report_and_assert_coverage();
}
- #[tokio::test]
- #[dpdk::with_eal]
- async fn a_configuration_carries_nothing_it_denies() {
- static SENT: LazyLock = LazyLock::new(|| AtomicU64::new(0));
- static BY_ACL: LazyLock = LazyLock::new(|| AtomicU64::new(0));
- static NARROWED: LazyLock = LazyLock::new(|| AtomicU64::new(0));
- static CONFIGS: LazyLock = LazyLock::new(|| AtomicU64::new(0));
+ pub(super) struct OverTime;
+
+ impl ValueGenerator for OverTime {
+ type Output = (Vec, Vec, Vec, Vec);
+
+ fn generate(&self, driver: &mut D) -> Option {
+ let (ops, vary, schedule) = Generated.generate(driver)?;
+ let cap =
+ (Masquerade::MASQUERADE_ONEWAY_TIMEOUT / 2).as_millis() / (POLLS as u128).max(1);
+ let cap = u64::try_from(cap).unwrap_or(u64::MAX).max(1);
+ let waits = (0..POLLS)
+ .map(|_| {
+ Some(Duration::from_millis(
+ driver.gen_u64(Included(&0), Included(&cap))?,
+ ))
+ })
+ .collect::