1use crate::consensus::{
4 agree_deadline_candidates_bounded, agree_disposition_status_bounded,
5 agree_submission_status_bounded, resolve_output_completions_bounded, validate_schedule_bounded,
6 BoundedConsensusTransport, CompletionObservation, CompletionResolution, ScheduledWork,
7};
8use crate::BoundedCompletionWait;
9use serde::{Deserialize, Serialize};
10use sha2::{Digest, Sha256};
11use std::{
12 collections::{BTreeMap, BTreeSet, VecDeque},
13 time::{Duration, Instant},
14};
15
16#[derive(Debug, Clone, Copy, Eq, Ord, PartialEq, PartialOrd, Hash, Serialize, Deserialize)]
18pub struct RequestId(u64);
19
20impl RequestId {
21 pub const fn new(value: u64) -> Self {
23 Self(value)
24 }
25
26 pub const fn value(self) -> u64 {
28 self.0
29 }
30}
31
32#[derive(Debug, Clone, Copy, Eq, Ord, PartialEq, PartialOrd, Hash, Serialize, Deserialize)]
34pub struct WorkId {
35 request: RequestId,
36 sequence: u64,
37}
38
39impl WorkId {
40 pub const fn new(request: RequestId, sequence: u64) -> Self {
42 Self { request, sequence }
43 }
44
45 pub const fn request(self) -> RequestId {
47 self.request
48 }
49
50 pub const fn sequence(self) -> u64 {
52 self.sequence
53 }
54}
55
56#[derive(Debug, Clone, Copy, Eq, PartialEq, Serialize, Deserialize)]
58#[serde(rename_all = "snake_case")]
59pub enum WorkLifecycle {
60 Queued,
62 Prepared,
64 Submitted,
66 Completing,
68 Committed,
70 Abandoned,
72 Failed,
74}
75
76#[derive(Debug, Clone, Copy, Eq, PartialEq, Serialize, Deserialize)]
78#[serde(rename_all = "snake_case")]
79pub enum RequestStatus {
80 Active,
82 Finished,
84 Cancelled,
86 DeadlineExceeded,
88 Failed,
90}
91
92#[derive(Debug, Clone, Copy, Eq, PartialEq, Serialize, Deserialize)]
94#[serde(rename_all = "snake_case")]
95pub enum CancellationCause {
96 Explicit,
98 Deadline,
100}
101
102#[derive(Debug, Clone, Copy, Eq, PartialEq, Serialize, Deserialize)]
104pub struct SchedulerLimits {
105 pub max_active_requests: usize,
107 pub max_queued_work: usize,
109 pub max_new_submissions_per_turn: usize,
111 pub max_in_flight_global: usize,
113 pub max_in_flight_per_request: usize,
115 pub execution_slice: usize,
117}
118
119impl SchedulerLimits {
120 pub fn new(max_active_requests: usize, max_queued_work: usize) -> Result<Self, SchedulerError> {
122 Self::with_execution_bounds(
123 max_active_requests,
124 max_queued_work,
125 1,
126 max_active_requests,
127 1,
128 usize::MAX,
129 )
130 }
131
132 pub fn with_execution_bounds(
134 max_active_requests: usize,
135 max_queued_work: usize,
136 max_new_submissions_per_turn: usize,
137 max_in_flight_global: usize,
138 max_in_flight_per_request: usize,
139 execution_slice: usize,
140 ) -> Result<Self, SchedulerError> {
141 let values = [
142 max_active_requests,
143 max_queued_work,
144 max_new_submissions_per_turn,
145 max_in_flight_global,
146 max_in_flight_per_request,
147 execution_slice,
148 ];
149 if values.contains(&0) {
150 return Err(SchedulerError::InvalidLimits(values));
151 }
152 Ok(Self {
153 max_active_requests,
154 max_queued_work,
155 max_new_submissions_per_turn,
156 max_in_flight_global,
157 max_in_flight_per_request,
158 execution_slice,
159 })
160 }
161}
162
163impl Default for SchedulerLimits {
164 fn default() -> Self {
165 Self {
166 max_active_requests: 64,
167 max_queued_work: 256,
168 max_new_submissions_per_turn: 1,
169 max_in_flight_global: 64,
170 max_in_flight_per_request: 1,
171 execution_slice: usize::MAX,
172 }
173 }
174}
175
176pub trait WorkDescriptor {
178 type Error: std::error::Error;
180
181 fn encode_descriptor(&self, output: &mut Vec<u32>) -> Result<(), Self::Error>;
183
184 fn execution_slice_size(&self) -> usize {
186 1
187 }
188}
189
190pub trait SemanticStateTransaction {
192 type Branch;
194 type Error: std::error::Error;
196
197 fn branch(&self) -> Result<Self::Branch, Self::Error>;
199
200 fn commit_branch(&mut self, branch: Self::Branch) -> Result<(), Self::Error>;
202
203 fn discard_branch(branch: Self::Branch) -> Result<(), Self::Error> {
205 drop(branch);
206 Ok(())
207 }
208
209 fn permits_parallel_branches(&self) -> bool {
211 false
212 }
213}
214
215pub trait TransitionOutput {
217 type Error: std::error::Error;
219
220 fn is_complete(&self) -> Result<bool, Self::Error>;
222
223 fn backend_name(&self) -> Option<String> {
225 None
226 }
227
228 fn physically_preemptible(&self) -> bool {
230 false
231 }
232
233 fn retained_resources(&self) -> usize;
235}
236
237pub trait DistributedTransitionOutput: TransitionOutput {
239 fn encode_distributed_output(&self, output: &mut Vec<u32>) -> Result<(), String>;
241}
242
243#[derive(Debug)]
244struct Request<W, S> {
245 state: S,
246 next: u64,
247 pending: VecDeque<Queued<W>>,
248}
249
250#[derive(Debug)]
251struct Queued<W> {
252 id: WorkId,
253 work: W,
254 deadline: Option<Instant>,
255}
256
257#[derive(Debug)]
258struct Prepared<W, B> {
259 id: WorkId,
260 work: W,
261 descriptor: Vec<u32>,
262 branch: B,
263 deadline: Option<Instant>,
264}
265
266#[derive(Debug, Clone, Copy)]
267enum Disposition {
268 Publish,
269 Abandon { cancelled_at: Instant },
270 Fail,
271}
272
273#[derive(Debug)]
274struct Submitted<W, B, O> {
275 id: WorkId,
276 work: W,
277 branch: B,
278 output: O,
279 disposition: Disposition,
280}
281
282#[derive(Debug)]
284pub struct SchedulerProgress<W, O> {
285 pub newly_submitted: usize,
287 pub committed: Vec<(WorkId, W, O)>,
289 pub failed: Vec<(WorkId, SchedulerError)>,
291}
292
293impl<W, O> Default for SchedulerProgress<W, O> {
294 fn default() -> Self {
295 Self {
296 newly_submitted: 0,
297 committed: Vec::new(),
298 failed: Vec::new(),
299 }
300 }
301}
302
303#[derive(Debug, Clone, Eq, PartialEq, Serialize, Deserialize)]
305pub struct SchedulerCapabilities {
306 pub limits: SchedulerLimits,
308 pub observed_backends: Vec<String>,
310 pub executing_work_physically_preemptible: bool,
312 pub non_preemptible_interval: String,
314}
315
316#[derive(Debug, Clone, Copy, Default, Eq, PartialEq, Serialize, Deserialize)]
318pub struct SchedulerReport {
319 pub active_requests: usize,
321 pub queued_work: usize,
323 pub prepared_work: usize,
325 pub submitted_in_flight_work: usize,
327 pub completing_work: usize,
329 pub abandoned_in_flight_work: usize,
331 pub failed_in_flight_work: usize,
333 pub current_in_flight_work: usize,
335 pub peak_in_flight_work: usize,
337 pub peak_queued_work: usize,
339 pub submitted_work: u64,
341 pub completed_work: u64,
343 pub failed_work: u64,
345 pub discarded_work: u64,
347 pub cancellation_before_submission: u64,
349 pub cancellation_after_submission: u64,
351 pub abandoned_released_work: u64,
353 pub abandoned_retained_resources: usize,
355 pub peak_abandoned_retained_resources: usize,
357 pub last_cancellation_to_release_ns: Option<u128>,
359 pub max_cancellation_to_release_ns: Option<u128>,
361 pub finished_requests: u64,
363 pub cancelled_requests: u64,
365 pub deadline_expired_requests: u64,
367 pub drain_cycles: u64,
369 pub configured_submission_bound: usize,
371 pub configured_slice_bound: usize,
373 pub poisoned: bool,
375}
376
377#[derive(Debug)]
379pub struct Scheduler<W, S: SemanticStateTransaction, O: TransitionOutput> {
380 limits: SchedulerLimits,
381 requests: BTreeMap<RequestId, Request<W, S>>,
382 terminal: BTreeMap<RequestId, RequestStatus>,
383 ready: VecDeque<RequestId>,
384 prepared: VecDeque<Prepared<W, S::Branch>>,
385 submitted: Vec<Submitted<W, S::Branch, O>>,
386 lifecycle: BTreeMap<WorkId, WorkLifecycle>,
387 accepted_work: usize,
388 peak_accepted_work: usize,
389 peak_in_flight_work: usize,
390 submitted_work: u64,
391 completed_work: u64,
392 failed_work: u64,
393 discarded_work: u64,
394 cancellation_before_submission: u64,
395 cancellation_after_submission: u64,
396 abandoned_released_work: u64,
397 peak_abandoned_retained_resources: usize,
398 last_cancellation_to_release: Option<Duration>,
399 max_cancellation_to_release: Option<Duration>,
400 finished_requests: u64,
401 cancelled_requests: u64,
402 deadline_expired_requests: u64,
403 drain_cycles: u64,
404 observed_backends: BTreeSet<String>,
405 all_outputs_preemptible: bool,
406 poisoned: Option<String>,
407}
408
409impl<W, S: SemanticStateTransaction, O: TransitionOutput> Scheduler<W, S, O> {
410 pub fn new(limits: SchedulerLimits) -> Result<Self, SchedulerError> {
412 let limits = SchedulerLimits::with_execution_bounds(
413 limits.max_active_requests,
414 limits.max_queued_work,
415 limits.max_new_submissions_per_turn,
416 limits.max_in_flight_global,
417 limits.max_in_flight_per_request,
418 limits.execution_slice,
419 )?;
420 Ok(Self {
421 limits,
422 requests: BTreeMap::new(),
423 terminal: BTreeMap::new(),
424 ready: VecDeque::new(),
425 prepared: VecDeque::new(),
426 submitted: Vec::new(),
427 lifecycle: BTreeMap::new(),
428 accepted_work: 0,
429 peak_accepted_work: 0,
430 peak_in_flight_work: 0,
431 submitted_work: 0,
432 completed_work: 0,
433 failed_work: 0,
434 discarded_work: 0,
435 cancellation_before_submission: 0,
436 cancellation_after_submission: 0,
437 abandoned_released_work: 0,
438 peak_abandoned_retained_resources: 0,
439 last_cancellation_to_release: None,
440 max_cancellation_to_release: None,
441 finished_requests: 0,
442 cancelled_requests: 0,
443 deadline_expired_requests: 0,
444 drain_cycles: 0,
445 observed_backends: BTreeSet::new(),
446 all_outputs_preemptible: true,
447 poisoned: None,
448 })
449 }
450
451 pub fn validate_registration(&self, id: RequestId) -> Result<(), SchedulerError> {
453 self.ensure_ready()?;
454 if self.requests.contains_key(&id) || self.terminal.contains_key(&id) {
455 return Err(SchedulerError::DuplicateRequest(id));
456 }
457 if self.requests.len() >= self.limits.max_active_requests {
458 return Err(SchedulerError::Capacity(format!(
459 "scheduler active-request capacity {} is exhausted",
460 self.limits.max_active_requests
461 )));
462 }
463 Ok(())
464 }
465
466 pub fn register(&mut self, id: RequestId, state: S) -> Result<(), SchedulerError> {
468 self.validate_registration(id)?;
469 self.requests.insert(
470 id,
471 Request {
472 state,
473 next: 0,
474 pending: VecDeque::new(),
475 },
476 );
477 Ok(())
478 }
479
480 pub fn request_state(&self, id: RequestId) -> Option<&S> {
482 self.requests.get(&id).map(|entry| &entry.state)
483 }
484
485 pub fn request_state_mut(&mut self, id: RequestId) -> Result<&mut S, SchedulerError> {
487 self.ensure_ready()?;
488 if self.branch_count(id) != 0 {
489 return Err(SchedulerError::State(format!(
490 "request {} has prepared or submitted state branches",
491 id.value()
492 )));
493 }
494 self.requests
495 .get_mut(&id)
496 .map(|entry| &mut entry.state)
497 .ok_or(SchedulerError::UnknownRequest(id))
498 }
499
500 pub fn enqueue(&mut self, request: RequestId, work: W) -> Result<WorkId, SchedulerError> {
502 self.enqueue_with_deadline(request, work, None)
503 }
504
505 pub fn enqueue_with_deadline(
507 &mut self,
508 request: RequestId,
509 work: W,
510 deadline: Option<Instant>,
511 ) -> Result<WorkId, SchedulerError> {
512 Ok(self
513 .enqueue_batch_with_deadlines(request, vec![(work, deadline)])?
514 .pop()
515 .expect("one work item was supplied"))
516 }
517
518 pub fn enqueue_batch(
520 &mut self,
521 request: RequestId,
522 work: Vec<W>,
523 ) -> Result<Vec<WorkId>, SchedulerError> {
524 self.enqueue_batch_with_deadlines(
525 request,
526 work.into_iter().map(|work| (work, None)).collect(),
527 )
528 }
529
530 fn enqueue_batch_with_deadlines(
531 &mut self,
532 request: RequestId,
533 work: Vec<(W, Option<Instant>)>,
534 ) -> Result<Vec<WorkId>, SchedulerError> {
535 self.ensure_ready()?;
536 let requested = work.len();
537 let accepted_after = self
538 .accepted_work
539 .checked_add(requested)
540 .ok_or_else(|| SchedulerError::Capacity("scheduler occupancy overflow".into()))?;
541 if accepted_after > self.limits.max_queued_work {
542 return Err(SchedulerError::Capacity(format!(
543 "scheduler queue capacity {} cannot accept {requested} items with {} outstanding",
544 self.limits.max_queued_work, self.accepted_work
545 )));
546 }
547 let entry = self
548 .requests
549 .get_mut(&request)
550 .ok_or(SchedulerError::UnknownRequest(request))?;
551 let count = u64::try_from(requested)
552 .map_err(|_| SchedulerError::Capacity("work batch length exceeds u64".into()))?;
553 let next = entry
554 .next
555 .checked_add(count)
556 .ok_or_else(|| SchedulerError::Capacity("work identity space exhausted".into()))?;
557 let was_empty = entry.pending.is_empty();
558 let mut ids = Vec::with_capacity(requested);
559 for (offset, (work, deadline)) in work.into_iter().enumerate() {
560 let id = WorkId::new(request, entry.next + offset as u64);
561 entry.pending.push_back(Queued { id, work, deadline });
562 self.lifecycle.insert(id, WorkLifecycle::Queued);
563 ids.push(id);
564 }
565 entry.next = next;
566 if was_empty && requested != 0 {
567 self.ready.push_back(request);
568 }
569 self.accepted_work = accepted_after;
570 self.peak_accepted_work = self.peak_accepted_work.max(accepted_after);
571 self.submitted_work = self.submitted_work.saturating_add(count);
572 Ok(ids)
573 }
574
575 pub fn prepare_bounded(&mut self, limit: usize, now: Instant) -> Result<usize, SchedulerError>
577 where
578 W: WorkDescriptor,
579 {
580 self.ensure_ready()?;
581 if limit == 0 {
582 return Err(SchedulerError::Capacity(
583 "scheduler preparation bound must be positive".into(),
584 ));
585 }
586 self.expire_deadlines(now)?;
587 let mut count = 0;
588 let mut stalled = 0;
589 while count < limit && !self.ready.is_empty() {
590 let request = self.ready.pop_front().expect("ready queue is nonempty");
591 let branches = self.branch_count(request);
592 let can_branch = self
593 .requests
594 .get(&request)
595 .is_some_and(|entry| branches == 0 || entry.state.permits_parallel_branches());
596 if !can_branch || branches >= self.limits.max_in_flight_per_request {
597 self.ready.push_back(request);
598 stalled += 1;
599 if stalled >= self.ready.len() {
600 break;
601 }
602 continue;
603 }
604 stalled = 0;
605 let queued = self
606 .requests
607 .get_mut(&request)
608 .and_then(|entry| entry.pending.pop_front())
609 .expect("ready request owns queued work");
610 if self
611 .requests
612 .get(&request)
613 .is_some_and(|entry| !entry.pending.is_empty())
614 {
615 self.ready.push_back(request);
616 }
617 let slice = queued.work.execution_slice_size();
618 if slice == 0 || slice > self.limits.execution_slice {
619 self.fail_before_submission(queued.id);
620 return Err(SchedulerError::Descriptor(format!(
621 "work {:?} execution slice {slice} exceeds configured bound {}",
622 queued.id, self.limits.execution_slice
623 )));
624 }
625 let mut descriptor = Vec::new();
626 if let Err(error) = queued.work.encode_descriptor(&mut descriptor) {
627 self.fail_before_submission(queued.id);
628 return Err(SchedulerError::Descriptor(error.to_string()));
629 }
630 let branch = match self
631 .requests
632 .get(&request)
633 .expect("active request exists")
634 .state
635 .branch()
636 {
637 Ok(branch) => branch,
638 Err(error) => {
639 self.fail_before_submission(queued.id);
640 return Err(SchedulerError::State(error.to_string()));
641 }
642 };
643 self.lifecycle.insert(queued.id, WorkLifecycle::Prepared);
644 self.prepared.push_back(Prepared {
645 id: queued.id,
646 work: queued.work,
647 descriptor,
648 branch,
649 deadline: queued.deadline,
650 });
651 count += 1;
652 }
653 Ok(count)
654 }
655
656 pub fn submit_prepared<E>(
658 &mut self,
659 now: Instant,
660 mut execute: impl FnMut(WorkId, &W, &mut S::Branch) -> Result<O, E>,
661 ) -> Result<usize, SchedulerError>
662 where
663 E: std::error::Error,
664 {
665 self.ensure_ready()?;
666 self.expire_deadlines(now)?;
667 let capacity = self
668 .limits
669 .max_in_flight_global
670 .saturating_sub(self.submitted.len())
671 .min(self.limits.max_new_submissions_per_turn);
672 let mut count = 0;
673 while count < capacity {
674 let Some(mut prepared) = self.prepared.pop_front() else {
675 break;
676 };
677 if prepared.deadline.is_some_and(|deadline| deadline <= now) {
678 let request = prepared.id.request();
679 self.prepared.push_front(prepared);
680 self.cancel_internal(request, CancellationCause::Deadline, now)?;
681 continue;
682 }
683 let output = match execute(prepared.id, &prepared.work, &mut prepared.branch) {
684 Ok(output) => output,
685 Err(error) => {
686 let id = prepared.id;
687 let discard = S::discard_branch(prepared.branch).err();
688 self.lifecycle.insert(id, WorkLifecycle::Failed);
689 self.failed_work = self.failed_work.saturating_add(1);
690 self.accepted_work = self.accepted_work.saturating_sub(1);
691 self.fail_request(id.request());
692 let message = discard.map_or_else(
693 || error.to_string(),
694 |discard| format!("{error}; branch discard also failed: {discard}"),
695 );
696 return Err(SchedulerError::Submission(message));
697 }
698 };
699 if let Some(backend) = output.backend_name() {
700 self.observed_backends.insert(backend);
701 }
702 self.all_outputs_preemptible &= output.physically_preemptible();
703 self.lifecycle.insert(prepared.id, WorkLifecycle::Submitted);
704 self.submitted.push(Submitted {
705 id: prepared.id,
706 work: prepared.work,
707 branch: prepared.branch,
708 output,
709 disposition: Disposition::Publish,
710 });
711 count += 1;
712 self.peak_in_flight_work = self.peak_in_flight_work.max(self.submitted.len());
713 }
714 if count != 0 {
715 self.drain_cycles = self.drain_cycles.saturating_add(1);
716 }
717 self.update_abandoned_resource_peak();
718 Ok(count)
719 }
720
721 pub fn poll_completions(&mut self, now: Instant) -> SchedulerProgress<W, O> {
723 let mut progress = SchedulerProgress::default();
724 let mut retained = Vec::with_capacity(self.submitted.len());
725 for submitted in std::mem::take(&mut self.submitted) {
726 match submitted.output.is_complete() {
727 Ok(false) => retained.push(submitted),
728 Ok(true) => self.resolve_completed(submitted, now, &mut progress),
729 Err(error) => {
730 let id = submitted.id;
731 let already_failed = matches!(submitted.disposition, Disposition::Fail);
732 let discard = S::discard_branch(submitted.branch).err();
733 self.lifecycle.insert(id, WorkLifecycle::Failed);
734 if !already_failed {
735 self.failed_work = self.failed_work.saturating_add(1);
736 }
737 self.accepted_work = self.accepted_work.saturating_sub(1);
738 self.fail_request(id.request());
739 progress.failed.push((
740 id,
741 SchedulerError::Completion(discard.map_or_else(
742 || error.to_string(),
743 |discard| format!("{error}; branch discard also failed: {discard}"),
744 )),
745 ));
746 }
747 }
748 }
749 for submitted in &mut retained {
754 if !self.requests.contains_key(&submitted.id.request())
755 && matches!(submitted.disposition, Disposition::Publish)
756 {
757 submitted.disposition = Disposition::Abandon { cancelled_at: now };
758 self.lifecycle
759 .insert(submitted.id, WorkLifecycle::Abandoned);
760 }
761 }
762 self.submitted = retained;
763 self.update_abandoned_resource_peak();
764 progress
765 }
766
767 pub fn run_local_turn<E>(
769 &mut self,
770 now: Instant,
771 execute: impl FnMut(WorkId, &W, &mut S::Branch) -> Result<O, E>,
772 ) -> Result<SchedulerProgress<W, O>, SchedulerError>
773 where
774 W: WorkDescriptor,
775 E: std::error::Error,
776 {
777 self.ensure_ready()?;
778 let mut progress = self.poll_completions(now);
779 self.prepare_bounded(self.limits.max_new_submissions_per_turn, now)?;
780 progress.newly_submitted = self.submit_prepared(now, execute)?;
781 let after_submit = self.poll_completions(now);
782 progress.committed.extend(after_submit.committed);
783 progress.failed.extend(after_submit.failed);
784 Ok(progress)
785 }
786
787 pub fn run_distributed_turn<T, E>(
789 &mut self,
790 protocol: u64,
791 transport: &T,
792 wait: BoundedCompletionWait,
793 now: Instant,
794 execute: impl FnMut(WorkId, &W, &mut S::Branch) -> Result<O, E>,
795 ) -> Result<SchedulerProgress<W, O>, SchedulerError>
796 where
797 W: WorkDescriptor,
798 T: BoundedConsensusTransport,
799 <T::Completion as crate::Completion>::Error: std::fmt::Display,
800 E: std::error::Error,
801 O: DistributedTransitionOutput,
802 {
803 self.ensure_ready()?;
804 self.expire_deadlines_distributed(protocol, transport, wait, now)?;
805 let mut progress = self.poll_distributed(protocol, transport, wait, now)?;
806 self.prepare_bounded(self.limits.max_new_submissions_per_turn, now)?;
807 let plan = self
808 .prepared
809 .iter()
810 .take(
811 self.limits
812 .max_in_flight_global
813 .saturating_sub(self.submitted.len())
814 .min(self.limits.max_new_submissions_per_turn),
815 )
816 .map(|work| ScheduledWork {
817 id: work.id,
818 descriptor: &work.descriptor,
819 })
820 .collect::<Vec<_>>();
821 let submission_cycle = self.drain_cycles;
822 if let Err(error) =
823 validate_schedule_bounded(transport, &plan, submission_cycle, protocol, wait)
824 {
825 self.poison(error.to_string(), now);
826 return Err(SchedulerError::Consensus(error.to_string()));
827 }
828 let planned_count = plan.len();
829 drop(plan);
830 let before = self.submitted.len();
831 let submission = self.submit_prepared(now, execute);
832 let locally_submitted = self.submitted.len().saturating_sub(before);
833 let local_success = submission.is_ok();
834 let agreed = agree_submission_status_bounded(
835 transport,
836 protocol,
837 submission_cycle,
838 planned_count,
839 locally_submitted,
840 local_success,
841 wait,
842 );
843 match agreed {
844 Ok(true) => {
845 progress.newly_submitted = submission
846 .expect("topology-wide successful submission retains the local count");
847 }
848 Ok(false) => {
849 let reason = "backend submission failed or diverged on at least one rank";
850 self.poison(reason.into(), now);
851 return Err(SchedulerError::DistributedCompletion(reason.into()));
852 }
853 Err(error) => {
854 let reason = submission.err().map_or_else(
855 || error.to_string(),
856 |submission| format!("{error}; local submission also failed: {submission}"),
857 );
858 self.poison(reason.clone(), now);
859 return Err(SchedulerError::Consensus(reason));
860 }
861 }
862 Ok(progress)
863 }
864
865 pub fn cancel_distributed<T: BoundedConsensusTransport>(
867 &mut self,
868 protocol: u64,
869 request: RequestId,
870 transport: &T,
871 wait: BoundedCompletionWait,
872 now: Instant,
873 ) -> Result<(), SchedulerError>
874 where
875 <T::Completion as crate::Completion>::Error: std::fmt::Display,
876 {
877 self.ensure_ready()?;
878 let locally_ready =
879 self.requests.contains_key(&request) && !self.terminal.contains_key(&request);
880 let prepared = agree_disposition_status_bounded(
881 transport,
882 protocol,
883 request,
884 CancellationCause::Explicit,
885 1,
886 locally_ready,
887 wait,
888 );
889 let prepared = match prepared {
890 Ok(prepared) => prepared,
891 Err(error) => {
892 self.fence(error.to_string());
893 return Err(SchedulerError::Consensus(error.to_string()));
894 }
895 };
896 if !prepared {
897 let error = "distributed cancellation preparation failed on at least one rank";
898 self.fence(error.to_string());
899 return Err(SchedulerError::Consensus(error.into()));
900 }
901 let committed = agree_disposition_status_bounded(
902 transport,
903 protocol,
904 request,
905 CancellationCause::Explicit,
906 2,
907 true,
908 wait,
909 );
910 let committed = match committed {
911 Ok(committed) => committed,
912 Err(error) => {
913 self.fence(error.to_string());
914 return Err(SchedulerError::Consensus(error.to_string()));
915 }
916 };
917 if !committed {
918 let error = "distributed cancellation commit authorization failed";
919 self.fence(error.to_string());
920 return Err(SchedulerError::Consensus(error.into()));
921 }
922 let result = self.cancel_internal(request, CancellationCause::Explicit, now);
923 if let Err(error) = &result {
924 self.poison(error.to_string(), now);
928 }
929 result
930 }
931
932 pub fn finish(&mut self, request: RequestId) -> Result<(), SchedulerError> {
934 self.ensure_ready()?;
935 if self
936 .submitted
937 .iter()
938 .any(|work| work.id.request() == request)
939 {
940 return Err(SchedulerError::State(format!(
941 "request {} still has submitted work",
942 request.value()
943 )));
944 }
945 let entry = self
946 .requests
947 .remove(&request)
948 .ok_or(SchedulerError::UnknownRequest(request))?;
949 let queued = entry.pending.len();
950 for work in entry.pending {
951 self.lifecycle.insert(work.id, WorkLifecycle::Abandoned);
952 }
953 self.ready.retain(|candidate| *candidate != request);
954 let (prepared, discard_error) = self.discard_prepared_for_request(request);
955 let discarded = queued + prepared;
956 self.accepted_work = self.accepted_work.saturating_sub(discarded);
957 self.discarded_work = self.discarded_work.saturating_add(discarded as u64);
958 self.terminal.insert(
959 request,
960 if discard_error.is_some() {
961 RequestStatus::Failed
962 } else {
963 RequestStatus::Finished
964 },
965 );
966 self.finished_requests = self.finished_requests.saturating_add(1);
967 discard_error.map_or(Ok(()), |error| Err(SchedulerError::State(error)))
968 }
969
970 pub fn cancel(&mut self, request: RequestId) -> Result<(), SchedulerError> {
972 self.cancel_internal(request, CancellationCause::Explicit, Instant::now())
973 }
974
975 pub fn release(&mut self, request: RequestId) -> Result<S, SchedulerError> {
977 self.ensure_ready()?;
978 let entry = self
979 .requests
980 .get(&request)
981 .ok_or(SchedulerError::UnknownRequest(request))?;
982 if !entry.pending.is_empty() || self.branch_count(request) != 0 {
983 return Err(SchedulerError::State(format!(
984 "request {} still owns unpublished work",
985 request.value()
986 )));
987 }
988 Ok(self
989 .requests
990 .remove(&request)
991 .expect("checked active request")
992 .state)
993 }
994
995 pub fn forget_terminal(&mut self, request: RequestId) -> Result<RequestStatus, SchedulerError> {
997 self.ensure_ready()?;
998 if self
999 .submitted
1000 .iter()
1001 .any(|work| work.id.request() == request)
1002 {
1003 return Err(SchedulerError::State(format!(
1004 "request {} still has retained abandoned work",
1005 request.value()
1006 )));
1007 }
1008 self.terminal
1009 .remove(&request)
1010 .ok_or(SchedulerError::UnknownRequest(request))
1011 }
1012
1013 pub fn request_status(&self, id: RequestId) -> Option<RequestStatus> {
1015 self.requests
1016 .contains_key(&id)
1017 .then_some(RequestStatus::Active)
1018 .or_else(|| self.terminal.get(&id).copied())
1019 }
1020
1021 pub fn work_lifecycle(&self, id: WorkId) -> Option<WorkLifecycle> {
1023 self.lifecycle.get(&id).copied()
1024 }
1025
1026 pub fn queued_for_request(&self, request: RequestId) -> usize {
1028 self.requests
1029 .get(&request)
1030 .map_or(0, |entry| entry.pending.len())
1031 }
1032
1033 pub fn capabilities(&self) -> SchedulerCapabilities {
1035 SchedulerCapabilities {
1036 limits: self.limits,
1037 observed_backends: self.observed_backends.iter().cloned().collect(),
1038 executing_work_physically_preemptible: !self.observed_backends.is_empty()
1039 && self.all_outputs_preemptible,
1040 non_preemptible_interval:
1041 "from exact backend submission until that transition's completion resolves".into(),
1042 }
1043 }
1044
1045 pub fn report(&self) -> SchedulerReport {
1047 let abandoned = self
1048 .submitted
1049 .iter()
1050 .filter(|work| matches!(work.disposition, Disposition::Abandon { .. }))
1051 .count();
1052 let failed = self
1053 .submitted
1054 .iter()
1055 .filter(|work| matches!(work.disposition, Disposition::Fail))
1056 .count();
1057 SchedulerReport {
1058 active_requests: self.requests.len(),
1059 queued_work: self
1060 .requests
1061 .values()
1062 .map(|entry| entry.pending.len())
1063 .sum(),
1064 prepared_work: self.prepared.len(),
1065 submitted_in_flight_work: self.submitted.len() - abandoned - failed,
1066 completing_work: 0,
1067 abandoned_in_flight_work: abandoned,
1068 failed_in_flight_work: failed,
1069 current_in_flight_work: self.submitted.len(),
1070 peak_in_flight_work: self.peak_in_flight_work,
1071 peak_queued_work: self.peak_accepted_work,
1072 submitted_work: self.submitted_work,
1073 completed_work: self.completed_work,
1074 failed_work: self.failed_work,
1075 discarded_work: self.discarded_work,
1076 cancellation_before_submission: self.cancellation_before_submission,
1077 cancellation_after_submission: self.cancellation_after_submission,
1078 abandoned_released_work: self.abandoned_released_work,
1079 abandoned_retained_resources: self.abandoned_retained_resources(),
1080 peak_abandoned_retained_resources: self.peak_abandoned_retained_resources,
1081 last_cancellation_to_release_ns: self
1082 .last_cancellation_to_release
1083 .map(|duration| duration.as_nanos()),
1084 max_cancellation_to_release_ns: self
1085 .max_cancellation_to_release
1086 .map(|duration| duration.as_nanos()),
1087 finished_requests: self.finished_requests,
1088 cancelled_requests: self.cancelled_requests,
1089 deadline_expired_requests: self.deadline_expired_requests,
1090 drain_cycles: self.drain_cycles,
1091 configured_submission_bound: self.limits.max_new_submissions_per_turn,
1092 configured_slice_bound: self.limits.execution_slice,
1093 poisoned: self.poisoned.is_some(),
1094 }
1095 }
1096
1097 pub fn poison_reason(&self) -> Option<&str> {
1099 self.poisoned.as_deref()
1100 }
1101
1102 fn resolve_completed(
1103 &mut self,
1104 submitted: Submitted<W, S::Branch, O>,
1105 now: Instant,
1106 progress: &mut SchedulerProgress<W, O>,
1107 ) {
1108 match submitted.disposition {
1109 Disposition::Abandon { cancelled_at } => {
1110 if let Err(error) = S::discard_branch(submitted.branch) {
1111 self.lifecycle.insert(submitted.id, WorkLifecycle::Failed);
1112 self.failed_work = self.failed_work.saturating_add(1);
1113 progress.failed.push((
1114 submitted.id,
1115 SchedulerError::State(format!(
1116 "failed to discard abandoned state branch: {error}"
1117 )),
1118 ));
1119 }
1120 self.abandoned_released_work = self.abandoned_released_work.saturating_add(1);
1121 self.accepted_work = self.accepted_work.saturating_sub(1);
1122 let latency = now.saturating_duration_since(cancelled_at);
1123 self.last_cancellation_to_release = Some(latency);
1124 self.max_cancellation_to_release = Some(
1125 self.max_cancellation_to_release
1126 .map_or(latency, |previous| previous.max(latency)),
1127 );
1128 }
1129 Disposition::Publish => {
1130 self.lifecycle
1131 .insert(submitted.id, WorkLifecycle::Completing);
1132 let Some(request) = self.requests.get_mut(&submitted.id.request()) else {
1133 if let Err(error) = S::discard_branch(submitted.branch) {
1134 self.lifecycle.insert(submitted.id, WorkLifecycle::Failed);
1135 self.failed_work = self.failed_work.saturating_add(1);
1136 progress.failed.push((
1137 submitted.id,
1138 SchedulerError::State(format!(
1139 "failed to discard unpublished state branch: {error}"
1140 )),
1141 ));
1142 } else {
1143 self.lifecycle
1144 .insert(submitted.id, WorkLifecycle::Abandoned);
1145 }
1146 self.accepted_work = self.accepted_work.saturating_sub(1);
1147 return;
1148 };
1149 if let Err(error) = request.state.commit_branch(submitted.branch) {
1150 self.lifecycle.insert(submitted.id, WorkLifecycle::Failed);
1151 self.failed_work = self.failed_work.saturating_add(1);
1152 self.accepted_work = self.accepted_work.saturating_sub(1);
1153 self.fail_request(submitted.id.request());
1154 progress
1155 .failed
1156 .push((submitted.id, SchedulerError::State(error.to_string())));
1157 return;
1158 }
1159 self.lifecycle
1160 .insert(submitted.id, WorkLifecycle::Committed);
1161 self.completed_work = self.completed_work.saturating_add(1);
1162 self.accepted_work = self.accepted_work.saturating_sub(1);
1163 progress
1164 .committed
1165 .push((submitted.id, submitted.work, submitted.output));
1166 }
1167 Disposition::Fail => {
1168 let discard = S::discard_branch(submitted.branch).err();
1169 self.lifecycle.insert(submitted.id, WorkLifecycle::Failed);
1170 self.accepted_work = self.accepted_work.saturating_sub(1);
1171 self.fail_request(submitted.id.request());
1172 progress.failed.push((
1173 submitted.id,
1174 SchedulerError::DistributedCompletion(discard.map_or_else(
1175 || "backend completion failed on at least one rank".into(),
1176 |error| {
1177 format!(
1178 "backend completion failed on at least one rank; branch discard also failed: {error}"
1179 )
1180 },
1181 )),
1182 ));
1183 }
1184 }
1185 }
1186
1187 fn poll_distributed<T: BoundedConsensusTransport>(
1188 &mut self,
1189 protocol: u64,
1190 transport: &T,
1191 wait: BoundedCompletionWait,
1192 now: Instant,
1193 ) -> Result<SchedulerProgress<W, O>, SchedulerError>
1194 where
1195 <T::Completion as crate::Completion>::Error: std::fmt::Display,
1196 O: DistributedTransitionOutput,
1197 {
1198 let local = self
1199 .submitted
1200 .iter()
1201 .map(|work| {
1202 let (status, output) = match work.output.is_complete() {
1203 Ok(false) => (CompletionObservation::Incomplete, [0; 8]),
1204 Ok(true) => {
1205 let mut descriptor = Vec::new();
1206 match work.output.encode_distributed_output(&mut descriptor) {
1207 Ok(()) => {
1208 let mut digest = Sha256::new();
1209 for word in descriptor {
1210 digest.update(word.to_le_bytes());
1211 }
1212 let digest = digest.finalize();
1213 let output = std::array::from_fn(|index| {
1214 let start = index * 4;
1215 u32::from_le_bytes(
1216 digest[start..start + 4]
1217 .try_into()
1218 .expect("SHA-256 has eight complete u32 words"),
1219 )
1220 });
1221 (CompletionObservation::Complete, output)
1222 }
1223 Err(_) => (CompletionObservation::Failed, [0; 8]),
1224 }
1225 }
1226 Err(_) => (CompletionObservation::Failed, [0; 8]),
1227 };
1228 (work.id, status, output)
1229 })
1230 .collect::<Vec<_>>();
1231 let global = match resolve_output_completions_bounded(transport, protocol, &local, wait) {
1232 Ok(global) => global,
1233 Err(error) => {
1234 self.poison(error.to_string(), now);
1235 return Err(SchedulerError::Consensus(error.to_string()));
1236 }
1237 };
1238
1239 let mut progress = SchedulerProgress::default();
1240 let mut retained = Vec::with_capacity(self.submitted.len());
1241 for (mut work, status) in std::mem::take(&mut self.submitted).into_iter().zip(global) {
1242 match status {
1243 CompletionResolution::Incomplete => retained.push(work),
1244 CompletionResolution::Complete => {
1245 self.resolve_completed(work, now, &mut progress);
1246 }
1247 CompletionResolution::FailedPending => {
1248 if !matches!(work.disposition, Disposition::Fail) {
1249 work.disposition = Disposition::Fail;
1250 self.lifecycle.insert(work.id, WorkLifecycle::Failed);
1251 self.failed_work = self.failed_work.saturating_add(1);
1252 }
1253 retained.push(work);
1254 }
1255 CompletionResolution::FailedComplete => {
1256 if !matches!(work.disposition, Disposition::Fail) {
1257 work.disposition = Disposition::Fail;
1258 self.failed_work = self.failed_work.saturating_add(1);
1259 }
1260 self.resolve_completed(work, now, &mut progress);
1261 }
1262 }
1263 }
1264 for work in &mut retained {
1265 if !self.requests.contains_key(&work.id.request())
1266 && matches!(work.disposition, Disposition::Publish)
1267 {
1268 work.disposition = Disposition::Abandon { cancelled_at: now };
1269 self.lifecycle.insert(work.id, WorkLifecycle::Abandoned);
1270 }
1271 }
1272 self.submitted = retained;
1273 self.update_abandoned_resource_peak();
1274 Ok(progress)
1275 }
1276
1277 fn expire_deadlines(&mut self, now: Instant) -> Result<(), SchedulerError> {
1278 let mut expired = BTreeSet::new();
1279 for (request, entry) in &self.requests {
1280 if entry
1281 .pending
1282 .iter()
1283 .any(|work| work.deadline.is_some_and(|deadline| deadline <= now))
1284 {
1285 expired.insert(*request);
1286 }
1287 }
1288 for work in &self.prepared {
1289 if work.deadline.is_some_and(|deadline| deadline <= now) {
1290 expired.insert(work.id.request());
1291 }
1292 }
1293 for request in expired {
1294 self.cancel_internal(request, CancellationCause::Deadline, now)?;
1295 }
1296 Ok(())
1297 }
1298
1299 fn expire_deadlines_distributed<T: BoundedConsensusTransport>(
1300 &mut self,
1301 protocol: u64,
1302 transport: &T,
1303 wait: BoundedCompletionWait,
1304 now: Instant,
1305 ) -> Result<(), SchedulerError>
1306 where
1307 <T::Completion as crate::Completion>::Error: std::fmt::Display,
1308 {
1309 let mut expired = BTreeSet::new();
1310 for (request, entry) in &self.requests {
1311 if entry
1312 .pending
1313 .iter()
1314 .any(|work| work.deadline.is_some_and(|deadline| deadline <= now))
1315 {
1316 expired.insert(*request);
1317 }
1318 }
1319 for work in &self.prepared {
1320 if work.deadline.is_some_and(|deadline| deadline <= now) {
1321 expired.insert(work.id.request());
1322 }
1323 }
1324 let local = self
1325 .requests
1326 .keys()
1327 .copied()
1328 .map(|request| (request, expired.contains(&request)))
1329 .collect::<Vec<_>>();
1330 let globally_expired = match agree_deadline_candidates_bounded(
1331 transport,
1332 protocol,
1333 &local,
1334 self.limits.max_active_requests,
1335 wait,
1336 ) {
1337 Ok(expired) => expired,
1338 Err(error) => {
1339 self.fence(error.to_string());
1340 return Err(SchedulerError::Consensus(error.to_string()));
1341 }
1342 };
1343 for request in globally_expired {
1344 for phase in [1, 2] {
1345 let agreed = match agree_disposition_status_bounded(
1346 transport,
1347 protocol,
1348 request,
1349 CancellationCause::Deadline,
1350 phase,
1351 self.requests.contains_key(&request),
1352 wait,
1353 ) {
1354 Ok(agreed) => agreed,
1355 Err(error) => {
1356 self.fence(error.to_string());
1357 return Err(SchedulerError::Consensus(error.to_string()));
1358 }
1359 };
1360 if !agreed {
1361 let error = if phase == 1 {
1362 "distributed deadline preparation failed on at least one rank"
1363 } else {
1364 "distributed deadline commit authorization failed"
1365 };
1366 self.fence(error.into());
1367 return Err(SchedulerError::Consensus(error.into()));
1368 }
1369 }
1370 if let Err(error) = self.cancel_internal(request, CancellationCause::Deadline, now) {
1371 self.poison(error.to_string(), now);
1372 return Err(error);
1373 }
1374 }
1375 Ok(())
1376 }
1377
1378 fn cancel_internal(
1379 &mut self,
1380 request: RequestId,
1381 cause: CancellationCause,
1382 now: Instant,
1383 ) -> Result<(), SchedulerError> {
1384 self.ensure_ready()?;
1385 if self.terminal.contains_key(&request) {
1386 return Err(SchedulerError::State(format!(
1387 "request {} is already terminal",
1388 request.value()
1389 )));
1390 }
1391 let entry = self
1392 .requests
1393 .remove(&request)
1394 .ok_or(SchedulerError::UnknownRequest(request))?;
1395 let queued = entry.pending.len();
1396 for work in entry.pending {
1397 self.lifecycle.insert(work.id, WorkLifecycle::Abandoned);
1398 }
1399 self.ready.retain(|candidate| *candidate != request);
1400 let (prepared, discard_error) = self.discard_prepared_for_request(request);
1401 let before_submission = queued + prepared;
1402 self.accepted_work = self.accepted_work.saturating_sub(before_submission);
1403 self.discarded_work = self.discarded_work.saturating_add(before_submission as u64);
1404 self.cancellation_before_submission = self
1405 .cancellation_before_submission
1406 .saturating_add(before_submission as u64);
1407 let mut after_submission = 0u64;
1408 for work in &mut self.submitted {
1409 if work.id.request() == request && matches!(work.disposition, Disposition::Publish) {
1410 work.disposition = Disposition::Abandon { cancelled_at: now };
1411 self.lifecycle.insert(work.id, WorkLifecycle::Abandoned);
1412 after_submission += 1;
1413 }
1414 }
1415 self.cancellation_after_submission = self
1416 .cancellation_after_submission
1417 .saturating_add(after_submission);
1418 let status = match cause {
1419 CancellationCause::Explicit => {
1420 self.cancelled_requests = self.cancelled_requests.saturating_add(1);
1421 RequestStatus::Cancelled
1422 }
1423 CancellationCause::Deadline => {
1424 self.deadline_expired_requests = self.deadline_expired_requests.saturating_add(1);
1425 RequestStatus::DeadlineExceeded
1426 }
1427 };
1428 self.terminal.insert(
1429 request,
1430 if discard_error.is_some() {
1431 RequestStatus::Failed
1432 } else {
1433 status
1434 },
1435 );
1436 self.update_abandoned_resource_peak();
1437 discard_error.map_or(Ok(()), |error| Err(SchedulerError::State(error)))
1438 }
1439
1440 fn discard_prepared_for_request(&mut self, request: RequestId) -> (usize, Option<String>) {
1441 let mut retained = VecDeque::with_capacity(self.prepared.len());
1442 let mut discarded = 0;
1443 let mut errors = Vec::new();
1444 for work in std::mem::take(&mut self.prepared) {
1445 if work.id.request() == request {
1446 discarded += 1;
1447 match S::discard_branch(work.branch) {
1448 Ok(()) => {
1449 self.lifecycle.insert(work.id, WorkLifecycle::Abandoned);
1450 }
1451 Err(error) => {
1452 self.lifecycle.insert(work.id, WorkLifecycle::Failed);
1453 self.failed_work = self.failed_work.saturating_add(1);
1454 errors.push(format!("work {:?}: {error}", work.id));
1455 }
1456 }
1457 } else {
1458 retained.push_back(work);
1459 }
1460 }
1461 self.prepared = retained;
1462 let error = (!errors.is_empty()).then(|| errors.join("; "));
1463 (discarded, error)
1464 }
1465
1466 fn branch_count(&self, request: RequestId) -> usize {
1467 self.prepared
1468 .iter()
1469 .filter(|work| work.id.request() == request)
1470 .count()
1471 + self
1472 .submitted
1473 .iter()
1474 .filter(|work| work.id.request() == request)
1475 .count()
1476 }
1477
1478 fn fail_before_submission(&mut self, id: WorkId) {
1479 self.lifecycle.insert(id, WorkLifecycle::Failed);
1480 self.failed_work = self.failed_work.saturating_add(1);
1481 self.accepted_work = self.accepted_work.saturating_sub(1);
1482 self.fail_request(id.request());
1483 }
1484
1485 fn fail_request(&mut self, request: RequestId) {
1486 let Some(entry) = self.requests.remove(&request) else {
1487 return;
1488 };
1489 let queued = entry.pending.len();
1490 for work in entry.pending {
1491 self.lifecycle.insert(work.id, WorkLifecycle::Failed);
1492 }
1493 self.ready.retain(|candidate| *candidate != request);
1494 let (prepared, _) = self.discard_prepared_for_request(request);
1495 let discarded = queued + prepared;
1496 self.accepted_work = self.accepted_work.saturating_sub(discarded);
1497 self.discarded_work = self.discarded_work.saturating_add(discarded as u64);
1498 for work in &mut self.submitted {
1499 if work.id.request() == request && matches!(work.disposition, Disposition::Publish) {
1500 work.disposition = Disposition::Abandon {
1501 cancelled_at: Instant::now(),
1502 };
1503 self.lifecycle.insert(work.id, WorkLifecycle::Abandoned);
1504 }
1505 }
1506 self.terminal.insert(request, RequestStatus::Failed);
1507 }
1508
1509 fn ensure_ready(&self) -> Result<(), SchedulerError> {
1510 self.poisoned.as_ref().map_or(Ok(()), |reason| {
1511 Err(SchedulerError::Poisoned(reason.clone()))
1512 })
1513 }
1514
1515 fn fence(&mut self, reason: String) {
1516 if self.poisoned.is_none() {
1517 self.poisoned = Some(reason);
1518 }
1519 }
1520
1521 fn poison(&mut self, reason: String, now: Instant) {
1522 if self.poisoned.is_some() {
1523 return;
1524 }
1525 let mut discarded = 0usize;
1526 for (request, entry) in std::mem::take(&mut self.requests) {
1527 self.terminal.insert(request, RequestStatus::Failed);
1528 for work in entry.pending {
1529 self.lifecycle.insert(work.id, WorkLifecycle::Failed);
1530 discarded += 1;
1531 }
1532 }
1533 self.ready.clear();
1534 let mut cleanup_errors = Vec::new();
1535 for work in std::mem::take(&mut self.prepared) {
1536 let id = work.id;
1537 if let Err(error) = S::discard_branch(work.branch) {
1538 cleanup_errors.push(format!("work {id:?}: {error}"));
1539 self.failed_work = self.failed_work.saturating_add(1);
1540 }
1541 self.lifecycle.insert(id, WorkLifecycle::Failed);
1542 discarded += 1;
1543 }
1544 self.accepted_work = self.accepted_work.saturating_sub(discarded);
1545 self.discarded_work = self.discarded_work.saturating_add(discarded as u64);
1546 for work in &mut self.submitted {
1547 if matches!(work.disposition, Disposition::Publish) {
1548 work.disposition = Disposition::Abandon { cancelled_at: now };
1549 self.lifecycle.insert(work.id, WorkLifecycle::Abandoned);
1550 }
1551 }
1552 self.poisoned = Some(if cleanup_errors.is_empty() {
1553 reason
1554 } else {
1555 format!(
1556 "{reason}; prepared branch cleanup failed: {}",
1557 cleanup_errors.join("; ")
1558 )
1559 });
1560 self.update_abandoned_resource_peak();
1561 }
1562
1563 fn abandoned_retained_resources(&self) -> usize {
1564 self.submitted
1565 .iter()
1566 .filter(|work| matches!(work.disposition, Disposition::Abandon { .. }))
1567 .map(|work| work.output.retained_resources())
1568 .sum()
1569 }
1570
1571 fn update_abandoned_resource_peak(&mut self) {
1572 self.peak_abandoned_retained_resources = self
1573 .peak_abandoned_retained_resources
1574 .max(self.abandoned_retained_resources());
1575 }
1576}
1577
1578#[derive(Debug, Clone, Eq, PartialEq, thiserror::Error)]
1580pub enum SchedulerError {
1581 #[error("scheduler limits must be positive, got {0:?}")]
1583 InvalidLimits([usize; 6]),
1584 #[error("request {} is already registered", .0.value())]
1586 DuplicateRequest(RequestId),
1587 #[error("request {} is not active", .0.value())]
1589 UnknownRequest(RequestId),
1590 #[error("{0}")]
1592 Capacity(String),
1593 #[error("work descriptor failed: {0}")]
1595 Descriptor(String),
1596 #[error("semantic state transaction failed: {0}")]
1598 State(String),
1599 #[error("backend submission failed: {0}")]
1601 Submission(String),
1602 #[error("exact completion observation failed: {0}")]
1604 Completion(String),
1605 #[error("distributed scheduler consensus failed: {0}")]
1607 Consensus(String),
1608 #[error("distributed {0}")]
1610 DistributedCompletion(String),
1611 #[error("scheduler is poisoned after unsafe distributed ordering: {0}")]
1613 Poisoned(String),
1614}
1615
1616#[cfg(test)]
1617mod tests {
1618 use super::*;
1619 use crate::consensus::ConsensusTransport;
1620 use std::{
1621 cell::{Cell, RefCell},
1622 collections::VecDeque,
1623 convert::Infallible,
1624 rc::Rc,
1625 };
1626
1627 #[derive(Default)]
1628 struct State(u32);
1629
1630 impl SemanticStateTransaction for State {
1631 type Branch = u32;
1632 type Error = Infallible;
1633
1634 fn branch(&self) -> Result<u32, Infallible> {
1635 Ok(self.0 + 1)
1636 }
1637
1638 fn commit_branch(&mut self, branch: u32) -> Result<(), Infallible> {
1639 self.0 = branch;
1640 Ok(())
1641 }
1642
1643 fn permits_parallel_branches(&self) -> bool {
1644 true
1645 }
1646 }
1647
1648 impl WorkDescriptor for u32 {
1649 type Error = Infallible;
1650
1651 fn encode_descriptor(&self, output: &mut Vec<u32>) -> Result<(), Self::Error> {
1652 output.push(*self);
1653 Ok(())
1654 }
1655 }
1656
1657 #[derive(Debug)]
1658 struct Output {
1659 complete: Rc<Cell<bool>>,
1660 fail: bool,
1661 }
1662
1663 impl TransitionOutput for Output {
1664 type Error = std::io::Error;
1665
1666 fn is_complete(&self) -> Result<bool, Self::Error> {
1667 if self.fail {
1668 Err(std::io::Error::other("mock failure"))
1669 } else {
1670 Ok(self.complete.get())
1671 }
1672 }
1673
1674 fn backend_name(&self) -> Option<String> {
1675 Some("mock".into())
1676 }
1677
1678 fn retained_resources(&self) -> usize {
1679 2
1680 }
1681 }
1682
1683 impl DistributedTransitionOutput for Output {
1684 fn encode_distributed_output(&self, output: &mut Vec<u32>) -> Result<(), String> {
1685 output.push(0);
1686 Ok(())
1687 }
1688 }
1689
1690 #[derive(Default)]
1691 struct GatherStep {
1692 replacements: Vec<(usize, usize, u32)>,
1693 }
1694
1695 struct ScriptedTransport {
1696 participants: usize,
1697 steps: RefCell<VecDeque<GatherStep>>,
1698 }
1699
1700 impl ScriptedTransport {
1701 fn new(participants: usize, steps: Vec<GatherStep>) -> Self {
1702 Self {
1703 participants,
1704 steps: RefCell::new(steps.into()),
1705 }
1706 }
1707 }
1708
1709 impl ConsensusTransport for ScriptedTransport {
1710 type Error = Infallible;
1711
1712 fn participant_count(&self) -> usize {
1713 self.participants
1714 }
1715
1716 fn all_gather_words(&self, local: &[u32]) -> Result<Vec<u32>, Self::Error> {
1717 let mut gathered = local.repeat(self.participants);
1718 let step = self.steps.borrow_mut().pop_front().unwrap_or_default();
1719 for (rank, offset, value) in step.replacements {
1720 gathered[rank * local.len() + offset] = value;
1721 }
1722 Ok(gathered)
1723 }
1724 }
1725
1726 struct ReadyConsensusCompletion;
1727
1728 impl crate::Completion for ReadyConsensusCompletion {
1729 type Error = Infallible;
1730
1731 fn is_complete(&self) -> Result<bool, Self::Error> {
1732 Ok(true)
1733 }
1734
1735 fn wait(&self) -> Result<(), Self::Error> {
1736 Ok(())
1737 }
1738 }
1739
1740 impl crate::BoundedCompletion for ReadyConsensusCompletion {
1741 fn wait_bounded(
1742 self,
1743 _: crate::BoundedCompletionWait,
1744 ) -> Result<crate::BoundedCompletionOutcome, Self::Error> {
1745 Ok(crate::BoundedCompletionOutcome::Completed)
1746 }
1747 }
1748
1749 impl BoundedConsensusTransport for ScriptedTransport {
1750 type Completion = ReadyConsensusCompletion;
1751 type GatherOutput = Vec<u32>;
1752
1753 fn submit_all_gather_words(
1754 &self,
1755 local: &[u32],
1756 ) -> Result<crate::Submission<Self::GatherOutput, Self::Completion>, Self::Error> {
1757 Ok(crate::Submission {
1758 output: self.all_gather_words(local)?,
1759 completion: ReadyConsensusCompletion,
1760 })
1761 }
1762
1763 fn resolve_all_gather_words(
1764 &self,
1765 output: Self::GatherOutput,
1766 ) -> Result<Vec<u32>, Self::Error> {
1767 Ok(output)
1768 }
1769 }
1770
1771 fn consensus_wait() -> crate::BoundedCompletionWait {
1772 crate::BoundedCompletionWait::new(
1773 Duration::from_millis(10),
1774 crate::CompletionCancellationMode::QuarantineUntilComplete,
1775 )
1776 .unwrap()
1777 }
1778
1779 fn scheduler() -> Scheduler<u32, State, Output> {
1780 Scheduler::new(SchedulerLimits::default()).unwrap()
1781 }
1782
1783 #[test]
1784 fn scheduler_construction_revalidates_public_limit_fields() {
1785 let invalid = SchedulerLimits {
1786 max_active_requests: 0,
1787 ..SchedulerLimits::default()
1788 };
1789 assert!(matches!(
1790 Scheduler::<u32, State, Output>::new(invalid),
1791 Err(SchedulerError::InvalidLimits(_))
1792 ));
1793 }
1794
1795 #[test]
1796 fn queued_prepared_submitted_committed_exactly() {
1797 let done = Rc::new(Cell::new(false));
1798 let mut scheduler = scheduler();
1799 let request = RequestId::new(1);
1800 scheduler.register(request, State::default()).unwrap();
1801 let id = scheduler.enqueue(request, 7).unwrap();
1802 assert_eq!(scheduler.work_lifecycle(id), Some(WorkLifecycle::Queued));
1803 scheduler.prepare_bounded(1, Instant::now()).unwrap();
1804 assert_eq!(scheduler.work_lifecycle(id), Some(WorkLifecycle::Prepared));
1805 scheduler
1806 .submit_prepared(Instant::now(), |_, _, _| {
1807 Ok::<_, Infallible>(Output {
1808 complete: done.clone(),
1809 fail: false,
1810 })
1811 })
1812 .unwrap();
1813 assert!(scheduler
1814 .poll_completions(Instant::now())
1815 .committed
1816 .is_empty());
1817 done.set(true);
1818 assert_eq!(
1819 scheduler.poll_completions(Instant::now()).committed.len(),
1820 1
1821 );
1822 assert_eq!(scheduler.request_state(request).unwrap().0, 1);
1823 }
1824
1825 #[test]
1826 fn cancellation_retains_submitted_resources_until_completion() {
1827 let done = Rc::new(Cell::new(false));
1828 let mut scheduler = scheduler();
1829 let request = RequestId::new(2);
1830 scheduler.register(request, State::default()).unwrap();
1831 let id = scheduler.enqueue(request, 1).unwrap();
1832 scheduler.prepare_bounded(1, Instant::now()).unwrap();
1833 scheduler
1834 .submit_prepared(Instant::now(), |_, _, _| {
1835 Ok::<_, Infallible>(Output {
1836 complete: done.clone(),
1837 fail: false,
1838 })
1839 })
1840 .unwrap();
1841 scheduler.cancel(request).unwrap();
1842 assert_eq!(scheduler.work_lifecycle(id), Some(WorkLifecycle::Abandoned));
1843 assert_eq!(scheduler.report().abandoned_retained_resources, 2);
1844 assert_eq!(
1845 scheduler.request_status(request),
1846 Some(RequestStatus::Cancelled)
1847 );
1848 done.set(true);
1849 scheduler.poll_completions(Instant::now());
1850 assert_eq!(scheduler.report().current_in_flight_work, 0);
1851 assert_eq!(scheduler.report().abandoned_released_work, 1);
1852 }
1853
1854 #[test]
1855 fn exact_completion_failure_does_not_commit() {
1856 let mut scheduler = scheduler();
1857 let request = RequestId::new(3);
1858 scheduler.register(request, State::default()).unwrap();
1859 let id = scheduler.enqueue(request, 1).unwrap();
1860 scheduler.prepare_bounded(1, Instant::now()).unwrap();
1861 scheduler
1862 .submit_prepared(Instant::now(), |_, _, _| {
1863 Ok::<_, Infallible>(Output {
1864 complete: Rc::new(Cell::new(true)),
1865 fail: true,
1866 })
1867 })
1868 .unwrap();
1869 assert_eq!(scheduler.poll_completions(Instant::now()).failed.len(), 1);
1870 assert_eq!(scheduler.work_lifecycle(id), Some(WorkLifecycle::Failed));
1871 assert_eq!(
1872 scheduler.request_status(request),
1873 Some(RequestStatus::Failed)
1874 );
1875 }
1876
1877 #[test]
1878 fn request_failure_abandons_incomplete_sibling_until_exact_completion() {
1879 let waiting = Rc::new(Cell::new(false));
1880 let limits = SchedulerLimits::with_execution_bounds(1, 2, 2, 2, 2, 1).unwrap();
1881 let mut scheduler = Scheduler::new(limits).unwrap();
1882 let request = RequestId::new(31);
1883 scheduler.register(request, State::default()).unwrap();
1884 let ids = scheduler.enqueue_batch(request, vec![1, 2]).unwrap();
1885 scheduler.prepare_bounded(2, Instant::now()).unwrap();
1886 scheduler
1887 .submit_prepared(Instant::now(), |id, _, _| {
1888 Ok::<_, Infallible>(Output {
1889 complete: if id.sequence() == 0 {
1890 Rc::new(Cell::new(true))
1891 } else {
1892 waiting.clone()
1893 },
1894 fail: id.sequence() == 0,
1895 })
1896 })
1897 .unwrap();
1898
1899 let progress = scheduler.poll_completions(Instant::now());
1900 assert_eq!(progress.failed.len(), 1);
1901 assert_eq!(
1902 scheduler.work_lifecycle(ids[0]),
1903 Some(WorkLifecycle::Failed)
1904 );
1905 assert_eq!(
1906 scheduler.work_lifecycle(ids[1]),
1907 Some(WorkLifecycle::Abandoned)
1908 );
1909 assert_eq!(scheduler.report().abandoned_retained_resources, 2);
1910
1911 waiting.set(true);
1912 scheduler.poll_completions(Instant::now());
1913 assert_eq!(scheduler.report().current_in_flight_work, 0);
1914 assert_eq!(scheduler.report().abandoned_released_work, 1);
1915 }
1916
1917 #[test]
1918 fn submission_failure_marks_work_and_request_failed() {
1919 let mut scheduler = scheduler();
1920 let request = RequestId::new(4);
1921 scheduler.register(request, State::default()).unwrap();
1922 let id = scheduler.enqueue(request, 1).unwrap();
1923 scheduler.prepare_bounded(1, Instant::now()).unwrap();
1924 let error = scheduler
1925 .submit_prepared(Instant::now(), |_, _, _| {
1926 Err::<Output, _>(std::io::Error::other("mock submit"))
1927 })
1928 .unwrap_err();
1929 assert!(matches!(error, SchedulerError::Submission(_)));
1930 assert_eq!(scheduler.work_lifecycle(id), Some(WorkLifecycle::Failed));
1931 assert_eq!(
1932 scheduler.request_status(request),
1933 Some(RequestStatus::Failed)
1934 );
1935 }
1936
1937 #[test]
1938 fn distributed_submission_failure_is_agreed_before_any_rank_can_publish() {
1939 let transport = ScriptedTransport::new(
1940 2,
1941 vec![
1942 GatherStep::default(),
1943 GatherStep::default(),
1944 GatherStep::default(),
1945 GatherStep {
1946 replacements: vec![(1, 5, 1), (1, 6, 1)],
1947 },
1948 ],
1949 );
1950 let mut scheduler = scheduler();
1951 let request = RequestId::new(404);
1952 scheduler.register(request, State::default()).unwrap();
1953 let work = scheduler.enqueue(request, 1).unwrap();
1954 let error = scheduler
1955 .run_distributed_turn(
1956 0xFA11,
1957 &transport,
1958 consensus_wait(),
1959 Instant::now(),
1960 |_, _, _| Err::<Output, _>(std::io::Error::other("local submission failed")),
1961 )
1962 .unwrap_err();
1963 assert!(matches!(error, SchedulerError::DistributedCompletion(_)));
1964 assert!(scheduler.report().poisoned);
1965 assert_eq!(scheduler.report().completed_work, 0);
1966 assert_eq!(scheduler.work_lifecycle(work), Some(WorkLifecycle::Failed));
1967 }
1968
1969 #[test]
1970 fn deadlines_batches_and_fairness_are_backend_neutral() {
1971 let mut scheduler = scheduler();
1972 let first = RequestId::new(10);
1973 let second = RequestId::new(20);
1974 scheduler.register(first, State::default()).unwrap();
1975 scheduler.register(second, State::default()).unwrap();
1976 scheduler.enqueue_batch(first, vec![1, 2]).unwrap();
1977 scheduler.enqueue(second, 3).unwrap();
1978 scheduler.prepare_bounded(3, Instant::now()).unwrap();
1979 let mut order = Vec::new();
1980 scheduler
1981 .submit_prepared(Instant::now(), |id, _, _| {
1982 order.push(id.request());
1983 Ok::<_, Infallible>(Output {
1984 complete: Rc::new(Cell::new(true)),
1985 fail: false,
1986 })
1987 })
1988 .unwrap();
1989 assert_eq!(order, vec![first]);
1990
1991 let expired = RequestId::new(30);
1992 scheduler.register(expired, State::default()).unwrap();
1993 scheduler
1994 .enqueue_with_deadline(
1995 expired,
1996 4,
1997 Some(Instant::now().checked_sub(Duration::from_secs(1)).unwrap()),
1998 )
1999 .unwrap();
2000 scheduler.prepare_bounded(1, Instant::now()).unwrap();
2001 assert_eq!(
2002 scheduler.request_status(expired),
2003 Some(RequestStatus::DeadlineExceeded)
2004 );
2005 }
2006
2007 #[test]
2008 fn distributed_schedule_mismatch_poisons_the_canonical_machine() {
2009 let transport = ScriptedTransport::new(
2010 2,
2011 vec![
2012 GatherStep::default(),
2013 GatherStep::default(),
2014 GatherStep {
2015 replacements: vec![(1, 10, 99)],
2016 },
2017 ],
2018 );
2019 let mut scheduler = scheduler();
2020 let request = RequestId::new(50);
2021 scheduler.register(request, State::default()).unwrap();
2022 let work = scheduler.enqueue(request, 7).unwrap();
2023
2024 let error = scheduler
2025 .run_distributed_turn(
2026 0xAA,
2027 &transport,
2028 consensus_wait(),
2029 Instant::now(),
2030 |_, _, _| {
2031 Ok::<_, Infallible>(Output {
2032 complete: Rc::new(Cell::new(false)),
2033 fail: false,
2034 })
2035 },
2036 )
2037 .unwrap_err();
2038 assert!(matches!(error, SchedulerError::Consensus(_)));
2039 assert!(scheduler.poison_reason().is_some());
2040 assert!(scheduler.report().poisoned);
2041 assert_eq!(
2042 scheduler.request_status(request),
2043 Some(RequestStatus::Failed)
2044 );
2045 assert_eq!(scheduler.work_lifecycle(work), Some(WorkLifecycle::Failed));
2046 assert!(matches!(
2047 scheduler.enqueue(request, 8),
2048 Err(SchedulerError::Poisoned(_))
2049 ));
2050 }
2051
2052 #[test]
2053 fn poisoned_submissions_remain_retained_until_local_exact_completion() {
2054 let transport = ScriptedTransport::new(
2055 2,
2056 vec![
2057 GatherStep::default(),
2058 GatherStep::default(),
2059 GatherStep::default(),
2060 GatherStep::default(),
2061 GatherStep::default(),
2062 GatherStep {
2063 replacements: vec![(1, 4, 99)],
2064 },
2065 ],
2066 );
2067 let done = Rc::new(Cell::new(false));
2068 let mut scheduler = scheduler();
2069 let request = RequestId::new(54);
2070 scheduler.register(request, State::default()).unwrap();
2071 scheduler.enqueue(request, 1).unwrap();
2072 scheduler
2073 .run_distributed_turn(
2074 0xEE,
2075 &transport,
2076 consensus_wait(),
2077 Instant::now(),
2078 |_, _, _| {
2079 Ok::<_, Infallible>(Output {
2080 complete: done.clone(),
2081 fail: false,
2082 })
2083 },
2084 )
2085 .unwrap();
2086
2087 let error = scheduler
2088 .run_distributed_turn(
2089 0xEE,
2090 &transport,
2091 consensus_wait(),
2092 Instant::now(),
2093 |_, _, _| {
2094 Ok::<_, Infallible>(Output {
2095 complete: done.clone(),
2096 fail: false,
2097 })
2098 },
2099 )
2100 .unwrap_err();
2101 assert!(matches!(error, SchedulerError::Consensus(_)));
2102 assert_eq!(scheduler.report().abandoned_in_flight_work, 1);
2103 assert_eq!(scheduler.report().abandoned_retained_resources, 2);
2104
2105 assert!(scheduler
2106 .poll_completions(Instant::now())
2107 .committed
2108 .is_empty());
2109 assert_eq!(scheduler.report().current_in_flight_work, 1);
2110 done.set(true);
2111 assert!(scheduler
2112 .poll_completions(Instant::now())
2113 .committed
2114 .is_empty());
2115 assert_eq!(scheduler.report().current_in_flight_work, 0);
2116 assert_eq!(scheduler.report().abandoned_released_work, 1);
2117 }
2118
2119 #[test]
2120 fn distributed_failure_retains_output_until_every_rank_is_terminal() {
2121 let transport = ScriptedTransport::new(
2122 3,
2123 vec![
2124 GatherStep::default(),
2125 GatherStep::default(),
2126 GatherStep::default(),
2127 GatherStep::default(),
2128 GatherStep::default(),
2129 GatherStep {
2130 replacements: vec![(1, 7, 2), (2, 7, 0)],
2131 },
2132 GatherStep::default(),
2133 GatherStep::default(),
2134 GatherStep::default(),
2135 GatherStep {
2136 replacements: vec![(1, 7, 2), (2, 7, 1)],
2137 },
2138 GatherStep::default(),
2139 ],
2140 );
2141 let mut scheduler = scheduler();
2142 let request = RequestId::new(51);
2143 scheduler.register(request, State::default()).unwrap();
2144 scheduler.enqueue(request, 1).unwrap();
2145 let output = || Output {
2146 complete: Rc::new(Cell::new(true)),
2147 fail: false,
2148 };
2149
2150 scheduler
2151 .run_distributed_turn(
2152 0xBB,
2153 &transport,
2154 consensus_wait(),
2155 Instant::now(),
2156 |_, _, _| {
2157 Ok::<_, Infallible>(Output {
2158 complete: Rc::new(Cell::new(true)),
2159 fail: false,
2160 })
2161 },
2162 )
2163 .unwrap();
2164 let pending = scheduler
2165 .run_distributed_turn(
2166 0xBB,
2167 &transport,
2168 consensus_wait(),
2169 Instant::now(),
2170 |_, _, _| Ok::<_, Infallible>(output()),
2171 )
2172 .unwrap();
2173 assert!(pending.failed.is_empty());
2174 assert_eq!(scheduler.report().failed_in_flight_work, 1);
2175 assert_eq!(scheduler.report().current_in_flight_work, 1);
2176
2177 let terminal = scheduler
2178 .run_distributed_turn(
2179 0xBB,
2180 &transport,
2181 consensus_wait(),
2182 Instant::now(),
2183 |_, _, _| Ok::<_, Infallible>(output()),
2184 )
2185 .unwrap();
2186 assert_eq!(terminal.failed.len(), 1);
2187 assert_eq!(scheduler.report().current_in_flight_work, 0);
2188 assert_eq!(
2189 scheduler.request_status(request),
2190 Some(RequestStatus::Failed)
2191 );
2192 }
2193
2194 #[test]
2195 fn distributed_output_disagreement_publishes_nothing() {
2196 let transport = ScriptedTransport::new(
2197 2,
2198 vec![
2199 GatherStep::default(),
2200 GatherStep::default(),
2201 GatherStep::default(),
2202 GatherStep::default(),
2203 GatherStep::default(),
2204 GatherStep {
2205 replacements: vec![(1, 8, 99)],
2206 },
2207 ],
2208 );
2209 let mut scheduler = scheduler();
2210 let request = RequestId::new(405);
2211 scheduler.register(request, State::default()).unwrap();
2212 scheduler.enqueue(request, 1).unwrap();
2213 scheduler
2214 .run_distributed_turn(
2215 0x0A11,
2216 &transport,
2217 consensus_wait(),
2218 Instant::now(),
2219 |_, _, _| {
2220 Ok::<_, Infallible>(Output {
2221 complete: Rc::new(Cell::new(true)),
2222 fail: false,
2223 })
2224 },
2225 )
2226 .unwrap();
2227 let error = scheduler
2228 .run_distributed_turn(
2229 0x0A11,
2230 &transport,
2231 consensus_wait(),
2232 Instant::now(),
2233 |_, _, _| -> Result<Output, Infallible> {
2234 unreachable!("no additional work is queued")
2235 },
2236 )
2237 .unwrap_err();
2238 assert!(matches!(error, SchedulerError::Consensus(_)));
2239 assert_eq!(scheduler.report().completed_work, 0);
2240 assert!(scheduler.report().poisoned);
2241 }
2242
2243 #[test]
2244 fn distributed_cancellation_and_deadline_use_the_same_core_lifecycle() {
2245 let transport =
2246 ScriptedTransport::new(2, vec![GatherStep::default(), GatherStep::default()]);
2247 let mut cancelled = scheduler();
2248 let request = RequestId::new(52);
2249 cancelled.register(request, State::default()).unwrap();
2250 cancelled.enqueue(request, 1).unwrap();
2251 cancelled
2252 .cancel_distributed(0xCC, request, &transport, consensus_wait(), Instant::now())
2253 .unwrap();
2254 assert_eq!(
2255 cancelled.request_status(request),
2256 Some(RequestStatus::Cancelled)
2257 );
2258
2259 let transport = ScriptedTransport::new(
2260 2,
2261 vec![
2262 GatherStep::default(),
2263 GatherStep::default(),
2264 GatherStep::default(),
2265 ],
2266 );
2267 let mut scheduler = scheduler();
2268 let request = RequestId::new(53);
2269 scheduler.register(request, State::default()).unwrap();
2270 scheduler
2271 .enqueue_with_deadline(
2272 request,
2273 1,
2274 Some(Instant::now().checked_sub(Duration::from_secs(1)).unwrap()),
2275 )
2276 .unwrap();
2277 scheduler
2278 .run_distributed_turn(
2279 0xDD,
2280 &transport,
2281 consensus_wait(),
2282 Instant::now(),
2283 |_, _, _| {
2284 Ok::<_, Infallible>(Output {
2285 complete: Rc::new(Cell::new(true)),
2286 fail: false,
2287 })
2288 },
2289 )
2290 .unwrap();
2291 assert_eq!(
2292 scheduler.request_status(request),
2293 Some(RequestStatus::DeadlineExceeded)
2294 );
2295 }
2296
2297 #[test]
2298 fn bounded_distributed_cancellation_rolls_back_failed_phases_and_fences_retry() {
2299 for steps in [
2300 vec![
2301 GatherStep {
2302 replacements: vec![(1, 6, 0)],
2303 },
2304 GatherStep::default(),
2305 ],
2306 vec![
2307 GatherStep::default(),
2308 GatherStep {
2309 replacements: vec![(1, 6, 0)],
2310 },
2311 GatherStep::default(),
2312 ],
2313 ] {
2314 let transport = ScriptedTransport::new(2, steps);
2315 let mut scheduler = scheduler();
2316 let request = RequestId::new(53);
2317 scheduler.register(request, State::default()).unwrap();
2318 scheduler.enqueue(request, 1).unwrap();
2319
2320 assert!(scheduler
2321 .cancel_distributed(0xCD, request, &transport, consensus_wait(), Instant::now(),)
2322 .is_err());
2323 assert!(scheduler.request_state(request).is_some());
2324 assert_eq!(
2325 scheduler.request_status(request),
2326 Some(RequestStatus::Active)
2327 );
2328 let remaining = transport.steps.borrow().len();
2329 assert!(scheduler
2330 .cancel_distributed(0xCD, request, &transport, consensus_wait(), Instant::now(),)
2331 .is_err());
2332 assert_eq!(transport.steps.borrow().len(), remaining);
2333 assert!(scheduler.request_state(request).is_some());
2334 }
2335 }
2336
2337 #[test]
2338 fn bounded_distributed_deadline_uses_remote_status_and_fences_failed_phases() {
2339 let remote_expiry = GatherStep {
2340 replacements: vec![(1, 6, 1)],
2341 };
2342 let transport = ScriptedTransport::new(
2343 2,
2344 vec![
2345 remote_expiry,
2346 GatherStep::default(),
2347 GatherStep::default(),
2348 GatherStep::default(),
2349 GatherStep::default(),
2350 ],
2351 );
2352 let mut successful_scheduler = scheduler();
2353 let request = RequestId::new(54);
2354 successful_scheduler
2355 .register(request, State::default())
2356 .unwrap();
2357 successful_scheduler.enqueue(request, 1).unwrap();
2358 successful_scheduler
2359 .run_distributed_turn(
2360 0xCE,
2361 &transport,
2362 consensus_wait(),
2363 Instant::now(),
2364 |_, _, _| {
2365 Ok::<_, Infallible>(Output {
2366 complete: Rc::new(Cell::new(true)),
2367 fail: false,
2368 })
2369 },
2370 )
2371 .unwrap();
2372 assert_eq!(
2373 successful_scheduler.request_status(request),
2374 Some(RequestStatus::DeadlineExceeded)
2375 );
2376
2377 for steps in [
2378 vec![
2379 GatherStep {
2380 replacements: vec![(1, 6, 1)],
2381 },
2382 GatherStep {
2383 replacements: vec![(1, 6, 0)],
2384 },
2385 ],
2386 vec![
2387 GatherStep {
2388 replacements: vec![(1, 6, 1)],
2389 },
2390 GatherStep::default(),
2391 GatherStep {
2392 replacements: vec![(1, 6, 0)],
2393 },
2394 ],
2395 ] {
2396 let transport = ScriptedTransport::new(2, steps);
2397 let mut scheduler = scheduler();
2398 let request = RequestId::new(55);
2399 scheduler.register(request, State::default()).unwrap();
2400 scheduler.enqueue(request, 1).unwrap();
2401 assert!(scheduler
2402 .run_distributed_turn(
2403 0xCF,
2404 &transport,
2405 consensus_wait(),
2406 Instant::now(),
2407 |_, _, _| {
2408 Ok::<_, Infallible>(Output {
2409 complete: Rc::new(Cell::new(true)),
2410 fail: false,
2411 })
2412 },
2413 )
2414 .is_err());
2415 assert_eq!(
2416 scheduler.request_status(request),
2417 Some(RequestStatus::Active)
2418 );
2419 let remaining = transport.steps.borrow().len();
2420 assert!(scheduler
2421 .run_distributed_turn(
2422 0xCF,
2423 &transport,
2424 consensus_wait(),
2425 Instant::now(),
2426 |_, _, _| {
2427 Ok::<_, Infallible>(Output {
2428 complete: Rc::new(Cell::new(true)),
2429 fail: false,
2430 })
2431 },
2432 )
2433 .is_err());
2434 assert_eq!(transport.steps.borrow().len(), remaining);
2435 }
2436 }
2437
2438 #[test]
2439 fn production_telemetry_round_trips_without_backend_types() {
2440 let mut scheduler = scheduler();
2441 let request = RequestId::new(40);
2442 scheduler.register(request, State::default()).unwrap();
2443 scheduler.enqueue(request, 1).unwrap();
2444 scheduler.prepare_bounded(1, Instant::now()).unwrap();
2445 scheduler
2446 .submit_prepared(Instant::now(), |_, _, _| {
2447 Ok::<_, Infallible>(Output {
2448 complete: Rc::new(Cell::new(false)),
2449 fail: false,
2450 })
2451 })
2452 .unwrap();
2453
2454 let report = scheduler.report();
2455 let report_json = serde_json::to_string(&report).unwrap();
2456 assert_eq!(
2457 serde_json::from_str::<SchedulerReport>(&report_json).unwrap(),
2458 report
2459 );
2460
2461 let capabilities = scheduler.capabilities();
2462 assert_eq!(capabilities.observed_backends, ["mock"]);
2463 let capabilities_json = serde_json::to_string(&capabilities).unwrap();
2464 assert_eq!(
2465 serde_json::from_str::<SchedulerCapabilities>(&capabilities_json).unwrap(),
2466 capabilities
2467 );
2468 }
2469}