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