1use crate::{CancellationToken, Loader, Spawner, boxed_loader, cancellation::InFlight, increment};
4use derive_where::derive_where;
5use fusor::{
6 OwnerHandle, Registration,
7 coherence::{Attempt, BoundaryLifetime, ReadLease},
8 versions::Versions,
9};
10use futures_util::future::LocalBoxFuture;
11use std::{
12 cell::{Cell, RefCell},
13 future::Future,
14 rc::{Rc, Weak},
15};
16
17thread_local! { static NEXT: Cell<u64> = const { Cell::new(0) }; }
18
19pub enum AsyncRead<T> {
22 Pending,
23 Ready(Rc<T>),
24}
25
26enum State<T, E> {
27 Idle,
28 Pending,
29 Ready(Rc<T>),
30 Error(Rc<E>),
31}
32struct Inner<K, T, E> {
33 in_flight: RefCell<Option<InFlight>>,
36 id: u64,
37 owner: OwnerHandle,
38 boundary: RefCell<Option<(u64, BoundaryLifetime)>>,
39 key: Box<dyn Fn() -> K>,
40 selected: RefCell<Option<(K, Versions)>>,
41 state: RefCell<State<T, E>>,
42 generation: Cell<u64>,
43 retry: Cell<u64>,
44 cleanup: RefCell<Option<Registration>>,
45 load: Rc<Loader<K, Result<T, E>>>,
46 spawn: Box<Spawner>,
47}
48
49impl<K, T, E> Inner<K, T, E> {
50 fn take_in_flight(&self) -> Option<InFlight> {
52 increment(&self.generation, "read generation");
53 self.in_flight.take()
54 }
55 fn cancel_work(&self) {
57 let in_flight = self.take_in_flight();
58 if matches!(*self.state.borrow(), State::Pending) {
59 *self.state.borrow_mut() = State::Idle;
60 }
61 drop(in_flight);
62 }
63 fn reset(&self) -> State<T, E> {
66 let in_flight = self.take_in_flight();
67 let retired = self.state.replace(State::Idle);
68 drop(in_flight);
69 retired
70 }
71 fn is_latest(&self, generation: u64) -> bool {
74 !self.owner.is_disposed() && self.generation.get() == generation
75 }
76 fn adopt_boundary(&self, attempt: &Attempt) -> Result<(), String> {
77 let boundary = attempt.boundary_id();
78 let (changed, live) = self
79 .boundary
80 .borrow()
81 .as_ref()
82 .map_or((true, false), |(id, lifetime)| {
83 (*id != boundary, lifetime.is_live())
84 });
85 if changed && live {
86 return Err(
87 "one AsyncValue cannot participate in different live async boundaries".into(),
88 );
89 }
90 if changed {
91 let retired = self.reset();
95 self.selected.take();
96 self.boundary
97 .replace(Some((boundary, attempt.boundary_lifetime())));
98 drop(retired);
99 }
100 Ok(())
101 }
102 fn select_key(&self) -> K
104 where
105 K: Clone + PartialEq,
106 {
107 let (key, versions) = Versions::capture(|| (self.key)());
108 let compatible = self
109 .selected
110 .borrow()
111 .as_ref()
112 .is_some_and(|(old, inputs)| old == &key && inputs.same(&versions));
113 if !compatible {
114 let retired = self.reset();
115 self.selected.replace(Some((key.clone(), versions)));
116 drop(retired);
117 }
118 key
119 }
120 fn clear_error_on_retry(&self, attempt: &Attempt) {
122 let retry = attempt.retry_generation();
123 if self.retry.get() == retry {
124 return;
125 }
126 self.retry.set(retry);
127 if matches!(*self.state.borrow(), State::Error(_)) {
128 drop(self.state.replace(State::Idle));
129 }
130 }
131 fn start_load(self: &Rc<Self>, key: K, attempt: &Attempt)
132 where
133 K: 'static,
134 T: 'static,
135 E: 'static,
136 {
137 let generation = self.generation.get();
138 *self.state.borrow_mut() = State::Pending;
139 let weak = Rc::downgrade(self);
140 let load = self.load.clone();
141 let notify = attempt.notifier();
142 let (in_flight, work) = InFlight::start(move |token| async move {
143 let result = load(key, token).await;
144 let Some(inner) = weak.upgrade().filter(|inner| inner.is_latest(generation)) else {
145 return;
146 };
147 if let Some(in_flight) = inner.in_flight.take() {
148 in_flight.complete();
149 }
150 drop(inner.state.replace(match result {
151 Ok(value) => State::Ready(Rc::new(value)),
152 Err(error) => State::Error(Rc::new(error)),
153 }));
154 notify();
155 });
156 *self.in_flight.borrow_mut() = Some(in_flight);
157 (self.spawn)(work);
158 }
159 fn outcome(&self, attempt: &Attempt) -> Result<AsyncRead<T>, String>
160 where
161 E: std::fmt::Display,
162 {
163 let error = match &*self.state.borrow() {
165 State::Ready(value) => return Ok(AsyncRead::Ready(value.clone())),
166 State::Error(error) => error.clone(),
167 State::Idle | State::Pending => {
168 attempt.pending();
169 return Ok(AsyncRead::Pending);
170 }
171 };
172 Err(error.to_string())
173 }
174}
175
176struct Lease<K, T, E>(Weak<Inner<K, T, E>>);
177impl<K, T, E> ReadLease for Lease<K, T, E> {
178 fn cancel(&self) {
179 if let Some(inner) = self.0.upgrade() {
180 inner.cancel_work();
181 }
182 }
183}
184
185#[derive_where(Clone)]
189pub struct AsyncValue<K, T, E>(Rc<Inner<K, T, E>>);
190
191impl<K: Clone + PartialEq + 'static, T: 'static, E: std::fmt::Display + 'static>
192 AsyncValue<K, T, E>
193{
194 pub fn new<F: Future<Output = Result<T, E>> + 'static>(
195 owner: &OwnerHandle,
196 key: impl Fn() -> K + 'static,
197 load: impl Fn(K, CancellationToken) -> F + 'static,
198 spawn: impl Fn(LocalBoxFuture<'static, ()>) + 'static,
199 ) -> Self {
200 let inner = Rc::new(Inner {
201 in_flight: RefCell::new(None),
202 id: NEXT.with(|next| increment(next, "read id")),
203 owner: owner.clone(),
204 boundary: RefCell::new(None),
205 key: Box::new(key),
206 selected: RefCell::new(None),
207 state: RefCell::new(State::Idle),
208 generation: Cell::new(0),
209 retry: Cell::new(0),
210 cleanup: RefCell::new(None),
211 load: boxed_loader(load),
212 spawn: Box::new(spawn),
213 });
214 let weak = Rc::downgrade(&inner);
215 let cleanup = owner.on_cleanup(move || {
216 if let Some(inner) = weak.upgrade() {
217 drop(inner.reset());
218 }
219 });
220 *inner.cleanup.borrow_mut() = Some(cleanup);
221 Self(inner)
222 }
223
224 #[doc(hidden)]
226 pub fn read(&self, attempt: &Attempt) -> Result<AsyncRead<T>, String> {
227 let inner = &self.0;
228 if inner.owner.is_disposed() {
229 return Err("async read owner was disposed".into());
230 }
231 inner.adopt_boundary(attempt)?;
232 let key = inner.select_key();
233 inner.clear_error_on_retry(attempt);
234 attempt.register(inner.id, Rc::new(Lease(Rc::downgrade(inner))));
235 if matches!(*inner.state.borrow(), State::Idle) {
236 inner.start_load(key, attempt);
237 }
238 inner.outcome(attempt)
239 }
240}