Skip to main content

indicatrix_dispatch/item_queue/
mod.rs

1//! [`ItemQueue`]: whole-item work distribution over N lanes (batch previews, tilt-curve
2//! sets, anything that cannot be split into sample ranges).
3//!
4//! Generalised from the desktop batch's `WorkQueue` (one local pile, one remote lane,
5//! remote failures go to local only) to N equal lanes:
6//!
7//! - Items are handed out one at a time; whichever lane is free claims the next, so a
8//!   faster lane naturally does more and no throughput model is needed.
9//! - A failed item goes back to a retry pile, **never to a lane that already failed
10//!   it** -- the N-lane form of "a remote failure is retried locally only". An item that
11//!   fails for a persistent reason is therefore tried at most once per lane.
12//! - An item every still-running lane has failed is **abandoned** and handed back to
13//!   the caller to report as failed.
14//! - A lane with nothing claimable waits while any item is in flight (it may fail and
15//!   become claimable), and exits otherwise; exiting counts as retiring the lane.
16//!
17//! # No item is lost or duplicated
18//!
19//! Every item leaves the queue exactly once, through exactly one of:
20//! [`ItemQueue::complete`], [`FailOutcome::Abandoned`], the list [`ItemQueue::retire`]
21//! returns, [`ItemQueue::take_abandoned`] (items abandoned when a lane exited from
22//! [`ItemQueue::claim`]), or [`ItemQueue::close`]. All state lives under one mutex.
23
24use std::{
25    collections::VecDeque,
26    sync::{Condvar, Mutex, MutexGuard, PoisonError},
27};
28
29#[cfg(test)]
30mod tests;
31
32/// One claimed item. Hand it back through [`ItemQueue::complete`] or
33/// [`ItemQueue::fail`]; the queue counts it as in flight until then, so dropping a
34/// ticket instead would keep idle lanes waiting for it forever.
35#[derive(Debug)]
36#[must_use = "hand the ticket back through ItemQueue::complete or ItemQueue::fail"]
37pub struct ItemTicket<T> {
38    item: T,
39    /// Lanes that already failed this item.
40    failed_on: Vec<usize>,
41}
42
43impl<T> ItemTicket<T> {
44    /// The claimed item.
45    #[must_use]
46    pub const fn item(&self) -> &T {
47        &self.item
48    }
49
50    /// How many lanes failed this item before.
51    #[must_use]
52    pub const fn previous_failures(&self) -> usize {
53        self.failed_on.len()
54    }
55}
56
57/// What [`ItemQueue::fail`] did with the item.
58#[derive(Debug, Clone, PartialEq, Eq)]
59pub enum FailOutcome<T> {
60    /// Another running lane that has not failed it yet will get it.
61    Requeued,
62    /// Every running lane has failed it; report it as failed.
63    Abandoned(T),
64}
65
66/// A snapshot of the queue's bookkeeping.
67#[derive(Debug, Clone, Copy, PartialEq, Eq)]
68pub struct QueueCounts {
69    /// Items not yet claimed (fresh plus retry).
70    pub pending: usize,
71    /// Items claimed and not yet completed or failed.
72    pub in_flight: usize,
73    /// Items completed.
74    pub completed: usize,
75    /// Items abandoned (all of them, however they were handed back).
76    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[lane]`: the lane may still claim.
90    live: Vec<bool>,
91    in_flight: usize,
92    completed: usize,
93    abandoned: usize,
94    /// Items abandoned when a lane exited from `claim`, awaiting `take_abandoned`.
95    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    /// Whether no running lane is left that has not failed `entry`.
105    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    /// The next item for `lane`: a retry it has not failed first, then fresh work.
113    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    /// Marks `lane` as no longer claiming and removes every item nobody running can
123    /// take any more.
124    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/// See the module doc. Shared by reference between the lane threads.
146#[derive(Debug)]
147pub struct ItemQueue<T> {
148    state: Mutex<QueueState<T>>,
149    wake: Condvar,
150}
151
152impl<T> ItemQueue<T> {
153    /// A queue of `items` for lanes `0..lane_count`.
154    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    /// Claims the next item for `lane`, blocking while nothing is claimable for it but
181    /// another item is in flight. `None` means this lane is done for good (nothing it
182    /// can take is left, the lane was retired, or the queue was closed); items only it
183    /// could still have taken are then abandoned (see [`Self::take_abandoned`]).
184    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    /// The item was processed successfully; returns it.
212    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    /// `lane` failed the item: it goes back for a running lane that has not failed it,
222    /// or is abandoned when there is none.
223    #[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    /// Takes `lane` out of the rotation (it keeps failing, its worker went away) and
247    /// returns the items that no running lane can take any more.
248    #[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    /// Items abandoned when a lane exited from [`Self::claim`] (drained by this call).
256    pub fn take_abandoned(&self) -> Vec<T> {
257        std::mem::take(&mut self.lock().abandoned_unclaimed)
258    }
259
260    /// Stops the queue (cancellation): every later [`Self::claim`] returns `None`, and
261    /// every item not yet claimed -- plus any not yet taken with
262    /// [`Self::take_abandoned`] -- is returned. Items in flight still come back through
263    /// [`Self::complete`] / [`Self::fail`].
264    #[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    /// The queue's current bookkeeping.
277    #[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}