Skip to main content

fusor_async/
lib.rs

1//! Reactive, latest-key-wins data reads with explicit ownership.
2//!
3//! The core accepts a local task spawner. Enable `browser` for the browser
4//! executor and cancellable GET adapter. There is no implicit cache or retry.
5//! Only the synchronous key function tracks signals, never the async loader.
6mod cancellation;
7mod value;
8pub use cancellation::{CancelRegistration, CancellationSource, CancellationToken};
9pub use fusor::coherence::{AsyncBoundary, BoundaryStatus};
10pub use value::{AsyncRead, AsyncValue};
11#[cfg(feature = "browser")]
12pub mod browser;
13#[cfg(feature = "browser")]
14pub mod fetch;
15
16use cancellation::InFlight;
17use derive_where::derive_where;
18use fusor::{Effect, OwnerHandle, Registration, Signal, batch, effect, signal, untrack};
19use futures_util::future::LocalBoxFuture;
20use std::{
21    cell::{Cell, RefCell},
22    future::Future,
23    rc::Rc,
24};
25
26/// A successful value and the key that actually produced it.
27#[derive(Debug)]
28#[derive_where(Clone; K)]
29pub struct Data<K, T> {
30    pub key: K,
31    pub value: Rc<T>,
32}
33
34/// Previous data is explicitly labelled; it is never passed off as a new key's result.
35#[derive(Debug)]
36#[derive_where(Clone; K)]
37pub enum ResourceState<K, T, E> {
38    Idle,
39    Loading {
40        key: K,
41        previous: Option<Data<K, T>>,
42    },
43    Ready(Data<K, T>),
44    Error {
45        key: K,
46        error: Rc<E>,
47        previous: Option<Data<K, T>>,
48    },
49    Disposed,
50}
51impl<K, T, E> ResourceState<K, T, E> {
52    pub fn data(&self) -> Option<&Data<K, T>> {
53        match self {
54            Self::Ready(data) => Some(data),
55            Self::Loading { previous, .. } | Self::Error { previous, .. } => previous.as_ref(),
56            Self::Idle | Self::Disposed => None,
57        }
58    }
59    pub fn is_loading(&self) -> bool {
60        matches!(self, Self::Loading { .. })
61    }
62}
63
64/// A type-erased async operation: an input and its cancellation token in, a
65/// local future of `O` out. Reads use `Loader<K, Result<T, E>>`.
66pub type Loader<I, O> = dyn Fn(I, CancellationToken) -> LocalBoxFuture<'static, O>;
67/// Schedules a local future on the current thread's executor, without blocking.
68pub type Spawner = dyn Fn(LocalBoxFuture<'static, ()>);
69/// Erase a loader's future type, so one loader can serve many operations.
70pub fn boxed_loader<I, O, F: Future<Output = O> + 'static>(
71    load: impl Fn(I, CancellationToken) -> F + 'static,
72) -> Rc<Loader<I, O>> {
73    Rc::new(move |input, cancel| Box::pin(load(input, cancel)))
74}
75/// Increment and return `counter`. Overflow is a bug, never a wraparound.
76fn increment(counter: &Cell<u64>, name: &str) -> u64 {
77    let next = counter
78        .get()
79        .checked_add(1)
80        .unwrap_or_else(|| panic!("{name} overflow"));
81    counter.set(next);
82    next
83}
84
85struct Inner<K, T, E> {
86    // Declared first: dropping `Inner` aborts and cancels the in-flight read
87    // before any other field drops.
88    in_flight: RefCell<Option<InFlight>>,
89    owner: OwnerHandle,
90    state: Signal<ResourceState<K, T, E>>,
91    key: RefCell<Option<K>>,
92    generation: Cell<u64>,
93    disposed: Cell<bool>,
94    subscription: RefCell<Option<Effect>>,
95    registrations: RefCell<Vec<Registration>>,
96    load: Rc<Loader<K, Result<T, E>>>,
97    spawn: Rc<Spawner>,
98}
99
100/// A shared handle to one owned read, with no `Clone` bound on data or errors.
101/// Dropping its last handle cancels the read and detaches its key subscription.
102/// Cloned handles cannot extend the owner's lifetime.
103#[derive_where(Clone)]
104pub struct Resource<K, T, E>(Rc<Inner<K, T, E>>);
105impl<K: Clone + PartialEq + 'static, T: 'static, E: 'static> Resource<K, T, E> {
106    /// `spawn` must schedule the future on a local executor. The loader is first
107    /// called when the scheduled future is polled after owner activation.
108    /// Return `None` from `key` to disable loading and clear previous data.
109    pub fn new<F: Future<Output = Result<T, E>> + 'static>(
110        owner: &OwnerHandle,
111        key: impl Fn() -> Option<K> + 'static,
112        load: impl Fn(K, CancellationToken) -> F + 'static,
113        spawn: impl Fn(LocalBoxFuture<'static, ()>) + 'static,
114    ) -> Self {
115        Self::shared(owner, key, boxed_loader(load), Rc::new(spawn))
116    }
117
118    /// Like [`Self::new`], with a loader and spawner shared by other reads,
119    /// such as the entries of one query cache.
120    pub fn shared(
121        owner: &OwnerHandle,
122        key: impl Fn() -> Option<K> + 'static,
123        load: Rc<Loader<K, Result<T, E>>>,
124        spawn: Rc<Spawner>,
125    ) -> Self {
126        let inner = Rc::new(Inner {
127            in_flight: RefCell::new(None),
128            owner: owner.clone(),
129            state: signal(ResourceState::Idle),
130            key: RefCell::new(None),
131            generation: Cell::new(0),
132            disposed: Cell::new(false),
133            subscription: RefCell::new(None),
134            registrations: RefCell::new(Vec::new()),
135            load,
136            spawn,
137        });
138        let weak = Rc::downgrade(&inner);
139        let cleanup = owner.on_cleanup(move || {
140            if let Some(inner) = weak.upgrade() {
141                inner.dispose();
142            }
143        });
144        inner.registrations.borrow_mut().push(cleanup);
145        if inner.disposed.get() {
146            return Self(inner);
147        }
148        let weak = Rc::downgrade(&inner);
149        let subscription = effect(move || {
150            let next = key();
151            if let Some(inner) = weak.upgrade() {
152                untrack(|| inner.set_key(next));
153            }
154        });
155        if inner.disposed.get() {
156            subscription.dispose();
157        } else {
158            *inner.subscription.borrow_mut() = Some(subscription);
159        }
160        let weak = Rc::downgrade(&inner);
161        // Register only while prepared: an active owner already started in the effect.
162        if !owner.is_active() {
163            let activation = owner.on_activate(move || {
164                if let Some(inner) = weak.upgrade() {
165                    inner.reload();
166                }
167            });
168            inner.registrations.borrow_mut().push(activation);
169        }
170        Self(inner)
171    }
172    pub fn get(&self) -> ResourceState<K, T, E> {
173        self.0.state.get()
174    }
175    pub fn with<R>(&self, read: impl FnOnce(&ResourceState<K, T, E>) -> R) -> R {
176        self.0.state.with(read)
177    }
178    /// Reload the current key. No-op for a disabled or disposed resource.
179    pub fn refresh(&self) {
180        untrack(|| self.0.reload());
181    }
182    /// Permanently stop this resource, even while its owner remains alive.
183    pub fn dispose(&self) {
184        self.0.dispose();
185    }
186}
187
188impl<K, T, E> Inner<K, T, E> {
189    fn invalidate(&self) -> u64 {
190        let generation = increment(&self.generation, "resource generation");
191        // Dropping the in-flight read aborts it and cancels its token.
192        drop(self.in_flight.take());
193        generation
194    }
195    fn dispose(&self) {
196        if self.disposed.replace(true) {
197            return;
198        }
199        batch(|| {
200            self.invalidate();
201            drop(self.subscription.take());
202            self.publish(ResourceState::Disposed);
203        });
204    }
205    // Prepare the complete state before borrowing, notify (or queue in a batch),
206    // then retire payloads outside the borrow. Caller-written update closures keep
207    // their non-reentrant contract; framework-owned replacement does not drop there.
208    fn publish(&self, next: ResourceState<K, T, E>) {
209        drop(self.state.replace(next));
210    }
211    fn is_stopped(&self) -> bool {
212        self.disposed.get() || self.owner.is_disposed()
213    }
214    /// Only the latest generation publishes, and only while the owner is active.
215    fn is_current(&self, generation: u64) -> bool {
216        !self.disposed.get() && self.owner.is_active() && self.generation.get() == generation
217    }
218}
219
220impl<K: Clone + 'static, T: 'static, E: 'static> Inner<K, T, E> {
221    /// Load `next` unless it equals the current key.
222    fn set_key(self: &Rc<Self>, next: Option<K>)
223    where
224        K: PartialEq,
225    {
226        if self.is_stopped() || *self.key.borrow() == next {
227            return;
228        }
229        let retired_key = self.key.replace(next);
230        self.restart();
231        // Key destructors see the completed transition, including cancellation and
232        // scheduling. Cancellation callbacks retain their existing pre-publication order.
233        drop(retired_key);
234    }
235    /// Cancel any in-flight read and load the current key again, unless stopped.
236    fn reload(self: &Rc<Self>) {
237        if !self.is_stopped() {
238            self.restart();
239        }
240    }
241    fn restart(self: &Rc<Self>) {
242        let next = self.key.borrow().clone();
243        batch(|| {
244            let generation = self.invalidate();
245            // Cancellation callbacks may dispose this resource or start another load.
246            if self.disposed.get() || self.generation.get() != generation {
247                return;
248            }
249            let Some(key) = next else {
250                self.publish(ResourceState::Idle);
251                return;
252            };
253            if self.owner.is_active() {
254                self.spawn_load(key, generation);
255            }
256        });
257    }
258    fn spawn_load(self: &Rc<Self>, key: K, generation: u64) {
259        let previous = self.state.with_untracked(|state| state.data().cloned());
260        let loading = ResourceState::Loading {
261            key: key.clone(),
262            previous: previous.clone(),
263        };
264        let weak = Rc::downgrade(self);
265        let load = self.load.clone();
266        let (in_flight, work) = InFlight::start(move |token| async move {
267            let result = load(key.clone(), token).await;
268            let Some(inner) = weak.upgrade().filter(|i| i.is_current(generation)) else {
269                return;
270            };
271            // Remove completed handles before notifying subscribers. Notifications
272            // may immediately dispose this resource or start another generation.
273            if let Some(in_flight) = inner.in_flight.take() {
274                in_flight.complete();
275            }
276            let next = match result {
277                Ok(value) => ResourceState::Ready(Data {
278                    key,
279                    value: Rc::new(value),
280                }),
281                Err(error) => ResourceState::Error {
282                    key,
283                    error: Rc::new(error),
284                    previous,
285                },
286            };
287            inner.publish(next);
288            // On success, the unused previous data is retired here,
289            // after Ready has been published, never in an update closure.
290        });
291        *self.in_flight.borrow_mut() = Some(in_flight);
292        self.publish(loading);
293        (self.spawn)(work);
294    }
295}