1use super::{
2 defer_device_cleanup, invalid_resource, sequence_slot_is_poisoned, AdmissionDecision,
3 AdmissionDeferred, AdmissionFitPolicy, AdmissionPreflightDecision, AdmissionPressureAction,
4 AdmissionRejected, AllocationLifetime, Arc, AtomicU64, BTreeMap, BTreeSet,
5 BackingPrepareDecision, BatchStepId, CapacityClaimDecision, CapacityVector,
6 DeferredDeviceCleanupDisposition, DeferredDeviceCleanupTask, DeviceBufferRetention,
7 DeviceRuntime, Digest, DynamicBackingClaimScope, DynamicBackingDeferred,
8 DynamicDeferredMaintenanceOutcome, DynamicResourceShape, ExecutionFrameId,
9 InitialSequenceAdmissionDecision, LogicalAdmissionCoordinatorId, LogicalAdmissionLease,
10 LogicalBackingBufferView, LogicalBackingSliceAuthority, LogicalBackingSliceEvidence,
11 LogicalCapacityLease, LogicalRequestLease, ManuallyDrop, Mutex, NonZeroU64, Ordering,
12 ParticipantNodeKey, PlanBackingDeferral, PlanCapacityWaitRegistration, PreparedBackingClaim,
13 RequestAdmissionDecision, RequestAuthorityId, RequestIdentity, RequestResourceAdmissionRequest,
14 RequestStateHazardRegistration, ResourceId, ResourceWorkShape, RunId, SequenceAuthorityId,
15 SequenceRecoveryRegistry, SequenceResourceAdmissionRequest, Serialize, Sha256,
16 StaticProvisioningLease, TrustedPlanRuntimeBinding, TrustedPlanRuntimeEvidence, VNextError,
17 Weak, SEQUENCE_DISPATCH_POISONED_BIT,
18};
19use crate::vnext::CapacityAvailabilitySource;
20
21mod state_transfer;
22pub(crate) use state_transfer::*;
23mod completed_boundary;
24pub(crate) use completed_boundary::*;
25
26const SEQUENCE_DISPATCH_COUNT_MASK: u64 = SEQUENCE_DISPATCH_POISONED_BIT - 1;
27
28pub(super) fn sequence_dispatch_is_poisoned(gate: &AtomicU64) -> bool {
29 gate.load(Ordering::Acquire) & SEQUENCE_DISPATCH_POISONED_BIT != 0
30}
31
32pub(super) struct SequenceDispatchGuard<'a> {
33 gate: &'a AtomicU64,
34}
35
36impl Drop for SequenceDispatchGuard<'_> {
37 fn drop(&mut self) {
38 let previous = self.gate.fetch_sub(1, Ordering::AcqRel);
39 debug_assert!(previous & SEQUENCE_DISPATCH_COUNT_MASK > 0);
40 }
41}
42
43pub(super) fn enter_sequence_dispatch(
44 gate: &AtomicU64,
45) -> Result<SequenceDispatchGuard<'_>, VNextError> {
46 gate.fetch_update(Ordering::AcqRel, Ordering::Acquire, |state| {
47 if state & SEQUENCE_DISPATCH_POISONED_BIT != 0
48 || state & SEQUENCE_DISPATCH_COUNT_MASK == SEQUENCE_DISPATCH_COUNT_MASK
49 {
50 None
51 } else {
52 Some(state + 1)
53 }
54 })
55 .map_err(|state| {
56 if state & SEQUENCE_DISPATCH_POISONED_BIT != 0 {
57 invalid_resource("poisoned resource pool cannot dispatch another operation")
58 } else {
59 invalid_resource("resource pool dispatch counter is exhausted")
60 }
61 })?;
62 Ok(SequenceDispatchGuard { gate })
63}
64
65pub enum RequestResourceAdmissionDecision<R>
66where
67 R: DeviceRuntime,
68{
69 Admitted(Arc<AdmittedRequestResources<R>>),
70 Deferred(AdmissionDeferred),
71 BackingDeferred(RequestBackingDeferral<R>),
72 PermanentRejected(AdmissionRejected),
73}
74
75pub enum InitialSequenceResourceAdmissionDecision<R>
76where
77 R: DeviceRuntime,
78{
79 Admitted(Arc<AdmittedSequenceResources<R>>),
80 Deferred(AdmissionDeferred),
81 BackingDeferred(InitialSequenceBackingDeferral<R>),
82 PermanentRejected(AdmissionRejected),
83}
84
85#[must_use = "initial sequence backing deferral owns its exact bundle attempt"]
88pub struct InitialSequenceBackingDeferral<R>
89where
90 R: DeviceRuntime,
91{
92 evidence: DynamicBackingDeferred,
93 plan: TrustedPlanRuntimeBinding<R>,
94}
95
96impl<R> InitialSequenceBackingDeferral<R>
97where
98 R: DeviceRuntime,
99{
100 pub fn evidence(&self) -> &DynamicBackingDeferred {
101 &self.evidence
102 }
103
104 pub fn maintain(&self) -> Result<DynamicDeferredMaintenanceOutcome, VNextError> {
105 if self.evidence.scope() != DynamicBackingClaimScope::InitialSequenceBundle {
106 return Err(invalid_resource(
107 "initial sequence maintenance requires a bundle-scoped deferral",
108 ));
109 }
110 let _lifecycle = self
111 .plan
112 .resources
113 .read_lifecycle("maintain deferred initial sequence backing")?;
114 self.plan
115 .resources
116 .maintenance_controller
117 .maintain_for_live_deferred(&self.evidence)
118 }
119
120 pub fn register_waiter(&self) -> Result<PlanCapacityWaitRegistration<R>, VNextError> {
121 self.plan.register_backing_waiter(&self.evidence)
122 }
123}
124
125#[must_use = "request backing deferral owns its exact admission attempt"]
129pub struct RequestBackingDeferral<R>
130where
131 R: DeviceRuntime,
132{
133 evidence: DynamicBackingDeferred,
134 plan: TrustedPlanRuntimeBinding<R>,
135 work_shape: ResourceWorkShape,
136 run_id: RunId,
137 request_id: RequestIdentity,
138}
139
140impl<R> RequestBackingDeferral<R>
141where
142 R: DeviceRuntime,
143{
144 pub fn evidence(&self) -> &DynamicBackingDeferred {
145 &self.evidence
146 }
147
148 pub fn work_shape(&self) -> &ResourceWorkShape {
149 &self.work_shape
150 }
151
152 pub fn run_id(&self) -> &RunId {
153 &self.run_id
154 }
155
156 pub fn request_id(&self) -> &RequestIdentity {
157 &self.request_id
158 }
159
160 pub fn maintain(&self) -> Result<DynamicDeferredMaintenanceOutcome, VNextError> {
161 self.plan
162 .maintain_request_backing_for_deferred(&self.evidence)
163 }
164
165 pub fn register_waiter(&self) -> Result<PlanCapacityWaitRegistration<R>, VNextError> {
166 self.plan.register_backing_waiter(&self.evidence)
167 }
168}
169
170impl<R> TrustedPlanRuntimeBinding<R>
171where
172 R: DeviceRuntime,
173{
174 pub(super) fn maintain_request_backing_for_deferred(
177 &self,
178 deferred: &DynamicBackingDeferred,
179 ) -> Result<DynamicDeferredMaintenanceOutcome, VNextError> {
180 if deferred.scope() != DynamicBackingClaimScope::Request {
181 return Err(invalid_resource(
182 "request backing maintenance requires a request-lifetime deferral",
183 ));
184 }
185 let _lifecycle = self
186 .resources
187 .read_lifecycle("maintain deferred request backing")?;
188 self.resources
189 .maintenance_controller
190 .maintain_for_live_deferred(deferred)
191 }
192
193 pub fn try_admit_request(
196 &self,
197 request: RequestResourceAdmissionRequest,
198 run_id: RunId,
199 request_id: RequestIdentity,
200 ) -> Result<RequestResourceAdmissionDecision<R>, VNextError> {
201 let _lifecycle = self.resources.read_lifecycle("admit a request")?;
202 let RequestResourceAdmissionRequest {
203 work_shape,
204 fit_policy,
205 pressure_action,
206 } = request;
207 let request_shape = work_shape.fit_shape();
211 let (demand, requested_slices) = self.scoped_demand(
212 AllocationLifetime::Request,
213 None,
214 request_shape,
215 request_shape,
216 None,
217 fit_policy,
218 pressure_action,
219 )?;
220 if let Some(rejected) =
221 self.reject_impossible_plan_fit(demand.immediate_claim(), demand.fit_requirement())?
222 {
223 return Ok(RequestResourceAdmissionDecision::PermanentRejected(
224 rejected,
225 ));
226 }
227 let prepared = match self.prepare_backing_slices(requested_slices)? {
228 BackingPrepareDecision::Prepared(prepared) => prepared,
229 BackingPrepareDecision::Deferred(deferred) => {
230 return Ok(RequestResourceAdmissionDecision::BackingDeferred(
231 RequestBackingDeferral {
232 evidence: deferred,
233 plan: TrustedPlanRuntimeBinding {
234 resources: Arc::clone(&self.resources),
235 },
236 work_shape,
237 run_id,
238 request_id,
239 },
240 ));
241 }
242 };
243 match self.logical_admission().try_admit_request(&demand)? {
244 RequestAdmissionDecision::Admitted(logical_lease) => {
245 if !self.logical_admission().owns_request(&logical_lease) {
246 return Err(invalid_resource(
247 "request admission returned authority from another coordinator",
248 ));
249 }
250 let slices = prepared.commit();
251 Ok(RequestResourceAdmissionDecision::Admitted(Arc::new(
252 AdmittedRequestResources::new(
253 TrustedPlanRuntimeBinding {
254 resources: Arc::clone(&self.resources),
255 },
256 logical_lease,
257 slices,
258 work_shape,
259 run_id,
260 request_id,
261 )?,
262 )))
263 }
264 RequestAdmissionDecision::Deferred(deferred) => {
265 Ok(RequestResourceAdmissionDecision::Deferred(deferred))
266 }
267 RequestAdmissionDecision::PermanentRejected(rejected) => Ok(
268 RequestResourceAdmissionDecision::PermanentRejected(rejected),
269 ),
270 }
271 }
272
273 pub fn try_admit_initial_sequence(
277 &self,
278 request: RequestResourceAdmissionRequest,
279 sequence: SequenceResourceAdmissionRequest,
280 run_id: RunId,
281 request_id: RequestIdentity,
282 ) -> Result<InitialSequenceResourceAdmissionDecision<R>, VNextError> {
283 let _lifecycle = self
284 .resources
285 .read_lifecycle("admit an initial request/sequence bundle")?;
286 let RequestResourceAdmissionRequest {
287 work_shape: request_work_shape,
288 fit_policy: request_fit_policy,
289 pressure_action: request_pressure_action,
290 } = request;
291 let SequenceResourceAdmissionRequest {
292 work_shape: sequence_work_shape,
293 fit_policy: sequence_fit_policy,
294 pressure_action: sequence_pressure_action,
295 } = sequence;
296 if sequence_work_shape.fit_tokens() > request_work_shape.fit_tokens() {
297 return Err(invalid_resource(
298 "initial sequence token ceiling exceeds its request ceiling",
299 ));
300 }
301
302 let request_shape = request_work_shape.fit_shape();
303 let (request_demand, mut requested_slices) = self.scoped_demand(
304 AllocationLifetime::Request,
305 None,
306 request_shape,
307 request_shape,
308 None,
309 request_fit_policy,
310 request_pressure_action,
311 )?;
312 let sequence_immediate_shape = sequence_work_shape.immediate_shape();
313 let sequence_fit_shape = match sequence_fit_policy {
314 AdmissionFitPolicy::ImmediateOnly => sequence_immediate_shape,
315 AdmissionFitPolicy::FullInputMustFit => sequence_work_shape.fit_shape(),
316 };
317 let (sequence_demand, mut sequence_slices) = self.scoped_demand(
318 AllocationLifetime::Sequence,
319 None,
320 sequence_immediate_shape,
321 sequence_fit_shape,
322 None,
323 sequence_fit_policy,
324 sequence_pressure_action,
325 )?;
326
327 let request_resource_ids = requested_slices
328 .iter()
329 .flat_map(|request| request.projections.iter())
330 .map(|projection| projection.descriptor.base_resource_id().clone())
331 .collect::<BTreeSet<_>>();
332 let sequence_resource_ids = sequence_slices
333 .iter()
334 .flat_map(|request| request.projections.iter())
335 .map(|projection| projection.descriptor.base_resource_id().clone())
336 .collect::<BTreeSet<_>>();
337 if !request_resource_ids.is_disjoint(&sequence_resource_ids) {
338 return Err(invalid_resource(
339 "initial request and sequence backing resources are not disjoint",
340 ));
341 }
342 requested_slices.append(&mut sequence_slices);
343
344 let immediate = CapacityVector::checked_sum(
345 request_demand.immediate_claim(),
346 sequence_demand.immediate_claim(),
347 )?;
348 let fit = CapacityVector::checked_sum(
349 request_demand.fit_requirement(),
350 sequence_demand.fit_requirement(),
351 )?;
352 if let Some(rejected) = self.reject_impossible_plan_fit(&immediate, &fit)? {
353 return Ok(InitialSequenceResourceAdmissionDecision::PermanentRejected(
354 rejected,
355 ));
356 }
357
358 loop {
359 match self
360 .logical_admission()
361 .preflight_initial_sequence(&request_demand, &sequence_demand)?
362 {
363 AdmissionPreflightDecision::Eligible => {}
364 AdmissionPreflightDecision::Deferred(deferred) => {
365 return Ok(InitialSequenceResourceAdmissionDecision::Deferred(deferred));
366 }
367 AdmissionPreflightDecision::PermanentRejected(rejected) => {
368 return Ok(InitialSequenceResourceAdmissionDecision::PermanentRejected(
369 rejected,
370 ));
371 }
372 }
373
374 let prepared = match self.prepare_initial_sequence_backing_slices(&requested_slices)? {
375 BackingPrepareDecision::Prepared(prepared) => prepared,
376 BackingPrepareDecision::Deferred(evidence) => {
377 return Ok(InitialSequenceResourceAdmissionDecision::BackingDeferred(
378 InitialSequenceBackingDeferral {
379 evidence,
380 plan: TrustedPlanRuntimeBinding {
381 resources: Arc::clone(&self.resources),
382 },
383 },
384 ));
385 }
386 };
387
388 match self
389 .logical_admission()
390 .try_admit_initial_sequence(&request_demand, &sequence_demand)?
391 {
392 InitialSequenceAdmissionDecision::Admitted(logical) => {
393 let (request_logical, sequence_logical) = logical.into_parts();
394 if !self.logical_admission().owns_request(&request_logical)
395 || !self.logical_admission().owns(&sequence_logical)
396 || sequence_logical.request() != request_logical.request()
397 {
398 return Err(invalid_resource(
399 "initial bundle admission returned incoherent authority",
400 ));
401 }
402
403 let mut request_backing = Vec::new();
404 let mut sequence_backing = Vec::new();
405 let mut seen_request = BTreeSet::new();
406 let mut seen_sequence = BTreeSet::new();
407 for authority in prepared.commit() {
408 let resource_id = authority.resource_id().clone();
409 if request_resource_ids.contains(&resource_id) {
410 seen_request.insert(resource_id);
411 request_backing.push(authority);
412 } else if sequence_resource_ids.contains(&resource_id) {
413 seen_sequence.insert(resource_id);
414 sequence_backing.push(authority);
415 } else {
416 return Err(invalid_resource(
417 "initial bundle prepared an unrequested backing resource",
418 ));
419 }
420 }
421 if seen_request != request_resource_ids
422 || seen_sequence != sequence_resource_ids
423 {
424 return Err(invalid_resource(
425 "initial bundle backing did not cover every requested resource",
426 ));
427 }
428
429 let request = Arc::new(AdmittedRequestResources::new(
430 TrustedPlanRuntimeBinding {
431 resources: Arc::clone(&self.resources),
432 },
433 request_logical,
434 request_backing,
435 request_work_shape,
436 run_id,
437 request_id,
438 )?);
439 return Ok(InitialSequenceResourceAdmissionDecision::Admitted(
440 Arc::new(AdmittedSequenceResources::new(
441 request,
442 sequence_logical,
443 sequence_backing,
444 sequence_work_shape,
445 )?),
446 ));
447 }
448 InitialSequenceAdmissionDecision::Deferred => {
449 drop(prepared);
450 match self
451 .logical_admission()
452 .observe_initial_sequence(&request_demand, &sequence_demand)?
453 {
454 AdmissionPreflightDecision::Eligible => continue,
455 AdmissionPreflightDecision::Deferred(deferred) => {
456 return Ok(InitialSequenceResourceAdmissionDecision::Deferred(
457 deferred,
458 ));
459 }
460 AdmissionPreflightDecision::PermanentRejected(rejected) => {
461 return Ok(
462 InitialSequenceResourceAdmissionDecision::PermanentRejected(
463 rejected,
464 ),
465 );
466 }
467 }
468 }
469 InitialSequenceAdmissionDecision::PermanentRejected(rejected) => {
470 drop(prepared);
471 return Ok(InitialSequenceResourceAdmissionDecision::PermanentRejected(
472 rejected,
473 ));
474 }
475 }
476 }
477 }
478}
479
480#[must_use = "request resources release capacity after their last child sequence"]
484pub struct AdmittedRequestResources<R>
485where
486 R: DeviceRuntime,
487{
488 _request_state_hazard_registration: Option<RequestStateHazardRegistration>,
491 backing_slices: Vec<LogicalBackingSliceAuthority>,
492 logical_lease: LogicalRequestLease,
493 pub(super) plan: TrustedPlanRuntimeBinding<R>,
494 work_shape: ResourceWorkShape,
495 run_id: RunId,
496 request_id: RequestIdentity,
497}
498
499impl<R> AdmittedRequestResources<R>
500where
501 R: DeviceRuntime,
502{
503 pub(super) fn new(
504 plan: TrustedPlanRuntimeBinding<R>,
505 logical_lease: LogicalRequestLease,
506 backing_slices: Vec<LogicalBackingSliceAuthority>,
507 work_shape: ResourceWorkShape,
508 run_id: RunId,
509 request_id: RequestIdentity,
510 ) -> Result<Self, VNextError> {
511 if !plan.logical_admission().owns_request(&logical_lease) {
512 return Err(invalid_resource(
513 "logical request authority belongs to another coordinator",
514 ));
515 }
516 let request_state_hazard_registration = plan
517 .dynamic_pools()
518 .request_state_hazards
519 .register_request(logical_lease.request())?;
520 Ok(Self {
521 _request_state_hazard_registration: request_state_hazard_registration,
522 backing_slices,
523 logical_lease,
524 plan,
525 work_shape,
526 run_id,
527 request_id,
528 })
529 }
530
531 pub const fn request_authority(&self) -> RequestAuthorityId {
532 self.logical_lease.request()
533 }
534
535 pub fn coordinator_id(&self) -> LogicalAdmissionCoordinatorId {
536 self.logical_lease.coordinator_id()
537 }
538
539 pub fn run_id(&self) -> &RunId {
540 &self.run_id
541 }
542
543 pub fn request_id(&self) -> &RequestIdentity {
544 &self.request_id
545 }
546
547 pub fn backing_slices(&self) -> &[LogicalBackingSliceAuthority] {
548 &self.backing_slices
549 }
550
551 pub fn work_shape(&self) -> &ResourceWorkShape {
552 &self.work_shape
553 }
554
555 pub fn static_provisioning(&self) -> Option<&StaticProvisioningLease<R>> {
556 self.plan.static_provisioning()
557 }
558
559 pub(super) fn maintain_sequence_backing_for_deferred(
562 &self,
563 deferred: &DynamicBackingDeferred,
564 ) -> Result<DynamicDeferredMaintenanceOutcome, VNextError> {
565 if deferred.scope() != DynamicBackingClaimScope::Sequence {
566 return Err(invalid_resource(
567 "sequence backing maintenance requires a sequence-lifetime deferral",
568 ));
569 }
570 let _lifecycle = self
571 .plan
572 .resources
573 .read_lifecycle("maintain deferred sequence backing")?;
574 self.plan
575 .resources
576 .maintenance_controller
577 .maintain_for_live_deferred(deferred)
578 }
579
580 pub fn plan_evidence(&self) -> TrustedPlanRuntimeEvidence {
581 self.plan.evidence()
582 }
583
584 pub(super) fn backing_view(
585 &self,
586 resource_id: &ResourceId,
587 ) -> Result<LogicalBackingBufferView<'_, R::Buffer>, VNextError> {
588 let authority = self
589 .backing_slices
590 .iter()
591 .find(|authority| authority.resource_id() == resource_id)
592 .ok_or_else(|| invalid_resource("logical request does not own that backing slice"))?;
593 self.plan.dynamic_pools().view(authority)
594 }
595
596 pub fn try_admit_sequence(
598 self: &Arc<Self>,
599 request: SequenceResourceAdmissionRequest,
600 ) -> Result<SequenceResourceAdmissionDecision<R>, VNextError> {
601 let _lifecycle = self
602 .plan
603 .resources
604 .read_lifecycle("admit a child sequence")?;
605 let SequenceResourceAdmissionRequest {
606 work_shape,
607 fit_policy,
608 pressure_action,
609 } = request;
610 if work_shape.fit_tokens() > self.work_shape.fit_tokens() {
611 return Err(invalid_resource(
612 "sequence token ceiling exceeds its parent request ceiling",
613 ));
614 }
615 let immediate_shape = work_shape.immediate_shape();
616 let fit_shape = match fit_policy {
617 AdmissionFitPolicy::ImmediateOnly => immediate_shape,
618 AdmissionFitPolicy::FullInputMustFit => work_shape.fit_shape(),
619 };
620 let (demand, requested_slices) = self.plan.scoped_demand(
621 AllocationLifetime::Sequence,
622 None,
623 immediate_shape,
624 fit_shape,
625 None,
626 fit_policy,
627 pressure_action,
628 )?;
629 let immediate =
633 CapacityVector::checked_sum(self.logical_lease.claims(), demand.immediate_claim())?;
634 let fit =
635 CapacityVector::checked_sum(self.logical_lease.claims(), demand.fit_requirement())?;
636 if let Some(rejected) = self.plan.reject_impossible_plan_fit(&immediate, &fit)? {
637 return Ok(SequenceResourceAdmissionDecision::PermanentRejected(
638 rejected,
639 ));
640 }
641 match self
642 .plan
643 .logical_admission()
644 .preflight_sequence_ceiling_for_request(&self.logical_lease, &demand)?
645 {
646 AdmissionPreflightDecision::Eligible => {}
647 AdmissionPreflightDecision::Deferred(deferred) => {
648 return Ok(SequenceResourceAdmissionDecision::Deferred(deferred));
649 }
650 AdmissionPreflightDecision::PermanentRejected(rejected) => {
651 return Ok(SequenceResourceAdmissionDecision::PermanentRejected(
652 rejected,
653 ));
654 }
655 }
656 let prepared = match self.plan.prepare_backing_slices(requested_slices)? {
657 BackingPrepareDecision::Prepared(prepared) => prepared,
658 BackingPrepareDecision::Deferred(deferred) => {
659 return Ok(SequenceResourceAdmissionDecision::BackingDeferred(
660 SequenceAdmissionBackingDeferral {
661 evidence: deferred,
662 parent: Arc::clone(self),
663 },
664 ));
665 }
666 };
667 match self
668 .plan
669 .logical_admission()
670 .try_admit_sequence_for_request(&self.logical_lease, &demand)?
671 {
672 AdmissionDecision::Admitted(logical_lease) => {
673 if !self.plan.logical_admission().owns(&logical_lease)
674 || logical_lease.request() != self.request_authority()
675 {
676 return Err(invalid_resource(
677 "sequence admission returned authority from another request",
678 ));
679 }
680 let slices = prepared.commit();
681 Ok(SequenceResourceAdmissionDecision::Admitted(Arc::new(
682 AdmittedSequenceResources::new(
683 Arc::clone(self),
684 logical_lease,
685 slices,
686 work_shape,
687 )?,
688 )))
689 }
690 AdmissionDecision::Deferred(deferred) => {
691 Ok(SequenceResourceAdmissionDecision::Deferred(deferred))
692 }
693 AdmissionDecision::PermanentRejected(rejected) => Ok(
694 SequenceResourceAdmissionDecision::PermanentRejected(rejected),
695 ),
696 }
697 }
698}
699
700pub enum SequenceResourceAdmissionDecision<R>
701where
702 R: DeviceRuntime,
703{
704 Admitted(Arc<AdmittedSequenceResources<R>>),
705 Deferred(AdmissionDeferred),
706 BackingDeferred(SequenceAdmissionBackingDeferral<R>),
707 PermanentRejected(AdmissionRejected),
708}
709
710#[must_use = "sequence backing deferral owns its exact request parent"]
714pub struct SequenceAdmissionBackingDeferral<R>
715where
716 R: DeviceRuntime,
717{
718 evidence: DynamicBackingDeferred,
719 parent: Arc<AdmittedRequestResources<R>>,
720}
721
722impl<R> SequenceAdmissionBackingDeferral<R>
723where
724 R: DeviceRuntime,
725{
726 pub fn evidence(&self) -> &DynamicBackingDeferred {
727 &self.evidence
728 }
729
730 pub fn parent(&self) -> &Arc<AdmittedRequestResources<R>> {
731 &self.parent
732 }
733
734 pub fn into_parent(self) -> Arc<AdmittedRequestResources<R>> {
735 self.parent
736 }
737
738 pub fn maintain(&self) -> Result<DynamicDeferredMaintenanceOutcome, VNextError> {
739 self.parent
740 .maintain_sequence_backing_for_deferred(&self.evidence)
741 }
742
743 pub fn register_waiter(&self) -> Result<PlanCapacityWaitRegistration<R>, VNextError> {
744 self.parent.plan.register_backing_waiter(&self.evidence)
745 }
746}
747
748#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize)]
749pub struct SequenceSessionEpoch(pub(super) NonZeroU64);
750
751impl SequenceSessionEpoch {
752 pub const fn get(self) -> u64 {
753 self.0.get()
754 }
755}
756
757#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
758pub struct SequenceSessionFingerprint(pub(super) String);
759
760pub(super) const SEQUENCE_SESSION_FINGERPRINT_DOMAIN: &str = "sequence-session-v1";
761
762#[derive(Serialize)]
763pub(super) struct SequenceSessionFingerprintEnvelope<'a, T>
764where
765 T: Serialize + ?Sized,
766{
767 pub(super) domain: &'static str,
768 pub(super) payload: &'a T,
769}
770
771pub(super) fn sequence_session_fingerprint<T>(
772 payload: &T,
773) -> Result<SequenceSessionFingerprint, VNextError>
774where
775 T: Serialize + ?Sized,
776{
777 let envelope = SequenceSessionFingerprintEnvelope {
778 domain: SEQUENCE_SESSION_FINGERPRINT_DOMAIN,
779 payload,
780 };
781 let bytes = serde_json::to_vec(&envelope)
782 .map_err(|_| invalid_resource("trusted sequence session identity did not serialize"))?;
783 Ok(SequenceSessionFingerprint(format!(
784 "{:x}",
785 Sha256::digest(bytes)
786 )))
787}
788
789impl SequenceSessionFingerprint {
790 pub fn as_str(&self) -> &str {
791 &self.0
792 }
793}
794
795#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
796pub enum SequenceSessionTerminalDisposition {
797 Completed,
798 Aborted,
799}
800
801#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
802pub struct SequenceSessionTerminalReceipt {
803 pub(super) epoch: SequenceSessionEpoch,
804 pub(super) fingerprint: SequenceSessionFingerprint,
805 pub(super) disposition: SequenceSessionTerminalDisposition,
806 pub(super) retired_frames: u64,
807}
808
809impl SequenceSessionTerminalReceipt {
810 pub const fn epoch(&self) -> SequenceSessionEpoch {
811 self.epoch
812 }
813
814 pub fn fingerprint(&self) -> &SequenceSessionFingerprint {
815 &self.fingerprint
816 }
817
818 pub const fn disposition(&self) -> SequenceSessionTerminalDisposition {
819 self.disposition
820 }
821
822 pub const fn retired_frames(&self) -> u64 {
823 self.retired_frames
824 }
825}
826
827#[derive(Debug, Clone, Copy, PartialEq, Eq)]
828pub(super) enum SequenceSessionPhase {
829 Open,
830 CancelRequested,
831 Poisoned,
832}
833
834#[derive(Debug, Clone, Copy, PartialEq, Eq)]
835pub(super) struct ActiveSequenceFrame {
836 pub(super) frame_id: ExecutionFrameId,
837 pub(super) batch_step_id: BatchStepId,
838}
839
840#[derive(Debug, Clone, Copy, PartialEq, Eq)]
841pub(super) enum ParticipantFlightPhase {
842 Prepared,
843 InFlight,
844}
845
846#[derive(Debug, Clone)]
847pub(super) struct ActiveSequenceSessionState {
848 pub(super) epoch: SequenceSessionEpoch,
849 pub(super) fingerprint: SequenceSessionFingerprint,
850 pub(super) phase: SequenceSessionPhase,
851 pub(super) next_frame: Option<ExecutionFrameId>,
852 pub(super) active_frame: Option<ActiveSequenceFrame>,
853 pub(super) participant_flights: BTreeMap<ParticipantNodeKey, ParticipantFlightPhase>,
854 pub(super) submission_wave_flight: Option<ParticipantFlightPhase>,
855 pub(super) state_transfer: SequenceStateTransferSlot,
856 pub(super) completed_boundary: SequenceCompletedFrontier,
857 pub(super) retired_frames: u64,
858}
859
860impl ActiveSequenceSessionState {
861 pub(super) fn has_participant_flights(&self) -> bool {
862 !self.participant_flights.is_empty() || self.submission_wave_flight.is_some()
863 }
864
865 pub(super) fn participant_flight_count(&self) -> usize {
866 self.participant_flights.len() + usize::from(self.submission_wave_flight.is_some())
867 }
868}
869
870#[derive(Debug, Clone)]
871pub(super) enum SequenceSessionSlotState {
872 Dormant {
873 next_epoch: Option<SequenceSessionEpoch>,
874 },
875 Active(ActiveSequenceSessionState),
876 Terminal(SequenceSessionTerminalReceipt),
877 FailClosed,
878}
879
880pub(super) struct SequenceSessionSlot {
881 pub(super) state: Mutex<SequenceSessionSlotState>,
882}
883
884impl SequenceSessionSlot {
885 fn new() -> Self {
886 Self {
887 state: Mutex::new(SequenceSessionSlotState::Dormant {
888 next_epoch: Some(SequenceSessionEpoch(
889 NonZeroU64::new(1).expect("one is non-zero"),
890 )),
891 }),
892 }
893 }
894
895 fn is_poisoned(&self) -> bool {
896 match self.state.lock() {
897 Ok(state) => matches!(
898 &*state,
899 SequenceSessionSlotState::Active(ActiveSequenceSessionState {
900 phase: SequenceSessionPhase::Poisoned,
901 ..
902 }) | SequenceSessionSlotState::FailClosed
903 ),
904 Err(_) => true,
905 }
906 }
907}
908
909#[derive(Clone)]
914pub(crate) struct SequenceSessionLiveWitness {
915 pub(super) slot: Weak<SequenceSessionSlot>,
916 epoch: SequenceSessionEpoch,
917 fingerprint: SequenceSessionFingerprint,
918}
919
920impl SequenceSessionLiveWitness {
921 pub(crate) fn ensure_live(&self) -> Result<(), VNextError> {
922 let slot = self.slot.upgrade().ok_or_else(|| {
923 invalid_resource("sequence session live witness owner is no longer available")
924 })?;
925 let state = slot
926 .state
927 .lock()
928 .map_err(|_| invalid_resource("sequence session state mutex is poisoned"))?;
929 match &*state {
930 SequenceSessionSlotState::Active(active)
931 if active.epoch == self.epoch && active.fingerprint == self.fingerprint =>
932 {
933 Ok(())
934 }
935 _ => Err(invalid_resource(
936 "sequence session live witness is stale or no longer active",
937 )),
938 }
939 }
940
941 pub(crate) fn ensure_open(&self) -> Result<(), VNextError> {
942 let slot = self.slot.upgrade().ok_or_else(|| {
943 invalid_resource("sequence session live witness owner is no longer available")
944 })?;
945 let state = slot
946 .state
947 .lock()
948 .map_err(|_| invalid_resource("sequence session state mutex is poisoned"))?;
949 match &*state {
950 SequenceSessionSlotState::Active(active)
951 if active.epoch == self.epoch
952 && active.fingerprint == self.fingerprint
953 && active.phase == SequenceSessionPhase::Open =>
954 {
955 Ok(())
956 }
957 _ => Err(invalid_resource(
958 "sequence session live witness is stale or no longer open",
959 )),
960 }
961 }
962
963 pub(crate) fn ensure_identity(
964 &self,
965 epoch: SequenceSessionEpoch,
966 fingerprint: &SequenceSessionFingerprint,
967 ) -> Result<(), VNextError> {
968 if self.epoch != epoch || self.fingerprint != *fingerprint {
969 return Err(invalid_resource(
970 "sequence session live witness differs from the expected identity",
971 ));
972 }
973 self.ensure_open()
974 }
975
976 pub(crate) fn ensure_live_identity(
977 &self,
978 epoch: SequenceSessionEpoch,
979 fingerprint: &SequenceSessionFingerprint,
980 ) -> Result<(), VNextError> {
981 if self.epoch != epoch || self.fingerprint != *fingerprint {
982 return Err(invalid_resource(
983 "sequence session live witness differs from the expected identity",
984 ));
985 }
986 self.ensure_live()
987 }
988}
989
990#[must_use = "a sequence session must reach an explicit terminal disposition"]
993pub struct SequenceSession<R>
994where
995 R: DeviceRuntime,
996{
997 resources: Arc<AdmittedSequenceResources<R>>,
998 pub(super) slot: Arc<SequenceSessionSlot>,
999 pub(super) epoch: SequenceSessionEpoch,
1000 pub(super) fingerprint: SequenceSessionFingerprint,
1001}
1002
1003impl<R> SequenceSession<R>
1004where
1005 R: DeviceRuntime,
1006{
1007 pub fn resources(&self) -> &Arc<AdmittedSequenceResources<R>> {
1008 &self.resources
1009 }
1010
1011 pub fn write_release_capacity_sources(
1014 &self,
1015 sources: &mut Vec<CapacityAvailabilitySource>,
1016 ) -> Result<(), VNextError> {
1017 self.ensure_open_identity()?;
1018 let backing = self.resources.backing_snapshot()?;
1019 sources.clear();
1020 sources.push(CapacityAvailabilitySource::ActiveSequenceSlots);
1021 sources.extend(
1022 backing
1023 .backing_slices()
1024 .iter()
1025 .map(|slice| CapacityAvailabilitySource::Domain(slice.domain_id())),
1026 );
1027 sources.sort_unstable();
1028 sources.dedup();
1029 Ok(())
1030 }
1031
1032 pub const fn epoch(&self) -> SequenceSessionEpoch {
1033 self.epoch
1034 }
1035
1036 pub fn fingerprint(&self) -> &SequenceSessionFingerprint {
1037 &self.fingerprint
1038 }
1039
1040 pub(crate) fn live_witness(&self) -> Result<SequenceSessionLiveWitness, VNextError> {
1041 let witness = SequenceSessionLiveWitness {
1042 slot: Arc::downgrade(&self.slot),
1043 epoch: self.epoch,
1044 fingerprint: self.fingerprint.clone(),
1045 };
1046 witness.ensure_identity(self.epoch, &self.fingerprint)?;
1047 Ok(witness)
1048 }
1049
1050 pub(crate) fn ensure_open_identity(&self) -> Result<(), VNextError> {
1051 self.live_witness().map(|_| ())
1052 }
1053
1054 pub fn sequence_authority(&self) -> SequenceAuthorityId {
1055 self.resources.sequence_authority()
1056 }
1057
1058 pub fn request_authority(&self) -> RequestAuthorityId {
1059 self.resources.request_authority()
1060 }
1061
1062 pub fn register_admission_waiter(
1063 &self,
1064 deferred: &AdmissionDeferred,
1065 ) -> Result<PlanCapacityWaitRegistration<R>, VNextError> {
1066 self.resources
1067 .request
1068 .plan
1069 .register_admission_waiter(deferred)
1070 }
1071
1072 pub fn try_ensure_backing_covers(
1079 self: &Arc<Self>,
1080 request: SequenceResourceExtensionRequest,
1081 ) -> Result<SequenceResourceExtensionDecision<R>, VNextError> {
1082 let _lifecycle = self
1083 .resources
1084 .request
1085 .plan
1086 .resources
1087 .read_lifecycle("extend sequence backing")?;
1088 let target_fingerprint = request.target_work.fingerprint().to_owned();
1089 let target = request.target_work.fit_shape();
1090 if target.tokens() > self.resources.request.work_shape().fit_tokens() {
1091 return Err(invalid_resource(
1092 "sequence backing extension exceeds its parent request token ceiling",
1093 ));
1094 }
1095
1096 let expected = {
1097 let slot = self
1098 .slot
1099 .state
1100 .lock()
1101 .map_err(|_| invalid_resource("sequence session state mutex is poisoned"))?;
1102 let active = match &*slot {
1103 SequenceSessionSlotState::Active(active)
1104 if active.epoch == self.epoch
1105 && active.fingerprint == self.fingerprint
1106 && active.phase == SequenceSessionPhase::Open =>
1107 {
1108 active
1109 }
1110 SequenceSessionSlotState::Active(active)
1111 if active.epoch != self.epoch || active.fingerprint != self.fingerprint =>
1112 {
1113 return Err(invalid_resource("stale sequence session authority"));
1114 }
1115 SequenceSessionSlotState::Active(_) => {
1116 return Err(invalid_resource(
1117 "sequence backing can only grow for an open session",
1118 ));
1119 }
1120 _ => {
1121 return Err(invalid_resource(
1122 "inactive or terminal sequence session cannot grow backing",
1123 ));
1124 }
1125 };
1126 let backing = self.resources.lock_backing_state()?;
1127 let current = Arc::clone(&backing.current);
1128 if target.sequences() != 1 {
1129 return Err(invalid_resource(
1130 "sequence backing coverage target must contain exactly one sequence",
1131 ));
1132 }
1133 if current.committed_tokens() >= target.tokens()
1134 && current.committed_pages() >= target.pages()
1135 {
1136 return Ok(SequenceResourceExtensionDecision::Current(current));
1137 }
1138 if target.tokens() < current.committed_tokens()
1139 || target.pages() < current.committed_pages()
1140 {
1141 return Err(invalid_resource(
1142 "sequence backing coverage target is incomparable with committed work",
1143 ));
1144 }
1145 if active.active_frame.is_some()
1146 || active.has_participant_flights()
1147 || active.state_transfer.is_reserved()
1148 {
1149 return Ok(SequenceResourceExtensionDecision::RetryRequired(current));
1150 }
1151 current
1152 };
1153
1154 let plan = &self.resources.request.plan;
1155 let (demand, requested_slices) = plan.sequence_extension_demand(
1156 expected.committed_shape(),
1157 target,
1158 request.pressure_action,
1159 )?;
1160 let extension = if demand.immediate_claim().is_empty() {
1161 if !requested_slices.is_empty() {
1162 return Err(invalid_resource(
1163 "empty sequence extension demand produced physical backing requests",
1164 ));
1165 }
1166 PreparedSequenceExtension::empty()
1167 } else {
1168 let prepared = match plan.prepare_backing_slices(requested_slices)? {
1169 BackingPrepareDecision::Prepared(prepared) => prepared,
1170 BackingPrepareDecision::Deferred(deferred) => {
1171 return Ok(SequenceResourceExtensionDecision::BackingDeferred(
1172 SequenceExtensionBackingDeferral {
1173 backing: PlanBackingDeferral::new(
1174 Arc::clone(&plan.resources),
1175 deferred,
1176 )?,
1177 session: Arc::clone(self),
1178 expected_generation: expected.generation(),
1179 target_fingerprint,
1180 },
1181 ));
1182 }
1183 };
1184 let capacity = match plan
1185 .logical_admission()
1186 .try_claim_for_sequence(self.resources.logical_lease(), &demand)?
1187 {
1188 CapacityClaimDecision::Claimed(capacity) => capacity,
1189 CapacityClaimDecision::Deferred(deferred) => {
1190 return Ok(SequenceResourceExtensionDecision::Deferred(deferred));
1191 }
1192 CapacityClaimDecision::PermanentRejected(rejected) => {
1193 return Ok(SequenceResourceExtensionDecision::PermanentRejected(
1194 rejected,
1195 ));
1196 }
1197 };
1198 let extension = PreparedSequenceExtension::claimed(prepared, capacity);
1199 let capacity = extension
1200 .capacity()
1201 .expect("claimed sequence extension owns logical capacity");
1202 if !plan.logical_admission().owns_capacity_claim(capacity)
1203 || capacity.sequence() != self.sequence_authority()
1204 || capacity.request() != self.request_authority()
1205 || capacity.claims() != demand.immediate_claim()
1206 {
1207 return Err(invalid_resource(
1208 "sequence extension capacity belongs to another authority or demand",
1209 ));
1210 }
1211 extension
1212 };
1213
1214 let slot = self
1215 .slot
1216 .state
1217 .lock()
1218 .map_err(|_| invalid_resource("sequence session state mutex is poisoned"))?;
1219 let active = match &*slot {
1220 SequenceSessionSlotState::Active(active)
1221 if active.epoch == self.epoch
1222 && active.fingerprint == self.fingerprint
1223 && active.phase == SequenceSessionPhase::Open =>
1224 {
1225 active
1226 }
1227 SequenceSessionSlotState::Active(active)
1228 if active.epoch != self.epoch || active.fingerprint != self.fingerprint =>
1229 {
1230 return Err(invalid_resource("stale sequence session authority"));
1231 }
1232 SequenceSessionSlotState::Active(_) => {
1233 return Err(invalid_resource(
1234 "sequence backing can only grow for an open session",
1235 ));
1236 }
1237 _ => {
1238 return Err(invalid_resource(
1239 "inactive or terminal sequence session cannot grow backing",
1240 ));
1241 }
1242 };
1243 let mut backing = self.resources.lock_backing_state()?;
1244 if !Arc::ptr_eq(&backing.current, &expected) {
1245 let current = Arc::clone(&backing.current);
1246 if current.committed_tokens() >= target.tokens()
1247 && current.committed_pages() >= target.pages()
1248 {
1249 return Ok(SequenceResourceExtensionDecision::Current(current));
1250 }
1251 return Ok(SequenceResourceExtensionDecision::RetryRequired(current));
1252 }
1253 if active.active_frame.is_some()
1254 || active.has_participant_flights()
1255 || active.state_transfer.is_reserved()
1256 {
1257 return Ok(SequenceResourceExtensionDecision::RetryRequired(
1258 Arc::clone(&backing.current),
1259 ));
1260 }
1261
1262 let extension = extension.commit();
1263 let next = Arc::new(SequenceBackingSnapshot::advanced(
1264 &expected, extension, target,
1265 )?);
1266 backing.current = Arc::clone(&next);
1267 Ok(SequenceResourceExtensionDecision::Extended(next))
1268 }
1269
1270 pub fn request_cancel(&self) -> Result<SequenceSessionCancelSnapshot, VNextError> {
1271 let mut state = self
1272 .slot
1273 .state
1274 .lock()
1275 .map_err(|_| invalid_resource("sequence session state mutex is poisoned"))?;
1276 let active = match &mut *state {
1277 SequenceSessionSlotState::Active(active)
1278 if active.epoch == self.epoch && active.fingerprint == self.fingerprint =>
1279 {
1280 active
1281 }
1282 SequenceSessionSlotState::Active(_) => {
1283 return Err(invalid_resource("stale sequence session authority"));
1284 }
1285 SequenceSessionSlotState::Terminal(_) => {
1286 return Err(invalid_resource("sequence session is already terminal"));
1287 }
1288 SequenceSessionSlotState::Dormant { .. } => {
1289 return Err(invalid_resource("sequence session is not active"));
1290 }
1291 SequenceSessionSlotState::FailClosed => {
1292 return Err(invalid_resource("sequence session is fail-closed"));
1293 }
1294 };
1295 match active.phase {
1296 SequenceSessionPhase::Open => active.phase = SequenceSessionPhase::CancelRequested,
1297 SequenceSessionPhase::CancelRequested => {}
1298 SequenceSessionPhase::Poisoned => {
1299 return Err(invalid_resource(
1300 "poisoned sequence session cannot be cancelled",
1301 ));
1302 }
1303 }
1304 Ok(SequenceSessionCancelSnapshot {
1305 active_frame: active.active_frame.map(|frame| frame.frame_id),
1306 participant_flights: u64::try_from(active.participant_flight_count())
1307 .map_err(|_| invalid_resource("participant flight count exceeds u64"))?,
1308 state_transfer_pending: active.state_transfer.is_reserved(),
1309 })
1310 }
1311
1312 pub fn try_complete(&self) -> Result<SequenceSessionTerminalReceipt, VNextError> {
1313 self.terminalize(SequenceSessionTerminalDisposition::Completed)
1314 }
1315
1316 pub fn try_abort(&self) -> Result<SequenceSessionTerminalReceipt, VNextError> {
1317 self.terminalize(SequenceSessionTerminalDisposition::Aborted)
1318 }
1319
1320 pub fn try_abort_if_quiescent(&self) -> Result<SequenceSessionTerminalReceipt, VNextError> {
1329 let mut state = self
1330 .slot
1331 .state
1332 .lock()
1333 .map_err(|_| invalid_resource("sequence session state mutex is poisoned"))?;
1334 let active = match &*state {
1335 SequenceSessionSlotState::Active(active)
1336 if active.epoch == self.epoch && active.fingerprint == self.fingerprint =>
1337 {
1338 active
1339 }
1340 SequenceSessionSlotState::Active(_) => {
1341 return Err(invalid_resource("stale sequence session authority"));
1342 }
1343 SequenceSessionSlotState::Terminal(_) => {
1344 return Err(invalid_resource("sequence session is already terminal"));
1345 }
1346 SequenceSessionSlotState::Dormant { .. } => {
1347 return Err(invalid_resource("sequence session is not active"));
1348 }
1349 SequenceSessionSlotState::FailClosed => {
1350 return Err(invalid_resource("sequence session is fail-closed"));
1351 }
1352 };
1353 if active.phase == SequenceSessionPhase::Poisoned
1354 || active.active_frame.is_some()
1355 || active.has_participant_flights()
1356 || active.state_transfer.is_reserved()
1357 {
1358 return Err(invalid_resource(
1359 "quiescent sequence abort requires an open or cancel-requested phase with no active frame, participant flight, or state transfer",
1360 ));
1361 }
1362 let receipt = SequenceSessionTerminalReceipt {
1363 epoch: active.epoch,
1364 fingerprint: active.fingerprint.clone(),
1365 disposition: SequenceSessionTerminalDisposition::Aborted,
1366 retired_frames: active.retired_frames,
1367 };
1368 *state = SequenceSessionSlotState::Terminal(receipt.clone());
1369 Ok(receipt)
1370 }
1371
1372 fn terminalize(
1373 &self,
1374 disposition: SequenceSessionTerminalDisposition,
1375 ) -> Result<SequenceSessionTerminalReceipt, VNextError> {
1376 let mut state = self
1377 .slot
1378 .state
1379 .lock()
1380 .map_err(|_| invalid_resource("sequence session state mutex is poisoned"))?;
1381 let active = match &*state {
1382 SequenceSessionSlotState::Active(active)
1383 if active.epoch == self.epoch && active.fingerprint == self.fingerprint =>
1384 {
1385 active
1386 }
1387 SequenceSessionSlotState::Active(_) => {
1388 return Err(invalid_resource("stale sequence session authority"));
1389 }
1390 SequenceSessionSlotState::Terminal(_) => {
1391 return Err(invalid_resource("sequence session is already terminal"));
1392 }
1393 SequenceSessionSlotState::Dormant { .. } => {
1394 return Err(invalid_resource("sequence session is not active"));
1395 }
1396 SequenceSessionSlotState::FailClosed => {
1397 return Err(invalid_resource("sequence session is fail-closed"));
1398 }
1399 };
1400 let phase_matches = match disposition {
1401 SequenceSessionTerminalDisposition::Completed => {
1402 active.phase == SequenceSessionPhase::Open && active.retired_frames > 0
1403 }
1404 SequenceSessionTerminalDisposition::Aborted => matches!(
1405 active.phase,
1406 SequenceSessionPhase::CancelRequested | SequenceSessionPhase::Poisoned
1407 ),
1408 };
1409 if !phase_matches
1410 || active.active_frame.is_some()
1411 || active.has_participant_flights()
1412 || active.state_transfer.is_reserved()
1413 {
1414 return Err(invalid_resource(
1415 "sequence terminalization requires the matching phase with no active frame, participant flight, or state transfer",
1416 ));
1417 }
1418 let receipt = SequenceSessionTerminalReceipt {
1419 epoch: active.epoch,
1420 fingerprint: active.fingerprint.clone(),
1421 disposition,
1422 retired_frames: active.retired_frames,
1423 };
1424 *state = SequenceSessionSlotState::Terminal(receipt.clone());
1425 Ok(receipt)
1426 }
1427}
1428
1429#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1430pub struct SequenceSessionCancelSnapshot {
1431 active_frame: Option<ExecutionFrameId>,
1432 participant_flights: u64,
1433 state_transfer_pending: bool,
1434}
1435
1436impl SequenceSessionCancelSnapshot {
1437 pub const fn active_frame(self) -> Option<ExecutionFrameId> {
1438 self.active_frame
1439 }
1440
1441 pub const fn participant_flights(self) -> u64 {
1442 self.participant_flights
1443 }
1444
1445 pub const fn state_transfer_pending(self) -> bool {
1446 self.state_transfer_pending
1447 }
1448}
1449
1450impl<R> Drop for SequenceSession<R>
1451where
1452 R: DeviceRuntime,
1453{
1454 fn drop(&mut self) {
1455 let mut state = match self.slot.state.lock() {
1456 Ok(state) => state,
1457 Err(poisoned) => poisoned.into_inner(),
1458 };
1459 if matches!(
1460 &*state,
1461 SequenceSessionSlotState::Active(active)
1462 if active.epoch == self.epoch && active.fingerprint == self.fingerprint
1463 ) {
1464 *state = SequenceSessionSlotState::FailClosed;
1465 }
1466 }
1467}
1468
1469#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1470pub(super) enum SequenceExecutionAuthoritySource {
1471 Unselected,
1472 LegacyStream,
1473 SequenceSession,
1474 FailClosed,
1475}
1476
1477#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize)]
1478pub struct SequenceBackingGeneration(NonZeroU64);
1479
1480impl SequenceBackingGeneration {
1481 const INITIAL: Self = Self(NonZeroU64::MIN);
1482
1483 pub const fn get(self) -> u64 {
1484 self.0.get()
1485 }
1486}
1487
1488struct SequenceLogicalOwner<R>
1492where
1493 R: DeviceRuntime,
1494{
1495 logical_lease: LogicalAdmissionLease,
1496 _request: Arc<AdmittedRequestResources<R>>,
1497}
1498
1499struct PreparedSequenceExtension<R>
1502where
1503 R: DeviceRuntime,
1504{
1505 prepared: Option<PreparedBackingClaim<R>>,
1506 capacity: Option<LogicalCapacityLease>,
1507}
1508
1509impl<R> PreparedSequenceExtension<R>
1510where
1511 R: DeviceRuntime,
1512{
1513 fn empty() -> Self {
1514 Self {
1515 prepared: None,
1516 capacity: None,
1517 }
1518 }
1519
1520 fn claimed(prepared: PreparedBackingClaim<R>, capacity: LogicalCapacityLease) -> Self {
1521 Self {
1522 prepared: Some(prepared),
1523 capacity: Some(capacity),
1524 }
1525 }
1526
1527 fn capacity(&self) -> Option<&LogicalCapacityLease> {
1528 self.capacity.as_ref()
1529 }
1530
1531 fn commit(mut self) -> CommittedSequenceExtension {
1532 let backing_slices = self
1533 .prepared
1534 .take()
1535 .map(PreparedBackingClaim::commit)
1536 .unwrap_or_default();
1537 CommittedSequenceExtension {
1538 backing_slices,
1539 capacity: self.capacity.take(),
1540 }
1541 }
1542}
1543
1544struct CommittedSequenceExtension {
1547 backing_slices: Vec<LogicalBackingSliceAuthority>,
1548 capacity: Option<LogicalCapacityLease>,
1549}
1550
1551impl<R> SequenceLogicalOwner<R>
1552where
1553 R: DeviceRuntime,
1554{
1555 fn lease(&self) -> &LogicalAdmissionLease {
1556 &self.logical_lease
1557 }
1558}
1559
1560#[must_use = "a sequence backing snapshot owns its exact physical extents"]
1564pub struct SequenceBackingSnapshot<R>
1565where
1566 R: DeviceRuntime,
1567{
1568 generation: SequenceBackingGeneration,
1569 backing_slices: Vec<LogicalBackingSliceAuthority>,
1570 extension_capacity: Vec<Arc<LogicalCapacityLease>>,
1574 committed_shape: DynamicResourceShape,
1575 _logical_owner: Arc<SequenceLogicalOwner<R>>,
1578}
1579
1580impl<R> SequenceBackingSnapshot<R>
1581where
1582 R: DeviceRuntime,
1583{
1584 fn initial(
1585 backing_slices: Vec<LogicalBackingSliceAuthority>,
1586 work_shape: ResourceWorkShape,
1587 logical_owner: Arc<SequenceLogicalOwner<R>>,
1588 ) -> Result<Self, VNextError> {
1589 if work_shape.immediate_sequences() != 1
1590 || work_shape.fit_sequences() != 1
1591 || backing_slices
1592 .windows(2)
1593 .any(|pair| pair[0].resource_id() >= pair[1].resource_id())
1594 {
1595 return Err(invalid_resource(
1596 "initial sequence backing must be canonical and single-sequence",
1597 ));
1598 }
1599 Ok(Self {
1600 generation: SequenceBackingGeneration::INITIAL,
1601 backing_slices,
1602 extension_capacity: Vec::new(),
1603 committed_shape: work_shape.immediate_shape(),
1604 _logical_owner: logical_owner,
1605 })
1606 }
1607
1608 fn advanced(
1609 prior: &Self,
1610 mut extension: CommittedSequenceExtension,
1611 committed_shape: DynamicResourceShape,
1612 ) -> Result<Self, VNextError> {
1613 let generation = prior
1614 .generation
1615 .get()
1616 .checked_add(1)
1617 .and_then(NonZeroU64::new)
1618 .map(SequenceBackingGeneration)
1619 .ok_or_else(|| invalid_resource("sequence backing generation space is exhausted"))?;
1620 if committed_shape.sequences() != 1
1621 || committed_shape.tokens() < prior.committed_shape.tokens()
1622 || committed_shape.pages() < prior.committed_shape.pages()
1623 || (committed_shape.tokens() == prior.committed_shape.tokens()
1624 && committed_shape.pages() == prior.committed_shape.pages())
1625 || extension.backing_slices.is_empty() != extension.capacity.is_none()
1626 {
1627 return Err(invalid_resource(
1628 "sequence backing generation must monotonically advance with matching capacity",
1629 ));
1630 }
1631 extension
1632 .backing_slices
1633 .sort_by(|left, right| left.resource_id().cmp(right.resource_id()));
1634 if extension
1635 .backing_slices
1636 .windows(2)
1637 .any(|pair| pair[0].resource_id() >= pair[1].resource_id())
1638 {
1639 return Err(invalid_resource(
1640 "sequence backing extension slices must be canonical and unique",
1641 ));
1642 }
1643 let mut backing_slices = prior
1644 .backing_slices
1645 .iter()
1646 .map(LogicalBackingSliceAuthority::retained)
1647 .collect::<Vec<_>>();
1648 backing_slices.append(&mut extension.backing_slices);
1649 backing_slices.sort_by(|left, right| {
1650 left.resource_id().cmp(right.resource_id()).then_with(|| {
1651 left.evidence()
1652 .segment_generation()
1653 .cmp(&right.evidence().segment_generation())
1654 })
1655 });
1656 let mut extension_capacity = prior.extension_capacity.clone();
1657 if let Some(capacity) = extension.capacity.take() {
1658 extension_capacity.push(Arc::new(capacity));
1659 }
1660 Ok(Self {
1661 generation,
1662 backing_slices,
1663 extension_capacity,
1664 committed_shape,
1665 _logical_owner: Arc::clone(&prior._logical_owner),
1666 })
1667 }
1668
1669 pub const fn generation(&self) -> SequenceBackingGeneration {
1670 self.generation
1671 }
1672
1673 pub fn backing_slices(&self) -> &[LogicalBackingSliceAuthority] {
1674 &self.backing_slices
1675 }
1676
1677 pub const fn committed_tokens(&self) -> u64 {
1678 self.committed_shape.tokens()
1679 }
1680
1681 pub const fn committed_pages(&self) -> u64 {
1682 self.committed_shape.pages()
1683 }
1684
1685 pub(crate) const fn committed_shape(&self) -> DynamicResourceShape {
1686 self.committed_shape
1687 }
1688
1689 pub(super) fn backing_slices_for(
1690 &self,
1691 resource_id: &ResourceId,
1692 ) -> &[LogicalBackingSliceAuthority] {
1693 let start = self
1694 .backing_slices
1695 .partition_point(|slice| slice.resource_id() < resource_id);
1696 let end = self
1697 .backing_slices
1698 .partition_point(|slice| slice.resource_id() <= resource_id);
1699 &self.backing_slices[start..end]
1700 }
1701}
1702
1703pub(super) struct SequenceBackingState<R>
1704where
1705 R: DeviceRuntime,
1706{
1707 pub(super) current: Arc<SequenceBackingSnapshot<R>>,
1708}
1709
1710#[derive(Debug, Clone, PartialEq, Eq)]
1711pub struct SequenceResourceExtensionRequest {
1712 target_work: ResourceWorkShape,
1713 pressure_action: AdmissionPressureAction,
1714}
1715
1716impl SequenceResourceExtensionRequest {
1717 pub fn new(
1718 target_work: ResourceWorkShape,
1719 pressure_action: AdmissionPressureAction,
1720 ) -> Result<Self, VNextError> {
1721 if target_work.immediate_sequences() != 1 || target_work.fit_sequences() != 1 {
1722 return Err(invalid_resource(
1723 "sequence extension target requires single-sequence work evidence",
1724 ));
1725 }
1726 Ok(Self {
1727 target_work,
1728 pressure_action,
1729 })
1730 }
1731}
1732
1733pub enum SequenceResourceExtensionDecision<R>
1734where
1735 R: DeviceRuntime,
1736{
1737 Current(Arc<SequenceBackingSnapshot<R>>),
1738 Extended(Arc<SequenceBackingSnapshot<R>>),
1739 RetryRequired(Arc<SequenceBackingSnapshot<R>>),
1740 Deferred(AdmissionDeferred),
1741 BackingDeferred(SequenceExtensionBackingDeferral<R>),
1742 PermanentRejected(AdmissionRejected),
1743}
1744
1745#[must_use = "sequence extension backing must be maintained or explicitly dropped"]
1748pub struct SequenceExtensionBackingDeferral<R>
1749where
1750 R: DeviceRuntime,
1751{
1752 backing: PlanBackingDeferral<R>,
1753 session: Arc<SequenceSession<R>>,
1754 expected_generation: SequenceBackingGeneration,
1755 target_fingerprint: String,
1756}
1757
1758impl<R> SequenceExtensionBackingDeferral<R>
1759where
1760 R: DeviceRuntime,
1761{
1762 pub fn evidence(&self) -> &DynamicBackingDeferred {
1763 self.backing.evidence()
1764 }
1765
1766 pub fn expected_generation(&self) -> SequenceBackingGeneration {
1767 self.expected_generation
1768 }
1769
1770 pub fn target_fingerprint(&self) -> &str {
1771 &self.target_fingerprint
1772 }
1773
1774 pub fn maintain(&self) -> Result<DynamicDeferredMaintenanceOutcome, VNextError> {
1775 self.session.ensure_open_identity()?;
1776 let current_generation = self
1777 .session
1778 .resources
1779 .lock_backing_state()?
1780 .current
1781 .generation();
1782 if current_generation != self.expected_generation {
1783 return self.backing.retry_admission();
1784 }
1785 self.backing.maintain()
1786 }
1787
1788 pub fn register_waiter(&self) -> Result<PlanCapacityWaitRegistration<R>, VNextError> {
1789 self.backing.register_waiter()
1790 }
1791}
1792
1793#[must_use = "logical sequence resources release capacity when dropped"]
1797pub struct AdmittedSequenceResources<R>
1798where
1799 R: DeviceRuntime,
1800{
1801 pub(super) admitted_work: ResourceWorkShape,
1804 pub(super) sequence_recovery: ManuallyDrop<Arc<SequenceRecoveryRegistry<R>>>,
1807 backing_state: ManuallyDrop<Mutex<SequenceBackingState<R>>>,
1808 logical_owner: ManuallyDrop<Arc<SequenceLogicalOwner<R>>>,
1809 pub(super) request: ManuallyDrop<Arc<AdmittedRequestResources<R>>>,
1810 pub(super) authority_source: Mutex<SequenceExecutionAuthoritySource>,
1811 session_slot: Arc<SequenceSessionSlot>,
1812 pub(super) state: Arc<AtomicU64>,
1813 pub(super) sequence_dispatch_gate: Arc<AtomicU64>,
1814 pub(super) next_activation_epoch: AtomicU64,
1815}
1816
1817struct DeferredSequenceResourceCleanup<R>
1818where
1819 R: DeviceRuntime,
1820{
1821 sequence_recovery: ManuallyDrop<Arc<SequenceRecoveryRegistry<R>>>,
1822 backing_state: ManuallyDrop<Mutex<SequenceBackingState<R>>>,
1823 logical_owner: ManuallyDrop<Arc<SequenceLogicalOwner<R>>>,
1824 request: ManuallyDrop<Arc<AdmittedRequestResources<R>>>,
1825 completed: bool,
1826}
1827
1828impl<R> DeferredDeviceCleanupTask for DeferredSequenceResourceCleanup<R>
1829where
1830 R: DeviceRuntime,
1831{
1832 fn try_cleanup(&mut self) -> DeferredDeviceCleanupDisposition {
1833 if self.completed {
1834 return DeferredDeviceCleanupDisposition::Completed;
1835 }
1836 let recovery = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
1837 self.sequence_recovery
1838 .recover_all_for_owner_drop(self.request.plan.runtime())
1839 }));
1840 if !matches!(recovery, Ok(Ok(()))) {
1841 return DeferredDeviceCleanupDisposition::Retryable;
1842 }
1843
1844 unsafe {
1848 ManuallyDrop::drop(&mut self.sequence_recovery);
1849 ManuallyDrop::drop(&mut self.backing_state);
1850 ManuallyDrop::drop(&mut self.logical_owner);
1851 ManuallyDrop::drop(&mut self.request);
1852 }
1853 self.completed = true;
1854 DeferredDeviceCleanupDisposition::Completed
1855 }
1856}
1857
1858impl<R> AdmittedSequenceResources<R>
1859where
1860 R: DeviceRuntime,
1861{
1862 fn new(
1863 request: Arc<AdmittedRequestResources<R>>,
1864 logical_lease: LogicalAdmissionLease,
1865 backing_slices: Vec<LogicalBackingSliceAuthority>,
1866 work_shape: ResourceWorkShape,
1867 ) -> Result<Self, VNextError> {
1868 if !request.plan.logical_admission().owns(&logical_lease)
1869 || logical_lease.request() != request.request_authority()
1870 {
1871 return Err(invalid_resource(
1872 "logical sequence authority belongs to another request",
1873 ));
1874 }
1875 let plan_resources = Arc::clone(&request.plan.resources);
1876 let logical_owner = Arc::new(SequenceLogicalOwner {
1877 logical_lease,
1878 _request: Arc::clone(&request),
1879 });
1880 let backing_snapshot = Arc::new(SequenceBackingSnapshot::initial(
1881 backing_slices,
1882 work_shape.clone(),
1883 Arc::clone(&logical_owner),
1884 )?);
1885 Ok(Self {
1886 admitted_work: work_shape,
1887 sequence_recovery: ManuallyDrop::new(Arc::new(SequenceRecoveryRegistry::new(
1888 plan_resources,
1889 ))),
1890 backing_state: ManuallyDrop::new(Mutex::new(SequenceBackingState {
1891 current: backing_snapshot,
1892 })),
1893 logical_owner: ManuallyDrop::new(logical_owner),
1894 request: ManuallyDrop::new(request),
1895 authority_source: Mutex::new(SequenceExecutionAuthoritySource::Unselected),
1896 session_slot: Arc::new(SequenceSessionSlot::new()),
1897 state: Arc::new(AtomicU64::new(0)),
1898 sequence_dispatch_gate: Arc::new(AtomicU64::new(0)),
1899 next_activation_epoch: AtomicU64::new(1),
1900 })
1901 }
1902
1903 pub fn sequence_authority(&self) -> SequenceAuthorityId {
1904 self.logical_lease().sequence()
1905 }
1906
1907 pub fn request_authority(&self) -> RequestAuthorityId {
1908 self.logical_lease().request()
1909 }
1910
1911 pub fn coordinator_id(&self) -> LogicalAdmissionCoordinatorId {
1912 self.logical_lease().coordinator_id()
1913 }
1914
1915 pub fn run_id(&self) -> &RunId {
1916 self.request.run_id()
1917 }
1918
1919 pub fn request_id(&self) -> &RequestIdentity {
1920 self.request.request_id()
1921 }
1922
1923 pub(super) fn logical_lease(&self) -> &LogicalAdmissionLease {
1924 self.logical_owner.lease()
1925 }
1926
1927 pub(super) fn lock_backing_state(
1928 &self,
1929 ) -> Result<std::sync::MutexGuard<'_, SequenceBackingState<R>>, VNextError> {
1930 self.backing_state
1931 .lock()
1932 .map_err(|_| invalid_resource("sequence backing state mutex is poisoned"))
1933 }
1934
1935 pub(crate) fn backing_snapshot(&self) -> Result<Arc<SequenceBackingSnapshot<R>>, VNextError> {
1936 Ok(Arc::clone(&self.lock_backing_state()?.current))
1937 }
1938
1939 pub fn request_resources(&self) -> &Arc<AdmittedRequestResources<R>> {
1940 &self.request
1941 }
1942
1943 pub fn backing_generation(&self) -> Result<SequenceBackingGeneration, VNextError> {
1944 Ok(self.lock_backing_state()?.current.generation())
1945 }
1946
1947 pub fn static_provisioning(&self) -> Option<&StaticProvisioningLease<R>> {
1948 self.request.static_provisioning()
1949 }
1950
1951 pub fn plan_evidence(&self) -> TrustedPlanRuntimeEvidence {
1952 self.request.plan_evidence()
1953 }
1954
1955 pub(crate) fn device_buffer_retention(&self) -> DeviceBufferRetention {
1956 DeviceBufferRetention::plan(Arc::clone(&self.request.plan.resources))
1957 }
1958
1959 pub(super) fn lock_authority_source(
1960 &self,
1961 ) -> Result<std::sync::MutexGuard<'_, SequenceExecutionAuthoritySource>, VNextError> {
1962 match self.authority_source.lock() {
1963 Ok(source) => Ok(source),
1964 Err(poisoned) => {
1965 let mut source = poisoned.into_inner();
1966 *source = SequenceExecutionAuthoritySource::FailClosed;
1967 Err(invalid_resource(
1968 "logical sequence execution authority selector is fail-closed",
1969 ))
1970 }
1971 }
1972 }
1973
1974 fn authority_source_is_fail_closed(&self) -> bool {
1975 match self.authority_source.lock() {
1976 Ok(source) => *source == SequenceExecutionAuthoritySource::FailClosed,
1977 Err(poisoned) => {
1978 *poisoned.into_inner() = SequenceExecutionAuthoritySource::FailClosed;
1979 true
1980 }
1981 }
1982 }
1983
1984 pub fn open_session(self: &Arc<Self>) -> Result<Arc<SequenceSession<R>>, VNextError> {
1985 let _lifecycle = self
1986 .request
1987 .plan
1988 .resources
1989 .read_lifecycle("open a sequence session")?;
1990 #[derive(Serialize)]
1991 struct FingerprintInput<'a> {
1992 plan: &'a TrustedPlanRuntimeEvidence,
1993 coordinator_id: LogicalAdmissionCoordinatorId,
1994 request_authority: RequestAuthorityId,
1995 sequence_authority: SequenceAuthorityId,
1996 run_id: &'a RunId,
1997 request_id: &'a RequestIdentity,
1998 epoch: SequenceSessionEpoch,
1999 request_backing: Vec<&'a LogicalBackingSliceEvidence>,
2000 }
2001
2002 let mut authority_source = self.lock_authority_source()?;
2003 let selecting_session = match *authority_source {
2004 SequenceExecutionAuthoritySource::Unselected => true,
2005 SequenceExecutionAuthoritySource::SequenceSession => false,
2006 SequenceExecutionAuthoritySource::LegacyStream => {
2007 return Err(invalid_resource(
2008 "logical sequence execution authority is permanently selected for legacy streams",
2009 ));
2010 }
2011 SequenceExecutionAuthoritySource::FailClosed => {
2012 return Err(invalid_resource(
2013 "logical sequence execution authority selector is fail-closed",
2014 ));
2015 }
2016 };
2017 let mut state = match self.session_slot.state.lock() {
2018 Ok(state) => state,
2019 Err(_) => {
2020 *authority_source = SequenceExecutionAuthoritySource::FailClosed;
2021 return Err(invalid_resource("sequence session state mutex is poisoned"));
2022 }
2023 };
2024 let epoch = match &*state {
2025 SequenceSessionSlotState::Dormant {
2026 next_epoch: Some(epoch),
2027 } if selecting_session => *epoch,
2028 SequenceSessionSlotState::Dormant {
2029 next_epoch: Some(_),
2030 } => {
2031 *authority_source = SequenceExecutionAuthoritySource::FailClosed;
2032 return Err(invalid_resource(
2033 "sequence session authority source lost its active or terminal slot",
2034 ));
2035 }
2036 SequenceSessionSlotState::Dormant { next_epoch: None } => {
2037 *authority_source = SequenceExecutionAuthoritySource::FailClosed;
2038 return Err(invalid_resource(
2039 "sequence session epoch space is exhausted",
2040 ));
2041 }
2042 SequenceSessionSlotState::Active(_) => {
2043 return Err(invalid_resource(
2044 "logical sequence already has an active session",
2045 ));
2046 }
2047 SequenceSessionSlotState::Terminal(_) => {
2048 return Err(invalid_resource("logical sequence is already terminal"));
2049 }
2050 SequenceSessionSlotState::FailClosed => {
2051 return Err(invalid_resource(
2052 "logical sequence session slot is fail-closed",
2053 ));
2054 }
2055 };
2056 let plan = self.plan_evidence();
2057 let request_backing = self
2058 .request
2059 .backing_slices
2060 .iter()
2061 .map(LogicalBackingSliceAuthority::evidence)
2062 .collect();
2063 let input = FingerprintInput {
2064 plan: &plan,
2065 coordinator_id: self.coordinator_id(),
2066 request_authority: self.request_authority(),
2067 sequence_authority: self.sequence_authority(),
2068 run_id: self.run_id(),
2069 request_id: self.request_id(),
2070 epoch,
2071 request_backing,
2072 };
2073 let fingerprint = match sequence_session_fingerprint(&input) {
2074 Ok(fingerprint) => fingerprint,
2075 Err(error) => {
2076 *authority_source = SequenceExecutionAuthoritySource::FailClosed;
2077 return Err(error);
2078 }
2079 };
2080 let next_epoch = epoch
2081 .get()
2082 .checked_add(1)
2083 .and_then(NonZeroU64::new)
2084 .map(SequenceSessionEpoch);
2085 *state = SequenceSessionSlotState::Active(ActiveSequenceSessionState {
2086 epoch,
2087 fingerprint: fingerprint.clone(),
2088 phase: SequenceSessionPhase::Open,
2089 next_frame: Some(
2090 ExecutionFrameId::try_from(1_u64)
2091 .expect("the first execution frame id is non-zero"),
2092 ),
2093 active_frame: None,
2094 participant_flights: BTreeMap::new(),
2095 submission_wave_flight: None,
2096 state_transfer: SequenceStateTransferSlot::default(),
2097 completed_boundary: SequenceCompletedFrontier::default(),
2098 retired_frames: 0,
2099 });
2100 *authority_source = SequenceExecutionAuthoritySource::SequenceSession;
2101 let _ = next_epoch;
2104 drop(state);
2105 Ok(Arc::new(SequenceSession {
2106 resources: Arc::clone(self),
2107 slot: Arc::clone(&self.session_slot),
2108 epoch,
2109 fingerprint,
2110 }))
2111 }
2112
2113 pub fn is_poisoned(&self) -> bool {
2114 self.authority_source_is_fail_closed()
2115 || self.session_slot.is_poisoned()
2116 || sequence_dispatch_is_poisoned(&self.sequence_dispatch_gate)
2117 || sequence_slot_is_poisoned(self.state.load(Ordering::Acquire))
2118 }
2119}
2120
2121impl<R> Drop for AdmittedSequenceResources<R>
2122where
2123 R: DeviceRuntime,
2124{
2125 fn drop(&mut self) {
2126 if !self.sequence_recovery.is_empty() {
2127 self.sequence_dispatch_gate
2128 .fetch_or(SEQUENCE_DISPATCH_POISONED_BIT, Ordering::AcqRel);
2129 let cleanup_domain = self.request.plan.resources.deferred_cleanup_domain;
2130 let cleanup = unsafe {
2133 DeferredSequenceResourceCleanup {
2134 sequence_recovery: ManuallyDrop::new(ManuallyDrop::take(
2135 &mut self.sequence_recovery,
2136 )),
2137 backing_state: ManuallyDrop::new(ManuallyDrop::take(&mut self.backing_state)),
2138 logical_owner: ManuallyDrop::new(ManuallyDrop::take(&mut self.logical_owner)),
2139 request: ManuallyDrop::new(ManuallyDrop::take(&mut self.request)),
2140 completed: false,
2141 }
2142 };
2143 defer_device_cleanup(cleanup_domain, cleanup);
2144 return;
2145 }
2146
2147 unsafe {
2151 ManuallyDrop::drop(&mut self.sequence_recovery);
2152 ManuallyDrop::drop(&mut self.backing_state);
2153 ManuallyDrop::drop(&mut self.logical_owner);
2154 ManuallyDrop::drop(&mut self.request);
2155 }
2156 }
2157}