Skip to main content

sol_parser_sdk/grpc/
buffers.rs

1//! 事件缓冲区模块 - 用于有序模式下的事件排序和批次处理
2//!
3//! 提供多种缓冲策略:
4//! - `SlotBuffer`: 按 slot 缓冲,支持 Ordered 和 StreamingOrdered 模式
5//! - `MicroBatchBuffer`: 微秒级时间窗口批次,用于 MicroBatch 模式
6
7use crate::DexEvent;
8use std::collections::{BTreeMap, HashMap, HashSet};
9use tokio::time::Instant;
10
11// ==================== SlotBuffer ====================
12
13/// Slot 缓冲区,用于有序模式下缓存同一 slot 的事件
14#[derive(Default)]
15pub struct SlotBuffer {
16    /// slot -> Vec<(tx_index, event)>
17    slots: BTreeMap<u64, Vec<(u64, DexEvent)>>,
18    /// 当前处理的最大 slot
19    current_slot: u64,
20    /// 上次输出时间
21    last_flush_time: Option<Instant>,
22    /// 流式模式:每个 slot 已释放的最大连续 tx_index
23    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    /// 添加事件到缓冲区
44    #[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            // Preserve exact drop accounting without synchronous log I/O on every replay.
49            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    /// 输出所有小于 current_slot 的事件
62    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    /// 超时强制输出所有缓冲事件
85    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    /// Late Ordered events are dropped with a warning to preserve monotonic output.
106    pub fn ordered_late_events(&self) -> u64 { self.ordered_late_events }
107
108    /// 检查是否超时
109    #[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    /// Single-event compatibility API. Multi-event transactions must use
117    /// `push_streaming_batch` so their index advances only once.
118    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    /// Release complete transaction batches in contiguous transaction-index order.
123    /// Filtered streams with index gaps should use MicroBatch or the timeout fallback.
124    /// Once a slot is flushed by a newer slot, late batches for it are discarded.
125    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        // A streaming watermark must represent the next index without u64 overflow.
132        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            // Immediately emitted slots have no buffer entry, but still own a watermark.
143            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    /// 流式模式超时释放
178    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
200// ==================== MicroBatchBuffer ====================
201
202/// 微批次缓冲区,用于 MicroBatch 模式
203pub struct MicroBatchBuffer {
204    /// 当前窗口内的事件: (slot, tx_index, event)
205    events: Vec<(u64, u64, DexEvent)>,
206    /// 窗口开始时间(微秒)
207    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    /// 添加事件到窗口,返回是否需要刷新
217    #[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    /// 刷新窗口,返回排序后的事件
234    #[inline]
235    pub fn flush(&mut self) -> Vec<DexEvent> {
236        if self.events.is_empty() {
237            return Vec::new();
238        }
239
240        // Stable sort preserves parser order for multiple events from one transaction.
241        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    /// 检查是否需要刷新(窗口超时)
251    #[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                // Other parallel tests are isolated by this unique slot/index.
385                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}