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 EVALUATION_DEPTH: Cell<usize> = const { Cell::new(0) };
62 static PREPARING: RefCell<Option<OwnerHandle>> = const { RefCell::new(None) };
63}
64
65#[derive(Clone)]
68pub struct AsyncBoundary(Rc<Inner>);
69
70impl AsyncBoundary {
71 pub fn coherent() -> Self {
72 let id = NEXT.with(|next| {
73 let id = next.get().checked_add(1).expect("boundary id overflow");
74 next.set(id);
75 id
76 });
77 Self(Rc::new(Inner {
78 id,
79 status: signal(BoundaryStatus::Detached),
80 wake: signal(0),
81 retry: Cell::new(0),
82 epoch: Cell::new(0),
83 attached: Cell::new(false),
84 alive: Cell::new(true),
85 driving: Cell::new(false),
86 violation: RefCell::new(None),
87 inputs: RefCell::new(Versions::default()),
88 rejected_retry: Cell::new(None),
89 reads: RefCell::new(BTreeMap::new()),
90 }))
91 }
92
93 pub fn status(&self) -> BoundaryStatus {
94 let cycle = EVALUATING.with(|stack| {
95 if stack
96 .borrow()
97 .iter()
98 .filter_map(Weak::upgrade)
99 .any(|inner| inner.id == self.0.id)
100 {
101 *self.0.violation.borrow_mut() =
102 Some("a coherent region cannot read its own boundary status".into());
103 true
104 } else {
105 false
106 }
107 });
108 if cycle {
109 self.0.status.get_untracked()
110 } else {
111 self.0.status.get()
112 }
113 }
114
115 pub fn pending(&self) -> bool {
116 self.status() == BoundaryStatus::Pending
117 }
118 pub fn is_interactive(&self) -> bool {
119 self.status() == BoundaryStatus::Ready
120 }
121
122 pub fn retry(&self) {
124 if self.0.alive.get() {
125 self.0
126 .retry
127 .set(self.0.retry.get().checked_add(1).expect("retry overflow"));
128 self.0.wake.update(|n| *n += 1);
129 }
130 }
131
132 #[doc(hidden)]
133 pub fn attach(
134 &self,
135 owner: &OwnerHandle,
136 evaluate: impl Fn(&Attempt) -> Result<Box<dyn Publication>, String> + 'static,
137 ) -> Result<BoundaryMount, String> {
138 if !self.0.alive.get() || owner.is_disposed() {
139 return Err("cannot attach a disposed async boundary".into());
140 }
141 if self.0.attached.replace(true) {
142 return Err("an async boundary can only attach once".into());
143 }
144 let inner = Rc::downgrade(&self.0);
145 let cleanup = owner.on_cleanup(move || {
146 if let Some(inner) = inner.upgrade() {
147 dispose(&inner);
148 }
149 });
150 let weak = Rc::downgrade(&self.0);
151 let owner = owner.clone();
152 let evaluate: Rc<Evaluate> = Rc::new(evaluate);
153 let subscription = effect(move || {
154 if let Some(inner) = weak
155 .upgrade()
156 .filter(|inner| inner.alive.get() && !owner.is_disposed())
157 {
158 Versions::exclude(|| {
159 inner.wake.get();
160 });
161 drive(&inner, &*evaluate);
162 }
163 });
164 Ok(BoundaryMount {
165 inner: Rc::downgrade(&self.0),
166 subscription,
167 _cleanup: cleanup,
168 })
169 }
170}
171
172#[doc(hidden)]
174pub struct BoundaryMount {
175 inner: Weak<Inner>,
176 subscription: Effect,
177 _cleanup: Registration,
178}
179impl Drop for BoundaryMount {
180 fn drop(&mut self) {
181 self.subscription.dispose();
182 if let Some(inner) = self.inner.upgrade() {
183 dispose(&inner);
184 }
185 }
186}
187
188fn dispose(inner: &Inner) {
189 if !inner.alive.replace(false) {
190 return;
191 }
192 let reads = inner.reads.take();
193 for read in reads.into_values() {
194 read.cancel();
195 }
196 inner.status.set(BoundaryStatus::Disposed);
197}
198
199#[doc(hidden)]
201pub struct BoundaryLifetime(Weak<Inner>);
202impl BoundaryLifetime {
203 pub fn is_live(&self) -> bool {
204 self.0.upgrade().is_some_and(|inner| inner.alive.get())
205 }
206}
207
208#[doc(hidden)]
211pub struct Attempt {
212 inner: Weak<Inner>,
213 epoch: u64,
214 pending: Cell<bool>,
215 reads: RefCell<BTreeMap<u64, Rc<dyn ReadLease>>>,
216}
217impl Attempt {
218 #[doc(hidden)]
219 pub fn boundary_lifetime(&self) -> BoundaryLifetime {
220 BoundaryLifetime(self.inner.clone())
221 }
222 pub fn boundary_id(&self) -> u64 {
223 self.inner.upgrade().map_or(0, |inner| inner.id)
224 }
225 pub fn epoch(&self) -> u64 {
226 self.epoch
227 }
228 pub fn retry_generation(&self) -> u64 {
229 self.inner.upgrade().map_or(0, |inner| inner.retry.get())
230 }
231 pub fn pending(&self) {
232 self.pending.set(true);
233 }
234 pub fn register(&self, id: u64, lease: Rc<dyn ReadLease>) {
235 self.reads.borrow_mut().insert(id, lease);
236 }
237 pub fn notifier(&self) -> impl Fn() + 'static {
238 let weak = self.inner.clone();
239 let epoch = self.epoch;
240 move || {
241 if let Some(inner) = weak
242 .upgrade()
243 .filter(|inner| inner.alive.get() && inner.epoch.get() == epoch)
244 {
245 inner.wake.update(|n| *n += 1);
246 }
247 }
248 }
249}
250
251struct Evaluation;
252impl Drop for Evaluation {
253 fn drop(&mut self) {
254 EVALUATING.with(|stack| {
255 stack.borrow_mut().pop();
256 });
257 EVALUATION_DEPTH.set(EVALUATION_DEPTH.get() - 1);
258 }
259}
260struct Driving<'a>(&'a Cell<bool>);
261impl Drop for Driving<'_> {
262 fn drop(&mut self) {
263 self.0.set(false);
264 }
265}
266
267fn drive(inner: &Rc<Inner>, evaluate: &Evaluate) {
268 if inner.rejected_retry.get() == Some(inner.retry.get()) {
269 return;
270 }
271 if inner.driving.replace(true) {
272 return;
273 }
274 let _driving = Driving(&inner.driving);
275 inner.status.set(BoundaryStatus::Pending);
276 let versions = inner.inputs.borrow().clone();
277 if !untrack(|| versions.is_current()) || inner.epoch.get() == 0 {
278 inner
279 .epoch
280 .set(inner.epoch.get().checked_add(1).expect("attempt overflow"));
281 let reads = inner.reads.take();
282 for read in reads.into_values() {
283 read.cancel();
284 }
285 }
286 inner.violation.take();
287 let attempt = Attempt {
288 inner: Rc::downgrade(inner),
289 epoch: inner.epoch.get(),
290 pending: Cell::new(false),
291 reads: RefCell::new(BTreeMap::new()),
292 };
293 EVALUATING.with(|stack| stack.borrow_mut().push(Rc::downgrade(inner)));
294 EVALUATION_DEPTH.set(EVALUATION_DEPTH.get() + 1);
295 let guard = Evaluation;
296 let (result, inputs) = Versions::capture(|| evaluate(&attempt));
297 drop(guard);
298 let old = inner.reads.replace(attempt.reads.take());
299 let removed: Vec<_> = old
300 .into_iter()
301 .filter(|(id, _)| !inner.reads.borrow().contains_key(id))
302 .map(|(_, lease)| lease)
303 .collect();
304 for lease in removed {
305 lease.cancel();
306 }
307 *inner.inputs.borrow_mut() = inputs.clone();
308 if !inner.alive.get() {
309 return;
310 }
311 let result = match inner.violation.take() {
312 Some(error) => {
313 inner.rejected_retry.set(Some(inner.retry.get()));
314 Err(error)
315 }
316 None => result,
317 };
318 let mut publication = match result {
319 Ok(publication) => publication,
320 Err(error) => {
321 inner.status.set(BoundaryStatus::Error(error));
322 return;
323 }
324 };
325 if attempt.pending.get() || !untrack(|| inputs.is_current()) {
326 inner.status.set(BoundaryStatus::Pending);
327 return;
328 }
329 if let Err(error) = publication.validate() {
330 inner.status.set(BoundaryStatus::Error(error));
331 return;
332 }
333 if !inner.alive.get() || !untrack(|| inputs.is_current()) {
334 inner.status.set(BoundaryStatus::Pending);
335 return;
336 }
337 batch(|| {
338 if let Err(error) = publication.apply() {
339 inner.status.set(BoundaryStatus::Faulted(error));
340 return;
341 }
342 publication.finish();
345 if inner.alive.get() {
346 inner.status.set(if untrack(|| inputs.is_current()) {
347 BoundaryStatus::Ready
348 } else {
349 BoundaryStatus::Pending
350 });
351 }
352 });
353}
354
355#[inline]
358pub(crate) fn mutation(operation: &str) -> bool {
359 EVALUATION_DEPTH.get() != 0 && evaluating_mutation(operation)
360}
361
362#[cold]
363fn evaluating_mutation(operation: &str) -> bool {
364 EVALUATING.with(|stack| {
365 let inner = stack.borrow().last().and_then(Weak::upgrade);
366 if let Some(inner) = inner {
367 *inner.violation.borrow_mut() = Some(format!(
368 "coherent render evaluation must be pure: {operation}"
369 ));
370 true
371 } else {
372 false
373 }
374 })
375}
376
377pub(crate) fn preparing_owner() -> Option<OwnerHandle> {
378 PREPARING.with(|owner| owner.borrow().clone())
379}
380
381#[doc(hidden)]
382pub fn prepare_state<R>(owner: OwnerHandle, make: impl FnOnce(OwnerHandle) -> R) -> R {
383 struct Restore(Option<OwnerHandle>);
384 impl Drop for Restore {
385 fn drop(&mut self) {
386 PREPARING.with(|owner| *owner.borrow_mut() = self.0.take());
387 }
388 }
389 let _restore = Restore(PREPARING.with(|previous| previous.replace(Some(owner.clone()))));
390 make(owner)
391}