Compact the buffer: deduplicate by key field, keeping only the latest event per key value. DELETE events are retained as tombstones until they exceed `tombstone_grace_secs` age, then removed.
(&self, key_field: &str, tombstone_grace_secs: u64)
| 158 | /// event per key value. DELETE events are retained as tombstones until |
| 159 | /// they exceed `tombstone_grace_secs` age, then removed. |
| 160 | pub fn compact(&self, key_field: &str, tombstone_grace_secs: u64) -> u32 { |
| 161 | let mut events = self.events.write().unwrap_or_else(|p| { |
| 162 | tracing::warn!(stream = %self.name, "StreamBuffer RwLock poisoned during compact, recovering"); |
| 163 | p.into_inner() |
| 164 | }); |
| 165 | let before = events.len(); |
| 166 | |
| 167 | let now_ms = SystemTime::now() |
| 168 | .duration_since(UNIX_EPOCH) |
| 169 | .unwrap_or_default() |
| 170 | .as_millis() as u64; |
| 171 | let tombstone_cutoff_ms = now_ms.saturating_sub(tombstone_grace_secs * 1000); |
| 172 | |
| 173 | let mut latest: std::collections::HashMap<String, usize> = std::collections::HashMap::new(); |
| 174 | for (idx, event) in events.iter().enumerate() { |
| 175 | let key_value = extract_key_value(event, key_field); |
| 176 | latest.insert(key_value, idx); |
| 177 | } |
| 178 | |
| 179 | let mut keep = vec![false; events.len()]; |
| 180 | for (idx, event) in events.iter().enumerate() { |
| 181 | let key_value = extract_key_value(event, key_field); |
| 182 | let is_latest = latest.get(&key_value) == Some(&idx); |
| 183 | let is_tombstone = event.op == "DELETE"; |
| 184 | if is_latest && !(is_tombstone && event.event_time < tombstone_cutoff_ms) { |
| 185 | keep[idx] = true; |
| 186 | } |
| 187 | } |
| 188 | |
| 189 | let mut new_events = VecDeque::with_capacity(events.len()); |
| 190 | for (idx, event) in events.drain(..).enumerate() { |
| 191 | if keep[idx] { |
| 192 | new_events.push_back(event); |
| 193 | } |
| 194 | } |
| 195 | *events = new_events; |
| 196 | |
| 197 | let removed = (before - events.len()) as u32; |
| 198 | if removed > 0 { |
| 199 | self.total_evicted |
| 200 | .fetch_add(removed as u64, std::sync::atomic::Ordering::Relaxed); |
| 201 | } |
| 202 | removed |
| 203 | } |
| 204 | |
| 205 | pub fn len(&self) -> usize { |
| 206 | self.events.read().unwrap_or_else(|p| p.into_inner()).len() |