Skip to main content

ferrum_interfaces/vnext/
completion.rs

1use serde::{ser::SerializeSeq, Serialize, Serializer};
2use sha2::{Digest, Sha256};
3use std::collections::BTreeMap;
4use std::fmt;
5use std::num::NonZeroU64;
6use std::ops::Range;
7use std::panic::{catch_unwind, AssertUnwindSafe};
8use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
9use std::sync::{Arc, Mutex, MutexGuard, OnceLock, Weak};
10use std::time::Instant;
11
12use super::{
13    classify_device_error, defer_device_cleanup, translate_step_participant_readback_range,
14    AllocationKind, AllocationLifetime, BackingInitializationEncodeError, BatchOperationIdentity,
15    BatchParticipantAuthority, BufferUsage, CopyRegion, DeferredDeviceCleanupDisposition,
16    DeferredDeviceCleanupDomainId, DeferredDeviceCleanupTask, DefinitelyNotSubmittedRetryAuthority,
17    DefinitelyNotSubmittedWaveRetryAuthority, DeviceCommandBatch, DeviceDescriptor,
18    DeviceExecutionTiming, DeviceReusableExecutionPlan, DeviceReusableExecutionPreparation,
19    DeviceReusableExecutionPreparationState, DeviceReusableExecutionProgram, DeviceRuntime,
20    DeviceSubmissionExecutionTiming, DeviceSubmissionTimingSink, DeviceTerminal,
21    DeviceTerminalReceipt, DeviceTimingMeasurement, DeviceTimingMode,
22    DeviceTimingUnavailableReason, ExecutionIdentityEnvelope, ExecutionLaneId, FenceQuery,
23    HostTransferLayout, IdentifiedFailure, InvocationResourceLease, LogicalBackingBufferView,
24    NodeId, PreparedStepSubmissionWave, RequestStateHazardTerminalDisposition, ResourceId,
25    StreamState, VNextError,
26};
27
28mod readback_collection;
29pub use readback_collection::*;
30
31fn invalid_completion(reason: impl Into<String>) -> VNextError {
32    VNextError::InvalidExecutionPlan {
33        reason: reason.into(),
34    }
35}
36
37fn canonical_completion_fingerprint(value: &impl Serialize) -> String {
38    format!(
39        "{:x}",
40        Sha256::digest(
41            serde_json::to_vec(value).expect("trusted completion evidence must serialize")
42        )
43    )
44}
45
46#[derive(Debug)]
47pub enum ExecutionLaneCreationError<E> {
48    Contract(VNextError),
49    Device(E),
50}
51
52impl<E: fmt::Display> fmt::Display for ExecutionLaneCreationError<E> {
53    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
54        match self {
55            Self::Contract(error) => write!(formatter, "execution lane contract failed: {error}"),
56            Self::Device(error) => write!(formatter, "execution lane creation failed: {error}"),
57        }
58    }
59}
60
61impl<E: std::error::Error + 'static> std::error::Error for ExecutionLaneCreationError<E> {}
62
63struct ExecutionLaneState<S> {
64    stream: S,
65    in_flight: u64,
66    fail_closed: bool,
67}
68
69#[derive(Debug, Clone)]
70pub struct ExecutionLaneReusableExecutionCatalog {
71    epoch: u64,
72    programs: Vec<DeviceReusableExecutionProgram>,
73}
74
75impl ExecutionLaneReusableExecutionCatalog {
76    pub const fn epoch(&self) -> u64 {
77        self.epoch
78    }
79
80    pub fn programs(&self) -> &[DeviceReusableExecutionProgram] {
81        &self.programs
82    }
83
84    pub fn into_parts(self) -> (u64, Vec<DeviceReusableExecutionProgram>) {
85        (self.epoch, self.programs)
86    }
87}
88
89enum LaneReadbackError<E> {
90    Contract(VNextError),
91    Device(E),
92}
93
94struct LaneReadback {
95    bytes: Vec<u8>,
96    timing: DeviceTimingMeasurement<CompletionReadbackTiming>,
97}
98
99/// Scheduler-owned stream lane. It is intentionally not bound to any request
100/// or sequence and may enqueue multiple mixed-batch commands in stream order.
101#[must_use = "scheduler-owned lanes must outlive every fence enqueued on them"]
102pub struct ExecutionLane<R: DeviceRuntime> {
103    id: ExecutionLaneId,
104    runtime: Arc<R>,
105    descriptor: DeviceDescriptor,
106    fail_closed: AtomicBool,
107    reusable_execution_epoch: AtomicU64,
108    state: Mutex<ExecutionLaneState<R::Stream>>,
109}
110
111impl<R: DeviceRuntime> fmt::Debug for ExecutionLane<R> {
112    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
113        formatter
114            .debug_struct("ExecutionLane")
115            .field("id", &self.id)
116            .field("device_id", &self.descriptor.id)
117            .field(
118                "runtime_implementation_fingerprint",
119                &self.descriptor.runtime_implementation_fingerprint,
120            )
121            .finish_non_exhaustive()
122    }
123}
124
125impl<R: DeviceRuntime> ExecutionLane<R> {
126    pub fn create(runtime: Arc<R>) -> Result<Arc<Self>, ExecutionLaneCreationError<R::Error>> {
127        runtime
128            .descriptor()
129            .validate()
130            .map_err(ExecutionLaneCreationError::Contract)?;
131        let descriptor = runtime.descriptor().clone();
132        let stream = runtime
133            .create_stream()
134            .map_err(ExecutionLaneCreationError::Device)?;
135        if runtime.descriptor() != &descriptor
136            || runtime.stream_state(&stream) != StreamState::Ready
137        {
138            return Err(ExecutionLaneCreationError::Contract(invalid_completion(
139                "execution lane creation requires a stable runtime descriptor and ready stream",
140            )));
141        }
142        Ok(Arc::new(Self {
143            id: ExecutionLaneId::mint().map_err(ExecutionLaneCreationError::Contract)?,
144            runtime,
145            descriptor,
146            fail_closed: AtomicBool::new(false),
147            reusable_execution_epoch: AtomicU64::new(1),
148            state: Mutex::new(ExecutionLaneState {
149                stream,
150                in_flight: 0,
151                fail_closed: false,
152            }),
153        }))
154    }
155
156    pub fn descriptor(&self) -> &DeviceDescriptor {
157        &self.descriptor
158    }
159
160    pub const fn id(&self) -> ExecutionLaneId {
161        self.id
162    }
163
164    pub fn is_reusable(&self) -> bool {
165        if self.fail_closed.load(Ordering::Acquire) {
166            return false;
167        }
168        self.state
169            .lock()
170            .map(|state| !state.fail_closed && !self.fail_closed.load(Ordering::Acquire))
171            .unwrap_or(false)
172    }
173
174    pub fn is_fail_closed(&self) -> bool {
175        !self.is_reusable()
176    }
177
178    pub fn in_flight_count(&self) -> u64 {
179        self.state
180            .lock()
181            .map(|state| state.in_flight)
182            .unwrap_or(u64::MAX)
183    }
184
185    pub fn reusable_execution_epoch(&self) -> u64 {
186        self.reusable_execution_epoch.load(Ordering::Acquire)
187    }
188
189    fn with_quiescent_stream<T>(
190        &self,
191        operation: &'static str,
192        action: impl FnOnce(&R, &mut R::Stream) -> Result<T, R::Error>,
193    ) -> Result<T, VNextError> {
194        if self.fail_closed.load(Ordering::Acquire) {
195            return Err(invalid_completion(format!(
196                "fail-closed execution lane cannot {operation}"
197            )));
198        }
199        let mut state = self
200            .state
201            .lock()
202            .map_err(|_| invalid_completion("execution lane state mutex is poisoned"))?;
203        if state.fail_closed || self.fail_closed.load(Ordering::Acquire) {
204            return Err(invalid_completion(format!(
205                "fail-closed execution lane cannot {operation}"
206            )));
207        }
208        if state.in_flight != 0
209            || !self.current_descriptor_matches_snapshot()
210            || self.runtime.stream_state(&state.stream) != StreamState::Ready
211        {
212            return Err(invalid_completion(format!(
213                "{operation} requires a stable, quiescent execution lane"
214            )));
215        }
216        let result = catch_unwind(AssertUnwindSafe(|| {
217            action(self.runtime.as_ref(), &mut state.stream)
218        }));
219        match result {
220            Ok(Ok(value))
221                if self.current_descriptor_matches_snapshot()
222                    && self.runtime.stream_state(&state.stream) == StreamState::Ready =>
223            {
224                Ok(value)
225            }
226            Ok(Ok(_)) => {
227                state.fail_closed = true;
228                self.fail_closed.store(true, Ordering::Release);
229                Err(invalid_completion(format!(
230                    "{operation} changed the execution lane state"
231                )))
232            }
233            Ok(Err(error)) => {
234                state.fail_closed = true;
235                self.fail_closed.store(true, Ordering::Release);
236                Err(invalid_completion(format!("{operation} failed: {error}")))
237            }
238            Err(_) => {
239                state.fail_closed = true;
240                self.fail_closed.store(true, Ordering::Release);
241                Err(invalid_completion(format!(
242                    "device runtime panicked while attempting to {operation}"
243                )))
244            }
245        }
246    }
247
248    pub fn configure_reusable_executables(
249        &self,
250        plan: DeviceReusableExecutionPlan,
251    ) -> Result<DeviceReusableExecutionPreparation, VNextError> {
252        self.with_quiescent_stream("configure reusable executables", |runtime, stream| {
253            runtime.configure_reusable_executables(stream, plan)
254        })
255    }
256
257    pub fn seal_reusable_executables(
258        &self,
259    ) -> Result<DeviceReusableExecutionPreparation, VNextError> {
260        self.with_quiescent_stream("seal reusable executables", |runtime, stream| {
261            runtime.seal_reusable_executables(stream)
262        })
263    }
264
265    pub fn reusable_executable_preparation(
266        &self,
267    ) -> Result<DeviceReusableExecutionPreparation, VNextError> {
268        self.with_quiescent_stream(
269            "inspect reusable executable preparation",
270            |runtime, stream| runtime.reusable_executable_preparation(stream),
271        )
272    }
273
274    pub fn reusable_execution_catalog(
275        &self,
276    ) -> Result<ExecutionLaneReusableExecutionCatalog, VNextError> {
277        self.with_quiescent_stream("inspect reusable execution catalog", |runtime, stream| {
278            runtime.reusable_execution_catalog(stream).map(|programs| {
279                ExecutionLaneReusableExecutionCatalog {
280                    epoch: self.reusable_execution_epoch(),
281                    programs,
282                }
283            })
284        })
285    }
286
287    pub(crate) fn trim_reusable_executables_if_quiescent(&self) -> Result<bool, VNextError> {
288        if self.fail_closed.load(Ordering::Acquire) {
289            return Err(invalid_completion(
290                "fail-closed execution lane cannot trim reusable executables",
291            ));
292        }
293        let mut state = self
294            .state
295            .lock()
296            .map_err(|_| invalid_completion("execution lane state mutex is poisoned"))?;
297        if state.fail_closed || self.fail_closed.load(Ordering::Acquire) {
298            return Err(invalid_completion(
299                "fail-closed execution lane cannot trim reusable executables",
300            ));
301        }
302        if state.in_flight != 0 {
303            return Ok(false);
304        }
305        if !self.current_descriptor_matches_snapshot()
306            || self.runtime.stream_state(&state.stream) != StreamState::Ready
307        {
308            state.fail_closed = true;
309            self.fail_closed.store(true, Ordering::Release);
310            return Err(invalid_completion(
311                "reusable executable trim requires a stable, quiescent execution lane",
312            ));
313        }
314        let trimmed = catch_unwind(AssertUnwindSafe(|| {
315            let preparation = self
316                .runtime
317                .reusable_executable_preparation(&state.stream)?;
318            // A sealed resident catalog owns lane-stable backing that the
319            // immutable plan already budgets as reusable workspace. Releasing
320            // it to reclaim one idle slot invalidates the entire catalog; keep
321            // that inventory resident until the lane is torn down.
322            if preparation.state() == DeviceReusableExecutionPreparationState::Ready
323                && preparation.resident_executables() != 0
324            {
325                return Ok(None);
326            }
327            self.runtime
328                .trim_reusable_executables(&mut state.stream)
329                .map(Some)
330        }));
331        match trimmed {
332            Ok(Ok(None))
333                if self.current_descriptor_matches_snapshot()
334                    && self.runtime.stream_state(&state.stream) == StreamState::Ready =>
335            {
336                Ok(false)
337            }
338            Ok(Ok(Some(trim)))
339                if self.current_descriptor_matches_snapshot()
340                    && self.runtime.stream_state(&state.stream) == StreamState::Ready =>
341            {
342                if trim.released_executables() != 0 {
343                    let epoch = self.reusable_execution_epoch.load(Ordering::Relaxed);
344                    let Some(next_epoch) = epoch.checked_add(1) else {
345                        state.fail_closed = true;
346                        self.fail_closed.store(true, Ordering::Release);
347                        return Err(invalid_completion(
348                            "reusable execution catalog epoch is exhausted",
349                        ));
350                    };
351                    self.reusable_execution_epoch
352                        .store(next_epoch, Ordering::Release);
353                }
354                Ok(true)
355            }
356            Ok(Ok(_)) => {
357                state.fail_closed = true;
358                self.fail_closed.store(true, Ordering::Release);
359                Err(invalid_completion(
360                    "reusable executable trim changed the execution lane state",
361                ))
362            }
363            Ok(Err(error)) => {
364                state.fail_closed = true;
365                self.fail_closed.store(true, Ordering::Release);
366                Err(invalid_completion(format!(
367                    "device reusable executable inspection or trim failed: {error}"
368                )))
369            }
370            Err(_) => {
371                state.fail_closed = true;
372                self.fail_closed.store(true, Ordering::Release);
373                Err(invalid_completion(
374                    "device runtime panicked while trimming reusable executables",
375                ))
376            }
377        }
378    }
379
380    pub(crate) fn runtime(&self) -> &R {
381        &self.runtime
382    }
383
384    pub(crate) fn runtime_arc(&self) -> &Arc<R> {
385        &self.runtime
386    }
387
388    pub(crate) fn current_descriptor_matches_snapshot(&self) -> bool {
389        self.runtime.descriptor() == &self.descriptor
390    }
391
392    /// Locks only the enqueue critical section. Provider encode must complete
393    /// before this reservation is acquired.
394    pub(crate) fn reserve_enqueue(&self) -> Result<ExecutionLaneEnqueue<'_, R>, VNextError> {
395        if self.fail_closed.load(Ordering::Acquire) {
396            return Err(invalid_completion("execution lane is fail-closed"));
397        }
398        if !self.current_descriptor_matches_snapshot() {
399            return Err(invalid_completion(
400                "execution lane runtime descriptor differs from its creation snapshot",
401            ));
402        }
403        let state = self
404            .state
405            .lock()
406            .map_err(|_| invalid_completion("execution lane state mutex is poisoned"))?;
407        if self.fail_closed.load(Ordering::Acquire)
408            || state.fail_closed
409            || !matches!(
410                self.runtime.stream_state(&state.stream),
411                StreamState::Ready | StreamState::Submitted
412            )
413        {
414            return Err(invalid_completion(
415                "execution lane is failed or not enqueue-capable",
416            ));
417        }
418        Ok(ExecutionLaneEnqueue { lane: self, state })
419    }
420
421    fn query_fence(&self, fence: &R::Fence) -> FenceQuery<R::Error> {
422        self.runtime.query_fence(fence)
423    }
424
425    fn wait_fence(
426        &self,
427        fence: &R::Fence,
428    ) -> Result<DeviceTerminalReceipt<R::Error>, super::FenceIndeterminate<R::Error>> {
429        self.runtime.wait_fence(fence)
430    }
431
432    fn finish_one_terminal(&self) -> Result<(), QuiescentCompletionContractFailure> {
433        let mut state = match self.state.lock() {
434            Ok(state) => state,
435            Err(poisoned) => poisoned.into_inner(),
436        };
437        if state.in_flight == 0 {
438            state.fail_closed = true;
439            self.fail_closed.store(true, Ordering::Release);
440            return Err(QuiescentCompletionContractFailure::new(
441                "completion lane in-flight accounting underflowed",
442            ));
443        }
444        state.in_flight -= 1;
445        let descriptor_matches = catch_unwind(AssertUnwindSafe(|| {
446            self.current_descriptor_matches_snapshot()
447        }))
448        .unwrap_or(false);
449        let stream_not_failed = catch_unwind(AssertUnwindSafe(|| {
450            !matches!(
451                self.runtime.stream_state(&state.stream),
452                StreamState::Failed
453            )
454        }))
455        .unwrap_or(false);
456        if !descriptor_matches || !stream_not_failed {
457            state.fail_closed = true;
458            self.fail_closed.store(true, Ordering::Release);
459        }
460        if !descriptor_matches {
461            Err(QuiescentCompletionContractFailure::new(
462                "completion runtime descriptor drifted at terminal accounting",
463            ))
464        } else if !stream_not_failed {
465            Err(QuiescentCompletionContractFailure::new(
466                "completion lane stream entered failed state at terminal accounting",
467            ))
468        } else {
469            Ok(())
470        }
471    }
472
473    fn readback_buffer(
474        &self,
475        backing: &LogicalBackingBufferView<'_, R::Buffer>,
476        expected_usage: BufferUsage,
477        logical_offset_bytes: u64,
478        output_layout: HostTransferLayout,
479        timing_mode: DeviceTimingMode,
480    ) -> Result<LaneReadback, LaneReadbackError<R::Error>> {
481        let timing_started = timing_mode.completion_enabled().then(Instant::now);
482        let output_bytes = output_layout
483            .byte_len()
484            .map_err(LaneReadbackError::Contract)?;
485        let logical_end = logical_offset_bytes
486            .checked_add(output_bytes)
487            .ok_or_else(|| {
488                LaneReadbackError::Contract(invalid_completion(
489                    "completion readback logical range overflows u64",
490                ))
491            })?;
492        let element_bytes = output_layout.element_type().size_bytes();
493        if backing.usage() != expected_usage
494            || backing.element_type() != output_layout.element_type()
495            || logical_end > backing.size_bytes()
496            || logical_offset_bytes % element_bytes != 0
497            || output_bytes % element_bytes != 0
498        {
499            return Err(LaneReadbackError::Contract(invalid_completion(
500                "completion readback must select an aligned typed backing range with matching usage and element type",
501            )));
502        }
503        if self.fail_closed.load(Ordering::Acquire) || !self.current_descriptor_matches_snapshot() {
504            return Err(LaneReadbackError::Contract(invalid_completion(
505                "completion readback requires its original reusable execution lane",
506            )));
507        }
508        let mut state = self.state.lock().map_err(|_| {
509            LaneReadbackError::Contract(invalid_completion(
510                "execution lane state mutex is poisoned during completion readback",
511            ))
512        })?;
513        if state.fail_closed
514            || !matches!(
515                self.runtime.stream_state(&state.stream),
516                StreamState::Ready | StreamState::Submitted
517            )
518        {
519            return Err(LaneReadbackError::Contract(invalid_completion(
520                "completion readback lane is failed or not readable",
521            )));
522        }
523
524        let output_capacity = usize::try_from(output_bytes).map_err(|_| {
525            LaneReadbackError::Contract(invalid_completion(
526                "completion readback output exceeds host address space",
527            ))
528        })?;
529        let mut output = Vec::with_capacity(output_capacity);
530        let mut readback_calls = 0_u32;
531        let mut logical_cursor = 0_u64;
532        for binding in backing.segment_bindings() {
533            let segment = binding.segment();
534            let segment_logical_end = logical_cursor
535                .checked_add(segment.length_bytes())
536                .ok_or_else(|| {
537                    LaneReadbackError::Contract(invalid_completion(
538                        "completion readback backing coverage overflows u64",
539                    ))
540                })?;
541            let overlap_start = logical_cursor.max(logical_offset_bytes);
542            let overlap_end = segment_logical_end.min(logical_end);
543            if overlap_start < overlap_end {
544                let within_segment = overlap_start - logical_cursor;
545                let source_offset = segment
546                    .offset_bytes()
547                    .checked_add(within_segment)
548                    .ok_or_else(|| {
549                        LaneReadbackError::Contract(invalid_completion(
550                            "completion readback physical offset overflows u64",
551                        ))
552                    })?;
553                let length = overlap_end - overlap_start;
554                if source_offset % element_bytes != 0 || length % element_bytes != 0 {
555                    return Err(LaneReadbackError::Contract(invalid_completion(
556                        "completion readback backing segments split an element",
557                    )));
558                }
559                let actual = self.runtime.buffer_descriptor(binding.buffer());
560                if &actual != binding.descriptor()
561                    || source_offset
562                        .checked_add(length)
563                        .is_none_or(|end| end > actual.size_bytes)
564                {
565                    return Err(LaneReadbackError::Contract(invalid_completion(
566                        "completion readback backing descriptor differs from its committed extent",
567                    )));
568                }
569                let piece_layout =
570                    HostTransferLayout::new(output_layout.element_type(), length / element_bytes)
571                        .map_err(LaneReadbackError::Contract)?;
572                let region = CopyRegion::new(source_offset, 0, length)
573                    .map_err(LaneReadbackError::Contract)?;
574                readback_calls = readback_calls.checked_add(1).ok_or_else(|| {
575                    LaneReadbackError::Contract(invalid_completion(
576                        "completion readback call count exceeds u32",
577                    ))
578                })?;
579                let piece = match self.runtime.readback(
580                    &mut state.stream,
581                    binding.buffer(),
582                    region,
583                    piece_layout,
584                ) {
585                    Ok(piece) => piece,
586                    Err(error) => {
587                        state.fail_closed = true;
588                        self.fail_closed.store(true, Ordering::Release);
589                        return Err(LaneReadbackError::Device(error));
590                    }
591                };
592                if piece.len() != usize::try_from(length).unwrap_or(usize::MAX) {
593                    state.fail_closed = true;
594                    self.fail_closed.store(true, Ordering::Release);
595                    return Err(LaneReadbackError::Contract(invalid_completion(
596                        "device runtime returned an invalid completion readback byte count",
597                    )));
598                }
599                output.extend_from_slice(&piece);
600            }
601            logical_cursor = segment_logical_end;
602        }
603        if output.len() != output_capacity || !self.current_descriptor_matches_snapshot() {
604            state.fail_closed = true;
605            self.fail_closed.store(true, Ordering::Release);
606            return Err(LaneReadbackError::Contract(invalid_completion(
607                "completion readback did not cover its exact logical output",
608            )));
609        }
610        let timing = match timing_started {
611            Some(started) => {
612                let elapsed_ns = u64::try_from(started.elapsed().as_nanos()).map_err(|_| {
613                    LaneReadbackError::Contract(invalid_completion(
614                        "completion readback host duration exceeds u64 nanoseconds",
615                    ))
616                })?;
617                let bytes = u64::try_from(output.len()).map_err(|_| {
618                    LaneReadbackError::Contract(invalid_completion(
619                        "completion readback output length exceeds u64",
620                    ))
621                })?;
622                DeviceTimingMeasurement::Measured(CompletionReadbackTiming::new(
623                    elapsed_ns,
624                    readback_calls,
625                    bytes,
626                ))
627            }
628            None => DeviceTimingMeasurement::NotRequested,
629        };
630        Ok(LaneReadback {
631            bytes: output,
632            timing,
633        })
634    }
635
636    fn drain(&self, retires_submitted_fence: bool) -> bool {
637        self.fail_closed.store(true, Ordering::Release);
638        let mut state = match self.state.lock() {
639            Ok(state) => state,
640            Err(poisoned) => poisoned.into_inner(),
641        };
642        state.fail_closed = true;
643        let drained = catch_unwind(AssertUnwindSafe(|| {
644            self.runtime.synchronize(&mut state.stream).is_ok()
645                && self.current_descriptor_matches_snapshot()
646                && self.runtime.stream_state(&state.stream) == StreamState::Ready
647        }))
648        .unwrap_or(false);
649        if !drained {
650            state.fail_closed = true;
651            self.fail_closed.store(true, Ordering::Release);
652            return false;
653        }
654        // Synchronization proves lane-wide quiescence, but this call retires
655        // ownership for one exact completion record. Sibling records remain in
656        // the reaper and must account for their own terminal observations.
657        if retires_submitted_fence {
658            let Some(remaining) = state.in_flight.checked_sub(1) else {
659                state.fail_closed = true;
660                self.fail_closed.store(true, Ordering::Release);
661                return false;
662            };
663            state.in_flight = remaining;
664        }
665        true
666    }
667
668    pub(crate) fn fail_closed(&self) {
669        self.fail_closed.store(true, Ordering::Release);
670    }
671}
672
673pub(crate) enum LaneSubmitOutcome<F, E> {
674    Submitted(F),
675    DefinitelyNotSubmitted(E),
676    PossiblySubmittedPanic,
677}
678
679pub(crate) struct ExecutionLaneEnqueue<'a, R: DeviceRuntime> {
680    lane: &'a ExecutionLane<R>,
681    state: MutexGuard<'a, ExecutionLaneState<R::Stream>>,
682}
683
684impl<R: DeviceRuntime> ExecutionLaneEnqueue<'_, R> {
685    fn submit_via(
686        &mut self,
687        submit: impl FnOnce(
688            &R,
689            &mut R::Stream,
690        ) -> Result<R::Fence, super::DefinitelyNotSubmitted<R::Error>>,
691    ) -> LaneSubmitOutcome<R::Fence, R::Error> {
692        let runtime = self.lane.runtime.as_ref();
693        let stream = &mut self.state.stream;
694        let submission = catch_unwind(AssertUnwindSafe(|| submit(runtime, stream)));
695        match submission {
696            Ok(Ok(fence)) => {
697                if let Some(next) = self.state.in_flight.checked_add(1) {
698                    self.state.in_flight = next;
699                } else {
700                    self.state.fail_closed = true;
701                    self.lane.fail_closed.store(true, Ordering::Release);
702                }
703                LaneSubmitOutcome::Submitted(fence)
704            }
705            Ok(Err(not_submitted)) => {
706                LaneSubmitOutcome::DefinitelyNotSubmitted(not_submitted.into_error())
707            }
708            Err(_) => {
709                self.state.fail_closed = true;
710                self.lane.fail_closed.store(true, Ordering::Release);
711                LaneSubmitOutcome::PossiblySubmittedPanic
712            }
713        }
714    }
715
716    pub(crate) fn submit(
717        &mut self,
718        commands: super::DeviceCommandBatch<R::Command>,
719    ) -> LaneSubmitOutcome<R::Fence, R::Error> {
720        self.submit_via(|runtime, stream| runtime.submit(stream, commands))
721    }
722
723    pub(crate) fn submit_with_timing<S>(
724        &mut self,
725        commands: DeviceCommandBatch<R::Command>,
726        timing_sink: &S,
727    ) -> LaneSubmitOutcome<R::Fence, R::Error>
728    where
729        S: DeviceSubmissionTimingSink,
730    {
731        self.submit_via(|runtime, stream| runtime.submit_with_timing(stream, commands, timing_sink))
732    }
733}
734
735#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize)]
736#[serde(transparent)]
737pub struct CompletionSlotId(NonZeroU64);
738
739impl CompletionSlotId {
740    pub const fn get(self) -> u64 {
741        self.0.get()
742    }
743}
744
745#[derive(Debug)]
746struct SubmittedOperationReceiptData {
747    slot_id: CompletionSlotId,
748    batch_identity: BatchOperationIdentity,
749    participants: OnceLock<Vec<SubmittedOperationParticipantReceipt>>,
750    fingerprint: String,
751}
752
753#[derive(Debug, Clone)]
754#[must_use = "physical submission evidence must be recorded"]
755pub struct SubmittedOperationReceipt {
756    data: Arc<SubmittedOperationReceiptData>,
757}
758
759impl PartialEq for SubmittedOperationReceipt {
760    fn eq(&self, other: &Self) -> bool {
761        self.data.slot_id == other.data.slot_id
762            && self.data.batch_identity == other.data.batch_identity
763            && self.data.fingerprint == other.data.fingerprint
764    }
765}
766
767impl Eq for SubmittedOperationReceipt {}
768
769impl Serialize for SubmittedOperationReceipt {
770    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
771    where
772        S: Serializer,
773    {
774        #[derive(Serialize)]
775        struct Wire<'a> {
776            slot_id: CompletionSlotId,
777            batch_identity: &'a BatchOperationIdentity,
778            participants: &'a [SubmittedOperationParticipantReceipt],
779            fingerprint: &'a str,
780        }
781
782        Wire {
783            slot_id: self.data.slot_id,
784            batch_identity: &self.data.batch_identity,
785            participants: self.participants(),
786            fingerprint: &self.data.fingerprint,
787        }
788        .serialize(serializer)
789    }
790}
791
792#[derive(Debug, PartialEq, Eq, Serialize)]
793struct SubmittedOperationParticipantReceiptData {
794    slot_id: CompletionSlotId,
795    participant_index: u32,
796    identity: ExecutionIdentityEnvelope,
797    batch_submission_fingerprint: String,
798}
799
800#[derive(Debug, Clone, PartialEq, Eq)]
801#[must_use = "participant submission projections must remain linked to the physical batch"]
802pub struct SubmittedOperationParticipantReceipt {
803    data: Arc<SubmittedOperationParticipantReceiptData>,
804}
805
806impl SubmittedOperationParticipantReceipt {
807    fn new(
808        slot_id: CompletionSlotId,
809        participant_index: u32,
810        identity: ExecutionIdentityEnvelope,
811        batch_submission_fingerprint: String,
812    ) -> Self {
813        Self {
814            data: Arc::new(SubmittedOperationParticipantReceiptData {
815                slot_id,
816                participant_index,
817                identity,
818                batch_submission_fingerprint,
819            }),
820        }
821    }
822}
823
824impl Serialize for SubmittedOperationParticipantReceipt {
825    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
826    where
827        S: Serializer,
828    {
829        self.data.as_ref().serialize(serializer)
830    }
831}
832
833impl SubmittedOperationReceipt {
834    fn new(
835        slot_id: CompletionSlotId,
836        batch_identity: BatchOperationIdentity,
837        fingerprint: String,
838    ) -> Self {
839        Self {
840            data: Arc::new(SubmittedOperationReceiptData {
841                slot_id,
842                batch_identity,
843                participants: OnceLock::new(),
844                fingerprint,
845            }),
846        }
847    }
848
849    pub fn slot_id(&self) -> CompletionSlotId {
850        self.data.slot_id
851    }
852
853    pub fn batch_identity(&self) -> &BatchOperationIdentity {
854        &self.data.batch_identity
855    }
856
857    pub fn participants(&self) -> &[SubmittedOperationParticipantReceipt] {
858        self.data.participants.get_or_init(|| {
859            self.data
860                .batch_identity
861                .participants()
862                .iter()
863                .map(|participant| {
864                    SubmittedOperationParticipantReceipt::new(
865                        self.data.slot_id,
866                        participant.participant_index(),
867                        participant.identity().clone(),
868                        self.data.fingerprint.clone(),
869                    )
870                })
871                .collect()
872        })
873    }
874
875    #[doc(hidden)]
876    pub fn has_materialized_participant_receipts(&self) -> bool {
877        self.data.participants.get().is_some()
878    }
879
880    pub fn fingerprint(&self) -> &str {
881        &self.data.fingerprint
882    }
883}
884
885impl SubmittedOperationParticipantReceipt {
886    pub fn slot_id(&self) -> CompletionSlotId {
887        self.data.slot_id
888    }
889
890    pub fn participant_index(&self) -> u32 {
891        self.data.participant_index
892    }
893
894    pub fn identity(&self) -> &ExecutionIdentityEnvelope {
895        &self.data.identity
896    }
897
898    pub fn batch_submission_fingerprint(&self) -> &str {
899        &self.data.batch_submission_fingerprint
900    }
901}
902
903#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
904#[serde(rename_all = "snake_case", tag = "status", content = "failure")]
905pub enum OperationCompletionDisposition {
906    Succeeded,
907    FailedButQuiescent(Vec<IdentifiedFailure>),
908    ContractFailedButQuiescent(QuiescentCompletionContractFailure),
909}
910
911#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
912#[serde(rename_all = "snake_case", tag = "status", content = "failure")]
913pub enum OperationParticipantCompletionDisposition {
914    Succeeded,
915    FailedButQuiescent(IdentifiedFailure),
916    ContractFailedButQuiescent(QuiescentCompletionContractFailure),
917}
918
919#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
920#[must_use = "participant completion projections must remain linked to the physical batch"]
921pub struct OperationParticipantCompletionReceipt {
922    submission: SubmittedOperationParticipantReceipt,
923    disposition: OperationParticipantCompletionDisposition,
924    batch_completion_fingerprint: String,
925}
926
927impl OperationParticipantCompletionReceipt {
928    pub fn submission(&self) -> &SubmittedOperationParticipantReceipt {
929        &self.submission
930    }
931
932    pub fn disposition(&self) -> &OperationParticipantCompletionDisposition {
933        &self.disposition
934    }
935
936    pub fn batch_completion_fingerprint(&self) -> &str {
937        &self.batch_completion_fingerprint
938    }
939}
940
941#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
942#[must_use = "a quiescent contract failure is a terminal operation outcome"]
943pub struct QuiescentCompletionContractFailure {
944    reason: String,
945}
946
947impl QuiescentCompletionContractFailure {
948    fn new(reason: impl Into<String>) -> Self {
949        Self {
950            reason: reason.into(),
951        }
952    }
953
954    pub fn reason(&self) -> &str {
955        &self.reason
956    }
957}
958
959/// Fence timing for one exact operation completion. Device execution and host
960/// wait use different clocks and may overlap; consumers must not add them.
961#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
962pub struct CompletionFenceTiming {
963    timing_mode: DeviceTimingMode,
964    device_execution: DeviceTimingMeasurement<DeviceExecutionTiming>,
965    blocking_wait_host_ns: DeviceTimingMeasurement<u64>,
966}
967
968impl CompletionFenceTiming {
969    fn new(
970        timing_mode: DeviceTimingMode,
971        device_execution: DeviceTimingMeasurement<DeviceExecutionTiming>,
972        blocking_wait_host_ns: DeviceTimingMeasurement<u64>,
973    ) -> Self {
974        let device_execution = match (timing_mode, device_execution) {
975            (DeviceTimingMode::Off, _) => DeviceTimingMeasurement::NotRequested,
976            (
977                DeviceTimingMode::Completion
978                | DeviceTimingMode::Replay
979                | DeviceTimingMode::Kernel
980                | DeviceTimingMode::Verification,
981                DeviceTimingMeasurement::NotRequested,
982            ) => DeviceTimingMeasurement::Unavailable(
983                DeviceTimingUnavailableReason::BackendUnsupported,
984            ),
985            (
986                DeviceTimingMode::Completion
987                | DeviceTimingMode::Replay
988                | DeviceTimingMode::Kernel
989                | DeviceTimingMode::Verification,
990                measurement,
991            ) => measurement,
992        };
993        Self {
994            timing_mode,
995            device_execution,
996            blocking_wait_host_ns,
997        }
998    }
999
1000    pub const fn device_execution(self) -> DeviceTimingMeasurement<DeviceExecutionTiming> {
1001        self.device_execution
1002    }
1003
1004    pub const fn blocking_wait_host_ns(self) -> DeviceTimingMeasurement<u64> {
1005        self.blocking_wait_host_ns
1006    }
1007
1008    pub const fn timing_mode(self) -> DeviceTimingMode {
1009        self.timing_mode
1010    }
1011}
1012
1013#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1014pub struct CompletionReadbackTiming {
1015    host_elapsed_ns: u64,
1016    calls: u32,
1017    bytes: u64,
1018}
1019
1020impl CompletionReadbackTiming {
1021    const fn new(host_elapsed_ns: u64, calls: u32, bytes: u64) -> Self {
1022        Self {
1023            host_elapsed_ns,
1024            calls,
1025            bytes,
1026        }
1027    }
1028
1029    pub const fn host_elapsed_ns(self) -> u64 {
1030        self.host_elapsed_ns
1031    }
1032
1033    pub const fn calls(self) -> u32 {
1034        self.calls
1035    }
1036
1037    pub const fn bytes(self) -> u64 {
1038        self.bytes
1039    }
1040}
1041
1042#[derive(Debug, Clone)]
1043#[must_use = "a terminal completion receipt releases one exact invocation"]
1044pub struct OperationCompletionReceipt {
1045    submission: SubmittedOperationReceipt,
1046    disposition: OperationCompletionDisposition,
1047    fence_timing: CompletionFenceTiming,
1048    submission_timing: DeviceTimingMeasurement<DeviceSubmissionExecutionTiming>,
1049    participants: OnceLock<Vec<OperationParticipantCompletionReceipt>>,
1050    fingerprint: String,
1051}
1052
1053impl PartialEq for OperationCompletionReceipt {
1054    fn eq(&self, other: &Self) -> bool {
1055        self.submission == other.submission
1056            && self.disposition == other.disposition
1057            && self.fence_timing == other.fence_timing
1058            && self.submission_timing == other.submission_timing
1059            && self.fingerprint == other.fingerprint
1060    }
1061}
1062
1063impl Eq for OperationCompletionReceipt {}
1064
1065impl Serialize for OperationCompletionReceipt {
1066    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1067    where
1068        S: Serializer,
1069    {
1070        #[derive(Serialize)]
1071        struct Wire<'a> {
1072            submission: &'a SubmittedOperationReceipt,
1073            disposition: &'a OperationCompletionDisposition,
1074            fence_timing: CompletionFenceTiming,
1075            submission_timing: &'a DeviceTimingMeasurement<DeviceSubmissionExecutionTiming>,
1076            participants: &'a [OperationParticipantCompletionReceipt],
1077            fingerprint: &'a str,
1078        }
1079
1080        Wire {
1081            submission: &self.submission,
1082            disposition: &self.disposition,
1083            fence_timing: self.fence_timing,
1084            submission_timing: &self.submission_timing,
1085            participants: self.participants(),
1086            fingerprint: &self.fingerprint,
1087        }
1088        .serialize(serializer)
1089    }
1090}
1091
1092impl OperationCompletionReceipt {
1093    fn new(
1094        submission: SubmittedOperationReceipt,
1095        disposition: OperationCompletionDisposition,
1096        fence_timing: CompletionFenceTiming,
1097        submission_timing: DeviceTimingMeasurement<DeviceSubmissionExecutionTiming>,
1098    ) -> Result<Self, VNextError> {
1099        match &disposition {
1100            OperationCompletionDisposition::Succeeded
1101            | OperationCompletionDisposition::ContractFailedButQuiescent(_) => {}
1102            OperationCompletionDisposition::FailedButQuiescent(failures) => {
1103                if failures.len() != submission.participants().len()
1104                    || failures
1105                        .iter()
1106                        .zip(submission.participants())
1107                        .any(|(failure, participant)| failure.identity() != participant.identity())
1108                {
1109                    return Err(invalid_completion(
1110                        "batch completion failures differ from participant submission projections",
1111                    ));
1112                }
1113            }
1114        }
1115        #[derive(Serialize)]
1116        struct CompletionFingerprintInput<'a> {
1117            domain: &'static str,
1118            submission_fingerprint: &'a str,
1119            disposition: &'a OperationCompletionDisposition,
1120        }
1121        let fingerprint = canonical_completion_fingerprint(&CompletionFingerprintInput {
1122            domain: "ferrum.runtime-vnext.batch-operation-completion.v1",
1123            submission_fingerprint: submission.fingerprint(),
1124            disposition: &disposition,
1125        });
1126        Ok(Self {
1127            submission,
1128            disposition,
1129            fence_timing,
1130            submission_timing,
1131            participants: OnceLock::new(),
1132            fingerprint,
1133        })
1134    }
1135
1136    pub fn submission(&self) -> &SubmittedOperationReceipt {
1137        &self.submission
1138    }
1139
1140    pub fn disposition(&self) -> &OperationCompletionDisposition {
1141        &self.disposition
1142    }
1143
1144    pub const fn fence_timing(&self) -> CompletionFenceTiming {
1145        self.fence_timing
1146    }
1147
1148    pub const fn submission_timing(
1149        &self,
1150    ) -> &DeviceTimingMeasurement<DeviceSubmissionExecutionTiming> {
1151        &self.submission_timing
1152    }
1153
1154    pub fn participants(&self) -> &[OperationParticipantCompletionReceipt] {
1155        self.participants.get_or_init(|| {
1156            self.submission
1157                .participants()
1158                .iter()
1159                .enumerate()
1160                .map(|(index, submission)| {
1161                    let disposition = match &self.disposition {
1162                        OperationCompletionDisposition::Succeeded => {
1163                            OperationParticipantCompletionDisposition::Succeeded
1164                        }
1165                        OperationCompletionDisposition::FailedButQuiescent(failures) => {
1166                            OperationParticipantCompletionDisposition::FailedButQuiescent(
1167                                failures[index].clone(),
1168                            )
1169                        }
1170                        OperationCompletionDisposition::ContractFailedButQuiescent(failure) => {
1171                            OperationParticipantCompletionDisposition::ContractFailedButQuiescent(
1172                                failure.clone(),
1173                            )
1174                        }
1175                    };
1176                    OperationParticipantCompletionReceipt {
1177                        submission: submission.clone(),
1178                        disposition,
1179                        batch_completion_fingerprint: self.fingerprint.clone(),
1180                    }
1181                })
1182                .collect()
1183        })
1184    }
1185
1186    pub fn fingerprint(&self) -> &str {
1187        &self.fingerprint
1188    }
1189}
1190
1191#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1192#[serde(rename_all = "snake_case")]
1193pub enum CompletionRecoveryCause {
1194    SubmissionIndeterminate,
1195    FenceIndeterminate,
1196    FenceObservationPanicked,
1197}
1198
1199#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
1200#[must_use = "a lane-drain receipt proves quiescence for one exact completion slot"]
1201pub struct CompletionDrainReceipt {
1202    slot_id: CompletionSlotId,
1203    batch_identity: BatchOperationIdentity,
1204    submission: Option<SubmittedOperationReceipt>,
1205    cause: CompletionRecoveryCause,
1206    had_submission_fence: bool,
1207}
1208
1209impl CompletionDrainReceipt {
1210    pub const fn slot_id(&self) -> CompletionSlotId {
1211        self.slot_id
1212    }
1213
1214    pub fn batch_identity(&self) -> &BatchOperationIdentity {
1215        &self.batch_identity
1216    }
1217
1218    pub fn submission(&self) -> Option<&SubmittedOperationReceipt> {
1219        self.submission.as_ref()
1220    }
1221
1222    pub const fn cause(&self) -> CompletionRecoveryCause {
1223        self.cause
1224    }
1225
1226    pub const fn had_submission_fence(&self) -> bool {
1227        self.had_submission_fence
1228    }
1229}
1230
1231#[derive(Debug, Clone)]
1232struct CompletionQuarantineFreshness(Arc<AtomicBool>);
1233
1234impl CompletionQuarantineFreshness {
1235    fn current() -> Self {
1236        Self(Arc::new(AtomicBool::new(true)))
1237    }
1238
1239    fn is_current(&self) -> bool {
1240        self.0.load(Ordering::Acquire)
1241    }
1242
1243    fn invalidate(&self) {
1244        self.0.store(false, Ordering::Release);
1245    }
1246}
1247
1248#[derive(Debug, Clone, Serialize)]
1249#[must_use = "a quarantine receipt identifies retained device-visible ownership"]
1250pub struct CompletionQuarantineReceipt {
1251    slot_id: CompletionSlotId,
1252    batch_identity: BatchOperationIdentity,
1253    submission: Option<SubmittedOperationReceipt>,
1254    cause: CompletionRecoveryCause,
1255    had_submission_fence: bool,
1256    device_id: super::DeviceId,
1257    runtime_implementation_fingerprint: String,
1258    #[serde(skip)]
1259    freshness: CompletionQuarantineFreshness,
1260}
1261
1262impl PartialEq for CompletionQuarantineReceipt {
1263    fn eq(&self, other: &Self) -> bool {
1264        self.slot_id == other.slot_id
1265            && self.batch_identity == other.batch_identity
1266            && self.submission == other.submission
1267            && self.cause == other.cause
1268            && self.had_submission_fence == other.had_submission_fence
1269            && self.device_id == other.device_id
1270            && self.runtime_implementation_fingerprint == other.runtime_implementation_fingerprint
1271    }
1272}
1273
1274impl Eq for CompletionQuarantineReceipt {}
1275
1276impl CompletionQuarantineReceipt {
1277    pub const fn slot_id(&self) -> CompletionSlotId {
1278        self.slot_id
1279    }
1280
1281    pub fn batch_identity(&self) -> &BatchOperationIdentity {
1282        &self.batch_identity
1283    }
1284
1285    pub fn submission(&self) -> Option<&SubmittedOperationReceipt> {
1286        self.submission.as_ref()
1287    }
1288
1289    pub const fn cause(&self) -> CompletionRecoveryCause {
1290        self.cause
1291    }
1292
1293    pub const fn had_submission_fence(&self) -> bool {
1294        self.had_submission_fence
1295    }
1296
1297    pub fn device_id(&self) -> &super::DeviceId {
1298        &self.device_id
1299    }
1300
1301    pub fn runtime_implementation_fingerprint(&self) -> &str {
1302        &self.runtime_implementation_fingerprint
1303    }
1304
1305    pub fn is_current(&self) -> bool {
1306        self.freshness.is_current()
1307    }
1308}
1309
1310#[derive(Debug)]
1311#[must_use = "completion recovery must be recorded as drained or quarantined"]
1312pub enum CompletionRecoveryOutcome {
1313    Drained(CompletionDrainReceipt),
1314    Quarantined(CompletionQuarantineReceipt),
1315}
1316
1317#[derive(Debug)]
1318pub enum CompletionSweepObservation {
1319    Observed(CompletionObservation),
1320    Failed(VNextError),
1321}
1322
1323#[derive(Debug)]
1324pub struct CompletionSweepEntry {
1325    slot_id: CompletionSlotId,
1326    observation: CompletionSweepObservation,
1327}
1328
1329impl CompletionSweepEntry {
1330    pub const fn slot_id(&self) -> CompletionSlotId {
1331        self.slot_id
1332    }
1333
1334    pub fn observation(&self) -> &CompletionSweepObservation {
1335        &self.observation
1336    }
1337
1338    pub fn into_observation(self) -> CompletionSweepObservation {
1339        self.observation
1340    }
1341}
1342
1343#[derive(Debug)]
1344#[must_use = "a bounded completion sweep contains scheduler-owned progress evidence"]
1345pub struct CompletionSweepReceipt {
1346    entries: Vec<CompletionSweepEntry>,
1347    retained_after: usize,
1348    quarantined_after: usize,
1349}
1350
1351impl CompletionSweepReceipt {
1352    pub fn entries(&self) -> &[CompletionSweepEntry] {
1353        &self.entries
1354    }
1355
1356    pub fn into_entries(self) -> Vec<CompletionSweepEntry> {
1357        self.entries
1358    }
1359
1360    pub const fn retained_after(&self) -> usize {
1361        self.retained_after
1362    }
1363
1364    pub const fn quarantined_after(&self) -> usize {
1365        self.quarantined_after
1366    }
1367}
1368
1369enum CompletionResourceLease<R: DeviceRuntime> {
1370    Invocation(InvocationResourceLease<R>),
1371    Wave(PreparedStepSubmissionWave<R>),
1372}
1373
1374impl<R: DeviceRuntime> CompletionResourceLease<R> {
1375    fn runtime(&self) -> &Arc<R> {
1376        match self {
1377            Self::Invocation(invocation) => invocation.runtime(),
1378            Self::Wave(wave) => wave.runtime(),
1379        }
1380    }
1381
1382    fn deferred_cleanup_domain(&self) -> DeferredDeviceCleanupDomainId {
1383        match self {
1384            Self::Invocation(invocation) => invocation.deferred_cleanup_domain(),
1385            Self::Wave(wave) => wave.deferred_cleanup_domain(),
1386        }
1387    }
1388
1389    fn mark_submission_fence_installed(&mut self) -> Result<(), VNextError> {
1390        match self {
1391            Self::Invocation(invocation) => invocation.mark_submission_fence_installed(),
1392            Self::Wave(wave) => wave.mark_submission_fence_installed(),
1393        }
1394    }
1395
1396    fn encode_backing_initializations(
1397        &self,
1398        runtime: &R,
1399        commands: &mut DeviceCommandBatch<R::Command>,
1400    ) -> Result<usize, BackingInitializationEncodeError<R::Error>> {
1401        match self {
1402            Self::Invocation(invocation) => {
1403                invocation.encode_backing_initializations(runtime, commands)
1404            }
1405            Self::Wave(wave) => wave.encode_backing_initializations(runtime, commands),
1406        }
1407    }
1408
1409    fn mark_submission_indeterminate(&mut self) {
1410        match self {
1411            Self::Invocation(invocation) => invocation.mark_submission_indeterminate(),
1412            Self::Wave(wave) => wave.mark_submission_indeterminate(),
1413        }
1414    }
1415
1416    fn finish_backing_initializations(&mut self, succeeded: bool) -> Result<(), VNextError> {
1417        match self {
1418            Self::Invocation(invocation) => invocation.finish_backing_initializations(succeeded),
1419            Self::Wave(wave) => wave.finish_backing_initializations(succeeded),
1420        }
1421    }
1422
1423    fn finish_request_state_hazards(
1424        &mut self,
1425        disposition: RequestStateHazardTerminalDisposition,
1426    ) -> Result<(), VNextError> {
1427        match self {
1428            Self::Invocation(invocation) => invocation.finish_request_state_hazards(disposition),
1429            Self::Wave(wave) => wave.finish_request_state_hazards(disposition),
1430        }
1431    }
1432
1433    fn backing_view(
1434        &self,
1435        node_id: &NodeId,
1436        participant_index: u32,
1437        resource_id: &ResourceId,
1438    ) -> Result<LogicalBackingBufferView<'_, R::Buffer>, VNextError> {
1439        let participant_index = usize::try_from(participant_index).map_err(|_| {
1440            invalid_completion("completion readback participant index exceeds host address space")
1441        })?;
1442        match self {
1443            Self::Invocation(invocation) => {
1444                if invocation.node_id() != node_id {
1445                    return Err(invalid_completion(
1446                        "completion readback node differs from its invocation",
1447                    ));
1448                }
1449                let is_shared = invocation
1450                    .backing_slices()
1451                    .iter()
1452                    .chain(invocation.step_resources().backing_slices())
1453                    .any(|authority| authority.resource_id() == resource_id);
1454                if is_shared {
1455                    return invocation.backing_view(resource_id);
1456                }
1457                let participant = invocation
1458                    .participants()
1459                    .nth(participant_index)
1460                    .ok_or_else(|| {
1461                        invalid_completion(
1462                            "completion readback participant is absent from its invocation",
1463                        )
1464                    })?;
1465                invocation.step_resources().participant_backing_view(
1466                    BatchParticipantAuthority::new(
1467                        participant.sequence_authority(),
1468                        participant.request_authority(),
1469                    ),
1470                    resource_id,
1471                )
1472            }
1473            Self::Wave(wave) => {
1474                let node_index = wave
1475                    .nodes()
1476                    .iter()
1477                    .position(|node| node.node_id() == node_id)
1478                    .ok_or_else(|| {
1479                        invalid_completion("completion readback node is absent from its wave")
1480                    })?;
1481                let is_shared = wave
1482                    .claimed_backing()
1483                    .backing_slices()
1484                    .iter()
1485                    .chain(wave.step_resources().backing_slices())
1486                    .any(|authority| authority.resource_id() == resource_id);
1487                if is_shared {
1488                    return wave.backing_view(node_index, resource_id);
1489                }
1490                let participant = wave.nodes()[node_index]
1491                    .participants()
1492                    .nth(participant_index)
1493                    .ok_or_else(|| {
1494                        invalid_completion(
1495                            "completion readback participant is absent from its wave node",
1496                        )
1497                    })?;
1498                wave.step_resources().participant_backing_view(
1499                    BatchParticipantAuthority::new(
1500                        participant.sequence_authority(),
1501                        participant.request_authority(),
1502                    ),
1503                    resource_id,
1504                )
1505            }
1506        }
1507    }
1508
1509    fn readback_range(
1510        &self,
1511        node_id: &NodeId,
1512        participant_index: u32,
1513        resource_id: &ResourceId,
1514        semantic_range: Range<u64>,
1515    ) -> Result<Range<u64>, VNextError> {
1516        let participant_index = usize::try_from(participant_index).map_err(|_| {
1517            invalid_completion("completion readback participant index exceeds host address space")
1518        })?;
1519        let (step, work_shape) = match self {
1520            Self::Invocation(invocation) => {
1521                if invocation.node_id() != node_id {
1522                    return Err(invalid_completion(
1523                        "completion readback node differs from its invocation",
1524                    ));
1525                }
1526                (invocation.step_resources(), invocation.work_shape())
1527            }
1528            Self::Wave(wave) => {
1529                let node = wave
1530                    .nodes()
1531                    .iter()
1532                    .find(|node| node.node_id() == node_id)
1533                    .ok_or_else(|| {
1534                        invalid_completion("completion readback node is absent from its wave")
1535                    })?;
1536                (wave.step_resources(), node.work_shape())
1537            }
1538        };
1539        let descriptor = step.dynamic_descriptor(resource_id)?;
1540        if descriptor.lifetime() == AllocationLifetime::Step
1541            && descriptor.kind() == &AllocationKind::Value
1542        {
1543            translate_step_participant_readback_range(
1544                descriptor.demand(),
1545                work_shape,
1546                participant_index,
1547                semantic_range,
1548            )
1549        } else {
1550            Ok(semantic_range)
1551        }
1552    }
1553}
1554
1555enum CompletionRecord<R: DeviceRuntime> {
1556    Reserved,
1557    InFlight {
1558        resources: CompletionResourceLease<R>,
1559        lane: Arc<ExecutionLane<R>>,
1560        fence: R::Fence,
1561        batch_identity: BatchOperationIdentity,
1562        receipt: SubmittedOperationReceipt,
1563        timing_mode: DeviceTimingMode,
1564        recovery_state: CompletionRecoveryState,
1565    },
1566    SubmissionIndeterminate {
1567        resources: CompletionResourceLease<R>,
1568        lane: Arc<ExecutionLane<R>>,
1569        batch_identity: BatchOperationIdentity,
1570    },
1571    Quarantined {
1572        ownership: CompletionQuarantineOwnership<R>,
1573        receipt: CompletionQuarantineReceipt,
1574    },
1575    Reaped,
1576}
1577
1578enum BoundCompletionObservation<T> {
1579    Pending,
1580    Terminal {
1581        completion: OperationCompletionReceipt,
1582        terminal: T,
1583    },
1584    Indeterminate(Vec<IdentifiedFailure>),
1585    SubmissionIndeterminate,
1586    ObservationPanicked,
1587    Quarantined(CompletionQuarantineReceipt),
1588}
1589
1590#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1591enum CompletionRecoveryState {
1592    Unobserved,
1593    QueryIndeterminate,
1594    DrainEligible(CompletionRecoveryCause),
1595}
1596
1597enum CompletionQuarantineOwnership<R: DeviceRuntime> {
1598    InFlight {
1599        resources: CompletionResourceLease<R>,
1600        lane: Arc<ExecutionLane<R>>,
1601        fence: R::Fence,
1602        batch_identity: BatchOperationIdentity,
1603        submission: SubmittedOperationReceipt,
1604    },
1605    SubmissionIndeterminate {
1606        resources: CompletionResourceLease<R>,
1607        lane: Arc<ExecutionLane<R>>,
1608        batch_identity: BatchOperationIdentity,
1609    },
1610}
1611
1612impl<R: DeviceRuntime> CompletionQuarantineOwnership<R> {
1613    fn lane(&self) -> &Arc<ExecutionLane<R>> {
1614        match self {
1615            Self::InFlight { lane, .. } | Self::SubmissionIndeterminate { lane, .. } => lane,
1616        }
1617    }
1618
1619    fn deferred_cleanup_domain(&self) -> DeferredDeviceCleanupDomainId {
1620        match self {
1621            Self::InFlight { resources, .. } | Self::SubmissionIndeterminate { resources, .. } => {
1622                resources.deferred_cleanup_domain()
1623            }
1624        }
1625    }
1626
1627    fn resources_mut(&mut self) -> &mut CompletionResourceLease<R> {
1628        match self {
1629            Self::InFlight { resources, .. } | Self::SubmissionIndeterminate { resources, .. } => {
1630                resources
1631            }
1632        }
1633    }
1634}
1635
1636impl<R: DeviceRuntime> CompletionRecord<R> {
1637    fn deferred_cleanup_domain(&self) -> Option<DeferredDeviceCleanupDomainId> {
1638        match self {
1639            Self::InFlight { resources, .. } | Self::SubmissionIndeterminate { resources, .. } => {
1640                Some(resources.deferred_cleanup_domain())
1641            }
1642            Self::Quarantined { ownership, .. } => Some(ownership.deferred_cleanup_domain()),
1643            Self::Reserved | Self::Reaped => None,
1644        }
1645    }
1646
1647    fn finish_request_state_hazards_after_drain(&mut self) -> Result<(), VNextError> {
1648        let resources = match self {
1649            Self::InFlight { resources, .. } | Self::SubmissionIndeterminate { resources, .. } => {
1650                resources
1651            }
1652            Self::Quarantined { ownership, .. } => ownership.resources_mut(),
1653            Self::Reserved | Self::Reaped => return Ok(()),
1654        };
1655        resources.finish_request_state_hazards(
1656            RequestStateHazardTerminalDisposition::IndeterminateAfterDrain,
1657        )
1658    }
1659}
1660
1661type SharedCompletionRecord<R> = Arc<Mutex<CompletionRecord<R>>>;
1662
1663struct CompletionReaperState<R: DeviceRuntime> {
1664    next_slot: Option<NonZeroU64>,
1665    sweep_cursor: Option<CompletionSlotId>,
1666    slots: BTreeMap<CompletionSlotId, SharedCompletionRecord<R>>,
1667}
1668
1669/// Scheduler-owned completion registry. The global map lock only resolves a
1670/// slot; each fence is queried or waited under its own record lock.
1671#[must_use = "the scheduler must retain its completion reaper"]
1672pub struct CompletionReaper<R: DeviceRuntime> {
1673    state: Mutex<CompletionReaperState<R>>,
1674}
1675
1676pub const MAX_COMPLETION_SWEEP_SLOTS: usize = 64;
1677
1678impl<R: DeviceRuntime> fmt::Debug for CompletionReaper<R> {
1679    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
1680        formatter
1681            .debug_struct("CompletionReaper")
1682            .field("retained", &self.retained_count())
1683            .finish_non_exhaustive()
1684    }
1685}
1686
1687impl<R: DeviceRuntime> CompletionReaper<R> {
1688    pub fn new() -> Arc<Self> {
1689        Arc::new(Self {
1690            state: Mutex::new(CompletionReaperState {
1691                next_slot: NonZeroU64::new(1),
1692                sweep_cursor: None,
1693                slots: BTreeMap::new(),
1694            }),
1695        })
1696    }
1697
1698    pub fn retained_count(&self) -> usize {
1699        self.state
1700            .lock()
1701            .map(|state| state.slots.len())
1702            .unwrap_or(usize::MAX)
1703    }
1704
1705    pub fn quarantined_count(&self) -> usize {
1706        let records = match self.state.lock() {
1707            Ok(state) => state.slots.values().cloned().collect::<Vec<_>>(),
1708            Err(_) => return usize::MAX,
1709        };
1710        records
1711            .iter()
1712            .filter(|record| {
1713                record
1714                    .lock()
1715                    .map(|record| matches!(&*record, CompletionRecord::Quarantined { .. }))
1716                    .unwrap_or(true)
1717            })
1718            .count()
1719    }
1720
1721    /// Polls a bounded scheduler-owned snapshot. This path remains available
1722    /// after every external completion handle has detached.
1723    pub fn poll_bounded(&self, maximum_slots: usize) -> Result<CompletionSweepReceipt, VNextError> {
1724        if maximum_slots == 0 || maximum_slots > MAX_COMPLETION_SWEEP_SLOTS {
1725            return Err(invalid_completion(format!(
1726                "completion sweep size must be in 1..={MAX_COMPLETION_SWEEP_SLOTS}"
1727            )));
1728        }
1729        let slot_ids = {
1730            let mut state = self
1731                .state
1732                .lock()
1733                .map_err(|_| invalid_completion("completion reaper state mutex is poisoned"))?;
1734            let keys = state.slots.keys().copied().collect::<Vec<_>>();
1735            let start = state
1736                .sweep_cursor
1737                .and_then(|cursor| keys.iter().position(|slot_id| *slot_id > cursor))
1738                .unwrap_or(0);
1739            let slot_ids = keys
1740                .iter()
1741                .cycle()
1742                .skip(start)
1743                .take(maximum_slots.min(keys.len()))
1744                .copied()
1745                .collect::<Vec<_>>();
1746            if let Some(last) = slot_ids.last() {
1747                state.sweep_cursor = Some(*last);
1748            }
1749            slot_ids
1750        };
1751        let entries = slot_ids
1752            .into_iter()
1753            .map(|slot_id| CompletionSweepEntry {
1754                slot_id,
1755                observation: match self.poll_bound(slot_id) {
1756                    Ok(observation) => CompletionSweepObservation::Observed(observation),
1757                    Err(error) => CompletionSweepObservation::Failed(error),
1758                },
1759            })
1760            .collect();
1761        let records = self
1762            .state
1763            .lock()
1764            .map_err(|_| invalid_completion("completion reaper state mutex is poisoned"))?
1765            .slots
1766            .values()
1767            .cloned()
1768            .collect::<Vec<_>>();
1769        let retained_after = records.len();
1770        let quarantined_after = records
1771            .iter()
1772            .filter(|record| {
1773                record
1774                    .lock()
1775                    .map(|record| matches!(&*record, CompletionRecord::Quarantined { .. }))
1776                    .unwrap_or(true)
1777            })
1778            .count();
1779        Ok(CompletionSweepReceipt {
1780            entries,
1781            retained_after,
1782            quarantined_after,
1783        })
1784    }
1785
1786    /// Runs the blocking recovery path for an exact slot after an indeterminate
1787    /// observation. A failed drain moves ownership into an auditable, retryable
1788    /// quarantine record instead of releasing or forgetting it.
1789    pub fn recover_slot_by_draining_lane(
1790        &self,
1791        slot_id: CompletionSlotId,
1792    ) -> Result<CompletionRecoveryOutcome, VNextError> {
1793        self.recover_bound(slot_id)
1794    }
1795
1796    /// Performs the blocking exact-fence observation required before lane
1797    /// recovery. Schedulers must call this from an independent recovery worker,
1798    /// never from a request or admission thread.
1799    pub fn wait_slot_for_recovery(
1800        &self,
1801        slot_id: CompletionSlotId,
1802    ) -> Result<CompletionObservation, VNextError> {
1803        self.wait_bound(slot_id)
1804    }
1805
1806    pub(crate) fn reserve(
1807        reaper: &Arc<Self>,
1808        invocation: InvocationResourceLease<R>,
1809        lane: Arc<ExecutionLane<R>>,
1810        batch_identity: BatchOperationIdentity,
1811    ) -> Result<CompletionReservation<R>, VNextError> {
1812        Self::reserve_resources(
1813            reaper,
1814            CompletionResourceLease::Invocation(invocation),
1815            lane,
1816            batch_identity,
1817        )
1818    }
1819
1820    pub(crate) fn reserve_wave(
1821        reaper: &Arc<Self>,
1822        wave: PreparedStepSubmissionWave<R>,
1823        lane: Arc<ExecutionLane<R>>,
1824        batch_identity: BatchOperationIdentity,
1825    ) -> Result<CompletionReservation<R>, VNextError> {
1826        Self::reserve_resources(
1827            reaper,
1828            CompletionResourceLease::Wave(wave),
1829            lane,
1830            batch_identity,
1831        )
1832    }
1833
1834    fn reserve_resources(
1835        reaper: &Arc<Self>,
1836        resources: CompletionResourceLease<R>,
1837        lane: Arc<ExecutionLane<R>>,
1838        batch_identity: BatchOperationIdentity,
1839    ) -> Result<CompletionReservation<R>, VNextError> {
1840        if !Arc::ptr_eq(resources.runtime(), lane.runtime_arc()) {
1841            return Err(invalid_completion(
1842                "completion lane runtime is not the submission resource runtime instance",
1843            ));
1844        }
1845        let mut state = reaper
1846            .state
1847            .lock()
1848            .map_err(|_| invalid_completion("completion reaper state mutex is poisoned"))?;
1849        let raw = state
1850            .next_slot
1851            .ok_or_else(|| invalid_completion("completion slot identity space is exhausted"))?;
1852        state.next_slot = raw.get().checked_add(1).and_then(NonZeroU64::new);
1853        let slot_id = CompletionSlotId(raw);
1854        let record = Arc::new(Mutex::new(CompletionRecord::Reserved));
1855        if state.slots.insert(slot_id, Arc::clone(&record)).is_some() {
1856            return Err(invalid_completion("completion slot identity was reused"));
1857        }
1858        drop(state);
1859        #[derive(Serialize)]
1860        struct SubmissionFingerprintInput<'a> {
1861            domain: &'static str,
1862            slot_id: CompletionSlotId,
1863            batch_identity_fingerprint: &'a str,
1864        }
1865        let fingerprint = canonical_completion_fingerprint(&SubmissionFingerprintInput {
1866            domain: "ferrum.runtime-vnext.batch-operation-submission.v1",
1867            slot_id,
1868            batch_identity_fingerprint: batch_identity.fingerprint(),
1869        });
1870        let receipt = SubmittedOperationReceipt::new(slot_id, batch_identity.clone(), fingerprint);
1871        let record_receipt = receipt.clone();
1872        Ok(CompletionReservation {
1873            reaper: Arc::clone(reaper),
1874            record,
1875            slot_id,
1876            resources: Some(resources),
1877            lane: Some(lane),
1878            batch_identity: Some(batch_identity),
1879            receipt: Some(receipt),
1880            record_receipt: Some(record_receipt),
1881            submission_may_have_happened: false,
1882            finished: false,
1883        })
1884    }
1885
1886    fn lookup(&self, slot_id: CompletionSlotId) -> Result<SharedCompletionRecord<R>, VNextError> {
1887        self.state
1888            .lock()
1889            .map_err(|_| invalid_completion("completion reaper state mutex is poisoned"))?
1890            .slots
1891            .get(&slot_id)
1892            .cloned()
1893            .ok_or_else(|| invalid_completion("completion slot is unknown or already reaped"))
1894    }
1895
1896    fn remove_exact(&self, slot_id: CompletionSlotId, record: &SharedCompletionRecord<R>) {
1897        let mut state = match self.state.lock() {
1898            Ok(state) => state,
1899            Err(poisoned) => poisoned.into_inner(),
1900        };
1901        if state
1902            .slots
1903            .get(&slot_id)
1904            .is_some_and(|current| Arc::ptr_eq(current, record))
1905        {
1906            state.slots.remove(&slot_id);
1907        }
1908    }
1909
1910    pub(crate) fn poll_bound(
1911        &self,
1912        slot_id: CompletionSlotId,
1913    ) -> Result<CompletionObservation, VNextError> {
1914        Self::map_plain_observation(self.observe_bound_with(slot_id, false, |_, _, _, _, _| ())?)
1915    }
1916
1917    pub(crate) fn wait_bound(
1918        &self,
1919        slot_id: CompletionSlotId,
1920    ) -> Result<CompletionObservation, VNextError> {
1921        Self::map_plain_observation(self.observe_bound_with(slot_id, true, |_, _, _, _, _| ())?)
1922    }
1923
1924    pub(crate) fn wait_bound_with_readback(
1925        &self,
1926        slot_id: CompletionSlotId,
1927        request: CompletionReadbackRequest,
1928    ) -> Result<CompletionReadbackObservation, VNextError> {
1929        let observation = self.observe_bound_with(
1930            slot_id,
1931            true,
1932            |resources, lane, batch_identity, disposition, timing_mode| {
1933                attempt_completion_readback(
1934                    resources,
1935                    lane,
1936                    batch_identity,
1937                    disposition,
1938                    timing_mode,
1939                    request,
1940                )
1941            },
1942        )?;
1943        Ok(match observation {
1944            BoundCompletionObservation::Pending => CompletionReadbackObservation::Pending,
1945            BoundCompletionObservation::Terminal {
1946                completion,
1947                terminal,
1948            } => CompletionReadbackObservation::Terminal(CompletionReadbackReceipt::new(
1949                completion, terminal,
1950            )),
1951            BoundCompletionObservation::Indeterminate(failures) => {
1952                CompletionReadbackObservation::Indeterminate(failures)
1953            }
1954            BoundCompletionObservation::SubmissionIndeterminate => {
1955                CompletionReadbackObservation::SubmissionIndeterminate
1956            }
1957            BoundCompletionObservation::ObservationPanicked => {
1958                CompletionReadbackObservation::ObservationPanicked
1959            }
1960            BoundCompletionObservation::Quarantined(receipt) => {
1961                CompletionReadbackObservation::Quarantined(receipt)
1962            }
1963        })
1964    }
1965
1966    pub(crate) fn wait_bound_with_readbacks(
1967        &self,
1968        slot_id: CompletionSlotId,
1969        request: CompletionReadbackBatchRequest,
1970    ) -> Result<CompletionReadbackBatchObservation, VNextError> {
1971        self.validate_bound_readback_batch(slot_id, &request)?;
1972        let observation = self.observe_bound_with(
1973            slot_id,
1974            true,
1975            |resources, lane, batch_identity, disposition, timing_mode| {
1976                attempt_completion_readbacks(
1977                    resources,
1978                    lane,
1979                    batch_identity,
1980                    disposition,
1981                    timing_mode,
1982                    request,
1983                )
1984            },
1985        )?;
1986        Ok(match observation {
1987            BoundCompletionObservation::Pending => CompletionReadbackBatchObservation::Pending,
1988            BoundCompletionObservation::Terminal {
1989                completion,
1990                terminal,
1991            } => CompletionReadbackBatchObservation::Terminal(CompletionReadbackBatchReceipt::new(
1992                completion, terminal,
1993            )),
1994            BoundCompletionObservation::Indeterminate(failures) => {
1995                CompletionReadbackBatchObservation::Indeterminate(failures)
1996            }
1997            BoundCompletionObservation::SubmissionIndeterminate => {
1998                CompletionReadbackBatchObservation::SubmissionIndeterminate
1999            }
2000            BoundCompletionObservation::ObservationPanicked => {
2001                CompletionReadbackBatchObservation::ObservationPanicked
2002            }
2003            BoundCompletionObservation::Quarantined(receipt) => {
2004                CompletionReadbackBatchObservation::Quarantined(receipt)
2005            }
2006        })
2007    }
2008
2009    pub(crate) fn wait_bound_with_readback_collection(
2010        &self,
2011        slot_id: CompletionSlotId,
2012        request: CompletionReadbackCollectionRequest,
2013    ) -> Result<CompletionReadbackCollectionObservation, VNextError> {
2014        self.validate_bound_readback_collection(slot_id, &request)?;
2015        let observation = self.observe_bound_with(
2016            slot_id,
2017            true,
2018            |resources, lane, batch_identity, disposition, timing_mode| {
2019                attempt_completion_readback_collection(
2020                    resources,
2021                    lane,
2022                    batch_identity,
2023                    disposition,
2024                    timing_mode,
2025                    request,
2026                )
2027            },
2028        )?;
2029        Ok(match observation {
2030            BoundCompletionObservation::Pending => CompletionReadbackBatchObservation::Pending,
2031            BoundCompletionObservation::Terminal {
2032                completion,
2033                terminal,
2034            } => CompletionReadbackBatchObservation::Terminal(CompletionReadbackBatchReceipt::new(
2035                completion, terminal,
2036            )),
2037            BoundCompletionObservation::Indeterminate(failures) => {
2038                CompletionReadbackBatchObservation::Indeterminate(failures)
2039            }
2040            BoundCompletionObservation::SubmissionIndeterminate => {
2041                CompletionReadbackBatchObservation::SubmissionIndeterminate
2042            }
2043            BoundCompletionObservation::ObservationPanicked => {
2044                CompletionReadbackBatchObservation::ObservationPanicked
2045            }
2046            BoundCompletionObservation::Quarantined(receipt) => {
2047                CompletionReadbackBatchObservation::Quarantined(receipt)
2048            }
2049        })
2050    }
2051
2052    fn validate_bound_readback_batch(
2053        &self,
2054        slot_id: CompletionSlotId,
2055        request: &CompletionReadbackBatchRequest,
2056    ) -> Result<(), VNextError> {
2057        let record = self.lookup(slot_id)?;
2058        let guard = record
2059            .lock()
2060            .map_err(|_| invalid_completion("completion slot mutex is poisoned"))?;
2061        match &*guard {
2062            CompletionRecord::InFlight { batch_identity, .. }
2063            | CompletionRecord::SubmissionIndeterminate { batch_identity, .. } => {
2064                request.validate_for(batch_identity)
2065            }
2066            CompletionRecord::Reserved => Err(invalid_completion(
2067                "completion slot has not reached submission",
2068            )),
2069            CompletionRecord::Quarantined { .. } => Ok(()),
2070            CompletionRecord::Reaped => {
2071                Err(invalid_completion("completion slot is already reaped"))
2072            }
2073        }
2074    }
2075
2076    fn validate_bound_readback_collection(
2077        &self,
2078        slot_id: CompletionSlotId,
2079        request: &CompletionReadbackCollectionRequest,
2080    ) -> Result<(), VNextError> {
2081        let record = self.lookup(slot_id)?;
2082        let guard = record
2083            .lock()
2084            .map_err(|_| invalid_completion("completion slot mutex is poisoned"))?;
2085        match &*guard {
2086            CompletionRecord::InFlight { batch_identity, .. }
2087            | CompletionRecord::SubmissionIndeterminate { batch_identity, .. } => {
2088                request.validate_for(batch_identity)
2089            }
2090            CompletionRecord::Reserved => Err(invalid_completion(
2091                "completion slot has not reached submission",
2092            )),
2093            CompletionRecord::Quarantined { .. } => Ok(()),
2094            CompletionRecord::Reaped => {
2095                Err(invalid_completion("completion slot is already reaped"))
2096            }
2097        }
2098    }
2099
2100    fn map_plain_observation(
2101        observation: BoundCompletionObservation<()>,
2102    ) -> Result<CompletionObservation, VNextError> {
2103        Ok(match observation {
2104            BoundCompletionObservation::Pending => CompletionObservation::Pending,
2105            BoundCompletionObservation::Terminal { completion, .. } => {
2106                CompletionObservation::Terminal(completion)
2107            }
2108            BoundCompletionObservation::Indeterminate(failures) => {
2109                CompletionObservation::Indeterminate(failures)
2110            }
2111            BoundCompletionObservation::SubmissionIndeterminate => {
2112                CompletionObservation::SubmissionIndeterminate
2113            }
2114            BoundCompletionObservation::ObservationPanicked => {
2115                CompletionObservation::ObservationPanicked
2116            }
2117            BoundCompletionObservation::Quarantined(receipt) => {
2118                CompletionObservation::Quarantined(receipt)
2119            }
2120        })
2121    }
2122
2123    fn observe_bound_with<T, F>(
2124        &self,
2125        slot_id: CompletionSlotId,
2126        blocking: bool,
2127        terminal_action: F,
2128    ) -> Result<BoundCompletionObservation<T>, VNextError>
2129    where
2130        F: FnOnce(
2131            &CompletionResourceLease<R>,
2132            &Arc<ExecutionLane<R>>,
2133            &BatchOperationIdentity,
2134            &OperationCompletionDisposition,
2135            DeviceTimingMode,
2136        ) -> T,
2137    {
2138        let record = self.lookup(slot_id)?;
2139        let mut guard = record
2140            .lock()
2141            .map_err(|_| invalid_completion("completion slot mutex is poisoned"))?;
2142        let observation = match &*guard {
2143            CompletionRecord::Reserved => {
2144                return Err(invalid_completion(
2145                    "completion slot has not reached submission",
2146                ));
2147            }
2148            CompletionRecord::SubmissionIndeterminate { .. } => {
2149                FenceObservation::SubmissionIndeterminate
2150            }
2151            CompletionRecord::Quarantined { receipt, .. } => {
2152                FenceObservation::Quarantined(receipt.clone())
2153            }
2154            CompletionRecord::Reaped => {
2155                return Err(invalid_completion("completion slot is already reaped"));
2156            }
2157            CompletionRecord::InFlight {
2158                lane,
2159                fence,
2160                batch_identity,
2161                timing_mode,
2162                ..
2163            } if blocking => {
2164                let wait_started = timing_mode.completion_enabled().then(Instant::now);
2165                match catch_unwind(AssertUnwindSafe(|| lane.wait_fence(fence))) {
2166                    Ok(Ok(terminal)) => {
2167                        let wait_timing = match wait_started {
2168                            Some(started) => u64::try_from(started.elapsed().as_nanos()).map_or(
2169                                DeviceTimingMeasurement::Unavailable(
2170                                    DeviceTimingUnavailableReason::DurationOverflow,
2171                                ),
2172                                DeviceTimingMeasurement::Measured,
2173                            ),
2174                            None => DeviceTimingMeasurement::NotRequested,
2175                        };
2176                        catch_unwind(AssertUnwindSafe(|| {
2177                            terminal_observation(
2178                                lane,
2179                                batch_identity,
2180                                terminal,
2181                                *timing_mode,
2182                                wait_timing,
2183                            )
2184                        }))
2185                        .unwrap_or(FenceObservation::ObservationPanicked)
2186                    }
2187                    Ok(Err(indeterminate)) => catch_unwind(AssertUnwindSafe(|| {
2188                        classify_batch_device_error(
2189                            lane.runtime(),
2190                            batch_identity,
2191                            indeterminate.error(),
2192                        )
2193                    }))
2194                    .map_or(
2195                        FenceObservation::ObservationPanicked,
2196                        |classified| match classified {
2197                            Ok(failure) => FenceObservation::Indeterminate(failure),
2198                            Err(error) => FenceObservation::ContractIndeterminate(error),
2199                        },
2200                    ),
2201                    Err(_) => FenceObservation::ObservationPanicked,
2202                }
2203            }
2204            CompletionRecord::InFlight {
2205                lane,
2206                fence,
2207                batch_identity,
2208                timing_mode,
2209                ..
2210            } => match catch_unwind(AssertUnwindSafe(|| lane.query_fence(fence))) {
2211                Ok(FenceQuery::Pending) => FenceObservation::Pending,
2212                Ok(FenceQuery::Terminal(terminal)) => catch_unwind(AssertUnwindSafe(|| {
2213                    terminal_observation(
2214                        lane,
2215                        batch_identity,
2216                        terminal,
2217                        *timing_mode,
2218                        if timing_mode.completion_enabled() {
2219                            DeviceTimingMeasurement::Measured(0)
2220                        } else {
2221                            DeviceTimingMeasurement::NotRequested
2222                        },
2223                    )
2224                }))
2225                .unwrap_or(FenceObservation::ObservationPanicked),
2226                Ok(FenceQuery::Indeterminate(error)) => catch_unwind(AssertUnwindSafe(|| {
2227                    classify_batch_device_error(lane.runtime(), batch_identity, &error)
2228                }))
2229                .map_or(
2230                    FenceObservation::ObservationPanicked,
2231                    |classified| match classified {
2232                        Ok(failure) => FenceObservation::Indeterminate(failure),
2233                        Err(error) => FenceObservation::ContractIndeterminate(error),
2234                    },
2235                ),
2236                Err(_) => FenceObservation::ObservationPanicked,
2237            },
2238        };
2239        let recovery_state = match &observation {
2240            FenceObservation::Indeterminate(_) | FenceObservation::ContractIndeterminate(_) => {
2241                Some(if blocking {
2242                    CompletionRecoveryState::DrainEligible(
2243                        CompletionRecoveryCause::FenceIndeterminate,
2244                    )
2245                } else {
2246                    CompletionRecoveryState::QueryIndeterminate
2247                })
2248            }
2249            FenceObservation::ObservationPanicked => Some(if blocking {
2250                CompletionRecoveryState::DrainEligible(
2251                    CompletionRecoveryCause::FenceObservationPanicked,
2252                )
2253            } else {
2254                CompletionRecoveryState::QueryIndeterminate
2255            }),
2256            _ => None,
2257        };
2258        if let Some(recovery_state) = recovery_state {
2259            if let CompletionRecord::InFlight {
2260                recovery_state: current,
2261                ..
2262            } = &mut *guard
2263            {
2264                if !matches!(current, CompletionRecoveryState::DrainEligible(_)) {
2265                    *current = recovery_state;
2266                }
2267            }
2268        }
2269        if let FenceObservation::Terminal(mut disposition, fence_timing, submission_timing) =
2270            observation
2271        {
2272            let old = std::mem::replace(&mut *guard, CompletionRecord::Reaped);
2273            let CompletionRecord::InFlight {
2274                resources,
2275                lane,
2276                batch_identity,
2277                receipt,
2278                timing_mode,
2279                ..
2280            } = old
2281            else {
2282                unreachable!("terminal observation came from an in-flight record")
2283            };
2284            let terminal = terminal_action(
2285                &resources,
2286                &lane,
2287                &batch_identity,
2288                &disposition,
2289                timing_mode,
2290            );
2291            if let Err(failure) = lane.finish_one_terminal() {
2292                disposition = OperationCompletionDisposition::ContractFailedButQuiescent(failure);
2293            }
2294            let initialization_succeeded =
2295                matches!(disposition, OperationCompletionDisposition::Succeeded);
2296            let mut resources = resources;
2297            if let Err(error) = resources.finish_backing_initializations(initialization_succeeded) {
2298                disposition = OperationCompletionDisposition::ContractFailedButQuiescent(
2299                    QuiescentCompletionContractFailure::new(format!(
2300                        "backing initialization terminal transition failed: {error}"
2301                    )),
2302                );
2303            }
2304            let hazard_disposition =
2305                if matches!(disposition, OperationCompletionDisposition::Succeeded) {
2306                    RequestStateHazardTerminalDisposition::Succeeded
2307                } else {
2308                    RequestStateHazardTerminalDisposition::FailedButQuiescent
2309                };
2310            if let Err(error) = resources.finish_request_state_hazards(hazard_disposition) {
2311                disposition = OperationCompletionDisposition::ContractFailedButQuiescent(
2312                    QuiescentCompletionContractFailure::new(format!(
2313                        "request-state hazard terminal transition failed: {error}"
2314                    )),
2315                );
2316            }
2317            drop(guard);
2318            self.remove_exact(slot_id, &record);
2319            return OperationCompletionReceipt::new(
2320                receipt,
2321                disposition,
2322                fence_timing,
2323                submission_timing,
2324            )
2325            .map(|completion| BoundCompletionObservation::Terminal {
2326                completion,
2327                terminal,
2328            });
2329        }
2330        Ok(match observation {
2331            FenceObservation::Pending => BoundCompletionObservation::Pending,
2332            FenceObservation::Indeterminate(failure) => {
2333                BoundCompletionObservation::Indeterminate(failure)
2334            }
2335            FenceObservation::SubmissionIndeterminate => {
2336                BoundCompletionObservation::SubmissionIndeterminate
2337            }
2338            FenceObservation::ObservationPanicked => {
2339                BoundCompletionObservation::ObservationPanicked
2340            }
2341            FenceObservation::ContractIndeterminate(error) => return Err(error),
2342            FenceObservation::Quarantined(receipt) => {
2343                BoundCompletionObservation::Quarantined(receipt)
2344            }
2345            FenceObservation::Terminal(_, _, _) => {
2346                unreachable!("terminal observation returned through the reaping branch")
2347            }
2348        })
2349    }
2350
2351    fn recover_bound(
2352        &self,
2353        slot_id: CompletionSlotId,
2354    ) -> Result<CompletionRecoveryOutcome, VNextError> {
2355        let record = self.lookup(slot_id)?;
2356        let mut guard = record
2357            .lock()
2358            .map_err(|_| invalid_completion("completion slot mutex is poisoned"))?;
2359        let (lane, batch_identity, submission, cause, had_submission_fence) = match &*guard {
2360            CompletionRecord::InFlight {
2361                lane,
2362                batch_identity,
2363                receipt,
2364                recovery_state: CompletionRecoveryState::DrainEligible(cause),
2365                ..
2366            } => (
2367                Arc::clone(lane),
2368                batch_identity.clone(),
2369                Some(receipt.clone()),
2370                *cause,
2371                true,
2372            ),
2373            CompletionRecord::SubmissionIndeterminate {
2374                lane,
2375                batch_identity,
2376                ..
2377            } => (
2378                Arc::clone(lane),
2379                batch_identity.clone(),
2380                None,
2381                CompletionRecoveryCause::SubmissionIndeterminate,
2382                false,
2383            ),
2384            CompletionRecord::Quarantined { ownership, receipt } => (
2385                Arc::clone(ownership.lane()),
2386                receipt.batch_identity.clone(),
2387                receipt.submission.clone(),
2388                receipt.cause,
2389                receipt.had_submission_fence,
2390            ),
2391            CompletionRecord::InFlight { .. } => {
2392                return Err(invalid_completion(
2393                    "completion slot has no blocking indeterminate fence observation",
2394                ));
2395            }
2396            CompletionRecord::Reserved => {
2397                return Err(invalid_completion(
2398                    "completion slot has not reached submission",
2399                ));
2400            }
2401            CompletionRecord::Reaped => {
2402                return Err(invalid_completion("completion slot is already reaped"));
2403            }
2404        };
2405        // A recovery drain proves lane-wide quiescence rather than one fence's
2406        // terminal state. Retire the lane before draining so older records
2407        // cannot later decrement accounting for newly submitted work.
2408        lane.fail_closed();
2409        if !lane.drain(had_submission_fence) {
2410            if let CompletionRecord::Quarantined { receipt, .. } = &*guard {
2411                return Ok(CompletionRecoveryOutcome::Quarantined(receipt.clone()));
2412            }
2413            let old = std::mem::replace(&mut *guard, CompletionRecord::Reaped);
2414            let ownership = match old {
2415                CompletionRecord::InFlight {
2416                    resources,
2417                    lane,
2418                    fence,
2419                    batch_identity,
2420                    receipt,
2421                    ..
2422                } => CompletionQuarantineOwnership::InFlight {
2423                    resources,
2424                    lane,
2425                    fence,
2426                    batch_identity,
2427                    submission: receipt,
2428                },
2429                CompletionRecord::SubmissionIndeterminate {
2430                    resources,
2431                    lane,
2432                    batch_identity,
2433                } => CompletionQuarantineOwnership::SubmissionIndeterminate {
2434                    resources,
2435                    lane,
2436                    batch_identity,
2437                },
2438                _ => unreachable!("recovery source was validated before lane drain"),
2439            };
2440            let receipt = CompletionQuarantineReceipt {
2441                slot_id,
2442                batch_identity,
2443                submission,
2444                cause,
2445                had_submission_fence,
2446                device_id: lane.descriptor.id.clone(),
2447                runtime_implementation_fingerprint: lane
2448                    .descriptor
2449                    .runtime_implementation_fingerprint
2450                    .clone(),
2451                freshness: CompletionQuarantineFreshness::current(),
2452            };
2453            *guard = CompletionRecord::Quarantined {
2454                ownership,
2455                receipt: receipt.clone(),
2456            };
2457            return Ok(CompletionRecoveryOutcome::Quarantined(receipt));
2458        }
2459        let mut old = std::mem::replace(&mut *guard, CompletionRecord::Reaped);
2460        if let CompletionRecord::Quarantined { receipt, .. } = &old {
2461            receipt.freshness.invalidate();
2462        }
2463        let hazard_result = old.finish_request_state_hazards_after_drain();
2464        drop(guard);
2465        self.remove_exact(slot_id, &record);
2466        drop(old);
2467        hazard_result?;
2468        Ok(CompletionRecoveryOutcome::Drained(CompletionDrainReceipt {
2469            slot_id,
2470            batch_identity,
2471            submission,
2472            cause,
2473            had_submission_fence,
2474        }))
2475    }
2476}
2477
2478enum FenceObservation {
2479    Pending,
2480    Terminal(
2481        OperationCompletionDisposition,
2482        CompletionFenceTiming,
2483        DeviceTimingMeasurement<DeviceSubmissionExecutionTiming>,
2484    ),
2485    Indeterminate(Vec<IdentifiedFailure>),
2486    ContractIndeterminate(VNextError),
2487    SubmissionIndeterminate,
2488    ObservationPanicked,
2489    Quarantined(CompletionQuarantineReceipt),
2490}
2491
2492fn terminal_observation<R: DeviceRuntime>(
2493    lane: &ExecutionLane<R>,
2494    batch_identity: &BatchOperationIdentity,
2495    terminal: DeviceTerminalReceipt<R::Error>,
2496    timing_mode: DeviceTimingMode,
2497    blocking_wait_host_ns: DeviceTimingMeasurement<u64>,
2498) -> FenceObservation {
2499    let (terminal, device_execution, submission_timing) = terminal.into_parts();
2500    let timing = CompletionFenceTiming::new(timing_mode, device_execution, blocking_wait_host_ns);
2501    let submission_timing = match (timing_mode, submission_timing) {
2502        (DeviceTimingMode::Off | DeviceTimingMode::Completion, _) => {
2503            DeviceTimingMeasurement::NotRequested
2504        }
2505        (
2506            DeviceTimingMode::Replay | DeviceTimingMode::Kernel | DeviceTimingMode::Verification,
2507            DeviceTimingMeasurement::NotRequested,
2508        ) => {
2509            DeviceTimingMeasurement::Unavailable(DeviceTimingUnavailableReason::BackendUnsupported)
2510        }
2511        (
2512            DeviceTimingMode::Replay | DeviceTimingMode::Kernel | DeviceTimingMode::Verification,
2513            measurement,
2514        ) => measurement,
2515    };
2516    let descriptor_is_stable = catch_unwind(AssertUnwindSafe(|| {
2517        lane.current_descriptor_matches_snapshot()
2518    }))
2519    .unwrap_or(false);
2520    if !descriptor_is_stable {
2521        return FenceObservation::Terminal(
2522            OperationCompletionDisposition::ContractFailedButQuiescent(
2523                QuiescentCompletionContractFailure::new(
2524                    "completion runtime descriptor differs from its execution lane snapshot",
2525                ),
2526            ),
2527            timing,
2528            submission_timing,
2529        );
2530    }
2531    FenceObservation::Terminal(
2532        match terminal {
2533            DeviceTerminal::Succeeded => OperationCompletionDisposition::Succeeded,
2534            DeviceTerminal::FailedButQuiescent(error) => {
2535                match classify_batch_device_error(lane.runtime(), batch_identity, &error) {
2536                    Ok(failures) => OperationCompletionDisposition::FailedButQuiescent(failures),
2537                    Err(error) => OperationCompletionDisposition::ContractFailedButQuiescent(
2538                        QuiescentCompletionContractFailure::new(error.to_string()),
2539                    ),
2540                }
2541            }
2542        },
2543        timing,
2544        submission_timing,
2545    )
2546}
2547
2548fn classify_batch_device_error<R: DeviceRuntime>(
2549    runtime: &R,
2550    batch_identity: &BatchOperationIdentity,
2551    error: &R::Error,
2552) -> Result<Vec<IdentifiedFailure>, VNextError> {
2553    batch_identity
2554        .participants()
2555        .iter()
2556        .map(|participant| classify_device_error(runtime, participant.identity().clone(), error))
2557        .collect()
2558}
2559
2560struct DeferredCompletionCleanup<R: DeviceRuntime> {
2561    records: BTreeMap<CompletionSlotId, SharedCompletionRecord<R>>,
2562}
2563
2564impl<R: DeviceRuntime> DeferredCompletionCleanup<R> {
2565    fn new(records: BTreeMap<CompletionSlotId, SharedCompletionRecord<R>>) -> Self {
2566        Self { records }
2567    }
2568}
2569
2570impl<R: DeviceRuntime> DeferredDeviceCleanupTask for DeferredCompletionCleanup<R> {
2571    fn try_cleanup(&mut self) -> DeferredDeviceCleanupDisposition {
2572        let slot_ids = self.records.keys().copied().collect::<Vec<_>>();
2573        let mut retryable = false;
2574        let mut quarantined = false;
2575        for slot_id in slot_ids {
2576            let Some(record) = self.records.get(&slot_id).cloned() else {
2577                continue;
2578            };
2579            let quiescent = catch_unwind(AssertUnwindSafe(|| {
2580                cleanup_dropped_completion_record(&record)
2581            }))
2582            .unwrap_or(false);
2583            if quiescent {
2584                if self
2585                    .records
2586                    .get(&slot_id)
2587                    .is_some_and(|current| Arc::ptr_eq(current, &record))
2588                {
2589                    self.records.remove(&slot_id);
2590                }
2591                continue;
2592            }
2593            retryable = true;
2594            quarantined |= record
2595                .lock()
2596                .map(|record| matches!(&*record, CompletionRecord::Quarantined { .. }))
2597                .unwrap_or(true);
2598        }
2599        if self.records.is_empty() {
2600            DeferredDeviceCleanupDisposition::Completed
2601        } else if quarantined {
2602            DeferredDeviceCleanupDisposition::Quarantined
2603        } else {
2604            debug_assert!(retryable);
2605            DeferredDeviceCleanupDisposition::Retryable
2606        }
2607    }
2608}
2609
2610fn fail_close_completion_record<R: DeviceRuntime>(record: &SharedCompletionRecord<R>) {
2611    let guard = match record.lock() {
2612        Ok(guard) => guard,
2613        Err(poisoned) => poisoned.into_inner(),
2614    };
2615    match &*guard {
2616        CompletionRecord::InFlight { lane, .. }
2617        | CompletionRecord::SubmissionIndeterminate { lane, .. } => lane.fail_closed(),
2618        CompletionRecord::Quarantined { ownership, .. } => ownership.lane().fail_closed(),
2619        CompletionRecord::Reserved | CompletionRecord::Reaped => {}
2620    }
2621}
2622
2623fn cleanup_dropped_completion_record<R: DeviceRuntime>(record: &SharedCompletionRecord<R>) -> bool {
2624    let mut guard = match record.lock() {
2625        Ok(guard) => guard,
2626        Err(poisoned) => poisoned.into_inner(),
2627    };
2628    let (quiescent, drained) = match &*guard {
2629        CompletionRecord::InFlight { lane, fence, .. } => {
2630            let queried = catch_unwind(AssertUnwindSafe(|| lane.query_fence(fence)))
2631                .is_ok_and(|query| matches!(query, FenceQuery::Terminal(_)));
2632            let waited = queried
2633                || catch_unwind(AssertUnwindSafe(|| lane.wait_fence(fence)))
2634                    .is_ok_and(|result| result.is_ok());
2635            if waited {
2636                (true, false)
2637            } else {
2638                let drained = lane.drain(true);
2639                (drained, drained)
2640            }
2641        }
2642        CompletionRecord::SubmissionIndeterminate { lane, .. } => {
2643            let drained = lane.drain(false);
2644            (drained, drained)
2645        }
2646        CompletionRecord::Quarantined { ownership, receipt } => {
2647            let drained = ownership.lane().drain(receipt.had_submission_fence);
2648            (drained, drained)
2649        }
2650        CompletionRecord::Reserved | CompletionRecord::Reaped => (true, false),
2651    };
2652    if !quiescent {
2653        match &*guard {
2654            CompletionRecord::InFlight { lane, .. }
2655            | CompletionRecord::SubmissionIndeterminate { lane, .. } => lane.fail_closed(),
2656            CompletionRecord::Quarantined { ownership, .. } => ownership.lane().fail_closed(),
2657            CompletionRecord::Reserved | CompletionRecord::Reaped => {}
2658        }
2659        return false;
2660    }
2661
2662    let old = std::mem::replace(&mut *guard, CompletionRecord::Reaped);
2663    if let CompletionRecord::Quarantined { receipt, .. } = &old {
2664        receipt.freshness.invalidate();
2665    }
2666    if !drained {
2667        if let CompletionRecord::InFlight { lane, .. } = &old {
2668            let _ = lane.finish_one_terminal();
2669        }
2670    }
2671    drop(guard);
2672    drop(old);
2673    true
2674}
2675
2676impl<R: DeviceRuntime> Drop for CompletionReaper<R> {
2677    fn drop(&mut self) {
2678        let records = match self.state.get_mut() {
2679            Ok(state) => std::mem::take(&mut state.slots),
2680            Err(poisoned) => std::mem::take(&mut poisoned.into_inner().slots),
2681        };
2682        if records.is_empty() {
2683            return;
2684        }
2685        // No active reservation or observer can coexist with the final reaper
2686        // Arc: both retain or first upgrade that Arc. Record locks are therefore
2687        // uncontended here, and fail-closing each lane is an atomic operation.
2688        for record in records.values() {
2689            fail_close_completion_record(record);
2690        }
2691        for (slot_id, record) in records {
2692            let domain = record
2693                .lock()
2694                .unwrap_or_else(std::sync::PoisonError::into_inner)
2695                .deferred_cleanup_domain();
2696            if let Some(domain) = domain {
2697                // One slot per task prevents a blocked backend lane from
2698                // withholding unrelated lanes in the same plan. The global
2699                // registry mutex is never held during either recovery call.
2700                defer_device_cleanup(
2701                    domain,
2702                    DeferredCompletionCleanup::new(BTreeMap::from([(slot_id, record)])),
2703                );
2704            }
2705        }
2706    }
2707}
2708
2709#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
2710#[must_use = "a completion readback request identifies one exact typed logical backing range"]
2711pub struct CompletionReadbackRequest {
2712    node_id: NodeId,
2713    participant_index: u32,
2714    resource_id: ResourceId,
2715    expected_usage: BufferUsage,
2716    logical_offset_bytes: u64,
2717    output_layout: HostTransferLayout,
2718}
2719
2720impl CompletionReadbackRequest {
2721    pub fn new(
2722        node_id: NodeId,
2723        participant_index: u32,
2724        resource_id: ResourceId,
2725        logical_offset_bytes: u64,
2726        output_layout: HostTransferLayout,
2727    ) -> Result<Self, VNextError> {
2728        Self::new_typed(
2729            node_id,
2730            participant_index,
2731            resource_id,
2732            BufferUsage::Activations,
2733            logical_offset_bytes,
2734            output_layout,
2735        )
2736    }
2737
2738    pub fn new_typed(
2739        node_id: NodeId,
2740        participant_index: u32,
2741        resource_id: ResourceId,
2742        expected_usage: BufferUsage,
2743        logical_offset_bytes: u64,
2744        output_layout: HostTransferLayout,
2745    ) -> Result<Self, VNextError> {
2746        let output_bytes = output_layout.byte_len()?;
2747        if logical_offset_bytes.checked_add(output_bytes).is_none() {
2748            return Err(invalid_completion(
2749                "completion readback logical range overflows u64",
2750            ));
2751        }
2752        Ok(Self {
2753            node_id,
2754            participant_index,
2755            resource_id,
2756            expected_usage,
2757            logical_offset_bytes,
2758            output_layout,
2759        })
2760    }
2761
2762    pub fn node_id(&self) -> &NodeId {
2763        &self.node_id
2764    }
2765
2766    pub fn resource_id(&self) -> &ResourceId {
2767        &self.resource_id
2768    }
2769
2770    pub const fn expected_usage(&self) -> BufferUsage {
2771        self.expected_usage
2772    }
2773
2774    pub const fn participant_index(&self) -> u32 {
2775        self.participant_index
2776    }
2777
2778    pub const fn logical_offset_bytes(&self) -> u64 {
2779        self.logical_offset_bytes
2780    }
2781
2782    pub const fn output_layout(&self) -> HostTransferLayout {
2783        self.output_layout
2784    }
2785}
2786
2787#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
2788#[must_use = "a completion readback batch must cover one exact plan-node participant set"]
2789pub struct CompletionReadbackBatchRequest {
2790    requests: Vec<CompletionReadbackRequest>,
2791}
2792
2793impl CompletionReadbackBatchRequest {
2794    pub fn new(requests: Vec<CompletionReadbackRequest>) -> Result<Self, VNextError> {
2795        let Some(first) = requests.first() else {
2796            return Err(invalid_completion(
2797                "completion readback batch cannot be empty",
2798            ));
2799        };
2800        if requests.iter().enumerate().any(|(index, request)| {
2801            usize::try_from(request.participant_index()).ok() != Some(index)
2802                || request.node_id() != first.node_id()
2803                || request.resource_id() != first.resource_id()
2804                || request.expected_usage() != first.expected_usage()
2805                || request.output_layout().element_type() != first.output_layout().element_type()
2806        }) {
2807            return Err(invalid_completion(
2808                "completion readback batch must use canonical participant order and one typed node/resource/element group",
2809            ));
2810        }
2811        Ok(Self { requests })
2812    }
2813
2814    pub fn requests(&self) -> &[CompletionReadbackRequest] {
2815        &self.requests
2816    }
2817
2818    pub fn len(&self) -> usize {
2819        self.requests.len()
2820    }
2821
2822    pub fn is_empty(&self) -> bool {
2823        self.requests.is_empty()
2824    }
2825
2826    fn validate_for(&self, batch_identity: &BatchOperationIdentity) -> Result<(), VNextError> {
2827        let node_id = self.requests[0].node_id();
2828        let node_index = batch_identity.node_index(node_id).ok_or_else(|| {
2829            invalid_completion("completion readback batch node is absent from its submission")
2830        })?;
2831        let participant_count = batch_identity
2832            .node_participant_count(node_index)
2833            .ok_or_else(|| {
2834                invalid_completion("completion readback batch node is absent from its submission")
2835            })?;
2836        if participant_count != self.requests.len() {
2837            return Err(invalid_completion(
2838                "completion readback batch must cover every submitted node participant exactly once",
2839            ));
2840        }
2841        Ok(())
2842    }
2843
2844    fn into_requests(self) -> Vec<CompletionReadbackRequest> {
2845        self.requests
2846    }
2847}
2848
2849#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
2850#[must_use = "successful completion output bytes are exact readback evidence"]
2851pub struct CompletionReadbackOutput {
2852    request: CompletionReadbackRequest,
2853    bytes: Vec<u8>,
2854    sha256: String,
2855    #[serde(skip)]
2856    timing: DeviceTimingMeasurement<CompletionReadbackTiming>,
2857}
2858
2859impl CompletionReadbackOutput {
2860    fn new(request: CompletionReadbackRequest, readback: LaneReadback) -> Result<Self, VNextError> {
2861        request.output_layout.validate_bytes(readback.bytes.len())?;
2862        let sha256 = format!("{:x}", Sha256::digest(&readback.bytes));
2863        Ok(Self {
2864            request,
2865            bytes: readback.bytes,
2866            sha256,
2867            timing: readback.timing,
2868        })
2869    }
2870
2871    pub fn request(&self) -> &CompletionReadbackRequest {
2872        &self.request
2873    }
2874
2875    pub fn bytes(&self) -> &[u8] {
2876        &self.bytes
2877    }
2878
2879    pub fn sha256(&self) -> &str {
2880        &self.sha256
2881    }
2882
2883    pub const fn timing(&self) -> DeviceTimingMeasurement<CompletionReadbackTiming> {
2884        self.timing
2885    }
2886
2887    pub fn into_bytes(self) -> Vec<u8> {
2888        self.bytes
2889    }
2890}
2891
2892#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
2893#[serde(rename_all = "snake_case", tag = "status", content = "detail")]
2894pub enum CompletionReadbackDisposition {
2895    Succeeded(CompletionReadbackOutput),
2896    NotAttempted(CompletionReadbackRequest),
2897    FailedButQuiescent {
2898        request: CompletionReadbackRequest,
2899        failures: Vec<IdentifiedFailure>,
2900    },
2901    ContractFailedButQuiescent {
2902        request: CompletionReadbackRequest,
2903        failure: QuiescentCompletionContractFailure,
2904    },
2905}
2906
2907#[derive(Serialize)]
2908#[serde(rename_all = "snake_case", tag = "status", content = "detail")]
2909enum CompletionReadbackDispositionFingerprint<'a> {
2910    Succeeded {
2911        request: &'a CompletionReadbackRequest,
2912        output_sha256: &'a str,
2913    },
2914    NotAttempted {
2915        request: &'a CompletionReadbackRequest,
2916    },
2917    FailedButQuiescent {
2918        request: &'a CompletionReadbackRequest,
2919        failures: &'a [IdentifiedFailure],
2920    },
2921    ContractFailedButQuiescent {
2922        request: &'a CompletionReadbackRequest,
2923        failure: &'a str,
2924    },
2925}
2926
2927impl<'a> From<&'a CompletionReadbackDisposition> for CompletionReadbackDispositionFingerprint<'a> {
2928    fn from(disposition: &'a CompletionReadbackDisposition) -> Self {
2929        match disposition {
2930            CompletionReadbackDisposition::Succeeded(output) => Self::Succeeded {
2931                request: output.request(),
2932                output_sha256: output.sha256(),
2933            },
2934            CompletionReadbackDisposition::NotAttempted(request) => Self::NotAttempted { request },
2935            CompletionReadbackDisposition::FailedButQuiescent { request, failures } => {
2936                Self::FailedButQuiescent { request, failures }
2937            }
2938            CompletionReadbackDisposition::ContractFailedButQuiescent { request, failure } => {
2939                Self::ContractFailedButQuiescent {
2940                    request,
2941                    failure: failure.reason(),
2942                }
2943            }
2944        }
2945    }
2946}
2947
2948struct CompletionReadbackDispositionFingerprints<'a>(&'a [CompletionReadbackDisposition]);
2949
2950impl Serialize for CompletionReadbackDispositionFingerprints<'_> {
2951    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
2952    where
2953        S: serde::Serializer,
2954    {
2955        let mut sequence = serializer.serialize_seq(Some(self.0.len()))?;
2956        for disposition in self.0 {
2957            sequence
2958                .serialize_element(&CompletionReadbackDispositionFingerprint::from(disposition))?;
2959        }
2960        sequence.end()
2961    }
2962}
2963
2964fn completion_readback_batch_fingerprint(
2965    completion_fingerprint: &str,
2966    dispositions: &[CompletionReadbackDisposition],
2967) -> String {
2968    #[derive(Serialize)]
2969    struct FingerprintInput<'a> {
2970        domain: &'static str,
2971        completion_fingerprint: &'a str,
2972        dispositions: CompletionReadbackDispositionFingerprints<'a>,
2973    }
2974    canonical_completion_fingerprint(&FingerprintInput {
2975        domain: "ferrum.runtime-vnext.completion-readback-batch.v2",
2976        completion_fingerprint,
2977        dispositions: CompletionReadbackDispositionFingerprints(dispositions),
2978    })
2979}
2980
2981#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
2982#[must_use = "a terminal readback receipt couples output evidence to its exact completion"]
2983pub struct CompletionReadbackReceipt {
2984    completion: OperationCompletionReceipt,
2985    disposition: CompletionReadbackDisposition,
2986    readback_timing: Option<DeviceTimingMeasurement<CompletionReadbackTiming>>,
2987    fingerprint: String,
2988}
2989
2990impl CompletionReadbackReceipt {
2991    fn new(
2992        completion: OperationCompletionReceipt,
2993        disposition: CompletionReadbackDisposition,
2994    ) -> Self {
2995        let readback_timing = completion
2996            .fence_timing()
2997            .timing_mode()
2998            .completion_enabled()
2999            .then(|| readback_timing_for_disposition(&disposition));
3000        #[derive(Serialize)]
3001        struct FingerprintInput<'a> {
3002            domain: &'static str,
3003            completion_fingerprint: &'a str,
3004            request: &'a CompletionReadbackRequest,
3005            output_sha256: Option<&'a str>,
3006            failures: Option<&'a [IdentifiedFailure]>,
3007            contract_failure: Option<&'a str>,
3008        }
3009        let (request, output_sha256, failures, contract_failure) = match &disposition {
3010            CompletionReadbackDisposition::Succeeded(output) => {
3011                (output.request(), Some(output.sha256()), None, None)
3012            }
3013            CompletionReadbackDisposition::NotAttempted(request) => (request, None, None, None),
3014            CompletionReadbackDisposition::FailedButQuiescent { request, failures } => {
3015                (request, None, Some(failures.as_slice()), None)
3016            }
3017            CompletionReadbackDisposition::ContractFailedButQuiescent { request, failure } => {
3018                (request, None, None, Some(failure.reason()))
3019            }
3020        };
3021        let fingerprint = canonical_completion_fingerprint(&FingerprintInput {
3022            domain: "ferrum.runtime-vnext.completion-readback.v1",
3023            completion_fingerprint: completion.fingerprint(),
3024            request,
3025            output_sha256,
3026            failures,
3027            contract_failure,
3028        });
3029        Self {
3030            completion,
3031            disposition,
3032            readback_timing,
3033            fingerprint,
3034        }
3035    }
3036
3037    pub fn completion(&self) -> &OperationCompletionReceipt {
3038        &self.completion
3039    }
3040
3041    pub fn disposition(&self) -> &CompletionReadbackDisposition {
3042        &self.disposition
3043    }
3044
3045    pub const fn readback_timing(
3046        &self,
3047    ) -> Option<DeviceTimingMeasurement<CompletionReadbackTiming>> {
3048        self.readback_timing
3049    }
3050
3051    pub fn fingerprint(&self) -> &str {
3052        &self.fingerprint
3053    }
3054}
3055
3056#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
3057#[must_use = "a terminal batch readback receipt owns all participant output evidence"]
3058pub struct CompletionReadbackBatchReceipt {
3059    completion: OperationCompletionReceipt,
3060    dispositions: Vec<CompletionReadbackDisposition>,
3061    readback_timings: Option<Vec<DeviceTimingMeasurement<CompletionReadbackTiming>>>,
3062    fingerprint: String,
3063}
3064
3065impl CompletionReadbackBatchReceipt {
3066    fn new(
3067        completion: OperationCompletionReceipt,
3068        dispositions: Vec<CompletionReadbackDisposition>,
3069    ) -> Self {
3070        let readback_timings = completion
3071            .fence_timing()
3072            .timing_mode()
3073            .completion_enabled()
3074            .then(|| {
3075                dispositions
3076                    .iter()
3077                    .map(readback_timing_for_disposition)
3078                    .collect()
3079            });
3080        let fingerprint =
3081            completion_readback_batch_fingerprint(completion.fingerprint(), &dispositions);
3082        Self {
3083            completion,
3084            dispositions,
3085            readback_timings,
3086            fingerprint,
3087        }
3088    }
3089
3090    pub fn completion(&self) -> &OperationCompletionReceipt {
3091        &self.completion
3092    }
3093
3094    pub fn dispositions(&self) -> &[CompletionReadbackDisposition] {
3095        &self.dispositions
3096    }
3097
3098    pub fn readback_timings(&self) -> Option<&[DeviceTimingMeasurement<CompletionReadbackTiming>]> {
3099        self.readback_timings.as_deref()
3100    }
3101
3102    pub fn fingerprint(&self) -> &str {
3103        &self.fingerprint
3104    }
3105}
3106
3107fn readback_timing_for_disposition(
3108    disposition: &CompletionReadbackDisposition,
3109) -> DeviceTimingMeasurement<CompletionReadbackTiming> {
3110    match disposition {
3111        CompletionReadbackDisposition::Succeeded(output) => output.timing(),
3112        CompletionReadbackDisposition::NotAttempted(_)
3113        | CompletionReadbackDisposition::FailedButQuiescent { .. }
3114        | CompletionReadbackDisposition::ContractFailedButQuiescent { .. } => {
3115            DeviceTimingMeasurement::NotRequested
3116        }
3117    }
3118}
3119
3120fn attempt_completion_readback<R: DeviceRuntime>(
3121    resources: &CompletionResourceLease<R>,
3122    lane: &Arc<ExecutionLane<R>>,
3123    batch_identity: &BatchOperationIdentity,
3124    completion_disposition: &OperationCompletionDisposition,
3125    timing_mode: DeviceTimingMode,
3126    request: CompletionReadbackRequest,
3127) -> CompletionReadbackDisposition {
3128    if !matches!(
3129        completion_disposition,
3130        OperationCompletionDisposition::Succeeded
3131    ) {
3132        return CompletionReadbackDisposition::NotAttempted(request);
3133    }
3134    let byte_len = match request.output_layout().byte_len() {
3135        Ok(byte_len) => byte_len,
3136        Err(error) => {
3137            return CompletionReadbackDisposition::ContractFailedButQuiescent {
3138                request,
3139                failure: QuiescentCompletionContractFailure::new(error.to_string()),
3140            };
3141        }
3142    };
3143    let semantic_end = match request.logical_offset_bytes().checked_add(byte_len) {
3144        Some(end) => end,
3145        None => {
3146            return CompletionReadbackDisposition::ContractFailedButQuiescent {
3147                request,
3148                failure: QuiescentCompletionContractFailure::new(
3149                    "completion readback semantic range overflows u64",
3150                ),
3151            };
3152        }
3153    };
3154    let readback_range = match resources.readback_range(
3155        request.node_id(),
3156        request.participant_index(),
3157        request.resource_id(),
3158        request.logical_offset_bytes()..semantic_end,
3159    ) {
3160        Ok(range) => range,
3161        Err(error) => {
3162            return CompletionReadbackDisposition::ContractFailedButQuiescent {
3163                request,
3164                failure: QuiescentCompletionContractFailure::new(error.to_string()),
3165            };
3166        }
3167    };
3168    let backing = match resources.backing_view(
3169        request.node_id(),
3170        request.participant_index(),
3171        request.resource_id(),
3172    ) {
3173        Ok(backing) => backing,
3174        Err(error) => {
3175            return CompletionReadbackDisposition::ContractFailedButQuiescent {
3176                request,
3177                failure: QuiescentCompletionContractFailure::new(error.to_string()),
3178            };
3179        }
3180    };
3181    let readback = catch_unwind(AssertUnwindSafe(|| {
3182        lane.readback_buffer(
3183            &backing,
3184            request.expected_usage(),
3185            readback_range.start,
3186            request.output_layout(),
3187            timing_mode,
3188        )
3189    }));
3190    match readback {
3191        Ok(Ok(readback)) => match CompletionReadbackOutput::new(request.clone(), readback) {
3192            Ok(output) => CompletionReadbackDisposition::Succeeded(output),
3193            Err(error) => CompletionReadbackDisposition::ContractFailedButQuiescent {
3194                request,
3195                failure: QuiescentCompletionContractFailure::new(error.to_string()),
3196            },
3197        },
3198        Ok(Err(LaneReadbackError::Contract(error))) => {
3199            CompletionReadbackDisposition::ContractFailedButQuiescent {
3200                request,
3201                failure: QuiescentCompletionContractFailure::new(error.to_string()),
3202            }
3203        }
3204        Ok(Err(LaneReadbackError::Device(error))) => {
3205            match classify_batch_device_error(lane.runtime(), batch_identity, &error) {
3206                Ok(failures) => {
3207                    CompletionReadbackDisposition::FailedButQuiescent { request, failures }
3208                }
3209                Err(error) => CompletionReadbackDisposition::ContractFailedButQuiescent {
3210                    request,
3211                    failure: QuiescentCompletionContractFailure::new(error.to_string()),
3212                },
3213            }
3214        }
3215        Err(_) => {
3216            lane.fail_closed();
3217            CompletionReadbackDisposition::ContractFailedButQuiescent {
3218                request,
3219                failure: QuiescentCompletionContractFailure::new(
3220                    "device runtime panicked during completion readback",
3221                ),
3222            }
3223        }
3224    }
3225}
3226
3227fn attempt_completion_readbacks<R: DeviceRuntime>(
3228    resources: &CompletionResourceLease<R>,
3229    lane: &Arc<ExecutionLane<R>>,
3230    batch_identity: &BatchOperationIdentity,
3231    completion_disposition: &OperationCompletionDisposition,
3232    timing_mode: DeviceTimingMode,
3233    request: CompletionReadbackBatchRequest,
3234) -> Vec<CompletionReadbackDisposition> {
3235    request
3236        .into_requests()
3237        .into_iter()
3238        .map(|request| {
3239            attempt_completion_readback(
3240                resources,
3241                lane,
3242                batch_identity,
3243                completion_disposition,
3244                timing_mode,
3245                request,
3246            )
3247        })
3248        .collect()
3249}
3250
3251fn attempt_completion_readback_collection<R: DeviceRuntime>(
3252    resources: &CompletionResourceLease<R>,
3253    lane: &Arc<ExecutionLane<R>>,
3254    batch_identity: &BatchOperationIdentity,
3255    completion_disposition: &OperationCompletionDisposition,
3256    timing_mode: DeviceTimingMode,
3257    request: CompletionReadbackCollectionRequest,
3258) -> Vec<CompletionReadbackDisposition> {
3259    request
3260        .into_requests()
3261        .into_iter()
3262        .map(|request| {
3263            attempt_completion_readback(
3264                resources,
3265                lane,
3266                batch_identity,
3267                completion_disposition,
3268                timing_mode,
3269                request,
3270            )
3271        })
3272        .collect()
3273}
3274
3275#[derive(Debug)]
3276#[must_use = "nonterminal completion readback observations retain invocation ownership"]
3277pub enum CompletionReadbackObservation {
3278    Pending,
3279    Terminal(CompletionReadbackReceipt),
3280    Indeterminate(Vec<IdentifiedFailure>),
3281    SubmissionIndeterminate,
3282    ObservationPanicked,
3283    Quarantined(CompletionQuarantineReceipt),
3284}
3285
3286#[derive(Debug)]
3287#[must_use = "nonterminal batch readback observations retain invocation ownership"]
3288pub enum CompletionReadbackBatchObservation {
3289    Pending,
3290    Terminal(CompletionReadbackBatchReceipt),
3291    Indeterminate(Vec<IdentifiedFailure>),
3292    SubmissionIndeterminate,
3293    ObservationPanicked,
3294    Quarantined(CompletionQuarantineReceipt),
3295}
3296
3297#[derive(Debug)]
3298#[must_use = "nonterminal completion observations retain invocation ownership"]
3299pub enum CompletionObservation {
3300    Pending,
3301    Terminal(OperationCompletionReceipt),
3302    Indeterminate(Vec<IdentifiedFailure>),
3303    SubmissionIndeterminate,
3304    ObservationPanicked,
3305    Quarantined(CompletionQuarantineReceipt),
3306}
3307
3308/// Weak recovery authority for a submit unwind where no fence was returned.
3309/// Only a successful lane-wide drain can release the retained invocation.
3310#[must_use = "an indeterminate submission must be drained or retained"]
3311pub struct IndeterminateSubmissionHandle<R: DeviceRuntime> {
3312    reaper: Weak<CompletionReaper<R>>,
3313    slot_id: CompletionSlotId,
3314}
3315
3316impl<R: DeviceRuntime> Clone for IndeterminateSubmissionHandle<R> {
3317    fn clone(&self) -> Self {
3318        Self {
3319            reaper: Weak::clone(&self.reaper),
3320            slot_id: self.slot_id,
3321        }
3322    }
3323}
3324
3325impl<R: DeviceRuntime> fmt::Debug for IndeterminateSubmissionHandle<R> {
3326    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
3327        formatter
3328            .debug_struct("IndeterminateSubmissionHandle")
3329            .field("slot_id", &self.slot_id)
3330            .finish_non_exhaustive()
3331    }
3332}
3333
3334impl<R: DeviceRuntime> IndeterminateSubmissionHandle<R> {
3335    pub const fn slot_id(&self) -> CompletionSlotId {
3336        self.slot_id
3337    }
3338
3339    pub fn recover_by_draining_lane(&self) -> Result<CompletionRecoveryOutcome, VNextError> {
3340        self.reaper
3341            .upgrade()
3342            .ok_or_else(|| invalid_completion("completion reaper owner was dropped"))?
3343            .recover_bound(self.slot_id)
3344    }
3345}
3346
3347/// Weakly bound observation authority. Dropping a handle cannot drop or reap
3348/// the scheduler-owned completion registry.
3349#[must_use = "a submitted operation handle observes its exact completion slot"]
3350pub struct CompletionHandle<R: DeviceRuntime> {
3351    reaper: Weak<CompletionReaper<R>>,
3352    receipt: SubmittedOperationReceipt,
3353}
3354
3355impl<R: DeviceRuntime> Clone for CompletionHandle<R> {
3356    fn clone(&self) -> Self {
3357        Self {
3358            reaper: Weak::clone(&self.reaper),
3359            receipt: self.receipt.clone(),
3360        }
3361    }
3362}
3363
3364impl<R: DeviceRuntime> fmt::Debug for CompletionHandle<R> {
3365    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
3366        formatter
3367            .debug_struct("CompletionHandle")
3368            .field("receipt", &self.receipt)
3369            .finish_non_exhaustive()
3370    }
3371}
3372
3373impl<R: DeviceRuntime> CompletionHandle<R> {
3374    pub fn receipt(&self) -> &SubmittedOperationReceipt {
3375        &self.receipt
3376    }
3377
3378    pub fn batch_identity(&self) -> &BatchOperationIdentity {
3379        self.receipt.batch_identity()
3380    }
3381
3382    pub fn slot_id(&self) -> CompletionSlotId {
3383        self.receipt.slot_id()
3384    }
3385
3386    pub fn poll(&self) -> Result<CompletionObservation, VNextError> {
3387        self.reaper
3388            .upgrade()
3389            .ok_or_else(|| invalid_completion("completion reaper owner was dropped"))?
3390            .poll_bound(self.slot_id())
3391    }
3392
3393    pub fn wait(&self) -> Result<CompletionObservation, VNextError> {
3394        self.reaper
3395            .upgrade()
3396            .ok_or_else(|| invalid_completion("completion reaper owner was dropped"))?
3397            .wait_bound(self.slot_id())
3398    }
3399
3400    pub fn wait_with_readback(
3401        &self,
3402        request: CompletionReadbackRequest,
3403    ) -> Result<CompletionReadbackObservation, VNextError> {
3404        self.reaper
3405            .upgrade()
3406            .ok_or_else(|| invalid_completion("completion reaper owner was dropped"))?
3407            .wait_bound_with_readback(self.slot_id(), request)
3408    }
3409
3410    pub fn wait_with_readbacks(
3411        &self,
3412        request: CompletionReadbackBatchRequest,
3413    ) -> Result<CompletionReadbackBatchObservation, VNextError> {
3414        self.reaper
3415            .upgrade()
3416            .ok_or_else(|| invalid_completion("completion reaper owner was dropped"))?
3417            .wait_bound_with_readbacks(self.slot_id(), request)
3418    }
3419
3420    pub fn wait_with_readback_collection(
3421        &self,
3422        request: CompletionReadbackCollectionRequest,
3423    ) -> Result<CompletionReadbackCollectionObservation, VNextError> {
3424        self.reaper
3425            .upgrade()
3426            .ok_or_else(|| invalid_completion("completion reaper owner was dropped"))?
3427            .wait_bound_with_readback_collection(self.slot_id(), request)
3428    }
3429}
3430
3431pub(crate) struct CompletionReservation<R: DeviceRuntime> {
3432    reaper: Arc<CompletionReaper<R>>,
3433    record: SharedCompletionRecord<R>,
3434    slot_id: CompletionSlotId,
3435    resources: Option<CompletionResourceLease<R>>,
3436    lane: Option<Arc<ExecutionLane<R>>>,
3437    batch_identity: Option<BatchOperationIdentity>,
3438    receipt: Option<SubmittedOperationReceipt>,
3439    record_receipt: Option<SubmittedOperationReceipt>,
3440    submission_may_have_happened: bool,
3441    finished: bool,
3442}
3443
3444impl<R: DeviceRuntime> CompletionReservation<R> {
3445    pub(crate) fn invocation(&self) -> &InvocationResourceLease<R> {
3446        match self
3447            .resources
3448            .as_ref()
3449            .expect("live completion reservation owns submission resources")
3450        {
3451            CompletionResourceLease::Invocation(invocation) => invocation,
3452            CompletionResourceLease::Wave(_) => {
3453                unreachable!("single-operation reservation cannot own a submission wave")
3454            }
3455        }
3456    }
3457
3458    pub(crate) fn wave(&self) -> &PreparedStepSubmissionWave<R> {
3459        match self
3460            .resources
3461            .as_ref()
3462            .expect("live completion reservation owns submission resources")
3463        {
3464            CompletionResourceLease::Wave(wave) => wave,
3465            CompletionResourceLease::Invocation(_) => {
3466                unreachable!("wave reservation cannot own single-operation resources")
3467            }
3468        }
3469    }
3470
3471    pub(crate) fn backing_view(
3472        &self,
3473        node_id: &NodeId,
3474        participant_index: u32,
3475        resource_id: &ResourceId,
3476    ) -> Result<LogicalBackingBufferView<'_, R::Buffer>, VNextError> {
3477        self.resources
3478            .as_ref()
3479            .expect("live completion reservation owns submission resources")
3480            .backing_view(node_id, participant_index, resource_id)
3481    }
3482
3483    pub(crate) fn encode_backing_initializations(
3484        &self,
3485        runtime: &R,
3486        commands: &mut DeviceCommandBatch<R::Command>,
3487    ) -> Result<usize, BackingInitializationEncodeError<R::Error>> {
3488        self.resources
3489            .as_ref()
3490            .expect("live completion reservation owns submission resources")
3491            .encode_backing_initializations(runtime, commands)
3492    }
3493
3494    pub(crate) fn mark_submission_started(&mut self) {
3495        self.submission_may_have_happened = true;
3496    }
3497
3498    pub(crate) fn definitely_not_submitted(
3499        mut self,
3500    ) -> Result<DefinitelyNotSubmittedRetryAuthority<R>, VNextError> {
3501        self.remove_reserved_slot();
3502        self.submission_may_have_happened = false;
3503        let resources = self
3504            .resources
3505            .take()
3506            .expect("reservation owns definitely-not-submitted resources");
3507        let CompletionResourceLease::Invocation(invocation) = resources else {
3508            return Err(invalid_completion(
3509                "single-operation retry requested from a submission wave reservation",
3510            ));
3511        };
3512        let retry = invocation.definitely_not_submitted()?;
3513        self.finished = true;
3514        Ok(retry)
3515    }
3516
3517    pub(crate) fn definitely_not_submitted_wave(
3518        mut self,
3519    ) -> Result<DefinitelyNotSubmittedWaveRetryAuthority<R>, VNextError> {
3520        self.remove_reserved_slot();
3521        self.submission_may_have_happened = false;
3522        let resources = self
3523            .resources
3524            .take()
3525            .expect("reservation owns definitely-not-submitted resources");
3526        let CompletionResourceLease::Wave(wave) = resources else {
3527            return Err(invalid_completion(
3528                "wave retry requested from a single-operation reservation",
3529            ));
3530        };
3531        let retry = wave.definitely_not_submitted()?;
3532        self.finished = true;
3533        Ok(retry)
3534    }
3535
3536    pub(crate) fn arm(
3537        mut self,
3538        fence: R::Fence,
3539        timing_mode: DeviceTimingMode,
3540    ) -> Result<CompletionHandle<R>, (VNextError, CompletionHandle<R>)> {
3541        let resources = self.resources.take().expect("reservation owns resources");
3542        let lane = self.lane.take().expect("reservation owns lane");
3543        let batch_identity = self
3544            .batch_identity
3545            .take()
3546            .expect("reservation owns batch identity");
3547        let receipt = self.receipt.take().expect("reservation owns receipt");
3548        let record_receipt = self
3549            .record_receipt
3550            .take()
3551            .expect("reservation owns record receipt");
3552        let mut record = match self.record.lock() {
3553            Ok(record) => record,
3554            Err(poisoned) => poisoned.into_inner(),
3555        };
3556        if !matches!(&*record, CompletionRecord::Reserved) {
3557            lane.fail_closed();
3558            std::mem::forget((resources, lane, fence));
3559            self.finished = true;
3560            panic!("completion reservation changed after submission");
3561        }
3562        *record = CompletionRecord::InFlight {
3563            resources,
3564            lane,
3565            fence,
3566            batch_identity,
3567            receipt: record_receipt,
3568            timing_mode,
3569            recovery_state: CompletionRecoveryState::Unobserved,
3570        };
3571        let transition = match &mut *record {
3572            CompletionRecord::InFlight {
3573                resources, lane, ..
3574            } => match resources.mark_submission_fence_installed() {
3575                Ok(()) => Ok(()),
3576                Err(error) => {
3577                    // Submission already returned a fence. Any partial ownership
3578                    // transition is therefore unknown, never cleanly unsubmitted.
3579                    resources.mark_submission_indeterminate();
3580                    lane.fail_closed();
3581                    Err(error)
3582                }
3583            },
3584            _ => unreachable!("completion record was just armed"),
3585        };
3586        drop(record);
3587        self.finished = true;
3588        let handle = CompletionHandle {
3589            reaper: Arc::downgrade(&self.reaper),
3590            receipt,
3591        };
3592        match transition {
3593            Ok(()) => Ok(handle),
3594            Err(error) => Err((error, handle)),
3595        }
3596    }
3597
3598    pub(crate) fn submission_indeterminate(mut self) -> IndeterminateSubmissionHandle<R> {
3599        let mut resources = self.resources.take().expect("reservation owns resources");
3600        resources.mark_submission_indeterminate();
3601        let lane = self.lane.take().expect("reservation owns lane");
3602        let batch_identity = self
3603            .batch_identity
3604            .take()
3605            .expect("reservation owns batch identity");
3606        let mut record = match self.record.lock() {
3607            Ok(record) => record,
3608            Err(poisoned) => poisoned.into_inner(),
3609        };
3610        if matches!(&*record, CompletionRecord::Reserved) {
3611            *record = CompletionRecord::SubmissionIndeterminate {
3612                resources,
3613                lane,
3614                batch_identity,
3615            };
3616        } else {
3617            lane.fail_closed();
3618            std::mem::forget((resources, lane, batch_identity));
3619        }
3620        drop(record);
3621        self.finished = true;
3622        IndeterminateSubmissionHandle {
3623            reaper: Arc::downgrade(&self.reaper),
3624            slot_id: self.slot_id,
3625        }
3626    }
3627
3628    fn remove_reserved_slot(&mut self) {
3629        self.reaper.remove_exact(self.slot_id, &self.record);
3630        let mut record = match self.record.lock() {
3631            Ok(record) => record,
3632            Err(poisoned) => poisoned.into_inner(),
3633        };
3634        if matches!(&*record, CompletionRecord::Reserved) {
3635            *record = CompletionRecord::Reaped;
3636        }
3637    }
3638}
3639
3640impl<R: DeviceRuntime> Drop for CompletionReservation<R> {
3641    fn drop(&mut self) {
3642        if self.finished {
3643            return;
3644        }
3645        if self.submission_may_have_happened {
3646            let Some(mut resources) = self.resources.take() else {
3647                return;
3648            };
3649            resources.mark_submission_indeterminate();
3650            let Some(lane) = self.lane.take() else {
3651                std::mem::forget(resources);
3652                return;
3653            };
3654            let Some(batch_identity) = self.batch_identity.take() else {
3655                std::mem::forget((resources, lane));
3656                return;
3657            };
3658            lane.fail_closed();
3659            let mut record = match self.record.lock() {
3660                Ok(record) => record,
3661                Err(poisoned) => poisoned.into_inner(),
3662            };
3663            if matches!(&*record, CompletionRecord::Reserved) {
3664                *record = CompletionRecord::SubmissionIndeterminate {
3665                    resources,
3666                    lane,
3667                    batch_identity,
3668                };
3669            } else {
3670                std::mem::forget((resources, lane, batch_identity));
3671            }
3672        } else {
3673            self.remove_reserved_slot();
3674        }
3675    }
3676}
3677
3678#[cfg(test)]
3679mod fingerprint_tests {
3680    use super::*;
3681    use crate::vnext::ElementType;
3682
3683    const LARGE_READBACK_BYTES: usize = 512 * 1024;
3684
3685    fn request() -> CompletionReadbackRequest {
3686        CompletionReadbackRequest::new(
3687            NodeId::try_from("node.fingerprint-test".to_owned()).unwrap(),
3688            0,
3689            ResourceId::try_from("resource.fingerprint-test".to_owned()).unwrap(),
3690            0,
3691            HostTransferLayout::new(ElementType::U8, LARGE_READBACK_BYTES as u64).unwrap(),
3692        )
3693        .unwrap()
3694    }
3695
3696    fn successful_output(fill: u8) -> CompletionReadbackOutput {
3697        let bytes = vec![fill; LARGE_READBACK_BYTES];
3698        let sha256 = format!("{:x}", Sha256::digest(&bytes));
3699        CompletionReadbackOutput {
3700            request: request(),
3701            bytes,
3702            sha256,
3703            timing: DeviceTimingMeasurement::NotRequested,
3704        }
3705    }
3706
3707    #[test]
3708    fn batch_fingerprint_is_bounded_by_digest_evidence_not_output_bytes() {
3709        let first = vec![CompletionReadbackDisposition::Succeeded(successful_output(
3710            0x5a,
3711        ))];
3712        let encoded = serde_json::to_vec(&CompletionReadbackDispositionFingerprints(&first))
3713            .expect("fingerprint evidence serializes");
3714        assert!(encoded.len() < 1024, "encoded {} bytes", encoded.len());
3715        let first_sha = match &first[0] {
3716            CompletionReadbackDisposition::Succeeded(output) => output.sha256(),
3717            _ => unreachable!(),
3718        };
3719        assert!(String::from_utf8(encoded).unwrap().contains(first_sha));
3720
3721        let completion_fingerprint = "a".repeat(64);
3722        let first_fingerprint =
3723            completion_readback_batch_fingerprint(&completion_fingerprint, &first);
3724        assert_eq!(
3725            first_fingerprint,
3726            completion_readback_batch_fingerprint(&completion_fingerprint, &first)
3727        );
3728
3729        let second = vec![CompletionReadbackDisposition::Succeeded(successful_output(
3730            0xa5,
3731        ))];
3732        assert_ne!(
3733            first_fingerprint,
3734            completion_readback_batch_fingerprint(&completion_fingerprint, &second)
3735        );
3736
3737        let not_attempted = vec![CompletionReadbackDisposition::NotAttempted(request())];
3738        assert_ne!(
3739            first_fingerprint,
3740            completion_readback_batch_fingerprint(&completion_fingerprint, &not_attempted)
3741        );
3742    }
3743}