diff --git a/notes/architecture/01-CORE.md b/notes/architecture/01-CORE.md index 3357bf552e..2819292837 100644 --- a/notes/architecture/01-CORE.md +++ b/notes/architecture/01-CORE.md @@ -190,6 +190,9 @@ Composite key: `(height: u32, id: ScopeId)` in `BTreeSet` ensures parent scopes - `pop_work()` - Get next dirty scope or task - `pop_effect()` - Get next pending effect +### Task Passes +The synchronous drains (`poll_tasks`, `render_immediate`) poll each task at most `TaskPass::MAX_POLLS_PER_TASK` (32) times per pass. Short chains of immediate wakeups still settle within one call; a task that keeps waking itself past the bound is set aside and queued for the next pass, and `wait_for_work()` yields to the executor between passes, so a task that is always ready cannot keep the executor from running. + ## Runtime Manages async/scope/task coordination: diff --git a/packages/core/src/scheduler.rs b/packages/core/src/scheduler.rs index 22ec11482a..e55b375b7f 100644 --- a/packages/core/src/scheduler.rs +++ b/packages/core/src/scheduler.rs @@ -71,6 +71,7 @@ //! 2. Tasks: //! Description: Futures spawned in the dioxus runtime each have an unique task id. When the waker for that future is called, the task is rerun. //! Priority: These are the second highest priority tasks. They are run after all other dirty scopes have been resolved because those dirty scopes may cause children (and the tasks those children own) to drop which should cancel the futures. +//! Fairness: A synchronous pass over the queue polls each task a bounded number of times. A task that keeps waking itself (a `yield_now` that wakes immediately, or a tokio resource once the task's cooperative budget is spent) is set aside for the next pass once it reaches that bound, and `wait_for_work` hands control back to the executor between passes. Otherwise such a task would be polled forever without the executor ever running again. //! //! 3. Effects: //! Description: Effects should always run after all changes to the DOM have been applied. @@ -80,6 +81,7 @@ use crate::ScopeId; use crate::Task; use crate::VirtualDom; use crate::innerlude::Effect; +use rustc_hash::FxHashMap; use std::hash::Hash; #[derive(Debug, Clone, Copy, Eq)] @@ -118,6 +120,35 @@ impl Hash for ScopeOrder { } } +/// The tasks polled by one synchronous pass over the scheduler queue. +/// +/// A task that wakes itself while it is being polled is queued again before the poll returns. +/// Polling it again in the same pass lets short chains of immediate wakeups settle within one +/// `render_immediate`, but a task that is always ready would be polled forever and the pass would +/// never return. A pass therefore polls each task at most [`TaskPass::MAX_POLLS_PER_TASK`] times. +/// A task that comes up again after that is set aside, and [`VirtualDom::end_task_pass`] queues it +/// again for the next pass. +#[derive(Default)] +pub(crate) struct TaskPass { + polls: FxHashMap, + set_aside: Vec, +} + +impl TaskPass { + const MAX_POLLS_PER_TASK: u32 = 32; + + /// Returns true if `task` may be polled again in this pass, and sets it aside otherwise. + fn claim(&mut self, task: Task) -> bool { + let polls = self.polls.entry(task).or_insert(0); + if *polls == Self::MAX_POLLS_PER_TASK { + self.set_aside.push(task); + return false; + } + *polls += 1; + true + } +} + impl VirtualDom { fn remove_empty_dirty_tasks(&mut self) { let mut dirty_tasks = self.runtime.dirty_tasks.borrow_mut(); @@ -194,6 +225,27 @@ impl VirtualDom { Some(task) } + /// Take the top task that `pass` may still poll + pub(crate) fn pop_task_in_pass(&mut self, pass: &mut TaskPass) -> Option { + std::iter::from_fn(|| self.pop_task()).find(|&task| pass.claim(task)) + } + + /// Check if any task is queued to be polled + pub(crate) fn has_dirty_tasks(&self) -> bool { + self.runtime + .dirty_tasks + .borrow() + .values() + .any(|tasks| !tasks.is_empty()) + } + + /// Queue the tasks that `pass` set aside, so the next pass polls them + pub(crate) fn end_task_pass(&mut self, pass: TaskPass) { + for task in pass.set_aside { + self.mark_task_dirty(task); + } + } + /// Take any effects from the highest scope. This should only be called if there are no pending scope reruns or tasks. pub(crate) fn pop_effect(&mut self) -> Option { let mut pending_effects = self.runtime.pending_effects.borrow_mut(); @@ -241,6 +293,14 @@ impl VirtualDom { (None, None) => None, } } + + /// Take the top work item, skipping tasks that `pass` may not poll again + pub(crate) fn pop_work_in_pass(&mut self, pass: &mut TaskPass) -> Option { + std::iter::from_fn(|| self.pop_work()).find(|work| match work { + Work::PollTask(task) => pass.claim(*task), + Work::RerunScope(_) => true, + }) + } } #[derive(Debug)] diff --git a/packages/core/src/virtual_dom.rs b/packages/core/src/virtual_dom.rs index 581c437497..470495b8f5 100644 --- a/packages/core/src/virtual_dom.rs +++ b/packages/core/src/virtual_dom.rs @@ -2,7 +2,7 @@ //! //! This module provides the primary mechanics to create a hook-based, concurrent VDOM for Rust. -use crate::innerlude::{VProps, Work}; +use crate::innerlude::{TaskPass, VProps, Work}; use crate::properties::RootProps; use crate::root_wrapper::RootScopeWrapper; use crate::{ @@ -405,7 +405,7 @@ impl VirtualDom { } /// Mark a task as dirty - fn mark_task_dirty(&mut self, task: Task) { + pub(crate) fn mark_task_dirty(&mut self, task: Task) { let Some(scope) = self.runtime.task_scope(task) else { return; }; @@ -454,6 +454,14 @@ impl VirtualDom { return; } + // A task that woke itself again during that pass is still queued, and its wakeup has + // already been taken off the channel. Hand control back to the executor before polling + // it again, so a task that is always ready cannot keep this future from ever returning. + if self.has_dirty_tasks() { + yield_now().await; + continue; + } + // There isn't any more work we can do synchronously. Wait for any new work to be ready self.wait_for_event().await; } @@ -500,13 +508,21 @@ impl VirtualDom { /// Poll any queued tasks, then drain effects if task wakeups did not make /// any scope dirty. In that case no render pass is coming to flush them. + /// + /// Each task is polled a bounded number of times per call; see [`TaskPass`]. #[instrument(skip(self), level = "trace", name = "VirtualDom::poll_tasks")] fn poll_tasks(&mut self) { // Make sure we set the runtime since we're running user code let _runtime = RuntimeGuard::new(self.runtime.clone()); + let mut pass = TaskPass::default(); + self.poll_tasks_in_pass(&mut pass); + self.end_task_pass(pass); + } + + fn poll_tasks_in_pass(&mut self, pass: &mut TaskPass) { while !self.has_dirty_scopes() { - let Some(task) = self.pop_task() else { + let Some(task) = self.pop_task_in_pass(pass) else { break; }; @@ -522,7 +538,7 @@ impl VirtualDom { // task polling queued effects without dirtying any scope, wait_for_work // will not return for a render pass, so drain them before going back to // sleep. - self.drain_remaining_effects(); + self.drain_remaining_effects(pass); } /// Rebuild the virtualdom without handling any of the mutations @@ -604,7 +620,8 @@ impl VirtualDom { } fn render_immediate_with_writer(&mut self, to: &mut dyn WriteMutations) { - while let Some(work) = self.pop_work() { + let mut pass = TaskPass::default(); + while let Some(work) = self.pop_work_in_pass(&mut pass) { match work { Work::PollTask(task) => { _ = self.runtime.handle_task_wakeup(task); @@ -623,9 +640,10 @@ impl VirtualDom { // render would stop before the DOM converged. self.queue_events(); } + self.end_task_pass(pass); } - fn drain_remaining_effects(&mut self) { + fn drain_remaining_effects(&mut self, pass: &mut TaskPass) { loop { let mut progressed = false; @@ -638,7 +656,7 @@ impl VirtualDom { } } - while let Some(task) = self.pop_task() { + while let Some(task) = self.pop_task_in_pass(pass) { progressed = true; _ = self.runtime.handle_task_wakeup(task); self.queue_events(); diff --git a/packages/core/tests/task.rs b/packages/core/tests/task.rs index 10de4bf536..ebd5b16a7b 100644 --- a/packages/core/tests/task.rs +++ b/packages/core/tests/task.rs @@ -1,9 +1,13 @@ //! Verify that tasks get polled by the virtualdom properly, and that we escape wait_for_work safely -use std::{sync::atomic::AtomicUsize, time::Duration}; +use std::{ + sync::atomic::{AtomicBool, AtomicUsize, Ordering}, + task::Poll, + time::Duration, +}; use dioxus::prelude::*; -use dioxus_core::{generation, needs_update, spawn_forever}; +use dioxus_core::{NoOpMutations, generation, needs_update, spawn_forever}; async fn run_vdom(app: fn() -> Element) { let mut dom = VirtualDom::new(app); @@ -127,3 +131,222 @@ async fn yield_now_works() { SEQUENCE.with(|s| assert_eq!(s.borrow().len(), 20)); } + +/// Run `drive` on its own thread and fail the test if it has not returned within `deadline`. +/// +/// A scheduler that re-polls a task forever never returns control, so the check has to live +/// outside the thread that drives the VirtualDom. +fn returns_within( + deadline: Duration, + drive: impl FnOnce() -> T + Send + 'static, +) -> T { + let (tx, rx) = std::sync::mpsc::channel(); + let driver = std::thread::spawn(move || { + let _ = tx.send(drive()); + }); + match rx.recv_timeout(deadline) { + Ok(value) => value, + Err(std::sync::mpsc::RecvTimeoutError::Timeout) => { + panic!("the VirtualDom did not hand control back within {deadline:?}") + } + Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => { + std::panic::resume_unwind(driver.join().unwrap_err()) + } + } +} + +/// A task that can always make progress: every poll wakes its own waker before returning +/// `Pending`, so it is queued again while it is still being polled. +async fn wake_self_forever(polls: &'static AtomicUsize) { + std::future::poll_fn(|cx| { + polls.fetch_add(1, Ordering::Relaxed); + cx.waker().wake_by_ref(); + Poll::<()>::Pending + }) + .await +} + +/// Wake our own waker and return `Pending` once, then complete. +async fn yield_once() { + let mut yielded = false; + std::future::poll_fn(|cx| { + if yielded { + return Poll::Ready(()); + } + yielded = true; + cx.waker().wake_by_ref(); + Poll::Pending + }) + .await +} + +/// `wait_for_work` must give the executor control back between polls of a task that keeps +/// waking itself, instead of polling it again forever. +#[test] +fn wait_for_work_yields_between_polls_of_a_self_waking_task() { + static POLLS: AtomicUsize = AtomicUsize::new(0); + + fn app() -> Element { + use_hook(|| spawn(wake_self_forever(&POLLS))); + rsx!({}) + } + + let [earlier, later] = returns_within(Duration::from_secs(10), || { + let rt = tokio::runtime::Builder::new_current_thread() + .enable_time() + .build() + .unwrap(); + rt.block_on(async { + let mut dom = VirtualDom::new(app); + dom.rebuild(&mut NoOpMutations); + let work = dom.wait_for_work(); + tokio::pin!(work); + let mut polls = [0; 2]; + for sample in &mut polls { + tokio::select! { + _ = &mut work => panic!("no scope is ever dirtied, so wait_for_work should not return"), + _ = tokio::time::sleep(Duration::from_millis(50)) => {} + } + *sample = POLLS.load(Ordering::Relaxed); + } + polls + }) + }); + + // The task kept being polled while wait_for_work was pending, rather than being parked. + assert!(later > earlier, "polls: {earlier} then {later}"); +} + +/// The synchronous drains poll a task a bounded number of times per call, so they return even when +/// a task always wakes itself, and the task still runs again on the next call. +#[test] +fn synchronous_drains_return_with_a_self_waking_task() { + static POLLS: AtomicUsize = AtomicUsize::new(0); + + fn app() -> Element { + use_hook(|| spawn(wake_self_forever(&POLLS))); + rsx!({}) + } + + let polls_after_each_call = returns_within(Duration::from_secs(10), || { + let mut dom = VirtualDom::new(app); + dom.rebuild(&mut NoOpMutations); + let mut polls = Vec::new(); + for _ in 0..2 { + dom.process_events(); + polls.push(POLLS.load(Ordering::Relaxed)); + dom.render_immediate(&mut NoOpMutations); + polls.push(POLLS.load(Ordering::Relaxed)); + } + polls + }); + + // Every call polled the task again: it was queued for the next call, not dropped. + assert!( + polls_after_each_call + .windows(2) + .all(|pair| pair[1] > pair[0]), + "polls after each call: {polls_after_each_call:?}" + ); +} + +/// Effects that run while the scheduler drains after polling tasks can wake tasks too; the drain +/// must return even when an effect starts a task that always wakes itself. +#[test] +fn effect_draining_returns_with_a_self_waking_task() { + static POLLS: AtomicUsize = AtomicUsize::new(0); + + fn app() -> Element { + use_effect(|| { + spawn(wake_self_forever(&POLLS)); + }); + rsx!({}) + } + + returns_within(Duration::from_secs(10), || { + let mut dom = VirtualDom::new(app); + dom.rebuild(&mut NoOpMutations); + dom.process_events(); + }); + + assert!(POLLS.load(Ordering::Relaxed) > 0); +} + +/// A task that always wakes itself in a parent scope is polled before tasks in child scopes. +/// It must not keep the child's task from ever being polled. +#[test] +fn self_waking_task_does_not_starve_a_child_task() { + static PARENT_POLLS: AtomicUsize = AtomicUsize::new(0); + static CHILD_DONE: AtomicBool = AtomicBool::new(false); + + fn app() -> Element { + use_hook(|| spawn(wake_self_forever(&PARENT_POLLS))); + rsx! { Child {} } + } + + #[component] + fn Child() -> Element { + use_hook(|| { + spawn(async { + for _ in 0..3 { + yield_once().await; + } + CHILD_DONE.store(true, Ordering::Relaxed); + }) + }); + rsx!({}) + } + + returns_within(Duration::from_secs(10), || { + let mut dom = VirtualDom::new(app); + dom.rebuild(&mut NoOpMutations); + for _ in 0..8 { + dom.process_events(); + } + }); + + assert!(CHILD_DONE.load(Ordering::Relaxed)); +} + +/// Tokio makes a task that has used up its cooperative budget return `Pending` from its next +/// tokio resource. Off the runtime's worker threads it also wakes that task immediately, and it +/// only refills the budget once the executor gets control back. A task that simply loops over a +/// tokio resource (an interval, a channel) must therefore keep running rather than freeze the +/// VirtualDom. +#[test] +fn exhausted_tokio_coop_budget_does_not_freeze_the_virtual_dom() { + static UNITS: AtomicUsize = AtomicUsize::new(0); + + fn app() -> Element { + use_hook(|| { + spawn(async { + loop { + tokio::task::coop::consume_budget().await; + UNITS.fetch_add(1, Ordering::Relaxed); + } + }) + }); + rsx!({}) + } + + returns_within(Duration::from_secs(10), || { + let rt = tokio::runtime::Builder::new_multi_thread() + .worker_threads(1) + .enable_time() + .build() + .unwrap(); + // `block_on` from a thread that is not one of the runtime's workers. + rt.block_on(async { + let mut dom = VirtualDom::new(app); + dom.rebuild(&mut NoOpMutations); + tokio::select! { + _ = dom.wait_for_work() => panic!("no scope is ever dirtied, so wait_for_work should not return"), + _ = tokio::time::sleep(Duration::from_millis(100)) => {} + } + }); + }); + + // The default budget is 128 units per poll. Getting past it means the budget was refilled, + // which only happens when the executor got control back. + assert!(UNITS.load(Ordering::Relaxed) > 128); +}