1use std::collections::{HashSet, VecDeque};
12
13use crate::batch::PrefetchPlan;
14use engramdb_core::layout::Layout;
15
16pub struct StreamingPlanner {
17 window_badges: usize,
18 queue: VecDeque<(u64, u64)>,
19 in_window: HashSet<(u64, u64)>,
20}
21
22impl StreamingPlanner {
23 pub fn new(window_badges: usize) -> Self {
24 Self {
25 window_badges,
26 queue: VecDeque::new(),
27 in_window: HashSet::new(),
28 }
29 }
30
31 pub fn advance(&mut self, rows: &[u64], layout: &Layout, plan: &mut PrefetchPlan) -> usize {
34 let mut n_new = 0;
35 for &r in rows {
36 let (shard, badge, _) = layout.locate(r);
37 let key = (shard, badge);
38 if !self.in_window.contains(&key) {
39 self.queue.push_back(key);
40 self.in_window.insert(key);
41 plan.entry(shard, badge);
42 n_new += 1;
43 }
44 }
45 while self.in_window.len() > self.window_badges {
46 if let Some(old) = self.queue.pop_front() {
47 self.in_window.remove(&old);
48 }
49 }
50 n_new
51 }
52
53 pub fn window_len(&self) -> usize {
54 self.in_window.len()
55 }
56}
57
58#[cfg(test)]
59mod tests {
60 use super::*;
61
62 fn lay() -> Layout {
63 Layout::new(2, 10_000, 160, 1)
64 }
65
66 #[test]
67 fn first_token_all_new() {
68 let l = lay();
69 let mut sp = StreamingPlanner::new(100);
70 let mut plan = PrefetchPlan::default();
71 let rows: Vec<u64> = (0..16).map(|i| i * 1000).collect();
72 let n = sp.advance(&rows, &l, &mut plan);
73 assert_eq!(n, 16);
74 assert_eq!(plan.n_badges, 16);
75 assert_eq!(sp.window_len(), 16);
76 }
77
78 #[test]
79 fn repeated_token_reuses_window() {
80 let l = lay();
81 let mut sp = StreamingPlanner::new(100);
82 let mut plan = PrefetchPlan::default();
83 let rows: Vec<u64> = (0..16).map(|i| i * 1000).collect();
84 let _ = sp.advance(&rows, &l, &mut plan);
85 let mut plan2 = PrefetchPlan::default();
86 let n = sp.advance(&rows, &l, &mut plan2);
87 assert_eq!(n, 0, "same token again -> all badges resident");
88 assert_eq!(plan2.n_badges, 0);
89 }
90
91 #[test]
92 fn window_eviction_republish() {
93 let l = lay();
94 let mut sp = StreamingPlanner::new(4);
95 for t in 0..6 {
97 let mut plan = PrefetchPlan::default();
98 let rows: Vec<u64> = vec![t * 1000, (t + 30) * 1000];
99 let n = sp.advance(&rows, &l, &mut plan);
100 assert_eq!(n, 2, "t={t} window len={}", sp.window_len());
102 }
103 assert!(sp.window_len() <= 4);
104 }
105}