Get a buffer for a stream (if it exists). Used by consumers to poll events.
(&self, tenant_id: u64, stream_name: &str)
| 245 | |
| 246 | /// Get a buffer for a stream (if it exists). Used by consumers to poll events. |
| 247 | pub fn get_buffer(&self, tenant_id: u64, stream_name: &str) -> Option<Arc<StreamBuffer>> { |
| 248 | let key = (tenant_id, stream_name.to_string()); |
| 249 | let buffers = self.buffers.read().unwrap_or_else(|p| p.into_inner()); |
| 250 | buffers.get(&key).cloned() |
| 251 | } |
| 252 | |
| 253 | /// Remove a buffer when a stream is dropped. |
| 254 | pub fn remove_buffer(&self, tenant_id: u64, stream_name: &str) { |