Skip to main content

sim_lib_control/
jobs.rs

1use std::collections::{BTreeMap, BTreeSet, VecDeque};
2
3/// Stable identity assigned to an admitted job.
4#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
5pub struct JobId(u64);
6
7/// Maximum jobs accepted by a queue collection.
8#[derive(Clone, Copy, Debug, PartialEq, Eq)]
9pub struct AdmissionLimit(pub usize);
10
11/// Maximum jobs executed by one drain operation.
12#[derive(Clone, Copy, Debug, PartialEq, Eq)]
13pub struct WorkLimit(pub usize);
14
15/// Runtime-owned queue classes which must never share a FIFO.
16#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
17pub enum RuntimeJobClass {
18    /// Collector-produced finalization work.
19    Finalization,
20    /// A language microtask class, isolated by its stable profile name.
21    LanguageMicrotask(&'static str),
22}
23
24/// Recorded lifecycle state for a job.
25#[derive(Clone, Copy, Debug, PartialEq, Eq)]
26pub enum JobStatus {
27    /// Waiting in its typed FIFO.
28    Queued,
29    /// Executed by an explicit drain.
30    Completed,
31    /// Cancelled before execution.
32    Cancelled,
33}
34
35/// Receipt for admission or cancellation.
36#[derive(Clone, Copy, Debug, PartialEq, Eq)]
37pub struct JobReceipt {
38    /// Job identity.
39    pub id: JobId,
40    /// Current status.
41    pub status: JobStatus,
42}
43
44/// Receipt for a bounded drain.
45#[derive(Clone, Debug, PartialEq, Eq)]
46pub struct DrainReceipt<K> {
47    /// Selected queue class.
48    pub class: K,
49    /// Jobs completed in FIFO order.
50    pub completed: Vec<JobId>,
51    /// Whether jobs remain in the class.
52    pub pending: bool,
53}
54
55/// Receipt for a drain-to-empty checkpoint.
56#[derive(Clone, Debug, PartialEq, Eq)]
57pub struct CheckpointReceipt<K> {
58    /// Selected queue class.
59    pub class: K,
60    /// Jobs completed, including reentrant jobs.
61    pub completed: Vec<JobId>,
62}
63
64/// Closed failures from bounded admission and checkpoint work.
65#[derive(Clone, Copy, Debug, PartialEq, Eq)]
66pub enum CheckpointError {
67    /// The collection's admission bound was reached.
68    AdmissionExhausted,
69    /// The checkpoint could not reach empty within its work bound.
70    WorkExhausted,
71}
72
73type Job<K> = Box<dyn FnOnce(&mut JobQueues<K>)>;
74struct Entry<K: Ord> {
75    id: JobId,
76    run: Job<K>,
77}
78
79/// Deterministic, typed FIFO queues driven only by explicit caller operations.
80pub struct JobQueues<K: Ord> {
81    queues: BTreeMap<K, VecDeque<Entry<K>>>,
82    cancelled: BTreeSet<JobId>,
83    next_id: u64,
84    admitted: usize,
85    admission: AdmissionLimit,
86}
87
88impl<K: Clone + Ord> JobQueues<K> {
89    /// Creates empty queues with a lifetime admission limit.
90    pub fn new(admission: AdmissionLimit) -> Self {
91        Self {
92            queues: BTreeMap::new(),
93            cancelled: BTreeSet::new(),
94            next_id: 0,
95            admitted: 0,
96            admission,
97        }
98    }
99
100    /// Returns capacity remaining under the lifetime admission bound.
101    pub fn remaining_admission(&self) -> usize {
102        self.admission.0.saturating_sub(self.admitted)
103    }
104
105    /// Enqueues one job at the tail of its class FIFO.
106    pub fn enqueue(
107        &mut self,
108        class: K,
109        job: impl FnOnce(&mut Self) + 'static,
110    ) -> Result<JobReceipt, CheckpointError> {
111        if self.admitted >= self.admission.0 {
112            return Err(CheckpointError::AdmissionExhausted);
113        }
114        let id = JobId(self.next_id);
115        self.next_id = self.next_id.saturating_add(1);
116        self.admitted += 1;
117        self.queues.entry(class).or_default().push_back(Entry {
118            id,
119            run: Box::new(job),
120        });
121        Ok(JobReceipt {
122            id,
123            status: JobStatus::Queued,
124        })
125    }
126
127    /// Cancels a queued job. Cancellation is idempotent and receipt-bearing.
128    pub fn cancel(&mut self, id: JobId) -> JobReceipt {
129        self.cancelled.insert(id);
130        JobReceipt {
131            id,
132            status: JobStatus::Cancelled,
133        }
134    }
135
136    /// Drains at most `work` jobs from one class.
137    pub fn drain(&mut self, class: K, work: WorkLimit) -> DrainReceipt<K> {
138        let completed = self.run_selected(&class, work.0);
139        let pending = self
140            .queues
141            .get(&class)
142            .is_some_and(|queue| !queue.is_empty());
143        DrainReceipt {
144            class,
145            completed,
146            pending,
147        }
148    }
149
150    /// Drains one class to empty, including jobs enqueued reentrantly into it.
151    ///
152    /// Fails closed if empty cannot be reached within `work`; unrelated job
153    /// classes remain untouched.
154    pub fn checkpoint(
155        &mut self,
156        class: K,
157        work: WorkLimit,
158    ) -> Result<CheckpointReceipt<K>, CheckpointError> {
159        let completed = self.run_selected(&class, work.0);
160        if self
161            .queues
162            .get(&class)
163            .is_some_and(|queue| !queue.is_empty())
164        {
165            return Err(CheckpointError::WorkExhausted);
166        }
167        Ok(CheckpointReceipt { class, completed })
168    }
169
170    fn run_selected(&mut self, class: &K, limit: usize) -> Vec<JobId> {
171        let mut completed = Vec::new();
172        let mut examined = 0;
173        while examined < limit {
174            let Some(entry) = self.queues.get_mut(class).and_then(VecDeque::pop_front) else {
175                break;
176            };
177            examined += 1;
178            if self.cancelled.remove(&entry.id) {
179                continue;
180            }
181            (entry.run)(self);
182            completed.push(entry.id);
183        }
184        completed
185    }
186}