(
&mut self,
t: String,
data: serde_json::Value,
ack: PendingAckType,
source: EventSource,
ts: DateTime<Utc>,
)
| 621 | } |
| 622 | |
| 623 | async fn dispatch_worker_evt( |
| 624 | &mut self, |
| 625 | t: String, |
| 626 | data: serde_json::Value, |
| 627 | ack: PendingAckType, |
| 628 | source: EventSource, |
| 629 | ts: DateTime<Utc>, |
| 630 | ) { |
| 631 | if self.scripts.is_empty() { |
| 632 | return; |
| 633 | } |
| 634 | |
| 635 | let mut ack = ack; |
| 636 | loop { |
| 637 | self.ensure_session().await; |
| 638 | |
| 639 | let session = self.session.as_mut().expect("ensure_session claims a worker"); |
| 640 | match session.dispatch(t.clone(), data.clone(), ack, source, ts) { |
| 641 | Ok(()) => { |
| 642 | crate::dispatch_metrics::record_stage( |
| 643 | "worker_sent", |
| 644 | crate::dispatch_metrics::event_source_label(source), |
| 645 | ts, |
| 646 | ); |
| 647 | return; |
| 648 | } |
| 649 | Err(returned_ack) => { |
| 650 | ack = returned_ack; |
| 651 | error!("worker died while trying to dispatch event, retrying in a second"); |
| 652 | self.discard_broken_session().await; |
| 653 | tokio::time::sleep(Duration::from_secs(1)).await; |
| 654 | } |
| 655 | } |
| 656 | } |
| 657 | } |
| 658 | |
| 659 | #[instrument(skip(self), fields(guild_id = self.guild_id.get()))] |
| 660 | pub(crate) async fn shutdown(&mut self) { |
no test coverage detected