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