1use crate::DexEvent;
8use std::collections::{BTreeMap, HashMap, HashSet};
9use tokio::time::Instant;
10
11#[derive(Default)]
15pub struct SlotBuffer {
16 slots: BTreeMap<u64, Vec<(u64, DexEvent)>>,
18 current_slot: u64,
20 last_flush_time: Option<Instant>,
22 streaming_watermarks: HashMap<u64, u64>,
24 streaming_pending_indexes: HashSet<u64>,
25 ordered_watermark: Option<(u64, u64)>,
26 ordered_late_events: u64,
27}
28
29impl SlotBuffer {
30 #[inline]
31 pub fn new() -> Self {
32 Self {
33 slots: BTreeMap::new(),
34 current_slot: 0,
35 last_flush_time: Some(Instant::now()),
36 streaming_watermarks: HashMap::new(),
37 streaming_pending_indexes: HashSet::new(),
38 ordered_watermark: None,
39 ordered_late_events: 0,
40 }
41 }
42
43 #[inline]
45 pub fn push(&mut self, slot: u64, tx_index: u64, event: DexEvent) {
46 if slot < self.current_slot || self.ordered_watermark.is_some_and(|last| (slot, tx_index) <= last) {
47 self.ordered_late_events = self.ordered_late_events.saturating_add(1);
48 let dropped = self.ordered_late_events;
50 if dropped <= 10 || dropped.is_power_of_two() {
51 log::warn!("Ordered continuity break: dropped late event ({slot},{tx_index}); total={dropped}");
52 }
53 return;
54 }
55 self.slots.entry(slot).or_default().push((tx_index, event));
56 if slot > self.current_slot {
57 self.current_slot = slot;
58 }
59 }
60
61 pub fn flush_before(&mut self, current_slot: u64) -> Vec<DexEvent> {
63 self.current_slot = self.current_slot.max(current_slot);
64 let slots_to_flush: Vec<u64> =
65 self.slots.keys().filter(|&&s| s < current_slot).copied().collect();
66
67 let mut result = Vec::with_capacity(slots_to_flush.len() * 4);
68 for slot in slots_to_flush {
69 if let Some(mut events) = self.slots.remove(&slot) {
70 events.sort_by_key(|(idx, _)| *idx);
71 if let Some((index, _)) = events.last() {
72 self.ordered_watermark = Some((slot, *index));
73 }
74 result.extend(events.into_iter().map(|(_, e)| e));
75 }
76 }
77
78 if !result.is_empty() {
79 self.last_flush_time = Some(Instant::now());
80 }
81 result
82 }
83
84 pub fn flush_all(&mut self) -> Vec<DexEvent> {
86 let all_slots: Vec<u64> = self.slots.keys().copied().collect();
87 let mut result = Vec::with_capacity(all_slots.len() * 4);
88
89 for slot in all_slots {
90 if let Some(mut events) = self.slots.remove(&slot) {
91 events.sort_by_key(|(idx, _)| *idx);
92 if let Some((index, _)) = events.last() {
93 self.ordered_watermark = Some((slot, *index));
94 }
95 result.extend(events.into_iter().map(|(_, e)| e));
96 }
97 }
98
99 if !result.is_empty() {
100 self.last_flush_time = Some(Instant::now());
101 }
102 result
103 }
104
105 pub fn ordered_late_events(&self) -> u64 { self.ordered_late_events }
107
108 #[inline]
110 pub fn should_timeout(&self, timeout_ms: u64) -> bool {
111 self.last_flush_time
112 .map(|t| !self.slots.is_empty() && t.elapsed().as_millis() as u64 > timeout_ms)
113 .unwrap_or(false)
114 }
115
116 pub fn push_streaming(&mut self, slot: u64, tx_index: u64, event: DexEvent) -> Vec<DexEvent> {
119 self.push_streaming_batch(slot, tx_index, [event])
120 }
121
122 pub fn push_streaming_batch(
126 &mut self,
127 slot: u64,
128 tx_index: u64,
129 events: impl IntoIterator<Item = DexEvent>,
130 ) -> Vec<DexEvent> {
131 if slot < self.current_slot || tx_index == u64::MAX {
133 return Vec::new();
134 }
135 let mut events = events.into_iter().peekable();
136 if events.peek().is_none() {
137 return Vec::new();
138 }
139 let mut result = Vec::new();
140 if slot > self.current_slot {
141 result = self.flush_before(slot);
142 self.streaming_watermarks.retain(|old_slot, _| *old_slot >= slot);
144 self.streaming_pending_indexes.clear();
145 self.current_slot = slot;
146 }
147
148 let next_expected = *self.streaming_watermarks.get(&slot).unwrap_or(&0);
149 if tx_index == next_expected {
150 result.extend(events);
151 let mut watermark = next_expected + 1;
152 if let Some(buffered) = self.slots.get_mut(&slot) {
153 buffered.sort_by_key(|(idx, _)| *idx);
154 let mut released = 0;
155 while released < buffered.len() && buffered[released].0 == watermark {
156 while released < buffered.len() && buffered[released].0 == watermark {
157 released += 1;
158 }
159 self.streaming_pending_indexes.remove(&watermark);
160 watermark += 1;
161 }
162 result.extend(buffered.drain(..released).map(|(_, event)| event));
163 if buffered.is_empty() {
164 self.slots.remove(&slot);
165 }
166 }
167 self.streaming_watermarks.insert(slot, watermark);
168 } else if tx_index > next_expected && self.streaming_pending_indexes.insert(tx_index) {
169 self.slots.entry(slot).or_default().extend(events.map(|event| (tx_index, event)));
170 }
171 if !result.is_empty() {
172 self.last_flush_time = Some(Instant::now());
173 }
174 result
175 }
176
177 pub fn flush_streaming_timeout(&mut self) -> Vec<DexEvent> {
179 let mut result = Vec::new();
180 for (slot, mut events) in std::mem::take(&mut self.slots) {
181 events.sort_by_key(|(idx, _)| *idx);
182 if slot == self.current_slot {
183 if let Some((index, _)) = events.last() {
184 let next = index.saturating_add(1);
185 let watermark = self.streaming_watermarks.entry(slot).or_default();
186 *watermark = (*watermark).max(next);
187 }
188 }
189 result.extend(events.into_iter().map(|(_, e)| e));
190 }
191 self.streaming_pending_indexes.clear();
192 self.streaming_watermarks.retain(|slot, _| *slot == self.current_slot);
193 if !result.is_empty() {
194 self.last_flush_time = Some(Instant::now());
195 }
196 result
197 }
198}
199
200pub struct MicroBatchBuffer {
204 events: Vec<(u64, u64, DexEvent)>,
206 window_start_us: i64,
208}
209
210impl MicroBatchBuffer {
211 #[inline]
212 pub fn new() -> Self {
213 Self { events: Vec::with_capacity(64), window_start_us: 0 }
214 }
215
216 #[inline]
218 pub fn push(
219 &mut self,
220 slot: u64,
221 tx_index: u64,
222 event: DexEvent,
223 now_us: i64,
224 window_us: u64,
225 ) -> bool {
226 if self.events.is_empty() {
227 self.window_start_us = now_us;
228 }
229 self.events.push((slot, tx_index, event));
230 (now_us - self.window_start_us) as u64 >= window_us
231 }
232
233 #[inline]
235 pub fn flush(&mut self) -> Vec<DexEvent> {
236 if self.events.is_empty() {
237 return Vec::new();
238 }
239
240 self.events.sort_by_key(|(slot, tx_index, _)| (*slot, *tx_index));
242
243 let mut result = Vec::with_capacity(self.events.len());
244 result.extend(self.events.drain(..).map(|(_, _, event)| event));
245
246 self.window_start_us = 0;
247 result
248 }
249
250 #[inline]
252 pub fn should_flush(&self, now_us: i64, window_us: u64) -> bool {
253 !self.events.is_empty() && (now_us - self.window_start_us) as u64 >= window_us
254 }
255
256 #[inline]
257 pub fn is_empty(&self) -> bool {
258 self.events.is_empty()
259 }
260}
261
262impl Default for MicroBatchBuffer {
263 fn default() -> Self {
264 Self::new()
265 }
266}
267
268#[cfg(test)]
269mod tests {
270 use super::*;
271 use crate::core::events::{BlockMetaEvent, EventMetadata};
272
273 fn event(id: u64) -> DexEvent {
274 DexEvent::BlockMeta(BlockMetaEvent {
275 metadata: EventMetadata { slot: id, ..Default::default() },
276 })
277 }
278 fn ids(events: Vec<DexEvent>) -> Vec<u64> {
279 events.into_iter().map(|event| event.metadata().slot).collect()
280 }
281
282 #[test]
283 fn streaming_releases_complete_immediate_and_buffered_transactions() {
284 let mut buffer = SlotBuffer::new();
285 assert!(buffer.push_streaming_batch(42, 1, [event(20), event(21)]).is_empty());
286 assert_eq!(
287 ids(buffer.push_streaming_batch(42, 0, [event(10), event(11)])),
288 [10, 11, 20, 21]
289 );
290 assert_eq!(ids(buffer.push_streaming_batch(42, 2, [event(30), event(31)])), [30, 31]);
291 assert!(buffer.push_streaming_batch(42, 0, [event(99)]).is_empty());
292 assert!(buffer.slots.is_empty());
293 }
294
295 #[test]
296 fn streaming_bounds_watermarks_and_rejects_late_or_replayed_batches() {
297 let mut buffer = SlotBuffer::new();
298 for slot in 1..=10_000 {
299 assert_eq!(buffer.push_streaming_batch(slot, 0, [event(slot)]).len(), 1);
300 assert_eq!(buffer.streaming_watermarks.len(), 1);
301 }
302 assert!(buffer.flush_streaming_timeout().is_empty());
303 assert_eq!(buffer.streaming_watermarks.len(), 1);
304 for old_slot in 1..10_000 {
305 assert!(buffer.push_streaming_batch(old_slot, 0, [event(99)]).is_empty());
306 }
307 assert_eq!(buffer.streaming_watermarks.len(), 1);
308 assert!(buffer.push_streaming_batch(10_000, 0, [event(99)]).is_empty());
309 }
310
311 #[test]
312 fn streaming_slot_and_timeout_flush_preserve_event_order_within_a_transaction() {
313 let mut buffer = SlotBuffer::new();
314 assert!(buffer.push_streaming_batch(1, 5, [event(10), event(11)]).is_empty());
315 assert_eq!(
316 ids(buffer.push_streaming_batch(2, 0, [event(20), event(21)])),
317 [10, 11, 20, 21]
318 );
319 assert!(buffer.push_streaming_batch(2, 3, [event(30), event(31)]).is_empty());
320 assert_eq!(ids(buffer.flush_streaming_timeout()), [30, 31]);
321 assert_eq!(buffer.streaming_watermarks.len(), 1);
322 assert!(buffer.push_streaming_batch(2, 3, [event(99)]).is_empty());
323 assert!(buffer.push_streaming_batch(2, 0, [event(99)]).is_empty());
324 assert_eq!(ids(buffer.push_streaming_batch(2, 4, [event(40)])), [40]);
325 }
326
327 #[test]
328 fn repeated_pending_and_old_indexes_remain_bounded() {
329 let mut buffer = SlotBuffer::new();
330 for _ in 0..10_000 {
331 assert!(buffer.push_streaming_batch(42, 2, [event(2)]).is_empty());
332 }
333 assert_eq!(buffer.slots[&42].len(), 1);
334 assert_eq!(buffer.streaming_pending_indexes.len(), 1);
335 assert_eq!(ids(buffer.flush_streaming_timeout()), [2]);
336 assert!(buffer.slots.is_empty());
337 assert!(buffer.streaming_pending_indexes.is_empty());
338 assert!(buffer.push_streaming_batch(42, u64::MAX, [event(99)]).is_empty());
339 assert!(buffer.flush_streaming_timeout().is_empty());
340 assert_eq!(ids(buffer.push_streaming_batch(43, 0, [event(3)])), [3]);
341 for _ in 0..10_000 {
342 assert!(buffer.push_streaming_batch(42, 2, [event(2)]).is_empty());
343 }
344 assert_eq!(buffer.streaming_watermarks.len(), 1);
345 assert!(buffer.slots.is_empty());
346 }
347
348 #[test]
349 fn duplicate_buffered_batch_is_not_emitted_twice() {
350 let mut buffer = SlotBuffer::new();
351 assert!(buffer.push_streaming_batch(42, 1, [event(20), event(21)]).is_empty());
352 assert!(buffer.push_streaming_batch(42, 1, [event(20), event(21)]).is_empty());
353 assert_eq!(ids(buffer.push_streaming_batch(42, 0, [event(10)])), [10, 20, 21]);
354 }
355 #[test]
356 fn ordered_rejects_closed_slots_and_retains_timeout_watermark() {
357 let mut buffer = SlotBuffer::new();
358 buffer.push(10, 2, event(1));
359 assert_eq!(ids(buffer.flush_before(11)), [1]);
360 buffer.push(10, 1, event(99));
361 buffer.push(10, 9, event(99));
362 buffer.push(11, 0, event(2));
363 assert_eq!(ids(buffer.flush_all()), [2]);
364 assert_eq!(buffer.ordered_late_events(), 2);
365
366 buffer.push(11, 3, event(3));
367 buffer.push(11, 3, event(4));
368 assert_eq!(ids(buffer.flush_all()), [3, 4]);
369 for index in [0, 1, 2, 3] { buffer.push(11, index, event(99)); }
370 buffer.push(11, 4, event(5));
371 assert_eq!(ids(buffer.flush_all()), [5]);
372 assert_eq!(buffer.ordered_late_events(), 6);
373 }
374
375 #[test]
376 fn ordered_late_replay_bounds_diagnostics() {
377 use std::sync::atomic::{AtomicUsize, Ordering};
378 struct Capture;
379 static COUNT: AtomicUsize = AtomicUsize::new(0);
380 static LOGGER: Capture = Capture;
381 impl log::Log for Capture {
382 fn enabled(&self, _: &log::Metadata) -> bool { true }
383 fn log(&self, record: &log::Record) {
384 if record.args().to_string().contains("late event (300001,0)") {
386 COUNT.fetch_add(1, Ordering::Relaxed);
387 }
388 }
389 fn flush(&self) {}
390 }
391 log::set_logger(&LOGGER).unwrap();
392 log::set_max_level(log::LevelFilter::Warn);
393 let mut buffer = SlotBuffer::new();
394 buffer.push(300001, 1, event(1));
395 assert_eq!(ids(buffer.flush_all()), [1]);
396 for _ in 0..10000 { buffer.push(300001, 0, event(99)); }
397 assert_eq!(buffer.ordered_late_events(), 10000);
398 assert_eq!(COUNT.load(Ordering::Relaxed), 20);
399 assert!(buffer.slots.is_empty());
400 buffer.push(300001, 2, event(2));
401 assert_eq!(ids(buffer.flush_all()), [2]);
402 }
403
404 #[test]
405 fn ordered_watermark_accepts_max_index_without_overflow() {
406 let mut buffer = SlotBuffer::new();
407 buffer.push(u64::MAX, u64::MAX, event(1));
408 assert_eq!(ids(buffer.flush_all()), [1]);
409 buffer.push(u64::MAX, u64::MAX, event(99));
410 assert!(buffer.flush_all().is_empty());
411 assert_eq!(buffer.ordered_late_events(), 1);
412 }
413
414}