use std::cell::RefCell;
use std::rc::Rc;
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::Arc;
use crate::base::FrameRequester;
use super::arena::Key;
use super::execute::update_if_necessary;
use super::node::{
add_edge, remove_observer_edges, remove_source_edges, Graph, Node, NodeKind, NodeState,
};
use super::scheduler::RemoteShared;
const MAX_FLUSH_RUNS: usize = 100_000;
const MAX_RUNS_PER_EFFECT_PER_FLUSH: u32 = 1_000;
const DRAW_VIOLATION_SAMPLE_CAP: usize = 8;
pub(crate) const MSG_DISPOSED: &str =
"abstracttui reactive: handle used after its node was disposed. FIX: keep the owning \
scope alive as long as the handle (state a Dyn rebuilds belongs OUTSIDE its closure — \
see dyn_view vs dyn_view_scoped), or use Signal::try_get_untracked where 'gone' is a \
valid answer";
pub(crate) const MSG_WRONG_THREAD: &str =
"abstracttui reactive: handle used on a thread that did not create it. FIX: the reactive \
graph is single-threaded by design — send data to the UI thread with spawn_worker/post \
(reactive::remote) and write signals from the posted closure, never from the worker";
pub(crate) const MSG_CYCLE: &str =
"abstracttui reactive: dependency cycle — a computation re-entered itself while running. \
FIX: a memo/effect (transitively) reads its own output; break the loop by reading the \
input with get_untracked or splitting the state into two signals";
pub(crate) const MSG_DRAW_READ: &str =
"tracked signal read inside a DRAW closure — the region will never repaint when this \
value changes (RT1-2). FIX: move the read into a dyn_view (re-renders on change) or \
capture the value before the closure; use get_untracked for a deliberate stale peek";
#[derive(Copy, Clone, Debug, PartialEq, Eq, Hash)]
pub(crate) struct RawId {
pub key: Key,
pub rt: u32,
}
pub(crate) struct Runtime {
pub graph: Graph,
pub rt_id: u32,
pub current_owner: Option<Key>,
pub current_observer: Option<Key>,
pub epoch_counter: u64,
pub order_counter: u64,
pub batch_depth: u32,
pub flushing: bool,
pub queue: Vec<Key>,
pub remote: Arc<RemoteShared>,
pub frame_requested: bool,
pub frame_requester: Option<Rc<dyn FrameRequester>>,
pub draw_depth: u32,
pub draw_read_violations: u64,
pub draw_read_samples: Vec<String>,
pub flush_epoch: u64,
pub worker_failures: Vec<String>,
pub frame_tasks: Vec<Box<dyn FnMut(std::time::Instant) -> bool>>,
pub timers: Vec<(std::time::Instant, Box<dyn FnOnce()>)>,
pub contexts: std::collections::HashMap<Key, ContextEntries>,
}
static NEXT_RT_ID: AtomicU32 = AtomicU32::new(1);
impl Runtime {
fn new() -> Self {
Runtime {
graph: Graph::new(),
rt_id: NEXT_RT_ID.fetch_add(1, Ordering::Relaxed),
current_owner: None,
current_observer: None,
epoch_counter: 0,
order_counter: 0,
batch_depth: 0,
flushing: false,
queue: Vec::new(),
remote: Arc::new(RemoteShared::new()),
frame_requested: false,
frame_requester: None,
draw_depth: 0,
draw_read_violations: 0,
draw_read_samples: Vec::new(),
flush_epoch: 0,
worker_failures: Vec::new(),
frame_tasks: Vec::new(),
timers: Vec::new(),
contexts: std::collections::HashMap::new(),
}
}
pub(crate) fn check_thread(&self, id: RawId) {
if id.rt != self.rt_id {
panic!("{MSG_WRONG_THREAD}");
}
}
pub(crate) fn create_node(&mut self, owner: Option<Key>, kind: NodeKind) -> Key {
self.order_counter += 1;
let mut node = Node::new(kind, self.order_counter);
if matches!(node.kind, NodeKind::Memo { .. }) {
node.state = NodeState::Dirty;
}
node.parent = owner;
let key = self.graph.insert(node);
if let Some(o) = owner {
if !self.graph.contains(o) {
self.graph.remove(key);
panic!(
"abstracttui reactive: node created under a disposed scope. FIX: the \
scope you captured died (a Dyn generation, a closed modal); create \
state on a scope that outlives the use — the mount scope for durable \
state, the generation scope (dyn_view_scoped) for per-render state"
);
}
let needs_sweep = {
let onode = self.graph.get_mut(o).expect("checked above");
onode.owned.push(key);
onode.owned.len() >= 32 && onode.owned.len().is_power_of_two()
};
if needs_sweep {
let owned = std::mem::take(&mut self.graph.get_mut(o).expect("owner").owned);
let filtered: Vec<Key> = owned
.into_iter()
.filter(|k| self.graph.contains(*k))
.collect();
self.graph.get_mut(o).expect("owner").owned = filtered;
}
}
key
}
fn report_draw_read(&mut self, source: Key) {
let who = self
.graph
.get(source)
.map(|n| n.describe(source))
.unwrap_or_else(|| "disposed node".to_string());
if cfg!(debug_assertions) {
panic!("abstracttui reactive: {MSG_DRAW_READ}; offending read: {who}");
}
self.draw_read_violations += 1;
if self.draw_read_samples.len() < DRAW_VIOLATION_SAMPLE_CAP {
self.draw_read_samples
.push(format!("#FALLBACK {MSG_DRAW_READ}; offending read: {who}"));
}
}
pub(crate) fn track_read(&mut self, source: Key) {
if self.draw_depth > 0 && self.current_observer.is_none() {
self.report_draw_read(source);
return; }
let Some(observer) = self.current_observer else {
return;
};
let obs_epoch = match self.graph.get(observer) {
Some(n) => n.run_epoch,
None => return, };
{
let Some(src) = self.graph.get_mut(source) else {
panic!("{MSG_DISPOSED}");
};
if src.seen_epoch == obs_epoch {
return;
}
src.seen_epoch = obs_epoch;
}
let duplicate = self
.graph
.get(observer)
.map(|o| o.sources.contains(&source))
.unwrap_or(true);
if !duplicate {
add_edge(&mut self.graph, source, observer);
}
}
pub(crate) fn mark_written(&mut self, source: Key) {
let mut stack: Vec<(Key, NodeState)> = match self.graph.get(source) {
Some(n) => n.observers.iter().map(|&o| (o, NodeState::Dirty)).collect(),
None => return,
};
while let Some((id, level)) = stack.pop() {
let Some(node) = self.graph.get_mut(id) else {
continue;
};
if node.state >= level {
continue;
}
let was_clean = node.state == NodeState::Clean;
node.state = level;
if node.kind.is_effect() && !node.queued {
node.queued = true;
self.queue.push(id);
}
if was_clean {
for &o in node.observers.clone().iter() {
stack.push((o, NodeState::Check));
}
}
}
}
pub(crate) fn mark_direct_observers_dirty(&mut self, source: Key) {
let observers: Vec<Key> = match self.graph.get(source) {
Some(n) => n.observers.clone(),
None => return,
};
for id in observers {
if let Some(node) = self.graph.get_mut(id) {
if node.state < NodeState::Dirty {
node.state = NodeState::Dirty;
}
if node.kind.is_effect() && !node.queued {
node.queued = true;
self.queue.push(id);
}
}
}
}
pub(crate) fn collect_dispose(&mut self, root: Key, out: &mut DisposeBundle) {
enum Phase {
Enter(Key),
Finish(Key),
}
let mut stack = vec![Phase::Enter(root)];
while let Some(phase) = stack.pop() {
match phase {
Phase::Enter(id) => {
let Some(node) = self.graph.get_mut(id) else {
continue;
};
let owned = std::mem::take(&mut node.owned);
stack.push(Phase::Finish(id));
for c in owned {
stack.push(Phase::Enter(c));
}
}
Phase::Finish(id) => {
if let Some(node) = self.graph.get_mut(id) {
let mut cleanups = std::mem::take(&mut node.cleanups);
cleanups.reverse(); out.cleanups.extend(cleanups);
}
remove_source_edges(&mut self.graph, id);
remove_observer_edges(&mut self.graph, id);
if let Some(ctx) = self.contexts.remove(&id) {
out.dropped_contexts.push(ctx);
}
if let Some(node) = self.graph.remove(id) {
out.dropped.push(node);
}
}
}
}
}
}
pub(crate) type ContextEntries = Vec<(std::any::TypeId, Rc<dyn std::any::Any>)>;
#[derive(Default)]
pub(crate) struct DisposeBundle {
pub cleanups: Vec<Box<dyn FnOnce()>>,
pub dropped: Vec<Node>,
pub dropped_contexts: Vec<ContextEntries>,
}
thread_local! {
static RT: RefCell<Runtime> = RefCell::new(Runtime::new());
}
pub(crate) fn with_rt<R>(f: impl FnOnce(&mut Runtime) -> R) -> R {
RT.with(|cell| f(&mut cell.borrow_mut()))
}
pub(crate) fn restore_owner(prev: Option<Key>) {
let _ = RT.try_with(|cell| {
if let Ok(mut rt) = cell.try_borrow_mut() {
rt.current_owner = prev;
}
});
}
pub(crate) fn restore_context(node: Key, prev_owner: Option<Key>, prev_observer: Option<Key>) {
let _ = RT.try_with(|cell| {
if let Ok(mut rt) = cell.try_borrow_mut() {
rt.current_owner = prev_owner;
rt.current_observer = prev_observer;
if let Some(n) = rt.graph.get_mut(node) {
n.running = false;
}
}
});
}
pub fn flush_effects() {
let already = with_rt(|rt| {
if rt.flushing {
true
} else {
rt.flushing = true;
false
}
});
if already {
return;
}
struct FlushGuard;
impl Drop for FlushGuard {
fn drop(&mut self) {
let _ = RT.try_with(|cell| {
if let Ok(mut rt) = cell.try_borrow_mut() {
rt.flushing = false;
}
});
}
}
let _guard = FlushGuard;
with_rt(|rt| rt.flush_epoch += 1);
let mut runs: usize = 0;
loop {
let mut batch: Vec<(u64, Key)> = with_rt(|rt| {
let queue = std::mem::take(&mut rt.queue);
queue
.into_iter()
.filter_map(|k| rt.graph.get(k).map(|n| (n.order, k)))
.collect()
});
if batch.is_empty() {
break;
}
batch.sort_unstable_by_key(|(order, _)| *order);
for (_, id) in batch {
let culprit = with_rt(|rt| {
let epoch = rt.flush_epoch;
let node = rt.graph.get_mut(id)?;
node.queued = false;
if node.flush_stamp != epoch {
node.flush_stamp = epoch;
node.flush_runs = 0;
}
node.flush_runs += 1;
(node.flush_runs > MAX_RUNS_PER_EFFECT_PER_FLUSH).then(|| node.describe(id))
});
if let Some(who) = culprit {
panic!(
"abstracttui reactive: {who} ran more than \
{MAX_RUNS_PER_EFFECT_PER_FLUSH} times in one flush — it (transitively) \
rewrites its own dependencies. FIX: read the rewritten signal with \
get_untracked inside the effect, or split read/write state; name the \
culprit with effect_labeled to trace it"
);
}
runs += 1;
if runs > MAX_FLUSH_RUNS {
panic!(
"abstracttui reactive: flush did not settle after {MAX_FLUSH_RUNS} effect \
runs — some effect chain keeps re-dirtying itself. FIX: find the writer \
(label effects with effect_labeled — the per-effect ceiling usually \
names it first) and cut its tracked read of what it writes"
);
}
update_if_necessary(id);
}
}
}
pub(crate) fn maybe_flush() {
let should = with_rt(|rt| rt.batch_depth == 0 && !rt.flushing && !rt.queue.is_empty());
if should {
flush_effects();
}
}
pub fn batch<R>(f: impl FnOnce() -> R) -> R {
with_rt(|rt| rt.batch_depth += 1);
struct BatchGuard;
impl Drop for BatchGuard {
fn drop(&mut self) {
let _ = RT.try_with(|cell| {
if let Ok(mut rt) = cell.try_borrow_mut() {
rt.batch_depth = rt.batch_depth.saturating_sub(1);
}
});
}
}
let result = {
let _guard = BatchGuard;
f()
};
maybe_flush();
result
}
pub fn untrack<R>(f: impl FnOnce() -> R) -> R {
let prev = with_rt(|rt| rt.current_observer.take());
struct UntrackGuard(Option<Key>);
impl Drop for UntrackGuard {
fn drop(&mut self) {
let prev = self.0;
let _ = RT.try_with(|cell| {
if let Ok(mut rt) = cell.try_borrow_mut() {
rt.current_observer = prev;
}
});
}
}
let _guard = UntrackGuard(prev);
f()
}
pub fn on_cleanup(f: impl FnOnce() + 'static) {
with_rt(|rt| {
let Some(owner) = rt.current_owner else {
panic!(
"abstracttui reactive: on_cleanup called outside any scope or computation. \
FIX: call it inside an effect body (cleanup-before-rerun) or use \
Scope::on_cleanup(cx, ..) to target a scope explicitly"
);
};
rt.graph
.get_mut(owner)
.expect("current owner is always live")
.cleanups
.push(Box::new(f));
});
}
pub(crate) fn dispose_node(id: RawId) {
let bundle = with_rt(|rt| {
rt.check_thread(id);
let mut bundle = DisposeBundle::default();
rt.collect_dispose(id.key, &mut bundle);
bundle
});
for c in bundle.cleanups {
c();
}
drop(bundle.dropped);
}
#[derive(Copy, Clone, Debug, PartialEq, Eq)]
pub struct RuntimeStats {
pub live_nodes: usize,
pub slot_capacity: usize,
pub queued_effects: usize,
}
pub fn stats() -> RuntimeStats {
with_rt(|rt| RuntimeStats {
live_nodes: rt.graph.live(),
slot_capacity: rt.graph.capacity_slots(),
queued_effects: rt.queue.len(),
})
}
pub(crate) fn exit_draw_phase() {
let _ = RT.try_with(|cell| {
if let Ok(mut rt) = cell.try_borrow_mut() {
rt.draw_depth = rt.draw_depth.saturating_sub(1);
}
});
}