steel_core/command/queue/
requests.rs1use std::collections::VecDeque;
2
3use steel_utils::locks::SyncMutex;
4
5use crate::command::sender::{CommandExecutionOwner, CommandSuggestionKey};
6
7const DEFAULT_COMMAND_REQUEST_CAPACITY: usize = 1024;
8const DEFAULT_SUGGESTION_REQUEST_CAPACITY: usize = 1024;
9
10pub(crate) const COMMAND_REQUESTS_PER_TICK: usize = 128;
12
13pub(crate) enum CommandRequest {
15 Execute {
16 owner: CommandExecutionOwner,
17 command: String,
18 },
19 Suggestions {
20 owner: CommandExecutionOwner,
21 transaction_id: i32,
22 input: String,
23 },
24}
25
26#[derive(Clone, Copy, Debug, PartialEq, Eq)]
28pub struct CommandQueueFull;
29
30enum PendingRequest<E, S> {
31 Execute(E),
32 Suggestions(S),
33}
34
35struct PendingRequestQueues<K, E, S> {
40 executions: VecDeque<E>,
41 suggestions: VecDeque<(K, S)>,
42 execution_capacity: usize,
43 suggestion_capacity: usize,
44 prefer_execution: bool,
45}
46
47impl<K: Eq, E, S> PendingRequestQueues<K, E, S> {
48 const fn new(execution_capacity: usize, suggestion_capacity: usize) -> Self {
49 Self {
50 executions: VecDeque::new(),
51 suggestions: VecDeque::new(),
52 execution_capacity,
53 suggestion_capacity,
54 prefer_execution: true,
55 }
56 }
57
58 fn submit_execution(&mut self, request: E) -> Result<(), CommandQueueFull> {
59 if self.executions.len() >= self.execution_capacity {
60 return Err(CommandQueueFull);
61 }
62 self.executions.push_back(request);
63 Ok(())
64 }
65
66 fn submit_suggestions(&mut self, key: K, request: S) -> Result<(), CommandQueueFull> {
67 if let Some((_, pending)) = self
68 .suggestions
69 .iter_mut()
70 .find(|(pending_key, _)| pending_key == &key)
71 {
72 *pending = request;
73 return Ok(());
74 }
75 if self.suggestions.len() >= self.suggestion_capacity {
76 return Err(CommandQueueFull);
77 }
78 self.suggestions.push_back((key, request));
79 Ok(())
80 }
81
82 #[cfg(test)]
83 fn pop_front(&mut self) -> Option<PendingRequest<E, S>> {
84 self.pop_front_where(|_| true)
85 }
86
87 fn pop_front_where(
88 &mut self,
89 mut execution_allowed: impl FnMut(&E) -> bool,
90 ) -> Option<PendingRequest<E, S>> {
91 if self.prefer_execution {
92 if let Some(request) = self.pop_allowed_execution(&mut execution_allowed) {
93 self.prefer_execution = false;
94 return Some(PendingRequest::Execute(request));
95 }
96 let (_, request) = self.suggestions.pop_front()?;
97 self.prefer_execution = true;
98 return Some(PendingRequest::Suggestions(request));
99 }
100
101 if let Some((_, request)) = self.suggestions.pop_front() {
102 self.prefer_execution = true;
103 return Some(PendingRequest::Suggestions(request));
104 }
105 let request = self.pop_allowed_execution(&mut execution_allowed)?;
106 self.prefer_execution = false;
107 Some(PendingRequest::Execute(request))
108 }
109
110 fn pop_allowed_execution(
111 &mut self,
112 execution_allowed: &mut impl FnMut(&E) -> bool,
113 ) -> Option<E> {
114 let index = self.executions.iter().position(execution_allowed)?;
115 self.executions.remove(index)
116 }
117
118 fn clear(&mut self) {
119 self.executions.clear();
120 self.suggestions.clear();
121 self.prefer_execution = true;
122 }
123}
124
125pub(crate) struct CommandRequestQueue {
127 queued: SyncMutex<PendingRequestQueues<CommandSuggestionKey, CommandRequest, CommandRequest>>,
128}
129
130impl CommandRequestQueue {
131 pub(crate) const fn new() -> Self {
132 Self {
133 queued: SyncMutex::new(PendingRequestQueues::new(
134 DEFAULT_COMMAND_REQUEST_CAPACITY,
135 DEFAULT_SUGGESTION_REQUEST_CAPACITY,
136 )),
137 }
138 }
139
140 pub(crate) fn submit(&self, request: CommandRequest) -> Result<(), CommandQueueFull> {
141 let mut queued = self.queued.lock();
142 match request {
143 request @ CommandRequest::Execute { .. } => queued.submit_execution(request),
144 CommandRequest::Suggestions {
145 owner,
146 transaction_id,
147 input,
148 } => queued.submit_suggestions(
149 owner.suggestion_key(),
150 CommandRequest::Suggestions {
151 owner,
152 transaction_id,
153 input,
154 },
155 ),
156 }
157 }
158
159 pub(crate) fn pop_front_runnable(
160 &self,
161 mut execution_allowed: impl FnMut(&CommandExecutionOwner) -> bool,
162 ) -> Option<CommandRequest> {
163 let request = self.queued.lock().pop_front_where(|request| {
164 let CommandRequest::Execute { owner, .. } = request else {
165 return false;
166 };
167 execution_allowed(owner)
168 })?;
169 match request {
170 PendingRequest::Execute(request) | PendingRequest::Suggestions(request) => {
171 Some(request)
172 }
173 }
174 }
175
176 pub(crate) fn clear(&self) {
177 self.queued.lock().clear();
178 }
179}
180
181impl Default for CommandRequestQueue {
182 fn default() -> Self {
183 Self::new()
184 }
185}
186
187#[cfg(test)]
188mod tests {
189 use super::{CommandQueueFull, PendingRequest, PendingRequestQueues};
190
191 fn queue_with_capacity(
192 execution_capacity: usize,
193 suggestion_capacity: usize,
194 ) -> PendingRequestQueues<u8, &'static str, &'static str> {
195 PendingRequestQueues::new(execution_capacity, suggestion_capacity)
196 }
197
198 #[test]
199 fn executions_are_dequeued_in_submission_order() {
200 let mut queue = queue_with_capacity(3, 3);
201
202 assert!(queue.submit_execution("first").is_ok());
203 assert!(queue.submit_execution("second").is_ok());
204
205 assert!(matches!(
206 queue.pop_front(),
207 Some(PendingRequest::Execute("first"))
208 ));
209 assert!(matches!(
210 queue.pop_front(),
211 Some(PendingRequest::Execute("second"))
212 ));
213 assert!(queue.pop_front().is_none());
214 }
215
216 #[test]
217 fn blocked_execution_sources_are_skipped_without_reordering_their_requests() {
218 let mut queue = queue_with_capacity(3, 1);
219 assert!(queue.submit_execution("blocked first").is_ok());
220 assert!(queue.submit_execution("ready").is_ok());
221 assert!(queue.submit_execution("blocked second").is_ok());
222
223 assert!(matches!(
224 queue.pop_front_where(|request| !request.starts_with("blocked")),
225 Some(PendingRequest::Execute("ready"))
226 ));
227 assert!(matches!(
228 queue.pop_front(),
229 Some(PendingRequest::Execute("blocked first"))
230 ));
231 assert!(matches!(
232 queue.pop_front(),
233 Some(PendingRequest::Execute("blocked second"))
234 ));
235 }
236
237 #[test]
238 fn full_execution_queue_rejects_without_dropping_pending_requests() {
239 let mut queue = queue_with_capacity(2, 2);
240
241 assert!(queue.submit_execution("first").is_ok());
242 assert!(queue.submit_execution("second").is_ok());
243 assert_eq!(queue.submit_execution("third"), Err(CommandQueueFull));
244
245 assert!(matches!(
246 queue.pop_front(),
247 Some(PendingRequest::Execute("first"))
248 ));
249 assert!(matches!(
250 queue.pop_front(),
251 Some(PendingRequest::Execute("second"))
252 ));
253 assert!(queue.pop_front().is_none());
254 }
255
256 #[test]
257 fn suggestion_capacity_cannot_starve_execution_capacity() {
258 let mut queue = queue_with_capacity(2, 2);
259
260 assert!(queue.submit_suggestions(1, "first suggestion").is_ok());
261 assert!(queue.submit_suggestions(2, "second suggestion").is_ok());
262 assert_eq!(
263 queue.submit_suggestions(3, "rejected suggestion"),
264 Err(CommandQueueFull)
265 );
266
267 assert!(queue.submit_execution("command").is_ok());
268 assert!(matches!(
269 queue.pop_front(),
270 Some(PendingRequest::Execute("command"))
271 ));
272 }
273
274 #[test]
275 fn suggestions_from_one_sender_are_coalesced() {
276 let mut queue = queue_with_capacity(1, 1);
277
278 assert!(queue.submit_suggestions(1, "old").is_ok());
279 assert!(queue.submit_suggestions(1, "latest").is_ok());
280 assert!(matches!(
281 queue.pop_front(),
282 Some(PendingRequest::Suggestions("latest"))
283 ));
284 assert!(queue.pop_front().is_none());
285 }
286
287 #[test]
288 fn busy_queues_are_drained_fairly() {
289 let mut queue = queue_with_capacity(2, 2);
290 assert!(queue.submit_execution("first command").is_ok());
291 assert!(queue.submit_execution("second command").is_ok());
292 assert!(queue.submit_suggestions(1, "first suggestion").is_ok());
293 assert!(queue.submit_suggestions(2, "second suggestion").is_ok());
294
295 assert!(matches!(
296 queue.pop_front(),
297 Some(PendingRequest::Execute(_))
298 ));
299 assert!(matches!(
300 queue.pop_front(),
301 Some(PendingRequest::Suggestions(_))
302 ));
303 assert!(matches!(
304 queue.pop_front(),
305 Some(PendingRequest::Execute(_))
306 ));
307 assert!(matches!(
308 queue.pop_front(),
309 Some(PendingRequest::Suggestions(_))
310 ));
311 }
312
313 #[test]
314 fn clear_discards_all_pending_requests() {
315 let mut queue = queue_with_capacity(2, 2);
316
317 assert!(queue.submit_execution("command").is_ok());
318 assert!(queue.submit_suggestions(1, "suggestion").is_ok());
319 queue.clear();
320
321 assert!(queue.pop_front().is_none());
322 }
323}