Skip to main content

fusor/
coherence.rs

1//! Coherent async views and their loading status.
2//!
3//! Supply an [`AsyncBoundary`] to one `<Async>` region to coordinate its reads.
4//! Its [`BoundaryStatus`] reports when the complete view is ready or a read has
5//! failed. Keep status displays and retry controls outside the region so they
6//! remain usable while it is pending.
7use crate::{
8    Effect, OwnerHandle, Registration, Signal, batch, effect, signal, untrack, versions::Versions,
9};
10use std::{
11    cell::{Cell, RefCell},
12    collections::BTreeMap,
13    rc::{Rc, Weak},
14};
15
16/// Version of the supported renderer preparation and coherent-read contract.
17///
18/// This covers this module, [`crate::versions::Versions`], and
19/// [`Signal::with_render_value`]. A change to their integration semantics bumps
20/// this version. It is independent of the browser template format and of any
21/// external renderer's generated-code protocol.
22#[doc(hidden)]
23pub const VERSION: u32 = 2;
24
25#[derive(Clone, Debug, PartialEq, Eq)]
26pub enum BoundaryStatus {
27    Detached,
28    Pending,
29    Ready,
30    Error(String),
31    Faulted(String),
32    Disposed,
33}
34
35/// A speculative read lifetime, distinct from a component's owner lifetime.
36///
37/// Reads may start before owner activation. The boundary cancels leases when
38/// their inputs become obsolete, they leave the read set, or it is disposed.
39#[doc(hidden)]
40pub trait ReadLease {
41    /// Invalidate delivery and request cancellation. Must tolerate repeated calls;
42    /// cancellation does not roll back external work or await its termination.
43    fn cancel(&self);
44}
45
46/// A renderer-owned candidate update. Dropping it releases unpublished patches.
47/// A renderer may retain prepared scopes across discovery passes in one attempt;
48/// use [`Attempt::on_invalidate`] to retire them when the input epoch ends.
49///
50/// Formatting, constructors, key comparisons and fallible target setup belong
51/// in evaluation or [`Self::validate`], before any visible mutation. The boundary
52/// validates source versions both before and after renderer validation.
53#[doc(hidden)]
54pub trait Publication {
55    /// Check that all prepared operations can still target the intended nodes.
56    /// An error preserves the previously published scene and reports `Error`.
57    fn validate(&self) -> Result<(), String>;
58    /// Apply the validated update synchronously, without application callbacks,
59    /// signal writes, constructors, formatting, or other fallible preparation.
60    /// An unexpected failure reports `Faulted`; the boundary cannot undo a
61    /// renderer's partial scene mutation. Successful application and finishing
62    /// run within one reactive batch.
63    fn apply(&mut self) -> Result<(), String>;
64    /// Finish a successful publication: activate adopted owners and retire old
65    /// scopes in the renderer's documented order. Application callbacks may run
66    /// here, so release mutable scene borrows first. Input changes during this
67    /// phase are detected before the boundary declares the scene ready.
68    fn finish(self: Box<Self>);
69}
70
71type Evaluate = dyn Fn(&Attempt) -> Result<Box<dyn Publication>, String>;
72
73struct Inner {
74    id: u64,
75    status: Signal<BoundaryStatus>,
76    wake: Signal<u64>,
77    retry: Cell<u64>,
78    epoch: Cell<u64>,
79    attached: Cell<bool>,
80    alive: Cell<bool>,
81    driving: Cell<bool>,
82    violation: RefCell<Option<String>>,
83    rejected_retry: Cell<Option<u64>>,
84    inputs: RefCell<Versions>,
85    reads: RefCell<BTreeMap<u64, Rc<dyn ReadLease>>>,
86    candidates: RefCell<Vec<Box<dyn FnOnce()>>>,
87}
88
89thread_local! {
90    static NEXT: Cell<u64> = const { Cell::new(0) };
91    static EVALUATING: RefCell<Vec<Weak<Inner>>> = const { RefCell::new(Vec::new()) };
92    // The length of EVALUATING, readable without a borrow on every write.
93    static EVALUATION_DEPTH: Cell<usize> = const { Cell::new(0) };
94    static PREPARING: RefCell<Option<OwnerHandle>> = const { RefCell::new(None) };
95}
96
97/// A stable handle. Attach it to exactly one live `Async` region. Read
98/// status and put selection/retry controls outside that region.
99#[derive(Clone)]
100pub struct AsyncBoundary(Rc<Inner>);
101
102impl AsyncBoundary {
103    pub fn coherent() -> Self {
104        let id = NEXT.with(|next| {
105            let id = next.get().checked_add(1).expect("boundary id overflow");
106            next.set(id);
107            id
108        });
109        Self(Rc::new(Inner {
110            id,
111            status: signal(BoundaryStatus::Detached),
112            wake: signal(0),
113            retry: Cell::new(0),
114            epoch: Cell::new(0),
115            attached: Cell::new(false),
116            alive: Cell::new(true),
117            driving: Cell::new(false),
118            violation: RefCell::new(None),
119            inputs: RefCell::new(Versions::default()),
120            rejected_retry: Cell::new(None),
121            reads: RefCell::new(BTreeMap::new()),
122            candidates: RefCell::new(Vec::new()),
123        }))
124    }
125
126    pub fn status(&self) -> BoundaryStatus {
127        let cycle = EVALUATING.with(|stack| {
128            if stack
129                .borrow()
130                .iter()
131                .filter_map(Weak::upgrade)
132                .any(|inner| inner.id == self.0.id)
133            {
134                *self.0.violation.borrow_mut() =
135                    Some("a coherent region cannot read its own boundary status".into());
136                true
137            } else {
138                false
139            }
140        });
141        if cycle {
142            self.0.status.get_untracked()
143        } else {
144            self.0.status.get()
145        }
146    }
147
148    pub fn pending(&self) -> bool {
149        self.status() == BoundaryStatus::Pending
150    }
151    pub fn is_interactive(&self) -> bool {
152        self.status() == BoundaryStatus::Ready
153    }
154
155    /// Explicitly retry failed reads. Compatible successes remain reusable.
156    pub fn retry(&self) {
157        if self.0.alive.get() {
158            self.0
159                .retry
160                .set(self.0.retry.get().checked_add(1).expect("retry overflow"));
161            self.0.wake.update(|n| *n += 1);
162        }
163    }
164
165    /// Attach one renderer to this boundary and retain the returned mount.
166    ///
167    /// Evaluation can begin immediately, even while `owner` is prepared. It
168    /// collects reactive dependencies and speculative reads, and must be pure
169    /// apart from read discovery. Use [`prepare_state`] for candidate constructor
170    /// effects. Pending reads or evaluation/validation errors discard the new
171    /// publication. The renderer remains responsible for staging nodes, guarding
172    /// interaction while pending, and retaining the previous scene.
173    ///
174    /// Disposal is permanent: a boundary cannot attach again after its mount or
175    /// owner has been disposed. This method does not commit the supplied owner.
176    #[doc(hidden)]
177    pub fn attach(
178        &self,
179        owner: &OwnerHandle,
180        evaluate: impl Fn(&Attempt) -> Result<Box<dyn Publication>, String> + 'static,
181    ) -> Result<BoundaryMount, String> {
182        if !self.0.alive.get() || owner.is_disposed() {
183            return Err("cannot attach a disposed async boundary".into());
184        }
185        if self.0.attached.replace(true) {
186            return Err("an async boundary can only attach once".into());
187        }
188        let inner = Rc::downgrade(&self.0);
189        let cleanup = owner.on_cleanup(move || {
190            if let Some(inner) = inner.upgrade() {
191                dispose(&inner);
192            }
193        });
194        let weak = Rc::downgrade(&self.0);
195        let owner = owner.clone();
196        let evaluate: Rc<Evaluate> = Rc::new(evaluate);
197        let subscription = effect(move || {
198            if let Some(inner) = weak
199                .upgrade()
200                .filter(|inner| inner.alive.get() && !owner.is_disposed())
201            {
202                Versions::exclude(|| {
203                    inner.wake.get();
204                });
205                drive(&inner, &*evaluate);
206            }
207        });
208        Ok(BoundaryMount {
209            inner: Rc::downgrade(&self.0),
210            subscription,
211            _cleanup: cleanup,
212        })
213    }
214}
215
216/// Retained by the renderer. Dropping it cancels candidate reads and subscriptions.
217#[must_use = "retain the mount for as long as the coherent region is mounted"]
218#[doc(hidden)]
219pub struct BoundaryMount {
220    inner: Weak<Inner>,
221    subscription: Effect,
222    _cleanup: Registration,
223}
224impl Drop for BoundaryMount {
225    fn drop(&mut self) {
226        self.subscription.dispose();
227        if let Some(inner) = self.inner.upgrade() {
228            dispose(&inner);
229        }
230    }
231}
232
233fn dispose(inner: &Inner) {
234    if !inner.alive.replace(false) {
235        return;
236    }
237    let reads = inner.reads.take();
238    for read in reads.into_values() {
239        read.cancel();
240    }
241    discard_candidates(inner);
242    inner.status.set(BoundaryStatus::Disposed);
243}
244
245fn discard_candidates(inner: &Inner) {
246    let candidates = inner.candidates.take();
247    untrack(|| {
248        for cleanup in candidates {
249            cleanup();
250        }
251    });
252}
253
254/// Weak lifetime witness used by read declarations that outlive a mounted view.
255#[doc(hidden)]
256pub struct BoundaryLifetime(Weak<Inner>);
257impl BoundaryLifetime {
258    /// Whether the boundary remains attached or available to attach.
259    pub fn is_live(&self) -> bool {
260        self.0.upgrade().is_some_and(|inner| inner.alive.get())
261    }
262}
263
264/// One discovery pass in the current attempt. An unresolved read records pending
265/// and returns normally; no exception, panic, or unwinding implements suspension.
266#[doc(hidden)]
267pub struct Attempt {
268    inner: Weak<Inner>,
269    epoch: u64,
270    pending: Cell<bool>,
271    reads: RefCell<BTreeMap<u64, Rc<dyn ReadLease>>>,
272}
273impl Attempt {
274    /// A weak witness for declarations retained beyond the current evaluation.
275    pub fn boundary_lifetime(&self) -> BoundaryLifetime {
276        BoundaryLifetime(self.inner.clone())
277    }
278    /// Stable identity of this boundary; zero only after its lifetime expires.
279    pub fn boundary_id(&self) -> u64 {
280        self.inner.upgrade().map_or(0, |inner| inner.id)
281    }
282    /// Candidate input generation. Changes when captured source versions change.
283    pub fn epoch(&self) -> u64 {
284        self.epoch
285    }
286    /// Retire prepared state when this input epoch ends or the boundary is
287    /// disposed, even if its renderer slot is never visited again. Register once
288    /// per slot and epoch, normally capturing a weak slot reference. The callback
289    /// runs untracked and without internal borrows held.
290    /// Pending passes and retries within the same epoch keep the candidate alive.
291    pub fn on_invalidate(&self, cleanup: impl FnOnce() + 'static) {
292        if let Some(inner) = self
293            .inner
294            .upgrade()
295            .filter(|inner| inner.alive.get() && inner.epoch.get() == self.epoch)
296        {
297            inner.candidates.borrow_mut().push(Box::new(cleanup));
298        } else {
299            untrack(cleanup);
300        }
301    }
302    /// Explicit retry generation, used by read adapters to retry failed reads.
303    pub fn retry_generation(&self) -> u64 {
304        self.inner.upgrade().map_or(0, |inner| inner.retry.get())
305    }
306    /// Record an unresolved read. Evaluation returns normally to discover other
307    /// reads, but this pass cannot publish.
308    pub fn pending(&self) {
309        self.pending.set(true);
310    }
311    /// Retain a read under an identity stable across discovery passes. Different
312    /// read declarations must use different IDs within a boundary. Re-registering
313    /// the same declaration in one pass replaces its lease.
314    pub fn register(&self, id: u64, lease: Rc<dyn ReadLease>) {
315        self.reads.borrow_mut().insert(id, lease);
316    }
317    /// Wake this boundary after a read completes. Notifications from obsolete
318    /// epochs or disposed boundaries are ignored. Adapters must separately guard
319    /// their result delivery using their own cancellation/request generations.
320    pub fn notifier(&self) -> impl Fn() + 'static {
321        let weak = self.inner.clone();
322        let epoch = self.epoch;
323        move || {
324            if let Some(inner) = weak
325                .upgrade()
326                .filter(|inner| inner.alive.get() && inner.epoch.get() == epoch)
327            {
328                inner.wake.update(|n| *n += 1);
329            }
330        }
331    }
332}
333
334struct Evaluation;
335impl Drop for Evaluation {
336    fn drop(&mut self) {
337        EVALUATING.with(|stack| {
338            stack.borrow_mut().pop();
339        });
340        EVALUATION_DEPTH.set(EVALUATION_DEPTH.get() - 1);
341    }
342}
343struct Driving<'a>(&'a Cell<bool>);
344impl Drop for Driving<'_> {
345    fn drop(&mut self) {
346        self.0.set(false);
347    }
348}
349
350fn drive(inner: &Rc<Inner>, evaluate: &Evaluate) {
351    if inner.rejected_retry.get() == Some(inner.retry.get()) {
352        return;
353    }
354    if inner.driving.replace(true) {
355        return;
356    }
357    let _driving = Driving(&inner.driving);
358    inner.status.set(BoundaryStatus::Pending);
359    let versions = inner.inputs.borrow().clone();
360    if !untrack(|| versions.is_current()) || inner.epoch.get() == 0 {
361        inner
362            .epoch
363            .set(inner.epoch.get().checked_add(1).expect("attempt overflow"));
364        let reads = inner.reads.take();
365        for read in reads.into_values() {
366            read.cancel();
367        }
368        discard_candidates(inner);
369    }
370    if !inner.alive.get() {
371        return;
372    }
373    inner.violation.take();
374    let attempt = Attempt {
375        inner: Rc::downgrade(inner),
376        epoch: inner.epoch.get(),
377        pending: Cell::new(false),
378        reads: RefCell::new(BTreeMap::new()),
379    };
380    EVALUATING.with(|stack| stack.borrow_mut().push(Rc::downgrade(inner)));
381    EVALUATION_DEPTH.set(EVALUATION_DEPTH.get() + 1);
382    let guard = Evaluation;
383    let (result, inputs) = Versions::capture(|| evaluate(&attempt));
384    drop(guard);
385    let old = inner.reads.replace(attempt.reads.take());
386    let removed: Vec<_> = old
387        .into_iter()
388        .filter(|(id, _)| !inner.reads.borrow().contains_key(id))
389        .map(|(_, lease)| lease)
390        .collect();
391    for lease in removed {
392        lease.cancel();
393    }
394    *inner.inputs.borrow_mut() = inputs.clone();
395    if !inner.alive.get() {
396        return;
397    }
398    publish_attempt(inner, &attempt, inputs, result);
399}
400
401fn publish_attempt(
402    inner: &Rc<Inner>,
403    attempt: &Attempt,
404    inputs: Versions,
405    result: Result<Box<dyn Publication>, String>,
406) {
407    let result = match inner.violation.take() {
408        Some(error) => {
409            inner.rejected_retry.set(Some(inner.retry.get()));
410            Err(error)
411        }
412        None => result,
413    };
414    let mut publication = match result {
415        Ok(publication) => publication,
416        Err(error) => {
417            inner.status.set(BoundaryStatus::Error(error));
418            return;
419        }
420    };
421    if attempt.pending.get() || !untrack(|| inputs.is_current()) {
422        inner.status.set(BoundaryStatus::Pending);
423        return;
424    }
425    if let Err(error) = publication.validate() {
426        inner.status.set(BoundaryStatus::Error(error));
427        return;
428    }
429    if !inner.alive.get() || !untrack(|| inputs.is_current()) {
430        inner.status.set(BoundaryStatus::Pending);
431        return;
432    }
433    batch(|| {
434        if let Err(error) = publication.apply() {
435            inner.status.set(BoundaryStatus::Faulted(error));
436            return;
437        }
438        // Keep the guard through activation. Synchronous activation effects may
439        // invalidate these inputs and must never be overwritten by Ready.
440        publication.finish();
441        if inner.alive.get() {
442            inner.status.set(if untrack(|| inputs.is_current()) {
443                BoundaryStatus::Ready
444            } else {
445                BoundaryStatus::Pending
446            });
447        }
448    });
449}
450
451/// Record a write inside a coherent evaluation. Every signal write checks this,
452/// so the idle check inlines into callers.
453#[inline]
454pub(crate) fn mutation(operation: &str) -> bool {
455    EVALUATION_DEPTH.get() != 0 && evaluating_mutation(operation)
456}
457
458#[cold]
459fn evaluating_mutation(operation: &str) -> bool {
460    EVALUATING.with(|stack| {
461        let inner = stack.borrow().last().and_then(Weak::upgrade);
462        if let Some(inner) = inner {
463            *inner.violation.borrow_mut() = Some(format!(
464                "coherent render evaluation must be pure: {operation}"
465            ));
466            true
467        } else {
468            false
469        }
470    })
471}
472
473pub(crate) fn preparing_owner() -> Option<OwnerHandle> {
474    PREPARING.with(|owner| owner.borrow().clone())
475}
476
477/// Construct candidate state with effects gated by the supplied owner.
478///
479/// If the owner is not active, effects created inside `make` first run when it
480/// activates and stop when it is disposed, even if their handles are retained.
481/// Retain those effect handles as usual. Nested calls restore the previous owner,
482/// including during unwinding.
483///
484/// An already active owner preserves ordinary immediate effect timing and does
485/// not adopt those effects for cleanup. This function neither commits the owner
486/// nor validates or rolls back renderer work. Ordinary constructors should not
487/// use it unless staged effects are intended: outside this explicit preparation
488/// context, [`effect`] runs immediately even when an unrelated owner is prepared.
489/// Speculative async reads have separate leases and may start before activation.
490#[doc(hidden)]
491pub fn prepare_state<R>(owner: OwnerHandle, make: impl FnOnce(OwnerHandle) -> R) -> R {
492    struct Restore(Option<OwnerHandle>);
493    impl Drop for Restore {
494        fn drop(&mut self) {
495            PREPARING.with(|owner| *owner.borrow_mut() = self.0.take());
496        }
497    }
498    let _restore = Restore(PREPARING.with(|previous| previous.replace(Some(owner.clone()))));
499    make(owner)
500}