Skip to main content

fusor/
coherence.rs

1//! The renderer/readiness contract, independent of executors and the browser.
2//!
3//! A boundary validates tracked source versions before synchronously publishing
4//! owned patches. Async libraries implement `ReadLease`; renderers implement
5//! `Publication`. This module does not execute or cache application requests.
6use crate::{
7    Effect, OwnerHandle, Registration, Signal, batch, effect, signal, untrack, versions::Versions,
8};
9use std::{
10    cell::{Cell, RefCell},
11    collections::BTreeMap,
12    rc::{Rc, Weak},
13};
14
15#[derive(Clone, Debug, PartialEq, Eq)]
16pub enum BoundaryStatus {
17    Detached,
18    Pending,
19    Ready,
20    Error(String),
21    Faulted(String),
22    Disposed,
23}
24
25/// A read preparation lifetime, distinct from a component's DOM lifetime.
26#[doc(hidden)]
27pub trait ReadLease {
28    fn cancel(&self);
29}
30
31/// Generated render work. Formatting, constructors, key comparisons and target
32/// validation belong in preparation/validate, never in `apply`.
33#[doc(hidden)]
34pub trait Publication {
35    fn validate(&self) -> Result<(), String>;
36    fn apply(&mut self) -> Result<(), String>;
37    fn finish(self: Box<Self>);
38}
39
40type Evaluate = dyn Fn(&Attempt) -> Result<Box<dyn Publication>, String>;
41
42struct Inner {
43    id: u64,
44    status: Signal<BoundaryStatus>,
45    wake: Signal<u64>,
46    retry: Cell<u64>,
47    epoch: Cell<u64>,
48    attached: Cell<bool>,
49    alive: Cell<bool>,
50    driving: Cell<bool>,
51    violation: RefCell<Option<String>>,
52    rejected_retry: Cell<Option<u64>>,
53    inputs: RefCell<Versions>,
54    reads: RefCell<BTreeMap<u64, Rc<dyn ReadLease>>>,
55}
56
57thread_local! {
58    static NEXT: Cell<u64> = const { Cell::new(0) };
59    static EVALUATING: RefCell<Vec<Weak<Inner>>> = const { RefCell::new(Vec::new()) };
60    // The length of EVALUATING, readable without a borrow on every write.
61    static EVALUATION_DEPTH: Cell<usize> = const { Cell::new(0) };
62    static PREPARING: RefCell<Option<OwnerHandle>> = const { RefCell::new(None) };
63}
64
65/// A stable handle. Attach it to exactly one live `Async` region. Read
66/// status and put selection/retry controls outside that region.
67#[derive(Clone)]
68pub struct AsyncBoundary(Rc<Inner>);
69
70impl AsyncBoundary {
71    pub fn coherent() -> Self {
72        let id = NEXT.with(|next| {
73            let id = next.get().checked_add(1).expect("boundary id overflow");
74            next.set(id);
75            id
76        });
77        Self(Rc::new(Inner {
78            id,
79            status: signal(BoundaryStatus::Detached),
80            wake: signal(0),
81            retry: Cell::new(0),
82            epoch: Cell::new(0),
83            attached: Cell::new(false),
84            alive: Cell::new(true),
85            driving: Cell::new(false),
86            violation: RefCell::new(None),
87            inputs: RefCell::new(Versions::default()),
88            rejected_retry: Cell::new(None),
89            reads: RefCell::new(BTreeMap::new()),
90        }))
91    }
92
93    pub fn status(&self) -> BoundaryStatus {
94        let cycle = EVALUATING.with(|stack| {
95            if stack
96                .borrow()
97                .iter()
98                .filter_map(Weak::upgrade)
99                .any(|inner| inner.id == self.0.id)
100            {
101                *self.0.violation.borrow_mut() =
102                    Some("a coherent region cannot read its own boundary status".into());
103                true
104            } else {
105                false
106            }
107        });
108        if cycle {
109            self.0.status.get_untracked()
110        } else {
111            self.0.status.get()
112        }
113    }
114
115    pub fn pending(&self) -> bool {
116        self.status() == BoundaryStatus::Pending
117    }
118    pub fn is_interactive(&self) -> bool {
119        self.status() == BoundaryStatus::Ready
120    }
121
122    /// Explicitly retry failed reads. Compatible successes remain reusable.
123    pub fn retry(&self) {
124        if self.0.alive.get() {
125            self.0
126                .retry
127                .set(self.0.retry.get().checked_add(1).expect("retry overflow"));
128            self.0.wake.update(|n| *n += 1);
129        }
130    }
131
132    #[doc(hidden)]
133    pub fn attach(
134        &self,
135        owner: &OwnerHandle,
136        evaluate: impl Fn(&Attempt) -> Result<Box<dyn Publication>, String> + 'static,
137    ) -> Result<BoundaryMount, String> {
138        if !self.0.alive.get() || owner.is_disposed() {
139            return Err("cannot attach a disposed async boundary".into());
140        }
141        if self.0.attached.replace(true) {
142            return Err("an async boundary can only attach once".into());
143        }
144        let inner = Rc::downgrade(&self.0);
145        let cleanup = owner.on_cleanup(move || {
146            if let Some(inner) = inner.upgrade() {
147                dispose(&inner);
148            }
149        });
150        let weak = Rc::downgrade(&self.0);
151        let owner = owner.clone();
152        let evaluate: Rc<Evaluate> = Rc::new(evaluate);
153        let subscription = effect(move || {
154            if let Some(inner) = weak
155                .upgrade()
156                .filter(|inner| inner.alive.get() && !owner.is_disposed())
157            {
158                Versions::exclude(|| {
159                    inner.wake.get();
160                });
161                drive(&inner, &*evaluate);
162            }
163        });
164        Ok(BoundaryMount {
165            inner: Rc::downgrade(&self.0),
166            subscription,
167            _cleanup: cleanup,
168        })
169    }
170}
171
172/// Retained by the renderer. Dropping it cancels candidate reads and subscriptions.
173#[doc(hidden)]
174pub struct BoundaryMount {
175    inner: Weak<Inner>,
176    subscription: Effect,
177    _cleanup: Registration,
178}
179impl Drop for BoundaryMount {
180    fn drop(&mut self) {
181        self.subscription.dispose();
182        if let Some(inner) = self.inner.upgrade() {
183            dispose(&inner);
184        }
185    }
186}
187
188fn dispose(inner: &Inner) {
189    if !inner.alive.replace(false) {
190        return;
191    }
192    let reads = inner.reads.take();
193    for read in reads.into_values() {
194        read.cancel();
195    }
196    inner.status.set(BoundaryStatus::Disposed);
197}
198
199/// Weak lifetime witness used by read declarations that outlive a mounted view.
200#[doc(hidden)]
201pub struct BoundaryLifetime(Weak<Inner>);
202impl BoundaryLifetime {
203    pub fn is_live(&self) -> bool {
204        self.0.upgrade().is_some_and(|inner| inner.alive.get())
205    }
206}
207
208/// One discovery pass in the current attempt. An unresolved read records pending
209/// and returns normally; no exception, panic, or unwinding implements suspension.
210#[doc(hidden)]
211pub struct Attempt {
212    inner: Weak<Inner>,
213    epoch: u64,
214    pending: Cell<bool>,
215    reads: RefCell<BTreeMap<u64, Rc<dyn ReadLease>>>,
216}
217impl Attempt {
218    #[doc(hidden)]
219    pub fn boundary_lifetime(&self) -> BoundaryLifetime {
220        BoundaryLifetime(self.inner.clone())
221    }
222    pub fn boundary_id(&self) -> u64 {
223        self.inner.upgrade().map_or(0, |inner| inner.id)
224    }
225    pub fn epoch(&self) -> u64 {
226        self.epoch
227    }
228    pub fn retry_generation(&self) -> u64 {
229        self.inner.upgrade().map_or(0, |inner| inner.retry.get())
230    }
231    pub fn pending(&self) {
232        self.pending.set(true);
233    }
234    pub fn register(&self, id: u64, lease: Rc<dyn ReadLease>) {
235        self.reads.borrow_mut().insert(id, lease);
236    }
237    pub fn notifier(&self) -> impl Fn() + 'static {
238        let weak = self.inner.clone();
239        let epoch = self.epoch;
240        move || {
241            if let Some(inner) = weak
242                .upgrade()
243                .filter(|inner| inner.alive.get() && inner.epoch.get() == epoch)
244            {
245                inner.wake.update(|n| *n += 1);
246            }
247        }
248    }
249}
250
251struct Evaluation;
252impl Drop for Evaluation {
253    fn drop(&mut self) {
254        EVALUATING.with(|stack| {
255            stack.borrow_mut().pop();
256        });
257        EVALUATION_DEPTH.set(EVALUATION_DEPTH.get() - 1);
258    }
259}
260struct Driving<'a>(&'a Cell<bool>);
261impl Drop for Driving<'_> {
262    fn drop(&mut self) {
263        self.0.set(false);
264    }
265}
266
267fn drive(inner: &Rc<Inner>, evaluate: &Evaluate) {
268    if inner.rejected_retry.get() == Some(inner.retry.get()) {
269        return;
270    }
271    if inner.driving.replace(true) {
272        return;
273    }
274    let _driving = Driving(&inner.driving);
275    inner.status.set(BoundaryStatus::Pending);
276    let versions = inner.inputs.borrow().clone();
277    if !untrack(|| versions.is_current()) || inner.epoch.get() == 0 {
278        inner
279            .epoch
280            .set(inner.epoch.get().checked_add(1).expect("attempt overflow"));
281        let reads = inner.reads.take();
282        for read in reads.into_values() {
283            read.cancel();
284        }
285    }
286    inner.violation.take();
287    let attempt = Attempt {
288        inner: Rc::downgrade(inner),
289        epoch: inner.epoch.get(),
290        pending: Cell::new(false),
291        reads: RefCell::new(BTreeMap::new()),
292    };
293    EVALUATING.with(|stack| stack.borrow_mut().push(Rc::downgrade(inner)));
294    EVALUATION_DEPTH.set(EVALUATION_DEPTH.get() + 1);
295    let guard = Evaluation;
296    let (result, inputs) = Versions::capture(|| evaluate(&attempt));
297    drop(guard);
298    let old = inner.reads.replace(attempt.reads.take());
299    let removed: Vec<_> = old
300        .into_iter()
301        .filter(|(id, _)| !inner.reads.borrow().contains_key(id))
302        .map(|(_, lease)| lease)
303        .collect();
304    for lease in removed {
305        lease.cancel();
306    }
307    *inner.inputs.borrow_mut() = inputs.clone();
308    if !inner.alive.get() {
309        return;
310    }
311    let result = match inner.violation.take() {
312        Some(error) => {
313            inner.rejected_retry.set(Some(inner.retry.get()));
314            Err(error)
315        }
316        None => result,
317    };
318    let mut publication = match result {
319        Ok(publication) => publication,
320        Err(error) => {
321            inner.status.set(BoundaryStatus::Error(error));
322            return;
323        }
324    };
325    if attempt.pending.get() || !untrack(|| inputs.is_current()) {
326        inner.status.set(BoundaryStatus::Pending);
327        return;
328    }
329    if let Err(error) = publication.validate() {
330        inner.status.set(BoundaryStatus::Error(error));
331        return;
332    }
333    if !inner.alive.get() || !untrack(|| inputs.is_current()) {
334        inner.status.set(BoundaryStatus::Pending);
335        return;
336    }
337    batch(|| {
338        if let Err(error) = publication.apply() {
339            inner.status.set(BoundaryStatus::Faulted(error));
340            return;
341        }
342        // Keep the guard through activation. Synchronous activation effects may
343        // invalidate these inputs and must never be overwritten by Ready.
344        publication.finish();
345        if inner.alive.get() {
346            inner.status.set(if untrack(|| inputs.is_current()) {
347                BoundaryStatus::Ready
348            } else {
349                BoundaryStatus::Pending
350            });
351        }
352    });
353}
354
355/// Record a write inside a coherent evaluation. Every signal write checks this,
356/// so the idle check inlines into callers.
357#[inline]
358pub(crate) fn mutation(operation: &str) -> bool {
359    EVALUATION_DEPTH.get() != 0 && evaluating_mutation(operation)
360}
361
362#[cold]
363fn evaluating_mutation(operation: &str) -> bool {
364    EVALUATING.with(|stack| {
365        let inner = stack.borrow().last().and_then(Weak::upgrade);
366        if let Some(inner) = inner {
367            *inner.violation.borrow_mut() = Some(format!(
368                "coherent render evaluation must be pure: {operation}"
369            ));
370            true
371        } else {
372            false
373        }
374    })
375}
376
377pub(crate) fn preparing_owner() -> Option<OwnerHandle> {
378    PREPARING.with(|owner| owner.borrow().clone())
379}
380
381#[doc(hidden)]
382pub fn prepare_state<R>(owner: OwnerHandle, make: impl FnOnce(OwnerHandle) -> R) -> R {
383    struct Restore(Option<OwnerHandle>);
384    impl Drop for Restore {
385        fn drop(&mut self) {
386            PREPARING.with(|owner| *owner.borrow_mut() = self.0.take());
387        }
388    }
389    let _restore = Restore(PREPARING.with(|previous| previous.replace(Some(owner.clone()))));
390    make(owner)
391}