1use std::cell::RefCell;
4use std::rc::{Rc, Weak};
5use std::sync::atomic::Ordering;
6use std::task::{Context, Poll};
7use std::time::{Duration, Instant};
8
9use io_uring::{EnterFlags, IoUring, opcode, types};
10
11use crate::metrics::Metrics;
12use crate::park::{FUTEX_BITSET_MATCH_ANY, FUTEX2_PRIVATE, FUTEX2_SIZE_U32, PARKED, RUNNING, Unpark};
13use crate::shared::{Cqe, Op, Shared, Task};
14use crate::{Error, timer, udp};
15
16const SQ_ENTRIES: u32 = 256;
20
21const CQ_ENTRIES: u32 = 2048;
34
35const TEARDOWN_CQE_BATCH: usize = 64;
37
38const TEARDOWN_EINTR_RETRIES: usize = 8;
40
41const TEARDOWN_TIMEOUT: Duration = Duration::from_millis(3200);
43
44#[derive(Debug, Default)]
49#[non_exhaustive]
50pub struct Config {
51 pub metrics: Metrics,
60}
61
62impl Clone for Config {
63 fn clone(&self) -> Self {
64 Self {
67 metrics: Metrics::default(),
68 }
69 }
70}
71
72pub struct Worker {
90 shared: Rc<Shared>,
91 tasks: kio::Tasks<Task>,
92 park: kio::Park,
93 cqes: Vec<Cqe>,
95 futex_armed: bool,
97}
98
99impl Worker {
100 pub fn new(config: Config) -> Result<Self, Error> {
106 let Config { metrics } = config;
107 let metrics = metrics.counters().clone();
108 let ring = IoUring::builder()
109 .setup_single_issuer()
110 .setup_defer_taskrun()
111 .setup_coop_taskrun()
112 .setup_cqsize(CQ_ENTRIES)
113 .build(SQ_ENTRIES)
114 .map_err(|err| match err.raw_os_error() {
115 Some(libc::ENOSYS) | Some(libc::EPERM) | Some(libc::EACCES) | Some(libc::EINVAL) => {
119 Error::Unsupported(format!(
120 "io_uring is unavailable ({err}); kernel {} (Linux 6.12+ required, and container seccomp \
121 policies such as Docker's default commonly block io_uring)",
122 kernel_release()
123 ))
124 }
125 _ => Error::ring(err),
126 })?;
127
128 if !ring.params().is_feature_min_timeout() {
131 return Err(Error::Unsupported(format!(
132 "kernel {} is too old: moq-uring requires Linux 6.12+ (io_uring MIN_TIMEOUT feature missing)",
133 kernel_release()
134 )));
135 }
136
137 Ok(Self {
138 shared: Rc::new(Shared {
139 ring: RefCell::new(ring),
140 ops: RefCell::new(slab::Slab::new()),
141 timers: Rc::new(RefCell::new(timer::Heap::new(metrics.clone()))),
142 spawns: RefCell::new(Vec::new()),
143 unpark: Unpark::new(metrics.clone()),
144 metrics,
145 next_bgid: std::cell::Cell::new(0),
146 stopped: std::cell::Cell::new(false),
147 spill: RefCell::new(std::collections::VecDeque::new()),
148 }),
149 tasks: kio::Tasks::new(),
150 park: kio::Park::default(),
151 cqes: Vec::new(),
152 futex_armed: false,
153 })
154 }
155
156 pub fn handle(&self) -> Handle {
158 Handle {
159 shared: self.shared.clone(),
160 }
161 }
162
163 pub fn block_on<F: Future>(&mut self, future: F) -> Result<F::Output, Error> {
169 let mut future = std::pin::pin!(future);
170 let waker = self.shared.unpark.waker();
171 loop {
172 let spawns = std::mem::take(&mut *self.shared.spawns.borrow_mut());
174 for task in spawns {
175 self.tasks.push(task);
176 }
177
178 let cx = Context::from_waker(&waker);
179 let waiter = self.park.hold(&cx);
180 if let Poll::Ready(value) = waiter.poll_future(future.as_mut()) {
181 return Ok(value);
182 }
183 let _ = self.tasks.poll(waiter);
186
187 self.shared.timers.borrow_mut().fire(Instant::now());
188 self.pump()?;
189 self.maybe_park()?;
190 }
191 }
192
193 fn pump(&mut self) -> Result<(), Error> {
195 self.pump_inner(None)
196 }
197
198 fn pump_until(&mut self, deadline: Instant) -> Result<(), Error> {
200 self.pump_inner(Some(deadline))
201 }
202
203 fn pump_inner(&mut self, deadline: Option<Instant>) -> Result<(), Error> {
204 self.submit()?;
205 loop {
206 if deadline.is_some_and(|deadline| Instant::now() >= deadline) {
207 return Ok(());
208 }
209 self.cqes.clear();
213 {
214 let mut ring = self.shared.ring.borrow_mut();
215 let mut spill = self.shared.spill.borrow_mut();
216 let limit = deadline.map_or(usize::MAX, |_| TEARDOWN_CQE_BATCH);
217 let spilled = spill.len().min(limit);
218 self.cqes.extend(spill.drain(..spilled));
219 self.cqes
220 .extend(ring.completion().take(limit - spilled).map(|entry| Cqe {
221 user_data: entry.user_data(),
222 result: entry.result(),
223 flags: entry.flags(),
224 }));
225 }
226 self.shared.metrics.completions.add(self.cqes.len() as u64);
227 if self.cqes.is_empty() || !self.dispatch_batch(deadline, Instant::now) {
228 return Ok(());
229 }
230 }
231 }
232
233 fn dispatch_batch(&mut self, deadline: Option<Instant>, mut now: impl FnMut() -> Instant) -> bool {
236 for index in 0..self.cqes.len() {
237 if deadline.is_some_and(|deadline| now() >= deadline) {
238 return false;
241 }
242 let cqe = self.cqes[index];
243 self.dispatch(cqe);
244 }
245 true
246 }
247
248 fn submit(&mut self) -> Result<(), Error> {
249 let mut ring = self.shared.ring.borrow_mut();
250 if ring.submission().is_empty() {
251 return Ok(());
252 }
253 self.shared.metrics.enters.add(1);
254 match ring.submit() {
255 Ok(count) => {
257 self.shared.metrics.submissions.add(count as u64);
258 Ok(())
259 }
260 Err(err) if err.raw_os_error() == Some(libc::EINTR) => Ok(()),
263 Err(err) if err.raw_os_error() == Some(libc::EBUSY) => Ok(()),
265 Err(err) => Err(err.into()),
266 }
267 }
268
269 fn submit_teardown(&mut self) -> Result<(), Error> {
271 let mut ring = self.shared.ring.borrow_mut();
272 let mut interruptions = 0;
273 loop {
274 if ring.submission().is_empty() {
275 return Ok(());
276 }
277 self.shared.metrics.enters.add(1);
278 match ring.submit() {
279 Ok(0) => {
282 return Err(std::io::Error::other("io_uring teardown submission made no progress").into());
283 }
284 Ok(count) => self.shared.metrics.submissions.add(count as u64),
285 Err(err) => retry_teardown_submit(&mut interruptions, err)?,
286 }
287 }
288 }
289
290 fn drain_teardown(&mut self, deadline: Instant) {
292 let submission_failed = self.submit_teardown().is_err();
295 if !submission_failed {
296 loop {
297 if self.shared.ops.borrow().is_empty() {
298 return;
299 }
300 if Instant::now() >= deadline || self.pump_until(deadline).is_err() {
301 break;
302 }
303 let Some(remaining) = deadline.checked_duration_since(Instant::now()) else {
304 break;
305 };
306 let ring = self.shared.ring.borrow_mut();
307 let wait = remaining.min(std::time::Duration::from_millis(50));
308 let ts = types::Timespec::from(wait);
309 let args = types::SubmitArgs::new().timespec(&ts);
310 self.shared.metrics.enters.add(1);
311 let _ = ring.submitter().submit_with_args(1, &args);
312 }
313 }
314 if !self.shared.ops.borrow().is_empty() {
315 tracing::error!("dropping an io_uring worker with operations stuck in flight; leaking them");
318 std::mem::forget(std::mem::take(&mut *self.shared.ops.borrow_mut()));
319 }
320 }
321
322 fn dispatch(&mut self, cqe: Cqe) {
324 let key = cqe.user_data as usize;
325
326 enum Route {
332 Live(Rc<udp::SockShared>),
333 Done(Op),
334 }
335
336 let route = {
337 let mut ops = self.shared.ops.borrow_mut();
338 let Some(op) = ops.get(key) else {
339 tracing::error!(key, "completion for an unknown operation");
340 return;
341 };
342 let terminal = match op {
343 Op::Recv { .. } => cqe.result < 0 || !io_uring::cqueue::more(cqe.flags),
344 _ => true,
345 };
346 if terminal {
347 Route::Done(ops.remove(key))
348 } else {
349 match op {
350 Op::Recv { sock, .. } => Route::Live(sock.clone()),
351 _ => unreachable!("only receives are non-terminal"),
352 }
353 }
354 };
355
356 match route {
357 Route::Live(sock) => udp::on_recv(&self.shared, &sock, None, cqe, false),
358 Route::Done(Op::Recv { sock, one }) => udp::on_recv(&self.shared, &sock, one, cqe, true),
359 Route::Done(Op::Send(op)) => udp::on_send(op, cqe),
360 Route::Done(Op::FutexWait) => self.futex_armed = false,
361 Route::Done(Op::Cancel) => {}
362 }
363 }
364
365 fn maybe_park(&mut self) -> Result<(), Error> {
368 let unpark = self.shared.unpark.clone();
369 if unpark
370 .word
371 .compare_exchange(RUNNING, PARKED, Ordering::AcqRel, Ordering::Acquire)
372 .is_err()
373 {
374 unpark.word.store(RUNNING, Ordering::Release);
376 return Ok(());
377 }
378
379 if !self.futex_armed {
384 let key = self.shared.insert(Op::FutexWait);
385 let entry = opcode::FutexWait::new(
386 unpark.word.as_ptr(),
387 PARKED as u64,
388 FUTEX_BITSET_MATCH_ANY,
389 FUTEX2_SIZE_U32 | FUTEX2_PRIVATE,
390 )
391 .build()
392 .user_data(key);
393 if let Err(err) = self.shared.push(&entry) {
394 self.shared.ops.borrow_mut().remove(key as usize);
395 unpark.word.store(RUNNING, Ordering::Release);
396 return Err(err.into());
397 }
398 self.futex_armed = true;
399 }
400
401 let deadline = self.shared.timers.borrow().next();
402 self.shared.metrics.parks.add(1);
403 self.shared.metrics.enters.add(1);
404 let result = {
405 let mut ring = self.shared.ring.borrow_mut();
406 let to_submit = ring.submission().len() as u32;
407 let submitter = ring.submitter();
408 match deadline {
409 None => submitter.submit_and_wait(1),
410 Some(at) => {
411 let ts = abs_timespec(at);
414 let args = types::SubmitArgs::new().timespec(&ts);
415 let flags = EnterFlags::GETEVENTS | EnterFlags::EXT_ARG | EnterFlags::ABS_TIMER;
416 unsafe { submitter.enter(to_submit, 1, flags.bits(), Some(&args)) }
419 }
420 }
421 };
422 unpark.word.store(RUNNING, Ordering::Release);
423
424 match result {
425 Ok(count) => {
426 self.shared.metrics.submissions.add(count as u64);
427 Ok(())
428 }
429 Err(err)
430 if matches!(
431 err.raw_os_error(),
432 Some(libc::ETIME) | Some(libc::EINTR) | Some(libc::EBUSY)
433 ) =>
434 {
435 Ok(())
436 }
437 Err(err) => Err(err.into()),
438 }
439 }
440}
441
442impl Drop for Worker {
443 fn drop(&mut self) {
444 self.shared.stopped.set(true);
447 let deadline = Instant::now() + TEARDOWN_TIMEOUT;
449 let cancel: Vec<u64> = self
455 .shared
456 .ops
457 .borrow()
458 .iter()
459 .filter_map(|(key, op)| matches!(op, Op::Recv { .. } | Op::FutexWait).then_some(key as u64))
460 .collect();
461 let mut cancellation_failed = false;
462 for key in cancel {
463 if Instant::now() >= deadline {
464 cancellation_failed = true;
465 break;
466 }
467 cancellation_failed |= self.shared.cancel_until(key, deadline).is_err();
468 }
469 if cancellation_failed {
470 tracing::error!("failed to queue one or more io_uring teardown cancellations");
471 }
472 self.drain_teardown(deadline);
473 }
474}
475
476impl std::fmt::Debug for Worker {
477 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
478 f.debug_struct("Worker").field("tasks", &self.tasks.len()).finish()
479 }
480}
481
482pub struct Handle {
488 shared: Rc<Shared>,
489}
490
491impl Clone for Handle {
492 fn clone(&self) -> Self {
493 Self {
494 shared: self.shared.clone(),
495 }
496 }
497}
498
499impl Handle {
500 pub fn metrics(&self) -> Metrics {
502 Metrics::from_counters(self.shared.metrics.clone())
503 }
504
505 pub fn spawn(&self, future: impl Future<Output = ()> + 'static) {
510 if self.shared.stopped.get() {
511 return;
512 }
513 let mut future = Box::pin(future);
514 self.shared
515 .spawns
516 .borrow_mut()
517 .push(Box::new(move |waiter: &kio::Waiter| {
518 waiter.poll_future(future.as_mut())
519 }));
520 self.shared.unpark.unpark();
523 }
524
525 pub fn udp(&self, socket: impl Into<udp::Bound>, config: udp::Config) -> Result<udp::Socket, Error> {
536 if self.shared.stopped.get() {
537 return Err(Shared::gone_error().into());
538 }
539 udp::Socket::bind(&self.shared, socket.into(), config)
540 }
541}
542
543#[derive(Clone)]
551pub(crate) struct Owner {
552 shared: Weak<Shared>,
553 timers: Rc<RefCell<timer::Heap>>,
557}
558
559impl Owner {
560 pub(crate) fn new(shared: &Rc<Shared>) -> Self {
561 Self {
562 shared: Rc::downgrade(shared),
563 timers: shared.timers.clone(),
564 }
565 }
566
567 pub fn upgrade(&self) -> Option<Rc<Shared>> {
569 self.shared.upgrade()
570 }
571
572 pub fn handle(&self) -> Option<Handle> {
574 let shared = self.shared.upgrade()?;
575 (!shared.stopped.get()).then_some(Handle { shared })
576 }
577
578 pub fn spawn(&self, future: impl Future<Output = ()> + 'static) {
580 if let Some(handle) = self.handle() {
581 handle.spawn(future);
582 }
583 }
584
585 pub fn timer(&self) -> crate::Timer {
587 crate::Timer::from_heap(self.timers.clone())
588 }
589
590 pub fn after(&self, duration: Duration) -> crate::Timer {
592 let mut timer = self.timer();
593 timer.set(Instant::now().checked_add(duration));
594 timer
595 }
596}
597
598impl std::fmt::Debug for Handle {
599 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
600 f.debug_struct("Handle").finish()
601 }
602}
603
604impl Handle {
605 pub fn timer(&self) -> crate::Timer {
607 crate::Timer::from_heap(self.shared.timers.clone())
608 }
609
610 pub async fn run<D: moq_net::time::Driver>(&self, mut driver: D) -> moq_net::Error {
613 let mut timer = self.timer();
614 kio::wait(|waiter| {
615 loop {
616 match driver.poll(Instant::now(), waiter) {
617 Ok(at) => timer.set(at),
618 Err(err) => return Poll::Ready(err),
619 }
620 if timer.poll(waiter).is_pending() {
621 return Poll::Pending;
622 }
623 }
624 })
625 .await
626 }
627}
628
629fn kernel_release() -> String {
631 let mut uts: libc::utsname = unsafe { std::mem::zeroed() };
633 if unsafe { libc::uname(&mut uts) } != 0 {
635 return "unknown".into();
636 }
637 unsafe { std::ffi::CStr::from_ptr(uts.release.as_ptr()) }
639 .to_string_lossy()
640 .into_owned()
641}
642
643fn abs_timespec(at: Instant) -> types::Timespec {
646 let delta = at.saturating_duration_since(Instant::now());
649 let mut now = libc::timespec { tv_sec: 0, tv_nsec: 0 };
650 unsafe { libc::clock_gettime(libc::CLOCK_MONOTONIC, &mut now) };
652 let nanos = now.tv_nsec as u64 + delta.subsec_nanos() as u64;
653 let secs = (now.tv_sec as u64)
654 .saturating_add(delta.as_secs())
655 .saturating_add(nanos / 1_000_000_000);
656 types::Timespec::new().sec(secs).nsec((nanos % 1_000_000_000) as u32)
657}
658
659fn retry_teardown_submit(interruptions: &mut usize, err: std::io::Error) -> std::io::Result<()> {
661 if err.raw_os_error() != Some(libc::EINTR) || *interruptions >= TEARDOWN_EINTR_RETRIES {
662 return Err(err);
663 }
664 *interruptions += 1;
665 Ok(())
666}
667
668#[cfg(test)]
669mod tests {
670 use super::*;
671
672 use crate::Timer as Deadline;
673 use std::time::Duration;
674
675 fn worker() -> Option<Worker> {
678 worker_with(Config::default())
679 }
680
681 fn worker_with(config: Config) -> Option<Worker> {
682 match Worker::new(config) {
683 Ok(worker) => Some(worker),
684 Err(Error::Unsupported(reason)) => {
685 eprintln!("skipping io_uring test: {reason}");
686 None
687 }
688 Err(err) => panic!("worker setup failed: {err}"),
689 }
690 }
691
692 #[test]
693 fn cloned_config_has_fresh_metrics() {
694 let config = Config::default();
695 let clone = config.clone();
696 assert!(!std::sync::Arc::ptr_eq(
697 config.metrics.counters(),
698 clone.metrics.counters()
699 ));
700 }
701
702 #[test]
703 fn ready_future() {
704 let Some(mut worker) = worker() else { return };
705 let value = worker.block_on(async { 7 }).unwrap();
706 assert_eq!(value, 7);
707 }
708
709 #[test]
710 fn spawned_tasks_run() {
711 let Some(mut worker) = worker() else { return };
712 let handle = worker.handle();
713 let flag = Rc::new(std::cell::Cell::new(0));
714
715 for index in 0..3 {
716 let flag = flag.clone();
717 handle.spawn(async move {
718 flag.set(flag.get() + index + 1);
719 });
720 }
721 let handle2 = handle.clone();
723 worker
724 .block_on(async move {
725 Deadline::after(&handle2, Duration::from_millis(10)).wait().await;
726 })
727 .unwrap();
728 assert_eq!(flag.get(), 6);
729 }
730
731 #[test]
732 fn deadline_fires_at_park() {
733 let Some(mut worker) = worker() else { return };
734 let handle = worker.handle();
735 let start = Instant::now();
736 worker
739 .block_on(async move {
740 Deadline::after(&handle, Duration::from_millis(50)).wait().await;
741 })
742 .unwrap();
743 let elapsed = start.elapsed();
744 assert!(elapsed >= Duration::from_millis(50), "woke early: {elapsed:?}");
745 assert!(elapsed < Duration::from_secs(5), "woke far too late: {elapsed:?}");
746 }
747
748 #[test]
749 fn timer_rearm_and_disarm() {
750 let Some(mut worker) = worker() else { return };
751 let handle = worker.handle();
752 let mut timer = handle.timer();
753
754 assert!(timer.poll(&kio::Waiter::noop()).is_pending());
756
757 timer.set(Some(Instant::now() - Duration::from_millis(1)));
760 assert!(timer.poll(&kio::Waiter::noop()).is_ready());
761 assert!(timer.poll(&kio::Waiter::noop()).is_ready());
762
763 timer.set(Some(Instant::now() + Duration::from_secs(60)));
765 assert!(timer.poll(&kio::Waiter::noop()).is_pending());
766 timer.set(None);
767 assert!(timer.poll(&kio::Waiter::noop()).is_pending());
768
769 let start = Instant::now();
771 worker
772 .block_on(async move {
773 timer.set(Some(Instant::now() + Duration::from_millis(20)));
774 kio::wait(|waiter| timer.poll(waiter)).await;
775 })
776 .unwrap();
777 assert!(start.elapsed() >= Duration::from_millis(20));
778 }
779
780 #[test]
781 fn dropped_worker_rejects_operations() {
782 let Some(worker) = worker() else { return };
783 let handle = worker.handle();
784 let bind = || std::net::UdpSocket::bind("127.0.0.1:0").expect("bind");
785 let sock = handle.udp(bind(), udp::Config::default()).expect("socket");
786 let shared = sock.downgrade();
787 let to = sock.local_addr().expect("addr");
788 let Poll::Ready(Ok(tx)) = sock.poll_acquire(&kio::Waiter::noop()) else {
789 panic!("no tx buffer");
790 };
791 drop(worker);
792
793 assert!(handle.udp(bind(), udp::Config::default()).is_err());
796 assert!(matches!(sock.poll_recv(&kio::Waiter::noop()), Poll::Ready(Err(_))));
797 assert!(matches!(sock.poll_acquire(&kio::Waiter::noop()), Poll::Ready(Err(_))));
798 assert!(
799 tx.send(udp::Transmit {
800 to,
801 len: 1200,
802 segment: 1200,
803 ecn: None,
804 })
805 .is_err()
806 );
807 handle.spawn(async {});
809 drop(sock);
810 assert!(shared.upgrade().is_none(), "the worker leaked its staged receive");
811 }
812
813 #[test]
814 fn teardown_stops_between_completions_at_the_deadline() {
815 let Some(mut worker) = worker() else { return };
816 let first = worker.shared.insert(Op::Cancel);
817 let second = worker.shared.insert(Op::Cancel);
818 let cqe = |user_data| Cqe {
819 user_data,
820 result: 0,
821 flags: 0,
822 };
823 let before = Instant::now();
824 let deadline = before + Duration::from_millis(1);
825 let mut now = [before, deadline].into_iter();
826
827 worker.cqes = vec![cqe(first), cqe(second)];
828 assert!(!worker.dispatch_batch(Some(deadline), || {
829 now.next().expect("one deadline check per completion")
830 }));
831 assert!(!worker.shared.ops.borrow().contains(first as usize));
832 assert!(worker.shared.ops.borrow().contains(second as usize));
833 worker.shared.ops.borrow_mut().remove(second as usize);
834 }
835
836 #[test]
837 fn expired_teardown_submits_residual_sqes() {
838 let Some(mut worker) = worker() else { return };
839 for _ in 0..SQ_ENTRIES {
841 worker.shared.push(&opcode::Nop::new().build()).expect("stage NOP");
842 }
843 assert_eq!(worker.shared.ring.borrow_mut().submission().len(), SQ_ENTRIES as usize);
844
845 worker.drain_teardown(Instant::now());
846 assert!(worker.shared.ring.borrow_mut().submission().is_empty());
847 }
848
849 #[test]
850 fn teardown_submit_interrupt_budget_is_finite() {
851 let interrupted = || std::io::Error::from_raw_os_error(libc::EINTR);
852 let mut interruptions = 0;
853 for _ in 0..TEARDOWN_EINTR_RETRIES {
854 retry_teardown_submit(&mut interruptions, interrupted()).expect("retry interrupted submit");
855 }
856 assert_eq!(interruptions, TEARDOWN_EINTR_RETRIES);
857 assert_eq!(
858 retry_teardown_submit(&mut interruptions, interrupted())
859 .expect_err("interrupt budget must be finite")
860 .raw_os_error(),
861 Some(libc::EINTR)
862 );
863 }
864
865 #[test]
866 fn dropped_worker_drains_more_receives_than_the_submission_queue() {
867 let Some(worker) = worker() else { return };
868 let handle = worker.handle();
869 let config = udp::Config {
870 gro: false,
871 gso: false,
872 multishot: false,
873 rx_buffers_max: 1,
874 rx_buffer_len: 2048,
875 tx_buffers_max: 1,
876 tx_buffer_len: 2048,
877 };
878 let mut sockets = Vec::new();
879 let mut shared = Vec::new();
880 for _ in 0..=SQ_ENTRIES {
881 let sock = handle
882 .udp(std::net::UdpSocket::bind("127.0.0.1:0").expect("bind"), config.clone())
883 .expect("socket");
884 shared.push(sock.downgrade());
885 sockets.push(sock);
886 }
887
888 drop(worker);
889 drop(sockets);
890 assert!(
891 shared.iter().all(|shared| shared.upgrade().is_none()),
892 "the worker leaked a receive staged across submission batches"
893 );
894 }
895
896 #[test]
897 fn cq_covers_the_default_pool_ceilings() {
898 let config = udp::Config::default();
904 let per_socket = u32::from(config.tx_buffers_max) + u32::from(config.rx_buffers_max);
905 assert!(CQ_ENTRIES > per_socket, "CQ_ENTRIES fell behind the pool defaults");
906 }
907
908 #[test]
909 fn the_ring_honors_the_requested_cq_depth() {
910 let Some(worker) = worker() else { return };
915 let cq = worker.shared.ring.borrow().params().cq_entries();
916 assert!(cq >= CQ_ENTRIES, "kernel granted a {cq}-entry CQ, wanted {CQ_ENTRIES}");
917 }
918
919 #[test]
920 fn completion_overflow_is_survivable() {
921 let Some(mut worker) = worker() else { return };
922 let handle = worker.handle();
923 let ceiling = (CQ_ENTRIES * 2) as u16;
928 let config = udp::Config {
929 tx_buffers_max: ceiling,
930 tx_buffer_len: 2048,
931 ..Default::default()
932 };
933 let sock = handle
934 .udp(std::net::UdpSocket::bind("127.0.0.1:0").expect("bind"), config)
935 .expect("socket");
936 let to = sock.local_addr().expect("addr");
937
938 let mut held = Vec::new();
939 while let Poll::Ready(Ok(tx)) = sock.poll_acquire(&kio::Waiter::noop()) {
940 held.push(tx);
941 }
942 assert_eq!(held.len(), usize::from(ceiling));
943 for tx in held.drain(..) {
944 tx.send(udp::Transmit {
945 to,
946 len: 1200,
947 segment: 1200,
948 ecn: None,
949 })
950 .expect("send");
951 }
952
953 let saw_overflow = |worker: &Worker| worker.shared.ring.borrow_mut().submission().cq_overflow();
957 let mut overflowed = saw_overflow(&worker);
958
959 let deadline = Instant::now() + Duration::from_secs(10);
962 loop {
963 overflowed = overflowed || saw_overflow(&worker);
964 let h = handle.clone();
965 worker
966 .block_on(async move {
967 Deadline::after(&h, Duration::from_millis(10)).wait().await;
968 })
969 .unwrap();
970 let mut free = Vec::new();
971 loop {
972 match sock.poll_acquire(&kio::Waiter::noop()) {
973 Poll::Ready(Ok(tx)) => free.push(tx),
974 Poll::Ready(Err(err)) => panic!("send path failed: {err}"),
975 Poll::Pending => break,
976 }
977 }
978 if free.len() == usize::from(ceiling) {
979 break;
980 }
981 assert!(
982 Instant::now() < deadline,
983 "buffers stuck in flight: {} of {ceiling} free",
984 free.len()
985 );
986 }
987 assert!(overflowed, "the burst never overflowed the CQ; it proves nothing");
988 while let Poll::Ready(result) = sock.poll_recv(&kio::Waiter::noop()) {
991 result.expect("receive path failed");
992 }
993 }
994
995 #[test]
996 fn oversized_receive_pool_is_rejected() {
997 let Some(worker) = worker() else { return };
998 let handle = worker.handle();
999 let config = udp::Config {
1002 rx_buffers_max: u16::MAX,
1003 ..Default::default()
1004 };
1005 let err = handle
1006 .udp(std::net::UdpSocket::bind("127.0.0.1:0").expect("bind"), config)
1007 .expect_err("oversized pool");
1008 assert!(matches!(err, Error::Io(err) if err.kind() == std::io::ErrorKind::InvalidInput));
1009 }
1010
1011 #[test]
1012 fn the_send_pool_grows_to_its_ceiling() {
1013 let Some(worker) = worker() else { return };
1014 let handle = worker.handle();
1015 let config = udp::Config {
1018 tx_buffers_max: u16::MAX,
1019 ..Default::default()
1020 };
1021 handle
1022 .udp(std::net::UdpSocket::bind("127.0.0.1:0").expect("bind"), config)
1023 .expect("socket");
1024
1025 let config = udp::Config {
1027 tx_buffers_max: 200,
1028 tx_buffer_len: 4096,
1029 ..Default::default()
1030 };
1031 let sock = handle
1032 .udp(std::net::UdpSocket::bind("127.0.0.1:0").expect("bind"), config)
1033 .expect("socket");
1034
1035 let mut held = Vec::new();
1039 while let Poll::Ready(Ok(tx)) = sock.poll_acquire(&kio::Waiter::noop()) {
1040 held.push(tx);
1041 }
1042 assert_eq!(held.len(), 200);
1043 drop(worker);
1044 }
1045
1046 #[test]
1047 fn ungso_send_is_not_capped_at_a_train() {
1048 let Some(worker) = worker() else { return };
1049 let handle = worker.handle();
1050 let config = udp::Config {
1053 gso: false,
1054 ..Default::default()
1055 };
1056 let sock = handle
1057 .udp(std::net::UdpSocket::bind("127.0.0.1:0").expect("bind"), config)
1058 .expect("socket");
1059 let to = sock.local_addr().expect("addr");
1060 let Poll::Ready(Ok(tx)) = sock.poll_acquire(&kio::Waiter::noop()) else {
1061 panic!("no tx buffer");
1062 };
1063 tx.send(udp::Transmit {
1064 to,
1065 len: 64 * 1024,
1066 segment: 1000,
1067 ecn: None,
1068 })
1069 .expect("send 66 datagrams");
1070 drop(worker);
1071 }
1072
1073 #[test]
1074 fn ungso_send_is_capped_by_the_ring() {
1075 let Some(worker) = worker() else { return };
1076 let handle = worker.handle();
1077 let config = udp::Config {
1081 gso: false,
1082 ..Default::default()
1083 };
1084 let sock = handle
1085 .udp(std::net::UdpSocket::bind("127.0.0.1:0").expect("bind"), config)
1086 .expect("socket");
1087 let to = sock.local_addr().expect("addr");
1088 let Poll::Ready(Ok(tx)) = sock.poll_acquire(&kio::Waiter::noop()) else {
1089 panic!("no tx buffer");
1090 };
1091 let err = tx
1092 .send(udp::Transmit {
1093 to,
1094 len: 64 * 1024,
1095 segment: 1,
1096 ecn: None,
1097 })
1098 .expect_err("65536 datagrams from one buffer");
1099 assert_eq!(err.kind(), std::io::ErrorKind::InvalidInput);
1100 drop(worker);
1101 }
1102
1103 #[test]
1104 fn oversized_gso_segment_is_rejected() {
1105 let Some(worker) = worker() else { return };
1106 let handle = worker.handle();
1107 let sock = handle
1108 .udp(
1109 std::net::UdpSocket::bind("127.0.0.1:0").expect("bind"),
1110 udp::Config::default(),
1111 )
1112 .expect("socket");
1113 let to = sock.local_addr().expect("addr");
1114 let Poll::Ready(Ok(tx)) = sock.poll_acquire(&kio::Waiter::noop()) else {
1115 panic!("no tx buffer");
1116 };
1117 let err = tx
1120 .send(udp::Transmit {
1121 to,
1122 len: 60_000,
1123 segment: usize::from(u16::MAX) + 2,
1124 ecn: None,
1125 })
1126 .expect_err("oversized segment");
1127 assert_eq!(err.kind(), std::io::ErrorKind::InvalidInput);
1128 drop(worker);
1129 }
1130
1131 #[test]
1136 fn metrics_record_ring_and_socket_activity() {
1137 let metrics = Metrics::default();
1138 let config = Config {
1139 metrics: metrics.clone(),
1140 ..Default::default()
1141 };
1142 let Some(mut worker) = worker_with(config) else { return };
1143 let handle = worker.handle();
1144 let sock = handle
1145 .udp(
1146 std::net::UdpSocket::bind("127.0.0.1:0").expect("bind"),
1147 udp::Config::default(),
1148 )
1149 .expect("socket");
1150 let to = sock.local_addr().expect("addr");
1151
1152 let Poll::Ready(Ok(mut tx)) = sock.poll_acquire(&kio::Waiter::noop()) else {
1154 panic!("no tx buffer");
1155 };
1156 tx[..4 * 1200].fill(7);
1157 tx.send(udp::Transmit {
1158 to,
1159 len: 4 * 1200,
1160 segment: 1200,
1161 ecn: None,
1162 })
1163 .expect("send");
1164
1165 let deadline = Instant::now() + Duration::from_secs(5);
1168 let mut received = 0;
1169 while received == 0 && Instant::now() < deadline {
1170 let handle = handle.clone();
1171 worker
1172 .block_on(async move {
1173 Deadline::after(&handle, Duration::from_millis(10)).wait().await;
1174 })
1175 .unwrap();
1176 while let Poll::Ready(packet) = sock.poll_recv(&kio::Waiter::noop()) {
1177 let packet = packet.expect("receive path failed");
1178 received += packet.payload().len();
1179 }
1180 }
1181 assert!(received > 0, "the loopback never delivered the send");
1182
1183 let snap = metrics.snapshot();
1184 assert_eq!(snap.tx_sends, 1, "one GSO train is one sendmsg: {snap:?}");
1185 assert_eq!(snap.tx_datagrams, 4, "four segments: {snap:?}");
1186 assert!(snap.rx_receives > 0, "no receive completions: {snap:?}");
1187 assert!(
1188 snap.rx_datagrams >= snap.rx_receives,
1189 "fewer datagrams than receives: {snap:?}"
1190 );
1191 assert!(snap.submissions > 0, "nothing was submitted: {snap:?}");
1192 assert!(snap.completions > 0, "nothing completed: {snap:?}");
1193 assert!(snap.enters > 0, "the ring was never entered: {snap:?}");
1194 assert!(snap.parks > 0, "the worker never parked: {snap:?}");
1195 assert!(snap.timers_fired > 0, "the park deadlines never fired: {snap:?}");
1196 let own = handle.metrics().snapshot();
1199 assert_eq!((own.tx_sends, own.tx_datagrams), (snap.tx_sends, snap.tx_datagrams));
1200 }
1201
1202 #[test]
1206 fn metrics_record_pool_backpressure() {
1207 let metrics = Metrics::default();
1208 let config = Config {
1209 metrics: metrics.clone(),
1210 ..Default::default()
1211 };
1212 let Some(mut worker) = worker_with(config) else { return };
1213 let handle = worker.handle();
1214 let sock = handle
1218 .udp(
1219 std::net::UdpSocket::bind("127.0.0.1:0").expect("bind"),
1220 udp::Config {
1221 gro: false,
1222 gso: false,
1223 multishot: false,
1224 rx_buffers_max: 1,
1225 rx_buffer_len: 2048,
1226 tx_buffers_max: 1,
1227 tx_buffer_len: 2048,
1228 },
1229 )
1230 .expect("socket");
1231 let to = sock.local_addr().expect("addr");
1232
1233 let Poll::Ready(Ok(tx)) = sock.poll_acquire(&kio::Waiter::noop()) else {
1234 panic!("no tx buffer");
1235 };
1236 assert!(sock.poll_acquire(&kio::Waiter::noop()).is_pending());
1238 assert!(sock.poll_acquire(&kio::Waiter::noop()).is_pending());
1239 assert_eq!(metrics.snapshot().tx_stalls, 1);
1240 tx.send(udp::Transmit {
1241 to,
1242 len: 1200,
1243 segment: 1200,
1244 ecn: None,
1245 })
1246 .expect("send");
1247
1248 let deadline = Instant::now() + Duration::from_secs(5);
1251 let mut held = None;
1252 while held.is_none() && Instant::now() < deadline {
1253 let handle = handle.clone();
1254 worker
1255 .block_on(async move {
1256 Deadline::after(&handle, Duration::from_millis(10)).wait().await;
1257 })
1258 .unwrap();
1259 if let Poll::Ready(packet) = sock.poll_recv(&kio::Waiter::noop()) {
1260 held = Some(packet.expect("receive path failed"));
1261 }
1262 }
1263 assert!(held.is_some(), "the loopback never delivered the send");
1264 assert!(
1265 metrics.snapshot().rx_exhausted > 0,
1266 "a re-arm with every buffer held went unreported: {:?}",
1267 metrics.snapshot()
1268 );
1269
1270 let Poll::Ready(Ok(_tx)) = sock.poll_acquire(&kio::Waiter::noop()) else {
1273 panic!("completed tx buffer was not released");
1274 };
1275 assert!(sock.poll_acquire(&kio::Waiter::noop()).is_pending());
1276 assert!(sock.poll_acquire(&kio::Waiter::noop()).is_pending());
1277 assert_eq!(metrics.snapshot().tx_stalls, 2);
1278 }
1279
1280 #[test]
1284 fn metrics_count_timer_churn() {
1285 let metrics = Metrics::default();
1286 let config = Config {
1287 metrics: metrics.clone(),
1288 ..Default::default()
1289 };
1290 let Some(worker) = worker_with(config) else { return };
1291 let handle = worker.handle();
1292 let mut timer = handle.timer();
1293
1294 timer.set(Some(Instant::now() + Duration::from_secs(60)));
1295 assert_eq!(metrics.snapshot().timers_active(), 1);
1296
1297 timer.set(Some(Instant::now() + Duration::from_secs(60)));
1299 let snap = metrics.snapshot();
1300 assert_eq!((snap.timers_armed, snap.timers_cancelled, snap.timers_fired), (2, 1, 0));
1301 assert_eq!(snap.timers_active(), 1);
1302
1303 timer.set(Some(Instant::now() - Duration::from_millis(1)));
1305 assert!(timer.poll(&kio::Waiter::noop()).is_ready());
1306 let snap = metrics.snapshot();
1307 assert_eq!((snap.timers_armed, snap.timers_cancelled, snap.timers_fired), (3, 2, 1));
1308 assert_eq!(snap.timers_active(), 0);
1309
1310 timer.set(Some(Instant::now() + Duration::from_secs(60)));
1312 assert_eq!(metrics.snapshot().timers_active(), 1);
1313 drop(timer);
1314 assert_eq!(metrics.snapshot().timers_active(), 0);
1315 }
1316
1317 #[test]
1318 fn remote_wake_unparks() {
1319 let metrics = Metrics::default();
1320 let config = Config {
1321 metrics: metrics.clone(),
1322 ..Default::default()
1323 };
1324 let Some(mut worker) = worker_with(config) else { return };
1325 let flag = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
1326
1327 let thread_flag = flag.clone();
1328 let waker_slot = std::sync::Arc::new(std::sync::Mutex::new(None::<std::task::Waker>));
1329 let thread_slot = waker_slot.clone();
1330 let thread = std::thread::spawn(move || {
1331 std::thread::sleep(Duration::from_millis(50));
1333 thread_flag.store(true, Ordering::Release);
1334 if let Some(waker) = thread_slot.lock().unwrap().take() {
1335 waker.wake();
1336 }
1337 });
1338
1339 let start = Instant::now();
1340 worker
1341 .block_on(std::future::poll_fn(move |cx| {
1342 if flag.load(Ordering::Acquire) {
1343 return Poll::Ready(());
1344 }
1345 *waker_slot.lock().unwrap() = Some(cx.waker().clone());
1346 Poll::Pending
1347 }))
1348 .unwrap();
1349 assert!(start.elapsed() >= Duration::from_millis(50));
1350 thread.join().unwrap();
1351 assert!(metrics.snapshot().wakes > 0, "the remote wake went unreported");
1354 }
1355}