Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions notes/architecture/01-CORE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
60 changes: 60 additions & 0 deletions packages/core/src/scheduler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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)]
Expand Down Expand Up @@ -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<Task, u32>,
set_aside: Vec<Task>,
}

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();
Expand Down Expand Up @@ -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<Task> {
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<Effect> {
let mut pending_effects = self.runtime.pending_effects.borrow_mut();
Expand Down Expand Up @@ -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<Work> {
std::iter::from_fn(|| self.pop_work()).find(|work| match work {
Work::PollTask(task) => pass.claim(*task),
Work::RerunScope(_) => true,
})
}
}

#[derive(Debug)]
Expand Down
32 changes: 25 additions & 7 deletions packages/core/src/virtual_dom.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::{
Expand Down Expand Up @@ -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;
};
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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;
};

Expand All @@ -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
Expand Down Expand Up @@ -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);
Expand All @@ -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;

Expand All @@ -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();
Expand Down
Loading