1use 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#[doc(hidden)]
27pub trait ReadLease {
28 fn cancel(&self);
29}
30
31#[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#[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 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#[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#[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#[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 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}