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