Enqueue event
(
&self,
event: AgenticEvent,
priority: Option<EventPriority>,
)
| 74 | |
| 75 | /// Enqueue event |
| 76 | pub async fn enqueue( |
| 77 | &self, |
| 78 | event: AgenticEvent, |
| 79 | priority: Option<EventPriority>, |
| 80 | ) -> EventBusResult<String> { |
| 81 | let priority = priority.unwrap_or_else(|| event.default_priority()); |
| 82 | let envelope = EventEnvelope::new(event, priority); |
| 83 | let event_id = envelope.id.clone(); |
| 84 | |
| 85 | // Check queue size |
| 86 | { |
| 87 | let queue = self.queue.lock().await; |
| 88 | if queue.len() >= self.config.max_queue_size { |
| 89 | warn!("Event queue full, dropping event: event_id={}", event_id); |
| 90 | return Ok(event_id); |
| 91 | } |
| 92 | } |
| 93 | |
| 94 | // Add to queue |
| 95 | { |
| 96 | let mut queue = self.queue.lock().await; |
| 97 | queue.push(std::cmp::Reverse(envelope.clone())); |
| 98 | } |
| 99 | |
| 100 | let _ = self.broadcast_tx.send(envelope); |
| 101 | |
| 102 | // Update statistics: get queue size first, then update statistics (avoid getting queue lock while holding stats lock) |
| 103 | let queue_len = self.queue.lock().await.len(); |
| 104 | { |
| 105 | let mut stats = self.stats.lock().await; |
| 106 | stats.total_enqueued += 1; |
| 107 | stats.pending_events = queue_len; |
| 108 | } |
| 109 | |
| 110 | // Notify waiting consumers |
| 111 | self.notify.notify_one(); |
| 112 | |
| 113 | trace!( |
| 114 | "Event enqueued: event_id={}, priority={:?}", |
| 115 | event_id, |
| 116 | priority |
| 117 | ); |
| 118 | |
| 119 | Ok(event_id) |
| 120 | } |
| 121 | |
| 122 | /// Dequeue batch of events |
| 123 | pub async fn dequeue_batch(&self, max_size: usize) -> Vec<EventEnvelope> { |
nothing calls this directly
no test coverage detected