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