Skip to main content

eredu_core/
scheduler.rs

1//! Fair transactional scheduling independent of event and tensor runtimes.
2
3use 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/// Stable caller-assigned request identity.
17#[derive(Debug, Clone, Copy, Eq, Ord, PartialEq, PartialOrd, Hash, Serialize, Deserialize)]
18pub struct RequestId(u64);
19
20impl RequestId {
21    /// Creates an identity. Zero is valid.
22    pub const fn new(value: u64) -> Self {
23        Self(value)
24    }
25
26    /// Returns its numeric value.
27    pub const fn value(self) -> u64 {
28        self.0
29    }
30}
31
32/// Scheduler-assigned ordered transition identity.
33#[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    /// Creates an ordered transition identity.
41    pub const fn new(request: RequestId, sequence: u64) -> Self {
42        Self { request, sequence }
43    }
44
45    /// Request updated by this work.
46    pub const fn request(self) -> RequestId {
47        self.request
48    }
49
50    /// Zero-based transition number.
51    pub const fn sequence(self) -> u64 {
52        self.sequence
53    }
54}
55
56/// Authoritative lifecycle of one transition.
57#[derive(Debug, Clone, Copy, Eq, PartialEq, Serialize, Deserialize)]
58#[serde(rename_all = "snake_case")]
59pub enum WorkLifecycle {
60    /// Accepted but not branched.
61    Queued,
62    /// Branch exists; nothing submitted.
63    Prepared,
64    /// Exact completion exists and is incomplete.
65    Submitted,
66    /// Completion succeeded and commit is resolving.
67    Completing,
68    /// Branch was published.
69    Committed,
70    /// Submitted work was cancelled and is retained until exact completion.
71    Abandoned,
72    /// Preparation, execution, completion, or commit failed.
73    Failed,
74}
75
76/// Observable lifecycle of a request.
77#[derive(Debug, Clone, Copy, Eq, PartialEq, Serialize, Deserialize)]
78#[serde(rename_all = "snake_case")]
79pub enum RequestStatus {
80    /// Owns canonical state.
81    Active,
82    /// Finished normally.
83    Finished,
84    /// Explicitly cancelled.
85    Cancelled,
86    /// Deadline expired.
87    DeadlineExceeded,
88    /// Request-local failure.
89    Failed,
90}
91
92/// Cause attached to cancellation.
93#[derive(Debug, Clone, Copy, Eq, PartialEq, Serialize, Deserialize)]
94#[serde(rename_all = "snake_case")]
95pub enum CancellationCause {
96    /// Caller cancellation.
97    Explicit,
98    /// Deadline expiry.
99    Deadline,
100}
101
102/// Capacity and cooperative-preemption controls.
103#[derive(Debug, Clone, Copy, Eq, PartialEq, Serialize, Deserialize)]
104pub struct SchedulerLimits {
105    /// Maximum active requests.
106    pub max_active_requests: usize,
107    /// Maximum accepted nonterminal work.
108    pub max_queued_work: usize,
109    /// Maximum submissions per scheduling turn.
110    pub max_new_submissions_per_turn: usize,
111    /// Maximum global in-flight work.
112    pub max_in_flight_global: usize,
113    /// Maximum per-request in-flight work.
114    pub max_in_flight_per_request: usize,
115    /// Program-defined operations per transition.
116    pub execution_slice: usize,
117}
118
119impl SchedulerLimits {
120    /// Creates conservative bounds.
121    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    /// Creates explicit bounds.
133    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
176/// Stable semantic description of one work item.
177pub trait WorkDescriptor {
178    /// Descriptor construction error.
179    type Error: std::error::Error;
180
181    /// Appends every semantic field in stable wire order.
182    fn encode_descriptor(&self, output: &mut Vec<u32>) -> Result<(), Self::Error>;
183
184    /// Number of program-defined non-preemptible operations in this transition.
185    fn execution_slice_size(&self) -> usize {
186        1
187    }
188}
189
190/// Transaction over canonical semantic state.
191pub trait SemanticStateTransaction {
192    /// Transition-local branch.
193    type Branch;
194    /// State error.
195    type Error: std::error::Error;
196
197    /// Creates an unpublished branch.
198    fn branch(&self) -> Result<Self::Branch, Self::Error>;
199
200    /// Publishes a completed branch.
201    fn commit_branch(&mut self, branch: Self::Branch) -> Result<(), Self::Error>;
202
203    /// Discards an unpublished branch.
204    fn discard_branch(branch: Self::Branch) -> Result<(), Self::Error> {
205        drop(branch);
206        Ok(())
207    }
208
209    /// Whether independent branches may share the current canonical base.
210    fn permits_parallel_branches(&self) -> bool {
211        false
212    }
213}
214
215/// Exact backend completion retaining submitted resources.
216pub trait TransitionOutput {
217    /// Completion observation error.
218    type Error: std::error::Error;
219
220    /// Nonblocking exact completion observation.
221    fn is_complete(&self) -> Result<bool, Self::Error>;
222
223    /// Stable backend name used only for capability telemetry.
224    fn backend_name(&self) -> Option<String> {
225        None
226    }
227
228    /// Whether already executing work can be physically interrupted.
229    fn physically_preemptible(&self) -> bool {
230        false
231    }
232
233    /// Explicitly retained resource count.
234    fn retained_resources(&self) -> usize;
235}
236
237/// Completed output that can prove rank agreement before distributed publication.
238pub trait DistributedTransitionOutput: TransitionOutput {
239    /// Appends every portable output field in stable semantic order.
240    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/// Result of one scheduling turn.
283#[derive(Debug)]
284pub struct SchedulerProgress<W, O> {
285    /// Newly submitted transitions.
286    pub newly_submitted: usize,
287    /// Successfully committed work and backend outputs.
288    pub committed: Vec<(WorkId, W, O)>,
289    /// Failed work with structured scheduler context.
290    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/// Static scheduler capabilities and configured bounds.
304#[derive(Debug, Clone, Eq, PartialEq, Serialize, Deserialize)]
305pub struct SchedulerCapabilities {
306    /// Configured scheduler limits.
307    pub limits: SchedulerLimits,
308    /// Backends observed on submitted transitions.
309    pub observed_backends: Vec<String>,
310    /// Whether every observed backend output supports physical preemption.
311    pub executing_work_physically_preemptible: bool,
312    /// Exact description of the unavoidable cancellation interval.
313    pub non_preemptible_interval: String,
314}
315
316/// Snapshot of scheduler occupancy and cumulative telemetry.
317#[derive(Debug, Clone, Copy, Default, Eq, PartialEq, Serialize, Deserialize)]
318pub struct SchedulerReport {
319    /// Active requests.
320    pub active_requests: usize,
321    /// Queued work.
322    pub queued_work: usize,
323    /// Prepared work.
324    pub prepared_work: usize,
325    /// Submitted work still eligible to publish.
326    pub submitted_in_flight_work: usize,
327    /// Work currently resolving publication.
328    pub completing_work: usize,
329    /// Cancelled submitted work awaiting exact completion.
330    pub abandoned_in_flight_work: usize,
331    /// Globally failed work retained until every rank reaches completion.
332    pub failed_in_flight_work: usize,
333    /// All submitted work, including abandoned work.
334    pub current_in_flight_work: usize,
335    /// Peak submitted work.
336    pub peak_in_flight_work: usize,
337    /// Peak accepted nonterminal work.
338    pub peak_queued_work: usize,
339    /// Total accepted work.
340    pub submitted_work: u64,
341    /// Total committed work.
342    pub completed_work: u64,
343    /// Total failed work.
344    pub failed_work: u64,
345    /// Work discarded before submission.
346    pub discarded_work: u64,
347    /// Work cancelled before submission.
348    pub cancellation_before_submission: u64,
349    /// Work cancelled after submission.
350    pub cancellation_after_submission: u64,
351    /// Abandoned work released after exact completion.
352    pub abandoned_released_work: u64,
353    /// Resources retained by abandoned work.
354    pub abandoned_retained_resources: usize,
355    /// Maximum resources retained by abandoned work.
356    pub peak_abandoned_retained_resources: usize,
357    /// Last cancellation-to-release latency in nanoseconds.
358    pub last_cancellation_to_release_ns: Option<u128>,
359    /// Maximum cancellation-to-release latency in nanoseconds.
360    pub max_cancellation_to_release_ns: Option<u128>,
361    /// Requests ended normally.
362    pub finished_requests: u64,
363    /// Requests cancelled explicitly.
364    pub cancelled_requests: u64,
365    /// Requests ended by deadline.
366    pub deadline_expired_requests: u64,
367    /// Turns which submitted at least one transition.
368    pub drain_cycles: u64,
369    /// Configured per-turn submission bound.
370    pub configured_submission_bound: usize,
371    /// Configured execution-slice bound.
372    pub configured_slice_bound: usize,
373    /// Whether unsafe distributed ordering poisoned the scheduler.
374    pub poisoned: bool,
375}
376
377/// Fair scheduler whose state transitions are independent of backend objects.
378#[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    /// Creates an empty scheduler under validated limits.
411    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    /// Validates identity and active-request capacity before allocating state.
452    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    /// Registers canonical request state.
467    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    /// Returns immutable canonical state for an active request.
481    pub fn request_state(&self, id: RequestId) -> Option<&S> {
482        self.requests.get(&id).map(|entry| &entry.state)
483    }
484
485    /// Returns mutable canonical state only while no branch exists.
486    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    /// Enqueues one transition without a deadline.
501    pub fn enqueue(&mut self, request: RequestId, work: W) -> Result<WorkId, SchedulerError> {
502        self.enqueue_with_deadline(request, work, None)
503    }
504
505    /// Enqueues one transition with an optional absolute deadline.
506    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    /// Atomically enqueues an ordered batch.
519    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    /// Prepares bounded branches in round-robin request order.
576    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    /// Submits a bounded prepared prefix through the injected backend adapter.
657    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    /// Polls exact completions and publishes only successful work.
722    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        // A completion failure can make its whole request terminal while sibling
750        // submissions from the same poll batch are still incomplete. Those
751        // siblings remain backend-owned until their exact completions resolve,
752        // but can no longer publish into the removed canonical state.
753        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    /// Runs one bounded local scheduling turn.
768    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    /// Runs one bounded turn with topology-wide schedule and completion consensus.
788    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    /// Reaches topology-wide consensus before cancelling a request.
866    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            // The disposition is already globally authorized and local
925            // cancellation marks every owned work item abandoned before
926            // reporting branch-cleanup failure. Fence later scheduler work.
927            self.poison(error.to_string(), now);
928        }
929        result
930    }
931
932    /// Marks a request finished and discards its queued work.
933    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    /// Cancels a request, retaining submitted resources until exact completion.
971    pub fn cancel(&mut self, request: RequestId) -> Result<(), SchedulerError> {
972        self.cancel_internal(request, CancellationCause::Explicit, Instant::now())
973    }
974
975    /// Releases an idle request and returns its canonical state.
976    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    /// Removes a terminal identity so it may be reused explicitly.
996    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    /// Returns active or terminal request status.
1014    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    /// Returns a work lifecycle state.
1022    pub fn work_lifecycle(&self, id: WorkId) -> Option<WorkLifecycle> {
1023        self.lifecycle.get(&id).copied()
1024    }
1025
1026    /// Returns queued transitions for one active request.
1027    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    /// Returns configured bounds and observed backend capabilities.
1034    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    /// Returns current occupancy and cumulative telemetry.
1046    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    /// Returns the unsafe distributed-ordering failure, if one occurred.
1098    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/// Structured neutral scheduler error.
1579#[derive(Debug, Clone, Eq, PartialEq, thiserror::Error)]
1580pub enum SchedulerError {
1581    /// One or more limits were zero.
1582    #[error("scheduler limits must be positive, got {0:?}")]
1583    InvalidLimits([usize; 6]),
1584    /// Duplicate request identity.
1585    #[error("request {} is already registered", .0.value())]
1586    DuplicateRequest(RequestId),
1587    /// Unknown request identity.
1588    #[error("request {} is not active", .0.value())]
1589    UnknownRequest(RequestId),
1590    /// Capacity exhausted with contextual detail.
1591    #[error("{0}")]
1592    Capacity(String),
1593    /// Work descriptor validation failed.
1594    #[error("work descriptor failed: {0}")]
1595    Descriptor(String),
1596    /// Semantic state operation failed.
1597    #[error("semantic state transaction failed: {0}")]
1598    State(String),
1599    /// Backend submission failed.
1600    #[error("backend submission failed: {0}")]
1601    Submission(String),
1602    /// Exact completion observation failed.
1603    #[error("exact completion observation failed: {0}")]
1604    Completion(String),
1605    /// Topology-wide scheduler agreement failed.
1606    #[error("distributed scheduler consensus failed: {0}")]
1607    Consensus(String),
1608    /// A distributed backend failed and every rank reached a terminal completion.
1609    #[error("distributed {0}")]
1610    DistributedCompletion(String),
1611    /// Unsafe distributed ordering prevents further scheduler mutation.
1612    #[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}