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#[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 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 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 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 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#[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#[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 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 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 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 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
2566fn 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 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 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#[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#[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 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, ¬_attempted)
3869 );
3870 }
3871}