abstracttui 0.2.20

A reactive, compositor-grade terminal UI engine: fine-grained signals, layered rendering with damage tracking, images (kitty/iTerm2/sixel/mosaic), software-rasterized 3D (GLB), themes and animation.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
//! The thread-local reactive runtime: node creation, dependency tracking,
//! two-phase invalidation, lazy recomputation, ownership disposal and the
//! effect flush loop.
//!
//! ## Why single-threaded + thread-local
//!
//! The graph lives in ONE `thread_local` and is never shared across
//! threads. Terminal output is inherently serialized (one byte stream),
//! so a `Send + Sync` graph would buy nothing except a lock acquisition
//! on *every signal read* — the hottest operation in the system. Instead,
//! timers/IO threads hand work to the UI thread through
//! [`super::scheduler::WakeHandle`] (posted closures + wakeup), which is
//! both cheaper and impossible to deadlock. Handles carry the id of the
//! runtime that minted them; using one on the wrong thread is a loud
//! panic, never silent aliasing.
//!
//! ## Borrow discipline (the one rule that keeps this sound)
//!
//! `with_rt` hands out `&mut Runtime` under a `RefCell` borrow. User code
//! (computations, cleanups, `Drop` impls of stored values) re-enters the
//! runtime, so it must NEVER run under that borrow. Every operation is
//! therefore structured as: borrow -> mutate graph -> collect closures /
//! `Rc`s -> release -> run user code -> repeat. Node payloads are `Rc`ed
//! precisely so they can be cloned out before the borrow is released.

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;

/// Hard ceiling on TOTAL effect executions in one flush — the backstop
/// against many effects ping-ponging collectively.
const MAX_FLUSH_RUNS: usize = 100_000;

/// Ceiling on ONE effect's executions within a single flush (RT1-15a).
/// This fires long before the global ceiling (1k runs instead of 100k),
/// and the panic names the effect's creation label — a runaway loop
/// should cost milliseconds and produce a culprit, not seconds and a
/// mystery.
const MAX_RUNS_PER_EFFECT_PER_FLUSH: u32 = 1_000;

/// How many draw-phase violation descriptions are retained for the
/// diagnostics getter in release builds (the count is unbounded).
const DRAW_VIOLATION_SAMPLE_CAP: usize = 8;

// Panic messages NAME THE FIX (cycle-8 audit): a user hitting one at
// 2am should know what to change without opening this file.
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";

