From a65d23afe97ab7ed13a59be94cd646679697e76e Mon Sep 17 00:00:00 2001 From: Carl Lerche Date: Wed, 14 Jun 2023 12:45:52 -0700 Subject: [PATCH] wip --- tokio/src/loom/std/unsafe_cell.rs | 2 ++ tokio/src/runtime/scheduler/multi_thread/counters.rs | 8 ++++++++ tokio/src/runtime/scheduler/multi_thread/idle.rs | 2 ++ tokio/src/runtime/scheduler/multi_thread/queue.rs | 2 ++ tokio/src/runtime/scheduler/multi_thread/worker.rs | 2 ++ tokio/src/runtime/task/list.rs | 2 +- 6 files changed, 17 insertions(+), 1 deletion(-) diff --git a/tokio/src/loom/std/unsafe_cell.rs b/tokio/src/loom/std/unsafe_cell.rs index 66c1d7943..3d6513b46 100644 --- a/tokio/src/loom/std/unsafe_cell.rs +++ b/tokio/src/loom/std/unsafe_cell.rs @@ -6,10 +6,12 @@ impl UnsafeCell { UnsafeCell(std::cell::UnsafeCell::new(data)) } + #[inline(always)] pub(crate) fn with(&self, f: impl FnOnce(*const T) -> R) -> R { f(self.0.get()) } + #[inline(always)] pub(crate) fn with_mut(&self, f: impl FnOnce(*mut T) -> R) -> R { f(self.0.get()) } diff --git a/tokio/src/runtime/scheduler/multi_thread/counters.rs b/tokio/src/runtime/scheduler/multi_thread/counters.rs index 50bcc1198..8bd550f9f 100644 --- a/tokio/src/runtime/scheduler/multi_thread/counters.rs +++ b/tokio/src/runtime/scheduler/multi_thread/counters.rs @@ -8,6 +8,7 @@ mod imp { static NUM_UNPARKS_LOCAL: AtomicUsize = AtomicUsize::new(0); static NUM_LIFO_SCHEDULES: AtomicUsize = AtomicUsize::new(0); static NUM_LIFO_CAPPED: AtomicUsize = AtomicUsize::new(0); + static NUM_STEALS: AtomicUsize = AtomicUsize::new(0); impl Drop for super::Counters { fn drop(&mut self) { @@ -16,6 +17,7 @@ mod imp { let maintenance = NUM_MAINTENANCE.load(Relaxed); let lifo_scheds = NUM_LIFO_SCHEDULES.load(Relaxed); let lifo_capped = NUM_LIFO_CAPPED.load(Relaxed); + let num_steals = NUM_STEALS.load(Relaxed); println!("---"); println!("notifies (local): {}", notifies_local); @@ -23,6 +25,7 @@ mod imp { println!(" maintenance: {}", maintenance); println!(" LIFO schedules: {}", lifo_scheds); println!(" LIFO capped: {}", lifo_capped); + println!(" steals: {}", num_steals); } } @@ -45,6 +48,10 @@ mod imp { pub(crate) fn inc_lifo_capped() { NUM_LIFO_CAPPED.fetch_add(1, Relaxed); } + + pub(crate) fn inc_num_steals() { + NUM_STEALS.fetch_add(1, Relaxed); + } } #[cfg(not(tokio_internal_mt_counters))] @@ -54,6 +61,7 @@ mod imp { pub(crate) fn inc_num_maintenance() {} pub(crate) fn inc_lifo_schedules() {} pub(crate) fn inc_lifo_capped() {} + pub(crate) fn inc_num_steals() {} } #[derive(Debug)] diff --git a/tokio/src/runtime/scheduler/multi_thread/idle.rs b/tokio/src/runtime/scheduler/multi_thread/idle.rs index 9dcee08bc..c41d80e9b 100644 --- a/tokio/src/runtime/scheduler/multi_thread/idle.rs +++ b/tokio/src/runtime/scheduler/multi_thread/idle.rs @@ -113,6 +113,8 @@ impl Idle { return; } + super::counters::inc_num_unparks_local(); + // Acquire the lock let synced = shared.synced.lock(); self.notify_synced(synced, shared, true); diff --git a/tokio/src/runtime/scheduler/multi_thread/queue.rs b/tokio/src/runtime/scheduler/multi_thread/queue.rs index 03406f94e..100a63361 100644 --- a/tokio/src/runtime/scheduler/multi_thread/queue.rs +++ b/tokio/src/runtime/scheduler/multi_thread/queue.rs @@ -425,6 +425,8 @@ impl Steal { return None; } + super::counters::inc_num_steals(); + dst_stats.incr_steal_count(n as u16); dst_stats.incr_steal_operations(); diff --git a/tokio/src/runtime/scheduler/multi_thread/worker.rs b/tokio/src/runtime/scheduler/multi_thread/worker.rs index a2b92bf8c..2007c16da 100644 --- a/tokio/src/runtime/scheduler/multi_thread/worker.rs +++ b/tokio/src/runtime/scheduler/multi_thread/worker.rs @@ -1366,6 +1366,8 @@ impl Shared { if let Some(prev) = prev { core.run_queue .push_back_or_overflow(prev, self, &mut core.stats); + } else { + return; } } else { core.run_queue diff --git a/tokio/src/runtime/task/list.rs b/tokio/src/runtime/task/list.rs index fb7dbdc1d..da9ea92a0 100644 --- a/tokio/src/runtime/task/list.rs +++ b/tokio/src/runtime/task/list.rs @@ -119,7 +119,7 @@ impl OwnedTasks { /// a LocalNotified, giving the thread permission to poll this task. #[inline] pub(crate) fn assert_owner(&self, task: Notified) -> LocalNotified { - assert_eq!(task.header().get_owner_id(), self.id); + debug_assert_eq!(task.header().get_owner_id(), self.id); // safety: All tasks bound to this OwnedTasks are Send, so it is safe // to poll it on this thread no matter what thread we are on.