Get or create a StreamBuffer for a topic.
(
state: &SharedState,
tenant_id: u64,
topic_name: &str,
retention: &RetentionConfig,
)
| 92 | |
| 93 | /// Get or create a StreamBuffer for a topic. |
| 94 | fn get_or_create_topic_buffer( |
| 95 | state: &SharedState, |
| 96 | tenant_id: u64, |
| 97 | topic_name: &str, |
| 98 | retention: &RetentionConfig, |
| 99 | ) -> Arc<StreamBuffer> { |
| 100 | // Topics use the CdcRouter's buffer pool with a "topic:" prefix |
| 101 | // to avoid name collisions with change streams. |
| 102 | let buffer_key = format!("topic:{topic_name}"); |
| 103 | |
| 104 | if let Some(buf) = state.cdc_router.get_buffer(tenant_id, &buffer_key) { |
| 105 | return buf; |
| 106 | } |
| 107 | |
| 108 | // Create a new buffer. Use the router's internal mechanism. |
| 109 | // Since CdcRouter.get_or_create_buffer is private, we route through |
| 110 | // a dummy event to force buffer creation, then return it. |
| 111 | // Instead, let's add a public create method to CdcRouter. |
| 112 | // For now, use the public get_buffer after forcing creation. |
| 113 | // |
| 114 | // Actually, we can just create the buffer directly and register it. |
| 115 | state |
| 116 | .cdc_router |
| 117 | .ensure_buffer(tenant_id, &buffer_key, retention) |
| 118 | } |
| 119 | |
| 120 | /// Determine the home node for a topic. |
| 121 | /// |
no test coverage detected