/// A `Key` stamped with the id of the runtime (thread) that owns it.
/// Handles are `Copy + Send` so cross-thread *transport* (e.g. inside a
/// posted closure) is allowed; cross-thread *use* fails the stamp check.
#[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,
    /// Node that owns whatever is created right now (scope or computation).
    pub current_owner: Option<Key>,
    /// Computation whose dependencies are being tracked right now.
    /// `None` inside `untrack` and outside any computation.
    pub current_observer: Option<Key>,
    /// Bumped at every computation run; used to dedupe repeated reads of
    /// the same source within one run in O(1).
    pub epoch_counter: u64,
    pub order_counter: u64,
    pub batch_depth: u32,
    pub flushing: bool,
    /// Effects awaiting execution (dirty or check — resolved at flush).
    pub queue: Vec<Key>,
    pub remote: Arc<RemoteShared>,
    pub frame_requested: bool,
    pub frame_requester: Option<Rc<dyn FrameRequester>>,
    /// Nonzero while the frame's DRAW phase runs (nesting-safe). Tracked
    /// reads while set are the RT1-2 stale-pixel bug: no computation
    /// owns the region, so nothing ever re-renders it.
    pub draw_depth: u32,
    /// Total draw-phase tracked-read violations (release builds count
    /// instead of panicking; debug builds never get here).
    pub draw_read_violations: u64,
    /// First few violation descriptions, for the diagnostics getter.
    pub draw_read_samples: Vec<String>,
    /// Bumped once per `flush_effects` entry; per-effect run counters
    /// reset lazily by comparing their stamp against this.
    pub flush_epoch: u64,
    /// Labeled panics reported by workers spawned via
    /// `scheduler::spawn_worker` (RT1-15b), delivered as posted jobs.
    pub worker_failures: Vec<String>,
    /// Per-frame callbacks (animations). Run once per frame in phase U;
    /// a task returning false is dropped. Empty list = zero idle cost.
    pub frame_tasks: Vec<Box<dyn FnMut(std::time::Instant) -> bool>>,
    /// One-shot timers (toast dismissal, debounce). Unlike frame tasks
    /// these do NOT keep frames coming: the loop sleeps until the
    /// earliest deadline (`next_timer_deadline`) — zero wakeups before.
    /// Entries carry an id so a holder can cancel BEFORE the deadline
    /// (`reactive::interval` rides this); `after` never exposes its id,
    /// so its contract is unchanged.
    pub timers: Vec<TimerEntry>,
    /// Id source for [`TimerEntry`]; monotonically increasing, never
    /// reused (u64: exhaustion is not a practical concern).
    pub next_timer_id: u64,
    /// The clock reading `run_due_timers` is firing with, readable by
    /// timer callbacks that re-arm (an interval's next deadline must
    /// derive from the LOOP's clock — injected in tests — not from a
    /// fresh `Instant::now`). `None` outside a timer-fire pass.
    pub timer_now: Option<std::time::Instant>,
    /// The loop's published TURN clock (`reactive::set_loop_clock`):
    /// the arm-time authority for `after`/`interval` deadlines while a
    /// turn runs, so one injected clock (`Driver::set_clock`) governs
    /// both when timers are ARMED and when they FIRE. Arming against a
    /// fresh `Instant::now` under an injected clock plants deadlines
    /// the injected timeline may never reach (the wave-11 drawer-feed
    /// flake: a loaded machine stretched real time between rig
    /// creation and the width probe's `after(0)` arm while the test
    /// clock advanced 12 ms — the probe never fired). `None` outside a
    /// driven turn (bare rigs keep real-time arming).
    pub clock_now: Option<std::time::Instant>,
    /// Scope-provided context values (`provide_context`/`use_context`):
    /// sparse side map (only providers pay), removed on dispose.
    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(),
            next_timer_id: 0,
            timer_now: None,
            clock_now: None,
            contexts: std::collections::HashMap::new(),
        }
    }

    pub(crate) fn check_thread(&self, id: RawId) {
        if id.rt != self.rt_id {
            panic!("{MSG_WRONG_THREAD}");
        }
    }

    /// Create a node, optionally attaching it to an owner scope.
    /// Memos are born `Dirty` (lazy: nothing is computed until observed).
    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) {
                // Creating under a dead scope would leak the node forever
                // (nobody left to dispose it) — refuse loudly.
                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);
                // Explicitly-disposed children leave stale keys behind; sweep
                // them amortized (every power-of-two growth past 32) so a
                // long-lived scope with heavy child churn stays bounded
                // without paying O(siblings) on every dispose.
                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
    }

    /// A tracked read arrived while phase D was running: debug builds
    /// panic naming the offending node (loud during development, when the
    /// widget author is looking); release builds count + keep a bounded
    /// sample so `diagnostics()` can surface the label without killing a
    /// shipped app over a stale region.
    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}"));
        }
    }

    /// Record `current_observer reads source`.
    ///
    /// Dedupe subtlety (REDTEAM: this is diamond-country): a single global
    /// epoch is NOT enough. If memo B is pulled in the middle of memo A's
    /// run, B's run overwrites `seen_epoch` on shared sources; back in A, a
    /// pure epoch check would either miss the dedupe or (worse, if the
    /// check were `seen == current_global`) silently skip adding A's edge.
    /// So: epoch hit => certainly already added this run, done. Epoch miss
    /// => fall back to scanning this run's source list before linking.
    pub(crate) fn track_read(&mut self, source: Key) {
        // RT1-2: a tracked read from a DRAW closure is the stale-pixel bug
        // — draw runs outside any computation, so the read subscribes
        // nothing and the region never repaints when the value changes.
        // `current_observer.is_none()` identifies exactly that case: a
        // memo legitimately recomputed during draw (via an untracked pull)
        // reads its own sources under `observer = the memo`, which is
        // graph maintenance, not a widget read. Untracked reads
        // (`get_untracked`) never reach this function and remain the
        // sanctioned way to peek at data captured for painting.
        if self.draw_depth > 0 && self.current_observer.is_none() {
            self.report_draw_read(source);
            return; // nothing to subscribe anyway (no observer)
        }
        let Some(observer) = self.current_observer else {
            return;
        };
        let obs_epoch = match self.graph.get(observer) {
            Some(n) => n.run_epoch,
            None => return, // observer disposed mid-run (pathological); drop the edge
        };
        {
            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);
        }
    }

    /// Down-phase of the two-phase marking: direct observers of a written
    /// value become `Dirty`, transitive observers become `Check`. Only a
    /// node's FIRST transition away from `Clean` descends — its downstream
    /// was already marked then, which is what makes repeated writes cheap
    /// and update storms linear in the affected subgraph.
    ///
    /// Iterative on purpose: fanout chains can be deep and the native
    /// stack is not ours to burn during a storm.
    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));
                }
            }
        }
    }

    /// After a memo recomputed to a DIFFERENT value: its direct observers
    /// are definitely stale. No descend — the original down-phase already
    /// marked everything transitively `Check`; this only upgrades the
    /// certainty of the first hop (the equality cut-off gate).
    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;
                }
                // Defensive: an effect that attached mid-flush may not be
                // queued yet.
                if node.kind.is_effect() && !node.queued {
                    node.queued = true;
                    self.queue.push(id);
                }
            }
        }
    }

    /// Detach + free `root` and its whole ownership subtree. Cleanups are
    /// COLLECTED here (under the borrow) and RUN by the caller (outside
    /// it); freed `Node`s ride along so their user payloads drop outside
    /// the borrow too (a stored value's `Drop` may re-enter the runtime).
    ///
    /// Order invariant: children before parents (reverse creation order
    /// among siblings), cleanups LIFO within a node. Children may hold
    /// references to parent-provided state, so they must die first — same
    /// order leptos's `Owner::cleanup` and Rust drop order use.
    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));
                    // Push in creation order so the LIFO stack visits the
                    // most recently created child first.
                    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(); // LIFO
                        out.cleanups.extend(cleanups);
                    }
                    remove_source_edges(&mut self.graph, id);
                    remove_observer_edges(&mut self.graph, id);
                    // Context values provided on this scope die with it
                    // (dropped OUTSIDE the borrow, with the node).
                    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);
                    }
                }
            }
        }
    }
}

