Handles a message received from the worker (`None` meaning the channel died), returning the events the guild handler needs to act on.
(&mut self, msg: Option<WorkerMessage>)
| 208 | /// Handles a message received from the worker (`None` meaning the channel |
| 209 | /// died), returning the events the guild handler needs to act on. |
| 210 | pub fn handle_worker_msg(&mut self, msg: Option<WorkerMessage>) -> Vec<SessionEvent> { |
| 211 | let Some(msg) = msg else { |
| 212 | return vec![SessionEvent::Broken]; |
| 213 | }; |
| 214 | |
| 215 | match msg { |
| 216 | WorkerMessage::Shutdown(shutdown) => self.handle_shutdown_msg(shutdown), |
| 217 | WorkerMessage::Ack(id) => { |
| 218 | let Some(item) = self.pending_acks.remove(&id) else { |
| 219 | return Vec::new(); |
| 220 | }; |
| 221 | |
| 222 | if let Some(meta) = &item.dispatched { |
| 223 | metrics::histogram!( |
| 224 | "dispatch_event_acked_latency", |
| 225 | "event_source" => crate::dispatch_metrics::event_source_label(meta.source), |
| 226 | "vm_cold" => if meta.vm_cold { "true" } else { "false" }, |
| 227 | ) |
| 228 | .record(crate::dispatch_metrics::elapsed_millis(meta.source_timestamp)); |
| 229 | } |
| 230 | |
| 231 | match item.kind { |
| 232 | PendingAckType::Dispatch(Some(resp)) => { |
| 233 | let _ = resp.send(()); |
| 234 | Vec::new() |
| 235 | } |
| 236 | PendingAckType::Dispatch(None) => Vec::new(), |
| 237 | PendingAckType::ScheduledTask(task_id) => { |
| 238 | vec![SessionEvent::TaskAcked(task_id)] |
| 239 | } |
| 240 | PendingAckType::IntervalTimer(timer) => vec![SessionEvent::TimerAcked(timer)], |
| 241 | PendingAckType::Restart => Vec::new(), |
| 242 | } |
| 243 | } |
| 244 | WorkerMessage::ScriptStarted(meta) => vec![SessionEvent::ScriptStarted(meta)], |
| 245 | WorkerMessage::ScriptsInit => todo!(), |
| 246 | WorkerMessage::NonePending => { |
| 247 | // the vm's event loop has drained, script init is done and |
| 248 | // future dispatches run against a warm vm |
| 249 | self.vm_cold = false; |
| 250 | |
| 251 | if self.pending_acks.is_empty() { |
| 252 | vec![SessionEvent::VmIdle] |
| 253 | } else { |
| 254 | Vec::new() |
| 255 | } |
| 256 | } |
| 257 | WorkerMessage::TaskScheduled => vec![SessionEvent::TaskScheduled], |
| 258 | WorkerMessage::GuildLog(entry) => { |
| 259 | self.logger.log_raw(entry); |
| 260 | Vec::new() |
| 261 | } |
| 262 | WorkerMessage::Hello(_) => { |
| 263 | // handled when the connection is established, not applicable here |
| 264 | unreachable!(); |
| 265 | } |
| 266 | WorkerMessage::Metric(name, m, labels) => { |
| 267 | self.handle_metric(name, m, labels); |
no test coverage detected