1use std::fmt;
7#[cfg(any(test, feature = "test-support"))]
8use std::sync::atomic::AtomicU64;
9use std::sync::atomic::{AtomicBool, Ordering};
10use std::sync::{Arc, OnceLock};
11
12use crate::resources::ResourceProfile;
13#[cfg(any(
14 not(all(target_arch = "wasm32", target_os = "unknown")),
15 feature = "operation-deadlines",
16 feature = "system-timing",
17 test
18))]
19use std::time::Duration;
20
21#[cfg(not(all(
22 target_arch = "wasm32",
23 target_os = "unknown",
24 any(feature = "operation-deadlines", feature = "system-timing")
25)))]
26use std::time::Instant;
27#[cfg(all(
28 target_arch = "wasm32",
29 target_os = "unknown",
30 any(feature = "operation-deadlines", feature = "system-timing")
31))]
32use web_time::Instant;
33
34#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
36#[non_exhaustive]
37pub enum OperationPhase {
38 Admission,
39 Parse,
40 Semantic,
41 Analysis,
42 Layout,
43 Emit,
44 Postprocess,
45 Export,
46 Unknown,
47}
48
49impl OperationPhase {
50 pub const fn as_str(self) -> &'static str {
51 match self {
52 Self::Admission => "admission",
53 Self::Parse => "parse",
54 Self::Semantic => "semantic",
55 Self::Analysis => "analysis",
56 Self::Layout => "layout",
57 Self::Emit => "emit",
58 Self::Postprocess => "postprocess",
59 Self::Export => "export",
60 Self::Unknown => "unknown",
61 }
62 }
63}
64
65impl fmt::Display for OperationPhase {
66 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
67 f.write_str(self.as_str())
68 }
69}
70
71#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
73pub enum CancelReason {
74 Requested,
75 DeadlineExceeded,
76}
77
78impl CancelReason {
79 pub const fn as_str(self) -> &'static str {
80 match self {
81 Self::Requested => "requested",
82 Self::DeadlineExceeded => "deadline_exceeded",
83 }
84 }
85}
86
87impl fmt::Display for CancelReason {
88 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
89 f.write_str(self.as_str())
90 }
91}
92
93#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
95#[error("operation cancelled during {phase}: {reason}")]
96pub struct OperationCancelled {
97 pub phase: OperationPhase,
98 pub reason: CancelReason,
99}
100
101pub type OperationControlResult<T> = std::result::Result<T, OperationCancelled>;
103
104type Clock = Arc<dyn Fn() -> Instant + Send + Sync + 'static>;
105
106struct OperationState {
107 attention_required: AtomicBool,
108 cancelled: AtomicBool,
109 terminal: OnceLock<OperationLedgerError>,
110 parent: Option<Arc<OperationState>>,
111 deadline: OnceLock<Instant>,
112 clock: Clock,
113 #[cfg(any(test, feature = "test-support"))]
114 successful_checkpoints_before_cancellation: AtomicU64,
115}
116
117impl fmt::Debug for OperationState {
118 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
119 f.debug_struct("OperationState")
120 .field("cancelled", &self.cancelled.load(Ordering::Acquire))
121 .field("terminal", &self.terminal.get())
122 .field("deadline", &self.deadline.get())
123 .finish_non_exhaustive()
124 }
125}
126
127impl OperationState {
128 fn new(parent: Option<Arc<OperationState>>, clock: Clock) -> Self {
129 let attention_required = parent.is_some();
130 Self {
131 attention_required: AtomicBool::new(attention_required),
132 cancelled: AtomicBool::new(false),
133 terminal: OnceLock::new(),
134 parent,
135 deadline: OnceLock::new(),
136 clock,
137 #[cfg(any(test, feature = "test-support"))]
138 successful_checkpoints_before_cancellation: AtomicU64::new(u64::MAX),
139 }
140 }
141
142 fn latch_terminal(&self, error: OperationLedgerError) -> OperationLedgerError {
143 self.attention_required.store(true, Ordering::Release);
144 self.terminal.get_or_init(|| error).clone()
145 }
146
147 fn attention_required(&self) -> bool {
148 self.attention_required.load(Ordering::Acquire)
149 }
150}
151
152#[derive(Clone, Debug)]
158pub struct OperationControl {
159 state: Arc<OperationState>,
160 default_phase: OperationPhase,
161}
162
163impl Default for OperationControl {
164 fn default() -> Self {
165 Self::new()
166 }
167}
168
169impl OperationControl {
170 pub fn new() -> Self {
172 Self::with_clock(default_clock())
173 }
174
175 fn with_clock(clock: Clock) -> Self {
180 Self {
181 state: Arc::new(OperationState::new(None, clock)),
182 default_phase: OperationPhase::Unknown,
183 }
184 }
185
186 pub fn child(&self) -> Self {
189 Self {
190 state: Arc::new(OperationState::new(
191 Some(Arc::clone(&self.state)),
192 Arc::clone(&self.state.clock),
193 )),
194 default_phase: self.default_phase,
195 }
196 }
197
198 pub fn for_phase(&self, phase: OperationPhase) -> Self {
200 Self {
201 state: Arc::clone(&self.state),
202 default_phase: phase,
203 }
204 }
205
206 #[cfg(any(
208 not(all(target_arch = "wasm32", target_os = "unknown")),
209 feature = "operation-deadlines",
210 feature = "system-timing"
211 ))]
212 pub fn with_deadline(self, timeout: Duration) -> Self {
213 self.set_deadline(timeout);
214 self
215 }
216
217 #[cfg(any(
219 not(all(target_arch = "wasm32", target_os = "unknown")),
220 feature = "operation-deadlines",
221 feature = "system-timing"
222 ))]
223 pub fn set_deadline(&self, timeout: Duration) -> bool {
224 let now = (self.state.clock)();
225 self.state.attention_required.store(true, Ordering::Release);
226 self.state
227 .deadline
228 .set(deadline_after(now, timeout))
229 .is_ok()
230 }
231
232 pub fn cancel(&self) {
234 self.state.attention_required.store(true, Ordering::Release);
235 self.state.cancelled.store(true, Ordering::Release);
236 }
237
238 pub fn is_cancelled(&self) -> bool {
240 let mut state = Some(self.state.as_ref());
241 while let Some(current) = state {
242 if current.cancelled.load(Ordering::Acquire) {
243 return true;
244 }
245 state = current.parent.as_deref();
246 }
247 false
248 }
249
250 pub fn checkpoint(&self) -> Result<(), OperationCancelled> {
252 self.checkpoint_at(self.default_phase)
253 }
254
255 pub fn checkpoint_at(&self, phase: OperationPhase) -> Result<(), OperationCancelled> {
262 if !self.state.attention_required() {
263 return Ok(());
264 }
265 self.observe_cancellation_at(phase).map_or(Ok(()), Err)
266 }
267
268 pub fn terminal_checkpoint_at(
274 &self,
275 phase: OperationPhase,
276 ) -> Result<(), OperationLedgerError> {
277 if !self.state.attention_required() {
278 return Ok(());
279 }
280 if let Some(error) = self.state.terminal.get() {
281 return Err(error.clone());
282 }
283 let Some(cancellation) = self.observe_cancellation_at(phase) else {
284 return self.state.terminal.get().cloned().map_or(Ok(()), Err);
285 };
286 Err(self.latch_terminal_error(OperationLedgerError::Cancelled(cancellation)))
287 }
288
289 fn observe_cancellation_at(&self, phase: OperationPhase) -> Option<OperationCancelled> {
290 #[cfg(any(test, feature = "test-support"))]
291 if self.consume_scheduled_checkpoint() {
292 self.cancel();
293 }
294
295 self.resolve_cancellation_at(phase)
296 }
297
298 fn resolve_cancellation_at(&self, phase: OperationPhase) -> Option<OperationCancelled> {
299 let mut state = Some(self.state.as_ref());
300 let mut inherited = false;
301 while let Some(current) = state {
302 let terminal = if let Some(error) = current.terminal.get() {
303 match error {
304 OperationLedgerError::Cancelled(error) if inherited => {
305 Some(OperationCancelled {
306 phase,
307 reason: error.reason,
308 })
309 }
310 OperationLedgerError::Cancelled(error) => Some(*error),
311 OperationLedgerError::ResourceLimitExceeded(_)
312 | OperationLedgerError::ArithmeticOverflow { .. }
313 if !inherited =>
314 {
315 return None;
316 }
317 OperationLedgerError::ResourceLimitExceeded(_)
318 | OperationLedgerError::ArithmeticOverflow { .. }
319 if current.cancelled.load(Ordering::Acquire) =>
320 {
321 Some(OperationCancelled {
322 phase,
323 reason: CancelReason::Requested,
324 })
325 }
326 OperationLedgerError::ResourceLimitExceeded(_)
327 | OperationLedgerError::ArithmeticOverflow { .. }
328 if current
329 .deadline
330 .get()
331 .is_some_and(|deadline| (current.clock)() >= *deadline) =>
332 {
333 Some(OperationCancelled {
334 phase,
335 reason: CancelReason::DeadlineExceeded,
336 })
337 }
338 OperationLedgerError::ResourceLimitExceeded(_)
339 | OperationLedgerError::ArithmeticOverflow { .. } => None,
340 }
341 } else if current.cancelled.load(Ordering::Acquire) {
342 let cancellation = OperationCancelled {
343 phase,
344 reason: CancelReason::Requested,
345 };
346 match Self::resolve_latched_cancellation(current, inherited, cancellation) {
347 Some(error) => Some(error),
348 None => return None,
349 }
350 } else if current
351 .deadline
352 .get()
353 .is_some_and(|deadline| (current.clock)() >= *deadline)
354 {
355 let cancellation = OperationCancelled {
356 phase,
357 reason: CancelReason::DeadlineExceeded,
358 };
359 match Self::resolve_latched_cancellation(current, inherited, cancellation) {
360 Some(error) => Some(error),
361 None => return None,
362 }
363 } else {
364 None
365 };
366
367 if let Some(error) = terminal {
368 if !inherited {
369 return Some(error);
370 }
371 return match self
372 .state
373 .latch_terminal(OperationLedgerError::Cancelled(error))
374 {
375 OperationLedgerError::Cancelled(error) => Some(error),
376 OperationLedgerError::ResourceLimitExceeded(_)
377 | OperationLedgerError::ArithmeticOverflow { .. } => None,
378 };
379 }
380
381 state = current.parent.as_deref();
382 inherited = true;
383 }
384 None
385 }
386
387 fn resolve_latched_cancellation(
388 state: &OperationState,
389 inherited: bool,
390 cancellation: OperationCancelled,
391 ) -> Option<OperationCancelled> {
392 match state.latch_terminal(OperationLedgerError::Cancelled(cancellation)) {
393 OperationLedgerError::Cancelled(error) => Some(error),
394 OperationLedgerError::ResourceLimitExceeded(_)
395 | OperationLedgerError::ArithmeticOverflow { .. }
396 if inherited =>
397 {
398 Some(cancellation)
399 }
400 OperationLedgerError::ResourceLimitExceeded(_)
401 | OperationLedgerError::ArithmeticOverflow { .. } => None,
402 }
403 }
404
405 fn latch_terminal_error(&self, error: OperationLedgerError) -> OperationLedgerError {
406 self.state.latch_terminal(error)
407 }
408
409 pub fn terminate_resource_limit(
411 &self,
412 error: OperationResourceLimitExceeded,
413 ) -> OperationLedgerError {
414 self.latch_terminal_error(OperationLedgerError::ResourceLimitExceeded(error))
415 }
416
417 pub fn terminate_resource_overflow(
419 &self,
420 id: &'static str,
421 phase: OperationPhase,
422 resource_phase: &'static str,
423 actual: u64,
424 maximum: u64,
425 provenance: OperationResourceProvenance,
426 ) -> OperationLedgerError {
427 self.latch_terminal_error(OperationLedgerError::ArithmeticOverflow {
428 id,
429 phase,
430 resource_phase,
431 actual,
432 maximum,
433 provenance,
434 })
435 }
436
437 #[cfg(any(test, feature = "test-support"))]
438 fn consume_scheduled_checkpoint(&self) -> bool {
439 self.state
440 .successful_checkpoints_before_cancellation
441 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |remaining| {
442 (remaining != u64::MAX).then(|| remaining.saturating_sub(1))
443 })
444 .is_ok_and(|remaining| remaining == 0)
445 }
446
447 #[cfg(any(test, feature = "test-support"))]
449 pub fn cancel_after_checkpoints(&self, successful_checkpoints: usize) {
450 self.state.attention_required.store(true, Ordering::Release);
451 self.state
452 .successful_checkpoints_before_cancellation
453 .store(successful_checkpoints as u64, Ordering::Relaxed);
454 }
455}
456
457#[cfg(any(
458 not(all(target_arch = "wasm32", target_os = "unknown")),
459 feature = "operation-deadlines",
460 feature = "system-timing"
461))]
462fn deadline_after(now: Instant, mut timeout: Duration) -> Instant {
463 loop {
464 if let Some(deadline) = now.checked_add(timeout) {
465 return deadline;
466 }
467 timeout /= 2;
468 }
469}
470
471fn default_clock() -> Clock {
472 #[cfg(not(all(
473 target_arch = "wasm32",
474 target_os = "unknown",
475 not(any(feature = "operation-deadlines", feature = "system-timing"))
476 )))]
477 {
478 Arc::new(Instant::now)
479 }
480
481 #[cfg(all(
482 target_arch = "wasm32",
483 target_os = "unknown",
484 not(any(feature = "operation-deadlines", feature = "system-timing"))
485 ))]
486 {
487 Arc::new(|| unreachable!("deadline clock is unavailable in this artifact"))
488 }
489}
490
491#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
493#[non_exhaustive]
494pub enum OperationResourceDomain {
495 Input,
496 Render,
497 Ascii,
498 Export,
499}
500
501impl OperationResourceDomain {
502 pub const fn as_str(self) -> &'static str {
503 match self {
504 Self::Input => "input",
505 Self::Render => "render",
506 Self::Ascii => "ascii",
507 Self::Export => "export",
508 }
509 }
510}
511
512impl fmt::Display for OperationResourceDomain {
513 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
514 formatter.write_str(self.as_str())
515 }
516}
517
518#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
520pub struct OperationResourceOverride {
521 pub id: &'static str,
522 pub value: u64,
523}
524
525#[derive(Debug, Clone, PartialEq, Eq)]
527pub struct OperationResourceProvenance {
528 pub domain: OperationResourceDomain,
529 pub profile: Option<ResourceProfile>,
530 pub explicit_overrides: Arc<[OperationResourceOverride]>,
531}
532
533impl OperationResourceProvenance {
534 pub fn new(
535 domain: OperationResourceDomain,
536 profile: Option<ResourceProfile>,
537 explicit_overrides: impl IntoIterator<Item = OperationResourceOverride>,
538 ) -> Self {
539 Self {
540 domain,
541 profile,
542 explicit_overrides: explicit_overrides.into_iter().collect::<Vec<_>>().into(),
543 }
544 }
545}
546
547#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
549#[error(
550 "operation resource limit `{id}` exceeded during {phase}: {consumed} + {requested} > {limit}"
551)]
552pub struct OperationResourceLimitExceeded {
553 pub id: &'static str,
554 pub phase: OperationPhase,
555 pub resource_phase: &'static str,
556 pub limit: u64,
557 pub consumed: u64,
558 pub requested: u64,
559 pub provenance: OperationResourceProvenance,
560}
561
562#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
564pub enum OperationLedgerError {
565 #[error(transparent)]
566 Cancelled(#[from] OperationCancelled),
567 #[error(transparent)]
568 ResourceLimitExceeded(#[from] OperationResourceLimitExceeded),
569 #[error(
570 "operation resource `{id}` arithmetic overflow during {phase} ({resource_phase}; actual {actual}, maximum {maximum})"
571 )]
572 ArithmeticOverflow {
573 id: &'static str,
574 phase: OperationPhase,
575 resource_phase: &'static str,
576 actual: u64,
577 maximum: u64,
578 provenance: OperationResourceProvenance,
579 },
580}
581
582#[cfg(test)]
583mod tests {
584 use super::*;
585 use std::sync::Barrier;
586 use std::sync::atomic::AtomicU64;
587 use std::thread;
588
589 fn test_resource_provenance() -> OperationResourceProvenance {
590 OperationResourceProvenance::new(
591 OperationResourceDomain::Render,
592 Some(ResourceProfile::Interactive),
593 [],
594 )
595 }
596
597 #[test]
598 fn clones_and_children_observe_shared_cancellation() {
599 let control = OperationControl::new();
600 let clone = control.clone();
601 let child = control.child();
602 child.cancel();
603 assert!(child.is_cancelled());
604 assert!(!control.is_cancelled());
605 control.cancel();
606 assert!(clone.is_cancelled());
607 assert!(child.checkpoint_at(OperationPhase::Parse).is_err());
608 }
609
610 #[test]
611 fn phase_views_share_state_and_label_unnamed_checkpoints() {
612 let control = OperationControl::new();
613 let parse = control.for_phase(OperationPhase::Parse);
614
615 control.cancel();
616
617 let error = parse.checkpoint().unwrap_err();
618 assert_eq!(error.phase, OperationPhase::Parse);
619 assert_eq!(error.reason, CancelReason::Requested);
620 }
621
622 #[test]
623 fn deadline_is_sticky_and_distinct_from_requested_cancellation() {
624 let now = Arc::new(AtomicU64::new(0));
625 let source = Arc::clone(&now);
626 let base = Instant::now();
627 let clock = Arc::new(move || base + Duration::from_millis(source.load(Ordering::Relaxed)));
628 let control = OperationControl::with_clock(clock).with_deadline(Duration::from_millis(0));
629 let error = control.checkpoint_at(OperationPhase::Layout).unwrap_err();
630 assert_eq!(error.reason, CancelReason::DeadlineExceeded);
631 assert_eq!(control.checkpoint().unwrap_err(), error);
632 assert_eq!(
633 control
634 .terminal_checkpoint_at(OperationPhase::Emit)
635 .unwrap_err(),
636 OperationLedgerError::Cancelled(error)
637 );
638 }
639
640 #[test]
641 fn child_latches_the_first_reason_observed_across_its_parent_chain() {
642 let now = Arc::new(AtomicU64::new(0));
643 let source = Arc::clone(&now);
644 let base = Instant::now();
645 let clock = Arc::new(move || base + Duration::from_millis(source.load(Ordering::Relaxed)));
646 let parent = OperationControl::with_clock(clock);
647 let child = parent.child().with_deadline(Duration::from_millis(1));
648
649 parent.cancel();
650 assert_eq!(
651 child
652 .checkpoint_at(OperationPhase::Parse)
653 .unwrap_err()
654 .reason,
655 CancelReason::Requested
656 );
657 now.store(2, Ordering::Relaxed);
658 assert_eq!(
659 child
660 .checkpoint_at(OperationPhase::Layout)
661 .unwrap_err()
662 .reason,
663 CancelReason::Requested
664 );
665 }
666
667 #[test]
668 fn child_rebinds_an_observed_parent_cancellation_to_its_checkpoint_phase() {
669 let parent = OperationControl::new();
670 parent.cancel();
671 let parent_error = parent
672 .checkpoint_at(OperationPhase::Parse)
673 .expect_err("the parent cancellation must be observed");
674 let child = parent.child();
675
676 let child_error = child
677 .checkpoint_at(OperationPhase::Layout)
678 .expect_err("the child must observe the parent cancellation");
679
680 assert_eq!(parent_error.reason, CancelReason::Requested);
681 assert_eq!(parent_error.phase, OperationPhase::Parse);
682 assert_eq!(child_error.reason, parent_error.reason);
683 assert_eq!(child_error.phase, OperationPhase::Layout);
684 assert_eq!(
685 parent.checkpoint_at(OperationPhase::Emit).unwrap_err(),
686 parent_error
687 );
688 assert_eq!(
689 child.checkpoint_at(OperationPhase::Emit).unwrap_err(),
690 child_error
691 );
692 }
693
694 #[test]
695 fn child_rebinds_an_observed_parent_deadline_to_its_checkpoint_phase() {
696 let parent = OperationControl::new().with_deadline(Duration::ZERO);
697 let parent_error = parent
698 .checkpoint_at(OperationPhase::Parse)
699 .expect_err("the parent deadline must be observed");
700 let child = parent.child();
701
702 let child_error = child
703 .checkpoint_at(OperationPhase::Layout)
704 .expect_err("the child must observe the parent deadline");
705
706 assert_eq!(parent_error.reason, CancelReason::DeadlineExceeded);
707 assert_eq!(parent_error.phase, OperationPhase::Parse);
708 assert_eq!(child_error.reason, parent_error.reason);
709 assert_eq!(child_error.phase, OperationPhase::Layout);
710 assert_eq!(
711 parent.checkpoint_at(OperationPhase::Emit).unwrap_err(),
712 parent_error
713 );
714 assert_eq!(
715 child.checkpoint_at(OperationPhase::Emit).unwrap_err(),
716 child_error
717 );
718 }
719
720 #[test]
721 fn child_keeps_its_requested_reason_after_parent_deadline_expires() {
722 let now = Arc::new(AtomicU64::new(0));
723 let source = Arc::clone(&now);
724 let base = Instant::now();
725 let clock = Arc::new(move || base + Duration::from_millis(source.load(Ordering::Relaxed)));
726 let parent = OperationControl::with_clock(clock).with_deadline(Duration::from_millis(1));
727 let child = parent.child();
728
729 child.cancel();
730 assert_eq!(
731 child
732 .checkpoint_at(OperationPhase::Parse)
733 .unwrap_err()
734 .reason,
735 CancelReason::Requested
736 );
737 now.store(2, Ordering::Relaxed);
738 assert_eq!(
739 child
740 .checkpoint_at(OperationPhase::Layout)
741 .unwrap_err()
742 .reason,
743 CancelReason::Requested
744 );
745 }
746
747 #[test]
748 fn child_latches_an_expired_parent_deadline() {
749 let now = Arc::new(AtomicU64::new(0));
750 let source = Arc::clone(&now);
751 let base = Instant::now();
752 let clock = Arc::new(move || base + Duration::from_millis(source.load(Ordering::Relaxed)));
753 let parent = OperationControl::with_clock(clock).with_deadline(Duration::from_millis(1));
754 let child = parent.child();
755
756 now.store(2, Ordering::Relaxed);
757 assert_eq!(
758 child
759 .checkpoint_at(OperationPhase::Parse)
760 .unwrap_err()
761 .reason,
762 CancelReason::DeadlineExceeded
763 );
764 parent.cancel();
765 assert_eq!(
766 child
767 .checkpoint_at(OperationPhase::Layout)
768 .unwrap_err()
769 .reason,
770 CancelReason::DeadlineExceeded
771 );
772 }
773
774 #[test]
775 fn child_keeps_an_observed_parent_cancellation_when_parent_resource_wins_the_latch() {
776 let parent = OperationControl::new();
777 let child = parent.child();
778 let cancellation = OperationCancelled {
779 phase: OperationPhase::Parse,
780 reason: CancelReason::Requested,
781 };
782 let parent_terminal = OperationLedgerError::ArithmeticOverflow {
783 id: "parent_work",
784 phase: OperationPhase::Layout,
785 resource_phase: "layout",
786 actual: u64::MAX,
787 maximum: u64::MAX,
788 provenance: test_resource_provenance(),
789 };
790
791 parent.cancel();
792 assert_eq!(
793 parent.latch_terminal_error(parent_terminal.clone()),
794 parent_terminal
795 );
796 let observed = OperationControl::resolve_latched_cancellation(
797 parent.state.as_ref(),
798 true,
799 cancellation,
800 )
801 .expect("an ancestor resource terminal must not hide an observed cancellation");
802 assert_eq!(observed, cancellation);
803 assert_eq!(
804 child.latch_terminal_error(OperationLedgerError::Cancelled(observed)),
805 OperationLedgerError::Cancelled(cancellation)
806 );
807 assert_eq!(
808 child
809 .terminal_checkpoint_at(OperationPhase::Emit)
810 .expect_err("the child must replay its cancellation"),
811 OperationLedgerError::Cancelled(cancellation)
812 );
813 assert_eq!(
814 parent
815 .terminal_checkpoint_at(OperationPhase::Emit)
816 .expect_err("the parent must retain its resource terminal"),
817 parent_terminal
818 );
819 }
820
821 #[test]
822 fn child_latches_parent_deadline_when_parent_resource_wins_during_clock_read() {
823 let now = Arc::new(AtomicU64::new(0));
824 let clock_now = Arc::clone(&now);
825 let parent_slot = Arc::new(OnceLock::<OperationControl>::new());
826 let clock_parent = Arc::clone(&parent_slot);
827 let base = Instant::now();
828 let parent_terminal = OperationLedgerError::ArithmeticOverflow {
829 id: "parent_work",
830 phase: OperationPhase::Layout,
831 resource_phase: "layout",
832 actual: u64::MAX,
833 maximum: u64::MAX,
834 provenance: test_resource_provenance(),
835 };
836 let clock_terminal = parent_terminal.clone();
837 let clock = Arc::new(move || {
838 if let Some(parent) = clock_parent.get() {
839 parent.latch_terminal_error(clock_terminal.clone());
840 }
841 base + Duration::from_millis(clock_now.load(Ordering::Relaxed))
842 });
843 let parent = OperationControl::with_clock(clock).with_deadline(Duration::from_millis(1));
844 assert!(parent_slot.set(parent.clone()).is_ok());
845 let child = parent.child();
846
847 now.store(2, Ordering::Relaxed);
848 let cancellation = child
849 .terminal_checkpoint_at(OperationPhase::Parse)
850 .expect_err("the child must observe the expired parent deadline");
851 assert_eq!(
852 cancellation,
853 OperationLedgerError::Cancelled(OperationCancelled {
854 phase: OperationPhase::Parse,
855 reason: CancelReason::DeadlineExceeded,
856 })
857 );
858 assert_eq!(
859 child
860 .terminal_checkpoint_at(OperationPhase::Emit)
861 .expect_err("the child must replay its cancellation"),
862 cancellation
863 );
864 assert_eq!(
865 parent
866 .terminal_checkpoint_at(OperationPhase::Emit)
867 .expect_err("the parent must retain its resource terminal"),
868 parent_terminal
869 );
870 }
871
872 #[test]
873 fn oversized_deadline_clamps_without_panicking_or_expiring_immediately() {
874 let control = OperationControl::new().with_deadline(Duration::MAX);
875 assert!(control.state.deadline.get().is_some());
876 assert_eq!(control.checkpoint(), Ok(()));
877 }
878
879 #[test]
880 fn concurrent_checkpoints_replay_one_complete_cancellation() {
881 const THREADS: usize = 16;
882 for _ in 0..32 {
883 let control = OperationControl::new();
884 control.cancel();
885 let barrier = Arc::new(Barrier::new(THREADS));
886 let handles = (0..THREADS)
887 .map(|index| {
888 let control = control.clone();
889 let barrier = Arc::clone(&barrier);
890 thread::spawn(move || {
891 barrier.wait();
892 let phase = if index % 2 == 0 {
893 OperationPhase::Parse
894 } else {
895 OperationPhase::Layout
896 };
897 control.checkpoint_at(phase)
898 })
899 })
900 .collect::<Vec<_>>();
901
902 let mut first = None;
903 for handle in handles {
904 let cancelled = handle
905 .join()
906 .expect("checkpoint thread should not panic")
907 .expect_err("every concurrent observer must see cancellation");
908 assert_eq!(cancelled.reason, CancelReason::Requested);
909 if let Some(first) = first {
910 assert_eq!(cancelled, first);
911 } else {
912 first = Some(cancelled);
913 }
914 }
915 }
916 }
917
918 #[test]
919 fn operation_replays_the_first_resource_rejection_after_later_cancellation() {
920 let control = OperationControl::new();
921 let first = control.terminate_resource_limit(OperationResourceLimitExceeded {
922 id: "work",
923 phase: OperationPhase::Layout,
924 resource_phase: "layout",
925 limit: 2,
926 consumed: 2,
927 requested: 1,
928 provenance: test_resource_provenance(),
929 });
930 assert_eq!(
931 first,
932 OperationLedgerError::ResourceLimitExceeded(OperationResourceLimitExceeded {
933 id: "work",
934 phase: OperationPhase::Layout,
935 resource_phase: "layout",
936 limit: 2,
937 consumed: 2,
938 requested: 1,
939 provenance: test_resource_provenance(),
940 })
941 );
942
943 let clone = control.clone();
944 let child = control.child();
945 control.cancel();
946 assert_eq!(clone.checkpoint_at(OperationPhase::Emit), Ok(()));
947 assert_eq!(
948 clone
949 .terminal_checkpoint_at(OperationPhase::Emit)
950 .expect_err("the complete checkpoint must replay the resource terminal"),
951 first
952 );
953 let cancellation = child
954 .checkpoint_at(OperationPhase::Postprocess)
955 .expect_err("a child has a distinct terminal scope and observes parent cancellation");
956 assert_eq!(cancellation.reason, CancelReason::Requested);
957 assert_eq!(cancellation.phase, OperationPhase::Postprocess);
958 assert_eq!(
959 child
960 .terminal_checkpoint_at(OperationPhase::Emit)
961 .expect_err("the child's complete checkpoint replays its cancellation"),
962 OperationLedgerError::Cancelled(cancellation)
963 );
964 assert_eq!(
965 control.terminate_resource_overflow(
966 "later_work",
967 OperationPhase::Emit,
968 "emit",
969 u64::MAX,
970 u64::MAX,
971 test_resource_provenance(),
972 ),
973 first
974 );
975 }
976
977 #[test]
978 fn child_control_starts_a_distinct_terminal_scope() {
979 let parent = OperationControl::new();
980 let parent_terminal = parent.terminate_resource_limit(OperationResourceLimitExceeded {
981 id: "parent_work",
982 phase: OperationPhase::Layout,
983 resource_phase: "layout",
984 limit: 0,
985 consumed: 0,
986 requested: 1,
987 provenance: test_resource_provenance(),
988 });
989
990 let child = parent.child();
991 assert_eq!(child.terminal_checkpoint_at(OperationPhase::Emit), Ok(()));
992 let child_terminal = child.terminate_resource_limit(OperationResourceLimitExceeded {
993 id: "child_work",
994 phase: OperationPhase::Emit,
995 resource_phase: "emit",
996 limit: 0,
997 consumed: 0,
998 requested: 1,
999 provenance: test_resource_provenance(),
1000 });
1001 assert_ne!(child_terminal, parent_terminal);
1002 assert_eq!(
1003 parent
1004 .terminal_checkpoint_at(OperationPhase::Layout)
1005 .unwrap_err(),
1006 parent_terminal
1007 );
1008 assert_eq!(
1009 child
1010 .terminal_checkpoint_at(OperationPhase::Emit)
1011 .unwrap_err(),
1012 child_terminal
1013 );
1014 }
1015
1016 #[test]
1017 fn operation_replays_cancellation_before_later_resource_failure() {
1018 let control = OperationControl::new();
1019 control.cancel();
1020
1021 let first = control
1022 .terminal_checkpoint_at(OperationPhase::Parse)
1023 .unwrap_err();
1024 assert_eq!(
1025 first,
1026 OperationLedgerError::Cancelled(OperationCancelled {
1027 phase: OperationPhase::Parse,
1028 reason: CancelReason::Requested,
1029 })
1030 );
1031 assert_eq!(
1032 control.terminate_resource_limit(OperationResourceLimitExceeded {
1033 id: "work",
1034 phase: OperationPhase::Emit,
1035 resource_phase: "emit",
1036 limit: 0,
1037 consumed: 0,
1038 requested: 1,
1039 provenance: test_resource_provenance(),
1040 }),
1041 first
1042 );
1043 assert_eq!(
1044 control
1045 .terminal_checkpoint_at(OperationPhase::Postprocess)
1046 .unwrap_err(),
1047 first
1048 );
1049 }
1050
1051 #[test]
1052 fn operation_replays_the_first_arithmetic_overflow() {
1053 let control = OperationControl::new();
1054 let first = control.terminate_resource_overflow(
1055 "layout_work",
1056 OperationPhase::Layout,
1057 "layout",
1058 u64::MAX,
1059 u64::MAX,
1060 test_resource_provenance(),
1061 );
1062 assert_eq!(
1063 first,
1064 OperationLedgerError::ArithmeticOverflow {
1065 id: "layout_work",
1066 phase: OperationPhase::Layout,
1067 resource_phase: "layout",
1068 actual: u64::MAX,
1069 maximum: u64::MAX,
1070 provenance: test_resource_provenance(),
1071 }
1072 );
1073 assert_eq!(
1074 control.terminate_resource_limit(OperationResourceLimitExceeded {
1075 id: "output_bytes",
1076 phase: OperationPhase::Emit,
1077 resource_phase: "emit",
1078 limit: 0,
1079 consumed: 0,
1080 requested: 1,
1081 provenance: test_resource_provenance(),
1082 }),
1083 first
1084 );
1085 assert_eq!(
1086 control
1087 .terminal_checkpoint_at(OperationPhase::Postprocess)
1088 .unwrap_err(),
1089 first
1090 );
1091 }
1092
1093 #[test]
1094 fn concurrent_resource_and_cancellation_observers_replay_one_terminal_error() {
1095 for _ in 0..32 {
1096 let control = OperationControl::new();
1097 let barrier = Arc::new(Barrier::new(2));
1098
1099 let resource = {
1100 let control = control.clone();
1101 let barrier = Arc::clone(&barrier);
1102 thread::spawn(move || {
1103 barrier.wait();
1104 control.terminate_resource_limit(OperationResourceLimitExceeded {
1105 id: "race_work",
1106 phase: OperationPhase::Layout,
1107 resource_phase: "layout",
1108 limit: 0,
1109 consumed: 0,
1110 requested: 1,
1111 provenance: test_resource_provenance(),
1112 })
1113 })
1114 };
1115 let cancellation = {
1116 let control = control.clone();
1117 let barrier = Arc::clone(&barrier);
1118 thread::spawn(move || {
1119 barrier.wait();
1120 control.cancel();
1121 control.terminal_checkpoint_at(OperationPhase::Emit)
1122 })
1123 };
1124
1125 let resource = resource.join().expect("resource thread should not panic");
1126 let cancellation = cancellation
1127 .join()
1128 .expect("cancellation thread should not panic")
1129 .unwrap_err();
1130 assert_eq!(resource, cancellation);
1131 assert_eq!(
1132 control
1133 .terminal_checkpoint_at(OperationPhase::Postprocess)
1134 .unwrap_err(),
1135 resource
1136 );
1137 }
1138 }
1139}