/// One armed one-shot timer: deadline, cancellation id, callback.
pub(crate) struct TimerEntry {
    pub deadline: std::time::Instant,
    pub id: u64,
    pub f: Box<dyn FnOnce()>,
}

/// One scope's provided context values: (type, boxed value) pairs.
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>,
    /// Context values from disposed scopes; dropped after the borrow
    /// releases (a value's Drop may re-enter the runtime).
    pub dropped_contexts: Vec<ContextEntries>,
}

thread_local! {
    static RT: RefCell<Runtime> = RefCell::new(Runtime::new());
}

/// The ONLY way to touch the runtime. Never run user code inside `f`.
pub(crate) fn with_rt<R>(f: impl FnOnce(&mut Runtime) -> R) -> R {
    RT.with(|cell| f(&mut cell.borrow_mut()))
}

/// Panic-safe owner restore for guards living outside this module.
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;
        }
    });
}

/// Panic-safe full context restore after a computation run (owner,
/// observer, `running` flag). Used by `execute::CtxGuard`; a poisoned
/// tracking context after a caught panic would corrupt every later
/// computation, so this must succeed even during unwinding.
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;
            }
        }
    });
}

/// Drain the effect queue in creation order until it stays empty.
/// Creation order runs outer effects before the inner effects they own —
/// an outer re-render disposes stale inner effects, which then skip here
/// via the generation check instead of running against dead state.
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 {
            // RT1-15a: per-effect run accounting. A ping-pong pair (A's
            // effect writes B's dependency and vice versa) trips this at
            // ~1k runs with a NAMED culprit, milliseconds into the storm —
            // long before the global backstop would fire after seconds of
            // frozen UI with no attribution.
            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);
        }
    }
}

/// Flush unless writes are being coalesced (batch) or a flush is already
/// draining (its loop will pick up new work).
pub(crate) fn maybe_flush() {
    let should = with_rt(|rt| rt.batch_depth == 0 && !rt.flushing && !rt.queue.is_empty());
    if should {
        flush_effects();
    }
}

/// Coalesce writes: effects observe only the final state, once, when the
/// outermost batch ends. Nesting is allowed.
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
}

/// Run `f` with dependency tracking suspended: reads inside do not
/// subscribe the current computation.
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()
}

/// Register a cleanup on the CURRENTLY RUNNING owner (computation or
/// scope body). Inside an effect this runs before the effect's next
/// re-run and at disposal — the Solid `onCleanup` contract.
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));
    });
}

/// Detach + free a node tree; runs cleanups after releasing the borrow.
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);
}

/// Observable counters for leak tests and diagnostics.
#[derive(Copy, Clone, Debug, PartialEq, Eq)]
pub struct RuntimeStats {
    /// Live nodes of any kind (signals, memos, effects, scopes).
    pub live_nodes: usize,
    /// Slots ever allocated — bounded by peak concurrency, not churn.
    pub slot_capacity: usize,
    /// Effects currently queued for the next flush.
    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(),
    })
}

/// Panic-safe `timer_now` clear for `run_due_timers`' fire guard.
pub(crate) fn clear_timer_now() {
    let _ = RT.try_with(|cell| {
        if let Ok(mut rt) = cell.try_borrow_mut() {
            rt.timer_now = None;
        }
    });
}

/// Panic-safe `clock_now` write for [`set_loop_clock`]'s turn guard —
/// the driver publishes at turn start and clears on every turn exit
/// path (RAII), so a panicking turn never strands a stale clock on the
/// thread's runtime (test binaries reuse threads across tests).
///
/// [`set_loop_clock`]: crate::reactive::set_loop_clock
pub(crate) fn write_clock_now(now: Option<std::time::Instant>) {
    let _ = RT.try_with(|cell| {
        if let Ok(mut rt) = cell.try_borrow_mut() {
            rt.clock_now = now;
        }
    });
}

/// Panic-safe draw-depth decrement for `diag::DrawPhase`'s Drop.
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);
        }
    });
}