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