Skip to main content

behavior_actors/routing/
work_queue.rs

1//! Bounded FIFO work admission over explicitly available workers.
2
3use std::collections::VecDeque;
4
5use behavior::{
6    Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, Protocol,
7    SendEffects, User,
8};
9
10use crate::DeliveryRoute;
11
12/// Complete observable [`WorkQueue`] state.
13pub struct WorkQueueState<WorkerRoute> {
14    /// Maximum values that may wait; zero permits only immediate dispatch.
15    pub capacity: usize,
16    available: Vec<WorkerRoute>,
17    queued: usize,
18}
19
20impl<WorkerRoute> WorkQueueState<WorkerRoute> {
21    /// Workers eligible for one dispatch, in availability order.
22    #[must_use]
23    pub fn available(&self) -> &[WorkerRoute] {
24        &self.available
25    }
26
27    /// Number of owned values waiting in FIFO order.
28    #[must_use]
29    pub fn queued(&self) -> usize {
30        self.queued
31    }
32}
33
34/// Exhaustive work-admission rejection reason.
35#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
36pub enum WorkQueueRejection {
37    /// No worker was available and waiting capacity was full.
38    Full,
39}
40
41/// Factual result of one submission.
42#[derive(Debug, PartialEq, Eq)]
43pub enum WorkQueueOutcome<T> {
44    /// Value is retained pending worker availability.
45    Queued {
46        /// Waiting depth after admission.
47        depth: usize,
48    },
49    /// Value was assigned to an available worker.
50    Dispatched {
51        /// Waiting depth after dispatch.
52        queued: usize,
53    },
54    /// Value was not admitted; ownership is returned.
55    Rejected {
56        /// Unaccepted value.
57        value: T,
58        /// Rejection reason.
59        reason: WorkQueueRejection,
60    },
61}
62
63/// Operations accepted by [`WorkQueue`].
64pub enum WorkQueueMessage<T, WorkerRoute, ReplyRoute> {
65    /// Submit one value and its typed outcome recipient.
66    Submit {
67        /// Work value.
68        value: T,
69        /// Outcome recipient retained with queued work.
70        reply_to: ReplyRoute,
71    },
72    /// Announce one worker as eligible for exactly one dispatch.
73    Available {
74        /// Worker recipient.
75        worker: WorkerRoute,
76    },
77    /// Withdraw one currently available worker, idempotently.
78    Withdraw {
79        /// Worker recipient.
80        worker: WorkerRoute,
81    },
82}
83
84/// Named effect lanes emitted by [`WorkQueue`].
85#[derive(behavior_macros::SendProduct)]
86pub struct WorkQueueSends<Assignments, OutcomeSends> {
87    /// Work assigned to workers.
88    pub assignments: Assignments,
89    /// Submission admission and dispatch facts.
90    pub outcomes: OutcomeSends,
91}
92
93struct Waiting<T, Route> {
94    value: T,
95    reply_to: Route,
96}
97
98/// Bounded FIFO admission and worker-availability behavior.
99///
100/// Availability is a one-dispatch capability. It immediately consumes the
101/// oldest waiting value or joins a unique FIFO of available workers. Submission
102/// consumes the oldest worker or enters the bounded FIFO; at capacity it
103/// returns the owned value. Duplicate availability and withdrawal are
104/// idempotent. Initialization is empty, no actors are created, and the host
105/// never terminates by policy. FIFO selection and bounded admission are Bombay
106/// policy. Worker execution, mailbox admission, and physical backpressure are
107/// runtime responsibilities. Naming its protocol requires only typed routes;
108/// running the queue requires cloneable, comparable worker routes for state
109/// inspection and duplicate availability checks. No transition has a semantic
110/// panic condition.
111pub struct WorkQueue<
112    A: Address,
113    T,
114    WorkerRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
115    ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = WorkQueueOutcome<T>>>,
116> {
117    capacity: usize,
118    available: VecDeque<WorkerRoute>,
119    waiting: VecDeque<Waiting<T, ReplyRoute>>,
120    marker: core::marker::PhantomData<fn() -> A>,
121}
122
123type QueueActions<A, Assignments, OutcomeSends> =
124    Actions<A, Never, WorkQueueSends<Assignments, OutcomeSends>, NoBirths>;
125
126impl<A, T, WorkerRoute, ReplyRoute> WorkQueue<A, T, WorkerRoute, ReplyRoute>
127where
128    A: Address,
129    WorkerRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>> + Clone + PartialEq,
130    ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = WorkQueueOutcome<T>>> + Clone,
131{
132    /// Construct an empty queue with explicit waiting capacity.
133    #[must_use]
134    pub fn new(capacity: usize) -> Self {
135        Self {
136            capacity,
137            available: VecDeque::new(),
138            waiting: VecDeque::with_capacity(capacity),
139            marker: core::marker::PhantomData,
140        }
141    }
142    /// Return complete observable queue state.
143    #[must_use]
144    pub fn state(&self) -> WorkQueueState<WorkerRoute> {
145        WorkQueueState {
146            capacity: self.capacity,
147            available: self.available.iter().cloned().collect(),
148            queued: self.waiting.len(),
149        }
150    }
151    fn sends(
152        assignments: WorkerRoute::Sends,
153        outcomes: ReplyRoute::Sends,
154    ) -> QueueActions<A, WorkerRoute::Sends, ReplyRoute::Sends> {
155        Actions::send(WorkQueueSends {
156            assignments,
157            outcomes,
158        })
159    }
160    fn submit(
161        &mut self,
162        value: T,
163        reply_to: ReplyRoute,
164    ) -> QueueActions<A, WorkerRoute::Sends, ReplyRoute::Sends> {
165        if let Some(worker) = self.available.pop_front() {
166            return Self::sends(
167                worker.deliver(value),
168                reply_to.deliver(WorkQueueOutcome::Dispatched {
169                    queued: self.waiting.len(),
170                }),
171            );
172        }
173        if self.waiting.len() == self.capacity {
174            return Self::sends(
175                WorkerRoute::Sends::empty(),
176                reply_to.deliver(WorkQueueOutcome::Rejected {
177                    value,
178                    reason: WorkQueueRejection::Full,
179                }),
180            );
181        }
182        self.waiting.push_back(Waiting {
183            value,
184            reply_to: reply_to.clone(),
185        });
186        Self::sends(
187            WorkerRoute::Sends::empty(),
188            reply_to.deliver(WorkQueueOutcome::Queued {
189                depth: self.waiting.len(),
190            }),
191        )
192    }
193    fn announce(
194        &mut self,
195        worker: WorkerRoute,
196    ) -> QueueActions<A, WorkerRoute::Sends, ReplyRoute::Sends> {
197        if let Some(waiting) = self.waiting.pop_front() {
198            return Self::sends(
199                worker.deliver(waiting.value),
200                waiting.reply_to.deliver(WorkQueueOutcome::Dispatched {
201                    queued: self.waiting.len(),
202                }),
203            );
204        }
205        if !self.available.contains(&worker) {
206            self.available.push_back(worker);
207        }
208        Actions::cont()
209    }
210}
211
212impl<A, T, WorkerRoute, ReplyRoute> BehaviorBase for WorkQueue<A, T, WorkerRoute, ReplyRoute>
213where
214    A: Address,
215    WorkerRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
216    ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = WorkQueueOutcome<T>>>,
217{
218    type Base = Self;
219    fn base(&self) -> &Self {
220        self
221    }
222}
223impl<A, T, WorkerRoute, ReplyRoute> behavior::Protocol for WorkQueue<A, T, WorkerRoute, ReplyRoute>
224where
225    A: Address,
226    WorkerRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
227    ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = WorkQueueOutcome<T>>>,
228{
229    type Addr = A;
230    type Msg = WorkQueueMessage<T, WorkerRoute, ReplyRoute>;
231}
232
233impl<A, T, WorkerRoute, ReplyRoute> Behavior for WorkQueue<A, T, WorkerRoute, ReplyRoute>
234where
235    A: Address,
236    WorkerRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>> + Clone + PartialEq,
237    ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = WorkQueueOutcome<T>>> + Clone,
238    WorkerRoute::Sends: behavior::SendsFor<User<A, WorkQueueMessage<T, WorkerRoute, ReplyRoute>>>,
239    ReplyRoute::Sends: behavior::SendsFor<User<A, WorkQueueMessage<T, WorkerRoute, ReplyRoute>>>,
240{
241    type Protocol = Self;
242    type Event = User<A, behavior::BehaviorMessage<Self>>;
243    type Sends = WorkQueueSends<WorkerRoute::Sends, ReplyRoute::Sends>;
244    type Ph = Never;
245    type Error = Never;
246    type Birth = NoBirths;
247    fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
248        Ok(match event.message {
249            WorkQueueMessage::Submit { value, reply_to } => self.submit(value, reply_to),
250            WorkQueueMessage::Available { worker } => self.announce(worker),
251            WorkQueueMessage::Withdraw { worker } => {
252                if let Some(index) = self
253                    .available
254                    .iter()
255                    .position(|candidate| candidate == &worker)
256                {
257                    self.available.remove(index);
258                }
259                Actions::cont()
260            }
261        })
262    }
263}
264
265#[cfg(test)]
266mod tests {
267    use super::*;
268    use crate::Activate as _;
269    use behavior::MailAddr;
270    use behavior::Recipient;
271    struct Worker;
272    struct Reply;
273    impl behavior::Protocol for Worker {
274        type Addr = MailAddr;
275        type Msg = u8;
276    }
277
278    impl behavior::Protocol for Reply {
279        type Addr = MailAddr;
280        type Msg = WorkQueueOutcome<u8>;
281    }
282    type Subject = WorkQueue<MailAddr, u8, Recipient<Worker>, Recipient<Reply>>;
283    fn worker(n: u64) -> Recipient<Worker> {
284        Recipient::global(MailAddr(n))
285    }
286    fn reply() -> Recipient<Reply> {
287        Recipient::global(MailAddr(9))
288    }
289    #[test]
290    fn availability_and_waiting_are_fifo() {
291        let mut s = (Subject::new(2)).initialize().unwrap().behavior;
292        for w in [worker(1), worker(2)] {
293            let available = s
294                .receive(MailAddr(0), WorkQueueMessage::Available { worker: w })
295                .unwrap();
296            assert!(available.sends.assignments.is_empty());
297            assert!(available.sends.outcomes.is_empty());
298            assert!(available.creates.is_empty());
299            assert_eq!(available.become_, behavior::Step::Continue);
300        }
301        for value in [10, 20, 30] {
302            let a = s
303                .receive(
304                    MailAddr(0),
305                    WorkQueueMessage::Submit {
306                        value,
307                        reply_to: reply(),
308                    },
309                )
310                .unwrap();
311            if value < 30 {
312                assert_eq!(a.sends.assignments[0].message, value);
313            }
314        }
315        assert_eq!(s.state().queued(), 1);
316        let a = s
317            .receive(
318                MailAddr(0),
319                WorkQueueMessage::Available { worker: worker(1) },
320            )
321            .unwrap();
322        assert_eq!(a.sends.assignments[0].message, 30);
323    }
324    #[test]
325    fn zero_capacity_returns_unaccepted_value() {
326        let mut s = (Subject::new(0)).initialize().unwrap().behavior;
327        let a = s
328            .receive(
329                MailAddr(0),
330                WorkQueueMessage::Submit {
331                    value: 7,
332                    reply_to: reply(),
333                },
334            )
335            .unwrap();
336        assert!(matches!(
337            a.sends.outcomes[0].message,
338            WorkQueueOutcome::Rejected {
339                value: 7,
340                reason: WorkQueueRejection::Full
341            }
342        ));
343        assert_eq!(s.state().queued(), 0);
344    }
345}