1use std::collections::{BTreeMap, BTreeSet, VecDeque};
2
3#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
5pub struct JobId(u64);
6
7#[derive(Clone, Copy, Debug, PartialEq, Eq)]
9pub struct AdmissionLimit(pub usize);
10
11#[derive(Clone, Copy, Debug, PartialEq, Eq)]
13pub struct WorkLimit(pub usize);
14
15#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
17pub enum RuntimeJobClass {
18 Finalization,
20 LanguageMicrotask(&'static str),
22}
23
24#[derive(Clone, Copy, Debug, PartialEq, Eq)]
26pub enum JobStatus {
27 Queued,
29 Completed,
31 Cancelled,
33}
34
35#[derive(Clone, Copy, Debug, PartialEq, Eq)]
37pub struct JobReceipt {
38 pub id: JobId,
40 pub status: JobStatus,
42}
43
44#[derive(Clone, Debug, PartialEq, Eq)]
46pub struct DrainReceipt<K> {
47 pub class: K,
49 pub completed: Vec<JobId>,
51 pub pending: bool,
53}
54
55#[derive(Clone, Debug, PartialEq, Eq)]
57pub struct CheckpointReceipt<K> {
58 pub class: K,
60 pub completed: Vec<JobId>,
62}
63
64#[derive(Clone, Copy, Debug, PartialEq, Eq)]
66pub enum CheckpointError {
67 AdmissionExhausted,
69 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
79pub 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 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 pub fn remaining_admission(&self) -> usize {
102 self.admission.0.saturating_sub(self.admitted)
103 }
104
105 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 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 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 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}