(
&mut self,
handler: Box<dyn QueueHandler>,
queue_builder: QueueBuilder,
)
| 2198 | } |
| 2199 | |
| 2200 | async fn add_queue( |
| 2201 | &mut self, |
| 2202 | handler: Box<dyn QueueHandler>, |
| 2203 | queue_builder: QueueBuilder, |
| 2204 | ) -> QueueId { |
| 2205 | let (mut params, rate_limiter) = queue_builder.build(); |
| 2206 | params.name = Some("Queue".to_string()); |
| 2207 | |
| 2208 | let (token, rx) = ResponseToken::new(); |
| 2209 | handle_message( |
| 2210 | &mut self.state, |
| 2211 | &self.senders.events, |
| 2212 | AutoAllocMessage::AddQueue { |
| 2213 | server_directory: PathBuf::from("test"), |
| 2214 | params, |
| 2215 | queue_id: None, |
| 2216 | worker_resources: None, |
| 2217 | response: token, |
| 2218 | }, |
| 2219 | ) |
| 2220 | .await; |
| 2221 | let queue_id = rx.await.unwrap().unwrap(); |
| 2222 | |
| 2223 | // A bit hacky, but easier than creating a mock for an external service, as we do in |
| 2224 | // Python. |
| 2225 | let queue = self.state.get_queue_mut(queue_id).unwrap(); |
| 2226 | queue.set_handler(handler); |
| 2227 | *queue.limiter_mut() = rate_limiter; |
| 2228 | |
| 2229 | queue_id |
| 2230 | } |
| 2231 | |
| 2232 | async fn try_submit(&mut self) { |
| 2233 | perform_submits(&mut self.state, &self.senders) |
no test coverage detected