Clear all events for a session
(&self, session_id: &str)
| 179 | |
| 180 | /// Clear all events for a session |
| 181 | pub async fn clear_session(&self, session_id: &str) -> EventBusResult<()> { |
| 182 | // Remove all events for this session from the queue |
| 183 | let queue_len = { |
| 184 | let mut queue = self.queue.lock().await; |
| 185 | let mut new_queue = BinaryHeap::new(); |
| 186 | |
| 187 | while let Some(std::cmp::Reverse(envelope)) = queue.pop() { |
| 188 | if envelope.event.session_id() != Some(session_id) { |
| 189 | new_queue.push(std::cmp::Reverse(envelope)); |
| 190 | } |
| 191 | } |
| 192 | |
| 193 | *queue = new_queue; |
| 194 | queue.len() // Get size before releasing queue lock |
| 195 | }; |
| 196 | |
| 197 | // Update statistics: use the size obtained earlier |
| 198 | { |
| 199 | let mut stats = self.stats.lock().await; |
| 200 | stats.pending_events = queue_len; |
| 201 | } |
| 202 | |
| 203 | debug!("Cleared all events for session: session_id={}", session_id); |
| 204 | |
| 205 | Ok(()) |
| 206 | } |
| 207 | |
| 208 | /// Get queue statistics |
| 209 | pub async fn stats(&self) -> QueueStats { |
nothing calls this directly
no test coverage detected