Skip to main content

merman_core/
operation.rs

1//! Target-neutral lifecycle primitives for CPU-bound Merman operations.
2//!
3//! Parsing, analysis, SVG, ASCII, and export all need the same operation-scoped cancellation and
4//! deadline semantics. Their layout and output budgets remain adapter-owned.
5
6use 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/// The broad phase in which an operation observes a terminal condition.
35#[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/// Why an operation stopped at a cooperative checkpoint.
72#[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/// Structured cooperative cancellation returned by an operation checkpoint.
94#[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
101/// Result channel for a controlled operation stage.
102pub 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/// Cloneable operation-scoped cooperative cancellation and deadline control.
153///
154/// A synchronous callback cannot be forcefully interrupted. It observes a request when it returns
155/// to a checkpoint. Cancellation, deadlines, resource ceilings, and resource arithmetic failures
156/// compete for one sticky terminal outcome per operation scope.
157#[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    /// Creates active control using the platform monotonic clock.
171    pub fn new() -> Self {
172        Self::with_clock(default_clock())
173    }
174
175    /// Creates active control with a supplied monotonic clock.
176    ///
177    /// This remains crate-private; public callers use [`Self::with_deadline`]. Tests use the hook
178    /// to advance a fake monotonic clock without sleeping.
179    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    /// Creates an independently cancellable child that observes this control's cancellation and
187    /// deadline state.
188    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    /// Creates a shared view whose unnamed checkpoints report the supplied operation phase.
199    pub fn for_phase(&self, phase: OperationPhase) -> Self {
200        Self {
201            state: Arc::clone(&self.state),
202            default_phase: phase,
203        }
204    }
205
206    /// Sets a relative monotonic deadline if one is not already configured.
207    #[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    /// Sets a relative monotonic deadline, returning whether this call installed it.
218    #[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    /// Requests cancellation for this control and its clones.
233    pub fn cancel(&self) {
234        self.state.attention_required.store(true, Ordering::Release);
235        self.state.cancelled.store(true, Ordering::Release);
236    }
237
238    /// Returns whether this control or an ancestor has received a cancellation request.
239    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    /// Checks cancellation/deadline at an unspecified phase.
251    pub fn checkpoint(&self) -> Result<(), OperationCancelled> {
252        self.checkpoint_at(self.default_phase)
253    }
254
255    /// Checks cancellation/deadline at a named phase.
256    ///
257    /// This cancellation-only projection is used by parsing and analysis code that cannot produce
258    /// resource terminals. If a target adapter has already recorded a non-cancellation terminal,
259    /// this method does not replace it with a later cancellation. Target adapters use
260    /// [`Self::terminal_checkpoint_at`] to replay the complete terminal value.
261    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    /// Checks cancellation, deadlines, and any previously observed operation terminal.
269    ///
270    /// Target adapters call this before formal charges and controlled work. Speculative estimates
271    /// that may be discarded may use [`Self::checkpoint_at`], but they must not record a resource
272    /// terminal.
273    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    /// Records a target-owned resource ceiling as this operation's first terminal outcome.
410    pub fn terminate_resource_limit(
411        &self,
412        error: OperationResourceLimitExceeded,
413    ) -> OperationLedgerError {
414        self.latch_terminal_error(OperationLedgerError::ResourceLimitExceeded(error))
415    }
416
417    /// Records target-owned resource accounting overflow as this operation's first terminal outcome.
418    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    /// Schedules deterministic cancellation for tests after the requested successful checkpoints.
448    #[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/// Stable adapter domain that owns an operation resource terminal.
492#[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/// One explicit policy override captured when a resource terminal is first recorded.
519#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
520pub struct OperationResourceOverride {
521    pub id: &'static str,
522    pub value: u64,
523}
524
525/// Immutable policy provenance attached to the first operation resource terminal.
526#[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/// Target-neutral description of a checked operation resource rejection.
548#[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/// Sticky terminal failure shared by operation controls and target-owned resource adapters.
563#[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}