indicatrix_dispatch/item_queue/
mod.rs1use std::{
25 collections::VecDeque,
26 sync::{Condvar, Mutex, MutexGuard, PoisonError},
27};
28
29#[cfg(test)]
30mod tests;
31
32#[derive(Debug)]
36#[must_use = "hand the ticket back through ItemQueue::complete or ItemQueue::fail"]
37pub struct ItemTicket<T> {
38 item: T,
39 failed_on: Vec<usize>,
41}
42
43impl<T> ItemTicket<T> {
44 #[must_use]
46 pub const fn item(&self) -> &T {
47 &self.item
48 }
49
50 #[must_use]
52 pub const fn previous_failures(&self) -> usize {
53 self.failed_on.len()
54 }
55}
56
57#[derive(Debug, Clone, PartialEq, Eq)]
59pub enum FailOutcome<T> {
60 Requeued,
62 Abandoned(T),
64}
65
66#[derive(Debug, Clone, Copy, PartialEq, Eq)]
68pub struct QueueCounts {
69 pub pending: usize,
71 pub in_flight: usize,
73 pub completed: usize,
75 pub abandoned: usize,
77}
78
79#[derive(Debug)]
80struct Entry<T> {
81 item: T,
82 failed_on: Vec<usize>,
83}
84
85#[derive(Debug)]
86struct QueueState<T> {
87 fresh: VecDeque<Entry<T>>,
88 retry: VecDeque<Entry<T>>,
89 live: Vec<bool>,
91 in_flight: usize,
92 completed: usize,
93 abandoned: usize,
94 abandoned_unclaimed: Vec<T>,
96 closed: bool,
97}
98
99impl<T> QueueState<T> {
100 fn is_live(&self, lane: usize) -> bool {
101 self.live.get(lane).copied().unwrap_or(false)
102 }
103
104 fn nobody_left_for(&self, entry: &Entry<T>) -> bool {
106 self.live
107 .iter()
108 .enumerate()
109 .all(|(lane, &live)| !live || entry.failed_on.contains(&lane))
110 }
111
112 fn take_for(&mut self, lane: usize) -> Option<Entry<T>> {
114 let retried = self
115 .retry
116 .iter()
117 .position(|entry| !entry.failed_on.contains(&lane))
118 .and_then(|index| self.retry.remove(index));
119 retried.or_else(|| self.fresh.pop_front())
120 }
121
122 fn retire(&mut self, lane: usize) -> Vec<T> {
125 if let Some(live) = self.live.get_mut(lane) {
126 *live = false;
127 }
128 let mut dropped = Vec::new();
129 let retry = std::mem::take(&mut self.retry);
130 for entry in retry {
131 if self.nobody_left_for(&entry) {
132 dropped.push(entry.item);
133 } else {
134 self.retry.push_back(entry);
135 }
136 }
137 if !self.live.iter().any(|&live| live) {
138 dropped.extend(self.fresh.drain(..).map(|entry| entry.item));
139 }
140 self.abandoned += dropped.len();
141 dropped
142 }
143}
144
145#[derive(Debug)]
147pub struct ItemQueue<T> {
148 state: Mutex<QueueState<T>>,
149 wake: Condvar,
150}
151
152impl<T> ItemQueue<T> {
153 pub fn new(items: impl IntoIterator<Item = T>, lane_count: usize) -> Self {
155 Self {
156 state: Mutex::new(QueueState {
157 fresh: items
158 .into_iter()
159 .map(|item| Entry {
160 item,
161 failed_on: Vec::new(),
162 })
163 .collect(),
164 retry: VecDeque::new(),
165 live: vec![true; lane_count],
166 in_flight: 0,
167 completed: 0,
168 abandoned: 0,
169 abandoned_unclaimed: Vec::new(),
170 closed: false,
171 }),
172 wake: Condvar::new(),
173 }
174 }
175
176 fn lock(&self) -> MutexGuard<'_, QueueState<T>> {
177 self.state.lock().unwrap_or_else(PoisonError::into_inner)
178 }
179
180 pub fn claim(&self, lane: usize) -> Option<ItemTicket<T>> {
185 let mut state = self.lock();
186 loop {
187 if state.closed || !state.is_live(lane) {
188 return None;
189 }
190 if let Some(entry) = state.take_for(lane) {
191 state.in_flight += 1;
192 return Some(ItemTicket {
193 item: entry.item,
194 failed_on: entry.failed_on,
195 });
196 }
197 if state.in_flight == 0 {
198 let dropped = state.retire(lane);
199 state.abandoned_unclaimed.extend(dropped);
200 drop(state);
201 self.wake.notify_all();
202 return None;
203 }
204 state = self
205 .wake
206 .wait(state)
207 .unwrap_or_else(PoisonError::into_inner);
208 }
209 }
210
211 pub fn complete(&self, ticket: ItemTicket<T>) -> T {
213 let mut state = self.lock();
214 state.in_flight -= 1;
215 state.completed += 1;
216 drop(state);
217 self.wake.notify_all();
218 ticket.item
219 }
220
221 #[must_use = "an abandoned item must be reported"]
224 pub fn fail(&self, lane: usize, ticket: ItemTicket<T>) -> FailOutcome<T> {
225 let mut entry = Entry {
226 item: ticket.item,
227 failed_on: ticket.failed_on,
228 };
229 if !entry.failed_on.contains(&lane) {
230 entry.failed_on.push(lane);
231 }
232 let mut state = self.lock();
233 state.in_flight -= 1;
234 let outcome = if state.nobody_left_for(&entry) {
235 state.abandoned += 1;
236 FailOutcome::Abandoned(entry.item)
237 } else {
238 state.retry.push_back(entry);
239 FailOutcome::Requeued
240 };
241 drop(state);
242 self.wake.notify_all();
243 outcome
244 }
245
246 #[must_use = "the returned items were abandoned and must be reported"]
249 pub fn retire(&self, lane: usize) -> Vec<T> {
250 let dropped = self.lock().retire(lane);
251 self.wake.notify_all();
252 dropped
253 }
254
255 pub fn take_abandoned(&self) -> Vec<T> {
257 std::mem::take(&mut self.lock().abandoned_unclaimed)
258 }
259
260 #[must_use = "the returned items were never processed"]
265 pub fn close(&self) -> Vec<T> {
266 let mut state = self.lock();
267 state.closed = true;
268 let mut left = std::mem::take(&mut state.abandoned_unclaimed);
269 left.extend(state.retry.drain(..).map(|entry| entry.item));
270 left.extend(state.fresh.drain(..).map(|entry| entry.item));
271 drop(state);
272 self.wake.notify_all();
273 left
274 }
275
276 #[must_use]
278 pub fn counts(&self) -> QueueCounts {
279 let state = self.lock();
280 let counts = QueueCounts {
281 pending: state.fresh.len() + state.retry.len(),
282 in_flight: state.in_flight,
283 completed: state.completed,
284 abandoned: state.abandoned,
285 };
286 drop(state);
287 counts
288 }
289}