1use std::collections::VecDeque;
4
5use behavior::{
6 Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, Protocol,
7 SendEffects, User,
8};
9
10use crate::DeliveryRoute;
11
12pub struct WorkQueueState<WorkerRoute> {
14 pub capacity: usize,
16 available: Vec<WorkerRoute>,
17 queued: usize,
18}
19
20impl<WorkerRoute> WorkQueueState<WorkerRoute> {
21 #[must_use]
23 pub fn available(&self) -> &[WorkerRoute] {
24 &self.available
25 }
26
27 #[must_use]
29 pub fn queued(&self) -> usize {
30 self.queued
31 }
32}
33
34#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
36pub enum WorkQueueRejection {
37 Full,
39}
40
41#[derive(Debug, PartialEq, Eq)]
43pub enum WorkQueueOutcome<T> {
44 Queued {
46 depth: usize,
48 },
49 Dispatched {
51 queued: usize,
53 },
54 Rejected {
56 value: T,
58 reason: WorkQueueRejection,
60 },
61}
62
63pub enum WorkQueueMessage<T, WorkerRoute, ReplyRoute> {
65 Submit {
67 value: T,
69 reply_to: ReplyRoute,
71 },
72 Available {
74 worker: WorkerRoute,
76 },
77 Withdraw {
79 worker: WorkerRoute,
81 },
82}
83
84#[derive(behavior_macros::SendProduct)]
86pub struct WorkQueueSends<Assignments, OutcomeSends> {
87 pub assignments: Assignments,
89 pub outcomes: OutcomeSends,
91}
92
93struct Waiting<T, Route> {
94 value: T,
95 reply_to: Route,
96}
97
98pub 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 #[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 #[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}