use crate::arena::ReactiveState;
use crate::arena::effect_arena::take_pending_effects_split;
use crate::arena::signal_arena::remove_signal_writer;
use crate::arena::with_signal_arena;
use crate::arena::{
EffectId, EffectMetadata, current_effect, destroy_children, effect_arena_insert,
effect_arena_remove, mark_effect_pending, remove_from_pending_set, set_effect_parent,
take_pending_effects,
};
use std::cell::Cell;
use std::time::Instant;
thread_local! {
static PROCESSING_SCHEDULED: Cell<bool> = const { Cell::new(false) };
}
pub fn schedule_effect_processing() {
PROCESSING_SCHEDULED.with(|scheduled| {
scheduled.set(true);
});
crate::executor::notify_effect_loop();
}
pub fn is_processing_scheduled() -> bool {
PROCESSING_SCHEDULED.with(std::cell::Cell::get)
}
pub(crate) fn flush_effects_with_budget(
budget: std::time::Duration,
must_run: &mut Vec<EffectId>,
skippable: &mut Vec<EffectId>,
) {
PROCESSING_SCHEDULED.with(|scheduled| {
scheduled.set(false);
});
let start = Instant::now();
take_pending_effects_split(must_run, skippable);
let mut i = 0;
let over_budget = loop {
if start.elapsed() > budget {
break true;
}
let Some(effect_id) = skippable.get(i) else {
break false;
};
if effect_id.state() != ReactiveState::Clean {
effect_id.update_if_necessary(false);
}
i += 1;
};
skippable.drain(..i);
for effect_id in must_run {
if effect_id.state() != ReactiveState::Clean {
effect_id.update_if_necessary(over_budget);
}
}
}
pub fn flush_effects() -> usize {
PROCESSING_SCHEDULED.with(|scheduled| {
scheduled.set(false);
});
Effect::process_all()
}
fn run_single_effect_internal(effect_id: EffectId, remove_from_pending: bool) {
use crate::arena::ReactiveState;
if !effect_id.needs_work() {
return;
}
effect_id.set_state(ReactiveState::Clean);
if remove_from_pending {
remove_from_pending_set(effect_id);
}
if !effect_id.has_callback() {
return;
}
effect_id.clear_signals(
|source_id| source_id.remove_subscriber(effect_id),
|output_id| remove_signal_writer(output_id, effect_id),
);
let _guard = crate::arena::CurrentEffectGuard::new(Some(effect_id));
effect_id.run_callback();
}
pub fn run_single_effect(effect_id: EffectId) {
run_single_effect_internal(effect_id, false);
}
pub fn untracked<F, R>(f: F) -> R
where
F: FnOnce() -> R,
{
let _guard = crate::arena::CurrentEffectGuard::new(None);
f()
}
pub struct Effect {
id: EffectId,
}
impl Effect {
pub fn new<F>(f: F) -> Self
where
F: FnMut() + Send + 'static,
{
Self::new_internal(f, false, false)
}
pub fn new_skippable<F>(f: F) -> Self
where
F: FnMut() + Send + 'static,
{
Self::new_internal(f, true, false)
}
pub fn new_ui_updates_only<F>(f: F) -> Self
where
F: FnMut() + Send + 'static,
{
Self::new_internal(f, false, true)
}
fn new_internal<F>(f: F, skippable: bool, ui_updates_only: bool) -> Self
where
F: FnMut() + Send + 'static,
{
let parent = current_effect();
let metadata = EffectMetadata::new_with_callback_parent_and_flags(
Box::new(f),
parent,
skippable,
ui_updates_only,
);
let id = effect_arena_insert(metadata);
if let Some(parent_id) = parent {
set_effect_parent(id, parent_id);
parent_id.add_child(id);
}
let effect = Self { id };
effect.run_now();
effect
}
pub(crate) fn run_now(&self) {
self.id.with_sources(|old_sources| {
for source_id in old_sources {
source_id.remove_subscriber(self.id);
}
});
self.id.clear_sources();
let _guard = crate::arena::CurrentEffectGuard::new(Some(self.id));
self.id.run_callback();
use crate::arena::ReactiveState;
self.id.set_state(ReactiveState::Clean);
}
pub fn invalidate(&self) {
mark_effect_pending(self.id);
}
pub(crate) fn id(&self) -> EffectId {
self.id
}
pub(crate) fn from_raw(id: EffectId) -> Self {
Self { id }
}
pub fn process_all() -> usize {
use crate::arena::ReactiveState;
let mut total = 0;
loop {
let pending_ids = take_pending_effects();
if pending_ids.is_empty() {
break;
}
for effect_id in pending_ids {
if effect_id.state() != ReactiveState::Clean {
effect_id.update_if_necessary(false);
total += 1;
}
}
}
total
}
}
impl Drop for Effect {
fn drop(&mut self) {
let has_parent = self.id.parent().is_some();
remove_from_pending_set(self.id);
let effect_id = self.id;
let exists = self.id.with_sources(|sources| {
with_signal_arena(|arena| {
for source_id in sources {
if let Some(node) = arena.get(source_id.index()) {
node.subscribers.write().retain(|&eid| eid != effect_id);
}
}
});
});
if exists.is_none() {
return; }
self.id.with_outputs(|outputs| {
for signal_id in outputs {
remove_signal_writer(signal_id, effect_id);
}
});
if has_parent {
return;
}
destroy_children(self.id);
effect_arena_remove(self.id);
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
#[test]
fn effect_debounces_rapid_invalidations() {
let run_count = Arc::new(AtomicUsize::new(0));
let run_count_clone = run_count.clone();
let effect = Effect::new(move || {
run_count_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(run_count.load(Ordering::Relaxed), 1);
for _ in 0..20 {
effect.invalidate();
}
Effect::process_all();
assert_eq!(run_count.load(Ordering::Relaxed), 2);
}
#[test]
fn effect_automatic_reactivity() {
use std::sync::atomic::AtomicI32;
let signal = crate::Signal::new();
let signal_node_id = signal.id();
let counter = Arc::new(AtomicI32::new(0));
let counter_clone = counter.clone();
let effect = Effect::new(move || {
signal_node_id.track_dependency();
counter_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(counter.load(Ordering::Relaxed), 1);
signal_node_id.with_subscribers(|subscribers| {
assert_eq!(subscribers.len(), 1);
assert!(subscribers.contains(&effect.id()));
});
effect.id().with_sources(|sources| {
assert_eq!(sources.count(), 1);
});
let has_signal = effect
.id()
.with_sources(|sources| {
for s in sources {
if s == signal_node_id {
return true;
}
}
false
})
.unwrap_or(false);
assert!(has_signal);
signal.emit();
flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 2);
signal.emit();
flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 3);
}
#[test]
fn effect_automatic_reactivity_via_track_dependency() {
use std::sync::atomic::AtomicI32;
let signal = crate::Signal::new();
let signal_node_id = signal.id();
let counter = Arc::new(AtomicI32::new(0));
let counter_clone = counter.clone();
let _effect = Effect::new(move || {
signal_node_id.track_dependency();
counter_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(counter.load(Ordering::Relaxed), 1);
signal.emit();
flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 2);
crate::Transaction::run(|| {
signal.emit();
signal.emit();
signal.emit();
});
assert_eq!(counter.load(Ordering::Relaxed), 3);
}
#[test]
fn effect_ui_origin_filtering() {
use std::sync::atomic::AtomicI32;
let signal = crate::Signal::new();
let signal_node_id = signal.id();
let counter = Arc::new(AtomicI32::new(0));
let counter_clone = counter.clone();
let _effect = Effect::new(move || {
signal_node_id.track_dependency();
counter_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(counter.load(Ordering::Relaxed), 1);
signal.emit_from_api();
flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 2);
signal.emit_from_ui();
flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 3);
}
#[test]
fn debounce_multiple_emissions() {
use std::sync::atomic::AtomicI32;
let signal = crate::Signal::new();
let signal_node_id = signal.id();
let counter = Arc::new(AtomicI32::new(0));
let counter_clone = counter.clone();
let _effect = Effect::new(move || {
signal_node_id.track_dependency();
counter_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(counter.load(Ordering::Relaxed), 1);
signal.emit();
signal.emit();
signal.emit();
assert_eq!(counter.load(Ordering::Relaxed), 1);
flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 2);
}
#[test]
fn flush_effects_processes_pending() {
use std::sync::atomic::AtomicI32;
let signal = crate::Signal::new();
let signal_node_id = signal.id();
let counter = Arc::new(AtomicI32::new(0));
let counter_clone = counter.clone();
let _effect = Effect::new(move || {
signal_node_id.track_dependency();
counter_clone.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(counter.load(Ordering::Relaxed), 1);
signal.emit();
flush_effects();
assert_eq!(counter.load(Ordering::Relaxed), 2);
}
}