MCPcopy Create free account
hub / github.com/NodeDB-Lab/nodedb / compact

Method compact

nodedb/src/event/cdc/buffer.rs:160–203  ·  view source on GitHub ↗

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)

Source from the content-addressed store, hash-verified

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()

Callers 15

build_csr_for_collectionFunction · 0.45
rebuild_csr_threadFunction · 0.45
run_compactionMethod · 0.45
social_graphFunction · 0.45
diameter_pathFunction · 0.45
diameter_triangleFunction · 0.45
diameter_approximateFunction · 0.45
diameter_single_nodeFunction · 0.45
label_prop_triangleFunction · 0.45

Calls 10

nowFunction · 0.85
extract_key_valueFunction · 0.85
duration_sinceMethod · 0.80
writeMethod · 0.45
lenMethod · 0.45
as_millisMethod · 0.45
iterMethod · 0.45
insertMethod · 0.45
getMethod · 0.45
drainMethod · 0.45

Tested by 15

diameter_pathFunction · 0.36
diameter_triangleFunction · 0.36
diameter_approximateFunction · 0.36
diameter_single_nodeFunction · 0.36
label_prop_triangleFunction · 0.36
label_prop_isolated_nodeFunction · 0.36
label_prop_deterministicFunction · 0.36
closeness_pathFunction · 0.36
closeness_complete_graphFunction · 0.36
closeness_disconnectedFunction · 0.36