MCPcopy Create free account
hub / github.com/Botloader/botloader / handle_worker_msg

Method handle_worker_msg

components/scheduler/src/vm_session.rs:210–271  ·  view source on GitHub ↗

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>)

Source from the content-addressed store, hash-verified

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);

Callers 3

handle_next_actionMethod · 0.80
shutdownMethod · 0.80
destroy_brokenMethod · 0.80

Calls 7

elapsed_millisFunction · 0.85
handle_shutdown_msgMethod · 0.80
removeMethod · 0.80
recordMethod · 0.80
log_rawMethod · 0.80
handle_metricMethod · 0.80
sendMethod · 0.45

Tested by

no test coverage detected