Skip to main content

engramdb_io/
planner.rs

1//! 流式预取规划器(M1.5 首交付):连续的 token rowid 流 → 增量 badge 读取计划。
2//!
3//! 证据与设计(docs/design.md §3.4、§7.0):
4//! - 预取窗口 = 计算窗口:窗口随 token 推进滚动,IO 与计算重叠;
5//! - 窗口内 badge/BLOCK 只派发一次(下一 token 重复访问 = 已驻留,不再读盘);
6//! - 不做语料级热集:P2 证明大语料下唯一行占表 37-49%、top1K 覆盖 <6%。
7//!
8//! 粒度说明:本版本按 "badge"(行簇)为读取单元,等价于 Store-I 路径;
9//! Store-P(4KB 视图记录)接同一结构(badge 语义换成 view_page)。
10
11use 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    /// 推进一批 rowid(一个 token 的 16 个 head 行)。
32    /// 返回"新 badge"数量;`plan` 只追加本次需要读的 badge(升序语义由消费方 settle)。
33    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        // 推进 6 个互不相同的 token(各 2 行且不同 badge,窗口预算 4 badge)
96        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            // 窗口=4:前两 token 会被挤出 -> 后续每 token 仍 2 新 badge
101            assert_eq!(n, 2, "t={t} window len={}", sp.window_len());
102        }
103        assert!(sp.window_len() <= 4);
104    }
105}