Remove a buffer when a stream is dropped.
(&self, tenant_id: u64, stream_name: &str)
| 252 | |
| 253 | /// Remove a buffer when a stream is dropped. |
| 254 | pub fn remove_buffer(&self, tenant_id: u64, stream_name: &str) { |
| 255 | let key = (tenant_id, stream_name.to_string()); |
| 256 | let mut buffers = self.buffers.write().unwrap_or_else(|p| p.into_inner()); |
| 257 | buffers.remove(&key); |
| 258 | self.lag_warner.remove_stream(tenant_id, stream_name); |
| 259 | } |
| 260 | |
| 261 | /// Snapshot of all buffer stats (for SHOW CHANGE STREAMS). |
| 262 | pub fn buffer_stats(&self) -> Vec<BufferStats> { |