1use crate::{
8 Effect, OwnerHandle, Registration, Signal, batch, effect, signal, untrack, versions::Versions,
9};
10use std::{
11 cell::{Cell, RefCell},
12 collections::BTreeMap,
13 rc::{Rc, Weak},
14};
15
16#[doc(hidden)]
23pub const VERSION: u32 = 2;
24
25#[derive(Clone, Debug, PartialEq, Eq)]
26pub enum BoundaryStatus {
27 Detached,
28 Pending,
29 Ready,
30 Error(String),
31 Faulted(String),
32 Disposed,
33}
34
35#[doc(hidden)]
40pub trait ReadLease {
41 fn cancel(&self);
44}
45
46#[doc(hidden)]
54pub trait Publication {
55 fn validate(&self) -> Result<(), String>;
58 fn apply(&mut self) -> Result<(), String>;
64 fn finish(self: Box<Self>);
69}
70
71type Evaluate = dyn Fn(&Attempt) -> Result<Box<dyn Publication>, String>;
72
73struct Inner {
74 id: u64,
75 status: Signal<BoundaryStatus>,
76 wake: Signal<u64>,
77 retry: Cell<u64>,
78 epoch: Cell<u64>,
79 attached: Cell<bool>,
80 alive: Cell<bool>,
81 driving: Cell<bool>,
82 violation: RefCell<Option<String>>,
83 rejected_retry: Cell<Option<u64>>,
84 inputs: RefCell<Versions>,
85 reads: RefCell<BTreeMap<u64, Rc<dyn ReadLease>>>,
86 candidates: RefCell<Vec<Box<dyn FnOnce()>>>,
87}
88
89thread_local! {
90 static NEXT: Cell<u64> = const { Cell::new(0) };
91 static EVALUATING: RefCell<Vec<Weak<Inner>>> = const { RefCell::new(Vec::new()) };
92 static EVALUATION_DEPTH: Cell<usize> = const { Cell::new(0) };
94 static PREPARING: RefCell<Option<OwnerHandle>> = const { RefCell::new(None) };
95}
96
97#[derive(Clone)]
100pub struct AsyncBoundary(Rc<Inner>);
101
102impl AsyncBoundary {
103 pub fn coherent() -> Self {
104 let id = NEXT.with(|next| {
105 let id = next.get().checked_add(1).expect("boundary id overflow");
106 next.set(id);
107 id
108 });
109 Self(Rc::new(Inner {
110 id,
111 status: signal(BoundaryStatus::Detached),
112 wake: signal(0),
113 retry: Cell::new(0),
114 epoch: Cell::new(0),
115 attached: Cell::new(false),
116 alive: Cell::new(true),
117 driving: Cell::new(false),
118 violation: RefCell::new(None),
119 inputs: RefCell::new(Versions::default()),
120 rejected_retry: Cell::new(None),
121 reads: RefCell::new(BTreeMap::new()),
122 candidates: RefCell::new(Vec::new()),
123 }))
124 }
125
126 pub fn status(&self) -> BoundaryStatus {
127 let cycle = EVALUATING.with(|stack| {
128 if stack
129 .borrow()
130 .iter()
131 .filter_map(Weak::upgrade)
132 .any(|inner| inner.id == self.0.id)
133 {
134 *self.0.violation.borrow_mut() =
135 Some("a coherent region cannot read its own boundary status".into());
136 true
137 } else {
138 false
139 }
140 });
141 if cycle {
142 self.0.status.get_untracked()
143 } else {
144 self.0.status.get()
145 }
146 }
147
148 pub fn pending(&self) -> bool {
149 self.status() == BoundaryStatus::Pending
150 }
151 pub fn is_interactive(&self) -> bool {
152 self.status() == BoundaryStatus::Ready
153 }
154
155 pub fn retry(&self) {
157 if self.0.alive.get() {
158 self.0
159 .retry
160 .set(self.0.retry.get().checked_add(1).expect("retry overflow"));
161 self.0.wake.update(|n| *n += 1);
162 }
163 }
164
165 #[doc(hidden)]
177 pub fn attach(
178 &self,
179 owner: &OwnerHandle,
180 evaluate: impl Fn(&Attempt) -> Result<Box<dyn Publication>, String> + 'static,
181 ) -> Result<BoundaryMount, String> {
182 if !self.0.alive.get() || owner.is_disposed() {
183 return Err("cannot attach a disposed async boundary".into());
184 }
185 if self.0.attached.replace(true) {
186 return Err("an async boundary can only attach once".into());
187 }
188 let inner = Rc::downgrade(&self.0);
189 let cleanup = owner.on_cleanup(move || {
190 if let Some(inner) = inner.upgrade() {
191 dispose(&inner);
192 }
193 });
194 let weak = Rc::downgrade(&self.0);
195 let owner = owner.clone();
196 let evaluate: Rc<Evaluate> = Rc::new(evaluate);
197 let subscription = effect(move || {
198 if let Some(inner) = weak
199 .upgrade()
200 .filter(|inner| inner.alive.get() && !owner.is_disposed())
201 {
202 Versions::exclude(|| {
203 inner.wake.get();
204 });
205 drive(&inner, &*evaluate);
206 }
207 });
208 Ok(BoundaryMount {
209 inner: Rc::downgrade(&self.0),
210 subscription,
211 _cleanup: cleanup,
212 })
213 }
214}
215
216#[must_use = "retain the mount for as long as the coherent region is mounted"]
218#[doc(hidden)]
219pub struct BoundaryMount {
220 inner: Weak<Inner>,
221 subscription: Effect,
222 _cleanup: Registration,
223}
224impl Drop for BoundaryMount {
225 fn drop(&mut self) {
226 self.subscription.dispose();
227 if let Some(inner) = self.inner.upgrade() {
228 dispose(&inner);
229 }
230 }
231}
232
233fn dispose(inner: &Inner) {
234 if !inner.alive.replace(false) {
235 return;
236 }
237 let reads = inner.reads.take();
238 for read in reads.into_values() {
239 read.cancel();
240 }
241 discard_candidates(inner);
242 inner.status.set(BoundaryStatus::Disposed);
243}
244
245fn discard_candidates(inner: &Inner) {
246 let candidates = inner.candidates.take();
247 untrack(|| {
248 for cleanup in candidates {
249 cleanup();
250 }
251 });
252}
253
254#[doc(hidden)]
256pub struct BoundaryLifetime(Weak<Inner>);
257impl BoundaryLifetime {
258 pub fn is_live(&self) -> bool {
260 self.0.upgrade().is_some_and(|inner| inner.alive.get())
261 }
262}
263
264#[doc(hidden)]
267pub struct Attempt {
268 inner: Weak<Inner>,
269 epoch: u64,
270 pending: Cell<bool>,
271 reads: RefCell<BTreeMap<u64, Rc<dyn ReadLease>>>,
272}
273impl Attempt {
274 pub fn boundary_lifetime(&self) -> BoundaryLifetime {
276 BoundaryLifetime(self.inner.clone())
277 }
278 pub fn boundary_id(&self) -> u64 {
280 self.inner.upgrade().map_or(0, |inner| inner.id)
281 }
282 pub fn epoch(&self) -> u64 {
284 self.epoch
285 }
286 pub fn on_invalidate(&self, cleanup: impl FnOnce() + 'static) {
292 if let Some(inner) = self
293 .inner
294 .upgrade()
295 .filter(|inner| inner.alive.get() && inner.epoch.get() == self.epoch)
296 {
297 inner.candidates.borrow_mut().push(Box::new(cleanup));
298 } else {
299 untrack(cleanup);
300 }
301 }
302 pub fn retry_generation(&self) -> u64 {
304 self.inner.upgrade().map_or(0, |inner| inner.retry.get())
305 }
306 pub fn pending(&self) {
309 self.pending.set(true);
310 }
311 pub fn register(&self, id: u64, lease: Rc<dyn ReadLease>) {
315 self.reads.borrow_mut().insert(id, lease);
316 }
317 pub fn notifier(&self) -> impl Fn() + 'static {
321 let weak = self.inner.clone();
322 let epoch = self.epoch;
323 move || {
324 if let Some(inner) = weak
325 .upgrade()
326 .filter(|inner| inner.alive.get() && inner.epoch.get() == epoch)
327 {
328 inner.wake.update(|n| *n += 1);
329 }
330 }
331 }
332}
333
334struct Evaluation;
335impl Drop for Evaluation {
336 fn drop(&mut self) {
337 EVALUATING.with(|stack| {
338 stack.borrow_mut().pop();
339 });
340 EVALUATION_DEPTH.set(EVALUATION_DEPTH.get() - 1);
341 }
342}
343struct Driving<'a>(&'a Cell<bool>);
344impl Drop for Driving<'_> {
345 fn drop(&mut self) {
346 self.0.set(false);
347 }
348}
349
350fn drive(inner: &Rc<Inner>, evaluate: &Evaluate) {
351 if inner.rejected_retry.get() == Some(inner.retry.get()) {
352 return;
353 }
354 if inner.driving.replace(true) {
355 return;
356 }
357 let _driving = Driving(&inner.driving);
358 inner.status.set(BoundaryStatus::Pending);
359 let versions = inner.inputs.borrow().clone();
360 if !untrack(|| versions.is_current()) || inner.epoch.get() == 0 {
361 inner
362 .epoch
363 .set(inner.epoch.get().checked_add(1).expect("attempt overflow"));
364 let reads = inner.reads.take();
365 for read in reads.into_values() {
366 read.cancel();
367 }
368 discard_candidates(inner);
369 }
370 if !inner.alive.get() {
371 return;
372 }
373 inner.violation.take();
374 let attempt = Attempt {
375 inner: Rc::downgrade(inner),
376 epoch: inner.epoch.get(),
377 pending: Cell::new(false),
378 reads: RefCell::new(BTreeMap::new()),
379 };
380 EVALUATING.with(|stack| stack.borrow_mut().push(Rc::downgrade(inner)));
381 EVALUATION_DEPTH.set(EVALUATION_DEPTH.get() + 1);
382 let guard = Evaluation;
383 let (result, inputs) = Versions::capture(|| evaluate(&attempt));
384 drop(guard);
385 let old = inner.reads.replace(attempt.reads.take());
386 let removed: Vec<_> = old
387 .into_iter()
388 .filter(|(id, _)| !inner.reads.borrow().contains_key(id))
389 .map(|(_, lease)| lease)
390 .collect();
391 for lease in removed {
392 lease.cancel();
393 }
394 *inner.inputs.borrow_mut() = inputs.clone();
395 if !inner.alive.get() {
396 return;
397 }
398 publish_attempt(inner, &attempt, inputs, result);
399}
400
401fn publish_attempt(
402 inner: &Rc<Inner>,
403 attempt: &Attempt,
404 inputs: Versions,
405 result: Result<Box<dyn Publication>, String>,
406) {
407 let result = match inner.violation.take() {
408 Some(error) => {
409 inner.rejected_retry.set(Some(inner.retry.get()));
410 Err(error)
411 }
412 None => result,
413 };
414 let mut publication = match result {
415 Ok(publication) => publication,
416 Err(error) => {
417 inner.status.set(BoundaryStatus::Error(error));
418 return;
419 }
420 };
421 if attempt.pending.get() || !untrack(|| inputs.is_current()) {
422 inner.status.set(BoundaryStatus::Pending);
423 return;
424 }
425 if let Err(error) = publication.validate() {
426 inner.status.set(BoundaryStatus::Error(error));
427 return;
428 }
429 if !inner.alive.get() || !untrack(|| inputs.is_current()) {
430 inner.status.set(BoundaryStatus::Pending);
431 return;
432 }
433 batch(|| {
434 if let Err(error) = publication.apply() {
435 inner.status.set(BoundaryStatus::Faulted(error));
436 return;
437 }
438 publication.finish();
441 if inner.alive.get() {
442 inner.status.set(if untrack(|| inputs.is_current()) {
443 BoundaryStatus::Ready
444 } else {
445 BoundaryStatus::Pending
446 });
447 }
448 });
449}
450
451#[inline]
454pub(crate) fn mutation(operation: &str) -> bool {
455 EVALUATION_DEPTH.get() != 0 && evaluating_mutation(operation)
456}
457
458#[cold]
459fn evaluating_mutation(operation: &str) -> bool {
460 EVALUATING.with(|stack| {
461 let inner = stack.borrow().last().and_then(Weak::upgrade);
462 if let Some(inner) = inner {
463 *inner.violation.borrow_mut() = Some(format!(
464 "coherent render evaluation must be pure: {operation}"
465 ));
466 true
467 } else {
468 false
469 }
470 })
471}
472
473pub(crate) fn preparing_owner() -> Option<OwnerHandle> {
474 PREPARING.with(|owner| owner.borrow().clone())
475}
476
477#[doc(hidden)]
491pub fn prepare_state<R>(owner: OwnerHandle, make: impl FnOnce(OwnerHandle) -> R) -> R {
492 struct Restore(Option<OwnerHandle>);
493 impl Drop for Restore {
494 fn drop(&mut self) {
495 PREPARING.with(|owner| *owner.borrow_mut() = self.0.take());
496 }
497 }
498 let _restore = Restore(PREPARING.with(|previous| previous.replace(Some(owner.clone()))));
499 make(owner)
500}