Get or create a buffer for a stream.
(
&self,
tenant_id: u64,
stream_name: &str,
retention: &super::stream_def::RetentionConfig,
)
| 205 | |
| 206 | /// Get or create a buffer for a stream. |
| 207 | fn get_or_create_buffer( |
| 208 | &self, |
| 209 | tenant_id: u64, |
| 210 | stream_name: &str, |
| 211 | retention: &super::stream_def::RetentionConfig, |
| 212 | ) -> Arc<StreamBuffer> { |
| 213 | let key = (tenant_id, stream_name.to_string()); |
| 214 | |
| 215 | // Fast path: read lock. |
| 216 | { |
| 217 | let buffers = self.buffers.read().unwrap_or_else(|p| p.into_inner()); |
| 218 | if let Some(buf) = buffers.get(&key) { |
| 219 | return Arc::clone(buf); |
| 220 | } |
| 221 | } |
| 222 | |
| 223 | // Slow path: write lock + create. |
| 224 | let mut buffers = self.buffers.write().unwrap_or_else(|p| p.into_inner()); |
| 225 | buffers |
| 226 | .entry(key) |
| 227 | .or_insert_with(|| { |
| 228 | Arc::new(StreamBuffer::new( |
| 229 | stream_name.to_string(), |
| 230 | retention.clone(), |
| 231 | )) |
| 232 | }) |
| 233 | .clone() |
| 234 | } |
| 235 | |
| 236 | /// Ensure a buffer exists for a given key (stream or topic). Creates if missing. |
| 237 | pub fn ensure_buffer( |