Skip to main content

engramdb_io/
view.rs

1//! Store-P 物化视图(P4):构建器 + 视图读写/基准/延迟原语。
2//!
3//! 视图 = 每 gram 16 头行(160B×16=2560B 紧凑槽,可选 4096 对齐槽)连续物化,
4//! 1 次定长读替代 16 路随机行读(IOPS 16:1、字节放大 1.00× vs scatter 20×)。
5//! build/bench/lat 原语与探针(p4view)及 CLI(`engramdb view`)共用;
6//! 输出格式保持探测基线(gate.sh + probes/*.csv 解析)不变。
7//!
8//! 规模:流式分块(500K grams/chunk)→ 16GB 内存机可构建全表;manifest.json
9//! 记录 grans/slot/耗时(lat/bench 以 manifest 为 n 的缺省来源)。
10
11use std::fs::{File, OpenOptions};
12use std::io::{BufRead, BufReader, BufWriter, Write};
13use std::path::Path;
14
15use crate::backend::{platform_read_at, platform_read_exact_at};
16use crate::batch::BadgeGather;
17use engramdb_core::layout::Layout;
18
19pub const HEAD_W: u64 = 16;
20pub const ROW_BYTES: u64 = 160;
21pub const RECORD_BYTES: u64 = HEAD_W * ROW_BYTES; // 2560B
22
23fn slot_of(view_file: &Path) -> (u64, u64) {
24    let mp = view_file.with_extension("manifest.json");
25    if mp.exists() {
26        let m: serde_json::Value = serde_json::from_slice(&std::fs::read(&mp).unwrap_or_default())
27            .unwrap_or(serde_json::Value::Null);
28        let slot = m["slot_bytes"].as_u64().unwrap_or(RECORD_BYTES);
29        let grans = m["grans"].as_u64().unwrap_or(0);
30        (slot, grans)
31    } else {
32        (RECORD_BYTES, 0)
33    }
34}
35
36/// 构建视图:固定 LCG seed 按 gram 序列生成 16 头 rowid,流式分块 gather 后写入
37/// 定长槽;manifest 写 `<view>.manifest.json`;keys(rowid 列表)写 `keys_out`
38/// (None = 不写;n 已在 manifest,读取方可 B-only)。
39pub fn build_view(
40    batch: &BadgeGather,
41    n: usize,
42    slot_bytes: u64,
43    view_out: &Path,
44    keys_out: Option<&Path>,
45) -> std::io::Result<f64> {
46    const CHUNK_G: usize = 500_000;
47    let spec = engramdb_keygen::PleSpec::real();
48    let mut rng_state: u64 = 0xDEAD_BEEF_1234_5678;
49    let mut keys = keys_out.map(|p| BufWriter::new(File::create(p).unwrap()));
50    let mut view = BufWriter::with_capacity(
51        64 << 20,
52        OpenOptions::new()
53            .write(true)
54            .create(true)
55            .truncate(true)
56            .open(view_out)?,
57    );
58    let t0 = std::time::Instant::now();
59    let mut done = 0usize;
60    while done < n {
61        let m = CHUNK_G.min(n - done);
62        let mut rowids = Vec::with_capacity(m * HEAD_W as usize);
63        for _ in 0..m {
64            rng_state = rng_state
65                .wrapping_mul(6364136223846793005)
66                .wrapping_add(1442695040888963407);
67            let a = (rng_state % 248_320) as u32;
68            rng_state = rng_state
69                .wrapping_mul(6364136223846793005)
70                .wrapping_add(1442695040888963407);
71            let b = (rng_state % 248_320) as u32;
72            rng_state = rng_state
73                .wrapping_mul(6364136223846793005)
74                .wrapping_add(1442695040888963407);
75            let c = (rng_state % 248_320) as u32;
76            let ids = spec.rowids_for_seq(&[a, b, c]);
77            for &r in &ids[0] {
78                rowids.push(r as u64);
79            }
80        }
81        let mut out = vec![0u8; rowids.len() * ROW_BYTES as usize];
82        batch
83            .gather_pp(&rowids, &mut out, 8)
84            .map_err(|e| std::io::Error::other(e.to_string()))?;
85        let mut s = vec![0u8; slot_bytes as usize];
86        let rec_len = (HEAD_W * ROW_BYTES) as usize;
87        for i in 0..m {
88            let rec = &out[i * rec_len..(i + 1) * rec_len];
89            s[..rec.len()].copy_from_slice(rec);
90            view.write_all(&s)?;
91        }
92        if let Some(w) = keys.as_mut() {
93            for &r in &rowids {
94                writeln!(w, "{r}")?;
95            }
96        }
97        done += m;
98    }
99    view.flush()?;
100    let build_s = t0.elapsed().as_secs_f64();
101    let m = serde_json::json!({
102        "grans": n,
103        "heads": HEAD_W,
104        "slot_bytes": slot_bytes,
105        "record_bytes": RECORD_BYTES,
106        "build_seconds": build_s,
107        "build_mb_s": (n as f64 * slot_bytes as f64 / 1e6) / build_s,
108        "rows": n as u64 * HEAD_W,
109        "source": format!("shards={}", batch.layout.shards),
110    });
111    std::fs::write(
112        view_out.with_extension("manifest.json"),
113        serde_json::to_vec_pretty(&m).unwrap(),
114    )?;
115    println!(
116        "view built: n={n} slot={slot_bytes}B view={} took={build_s:.1}s",
117        view_out.display()
118    );
119    Ok(build_s)
120}
121
122/// A 路径(16 行 scatter)唯一 4KiB 页数(字节放大的分母)。
123pub fn unique_pages(keys: &[u64], layout: &Layout) -> usize {
124    let mut set = std::collections::HashSet::with_capacity(keys.len() / 8);
125    for &k in keys {
126        set.insert(k * layout.row_bytes / 4096);
127    }
128    set.len()
129}
130
131/// 从 keys 文件读 rowid(trim + 空行跳过)。
132pub fn read_keys(p: &Path) -> std::io::Result<Vec<u64>> {
133    let mut k = Vec::new();
134    for l in BufReader::new(File::open(p)?).lines().map_while(Result::ok) {
135        let l = l.trim();
136        if !l.is_empty() {
137            k.push(
138                l.parse::<u64>()
139                    .map_err(|e: std::num::ParseIntError| std::io::Error::other(e.to_string()))?,
140            );
141        }
142    }
143    Ok(k)
144}
145
146/// 视图读取器:打开已经构建好的 Store-P 视图,按物理槽位读取记录。
147/// 这是面向产品面/PyO3 的读取 API(`build_view` 负责写;本类型负责读)。
148pub struct ViewReader {
149    file: File,
150    slot_bytes: u64,
151    count: usize,
152}
153
154impl ViewReader {
155    /// 打开视图文件;优先从 `.manifest.json` 读取槽宽与记录数,缺失时按文件大小推断。
156    pub fn open(view_file: &Path) -> std::io::Result<Self> {
157        let (slot_bytes, manifest_grans) = slot_of(view_file);
158        let file = File::open(view_file)?;
159        let meta = file.metadata()?;
160        let count_file = meta
161            .len()
162            .checked_div(slot_bytes)
163            .map(|n| n as usize)
164            .unwrap_or(0);
165        let count = if manifest_grans > 0 {
166            count_file.min(manifest_grans as usize)
167        } else {
168            count_file
169        };
170        Ok(Self {
171            file,
172            slot_bytes,
173            count,
174        })
175    }
176
177    pub fn len(&self) -> usize {
178        self.count
179    }
180
181    pub fn is_empty(&self) -> bool {
182        self.count == 0
183    }
184
185    pub fn slot_bytes(&self) -> u64 {
186        self.slot_bytes
187    }
188
189    /// 读取一个 gram 的完整 e_t 记录到 `buf`。返回实际读到的槽宽(通常 2560B)。
190    pub fn read_record(&self, index: usize, buf: &mut [u8]) -> std::io::Result<usize> {
191        if index >= self.count {
192            return Err(std::io::Error::new(
193                std::io::ErrorKind::InvalidInput,
194                format!(
195                    "view record index {index} out of range (count {})",
196                    self.count
197                ),
198            ));
199        }
200        let want = (self.slot_bytes as usize).min(buf.len());
201        if want == 0 {
202            return Ok(0);
203        }
204        platform_read_exact_at(&self.file, &mut buf[..want], index as u64 * self.slot_bytes)?;
205        Ok(want)
206    }
207
208    /// 按物理槽位读取多条记录;`out` 长度必须 >= indices.len() * slot_bytes。
209    pub fn read_records(&self, indices: &[usize], out: &mut [u8]) -> std::io::Result<()> {
210        let want = self.slot_bytes as usize;
211        if out.len() < indices.len() * want {
212            return Err(std::io::Error::new(
213                std::io::ErrorKind::InvalidInput,
214                "read_records: output buffer too small",
215            ));
216        }
217        for (j, &idx) in indices.iter().enumerate() {
218            self.read_record(idx, &mut out[j * want..(j + 1) * want])?;
219        }
220        Ok(())
221    }
222}
223
224/// 视图构建器:持有源表 gather 句柄与槽宽,提供随机采样构建和调用方访问序构建。
225pub struct ViewBuilder<'a> {
226    batch: &'a BadgeGather<'a>,
227    slot_bytes: u64,
228}
229
230impl<'a> ViewBuilder<'a> {
231    pub fn new(batch: &'a BadgeGather<'a>, slot_bytes: u64) -> Self {
232        Self { batch, slot_bytes }
233    }
234
235    /// 随机采样构建(原有 LCG 路径)。
236    pub fn build_random(
237        &self,
238        n: usize,
239        view_out: &Path,
240        keys_out: Option<&Path>,
241    ) -> std::io::Result<f64> {
242        build_view(self.batch, n, self.slot_bytes, view_out, keys_out)
243    }
244
245    /// 按调用方提供的 rowid 访问序构建。
246    pub fn build_from_keys(
247        &self,
248        keys: &[u64],
249        view_out: &Path,
250        keys_out: Option<&Path>,
251    ) -> std::io::Result<f64> {
252        build_view_from_keys(self.batch, keys, self.slot_bytes, view_out, keys_out)
253    }
254}
255
256/// 从调用方提供的 rowid 列表构建视图(`keys` 为 16 头平铺:每 gram 连续 16 行)。
257/// 槽位物理顺序 = 调用方给的顺序;因此可用来生成“按访问序排布”的视图:
258/// 把实际推理/训练访问的 gram 顺序直接作为 keys 顺序写入,顺序读取即连续 IO。
259pub fn build_view_from_keys(
260    batch: &BadgeGather,
261    keys: &[u64],
262    slot_bytes: u64,
263    view_out: &Path,
264    keys_out: Option<&Path>,
265) -> std::io::Result<f64> {
266    const CHUNK_ROWS: usize = 500_000 * HEAD_W as usize;
267    if !keys.len().is_multiple_of(HEAD_W as usize) {
268        return Err(std::io::Error::new(
269            std::io::ErrorKind::InvalidInput,
270            format!(
271                "keys length {} is not a multiple of heads {HEAD_W}",
272                keys.len()
273            ),
274        ));
275    }
276    let n = keys.len() / HEAD_W as usize;
277    let rec_len = (HEAD_W * ROW_BYTES) as usize;
278    let mut view = BufWriter::with_capacity(
279        64 << 20,
280        OpenOptions::new()
281            .write(true)
282            .create(true)
283            .truncate(true)
284            .open(view_out)?,
285    );
286    let mut keys_w = keys_out.map(|p| BufWriter::new(File::create(p).unwrap()));
287    let t0 = std::time::Instant::now();
288    for chunk in keys.chunks(CHUNK_ROWS) {
289        let m = chunk.len() / HEAD_W as usize;
290        let mut out = vec![0u8; chunk.len() * ROW_BYTES as usize];
291        batch
292            .gather_pp(chunk, &mut out, 8)
293            .map_err(|e| std::io::Error::other(e.to_string()))?;
294        let mut slot = vec![0u8; slot_bytes as usize];
295        for i in 0..m {
296            let rec = &out[i * rec_len..(i + 1) * rec_len];
297            slot[..rec.len()].copy_from_slice(rec);
298            view.write_all(&slot)?;
299        }
300        if let Some(w) = keys_w.as_mut() {
301            for &r in chunk {
302                writeln!(w, "{r}")?;
303            }
304        }
305    }
306    view.flush()?;
307    let build_s = t0.elapsed().as_secs_f64();
308    let manifest = serde_json::json!({
309        "grans": n,
310        "heads": HEAD_W,
311        "slot_bytes": slot_bytes,
312        "record_bytes": RECORD_BYTES,
313        "build_seconds": build_s,
314        "build_mb_s": (n as f64 * slot_bytes as f64 / 1e6) / build_s.max(1e-9),
315        "rows": n as u64 * HEAD_W,
316        "source": format!("provided-keys:{} shards={}", keys.len(), batch.layout.shards),
317        "layout": "access-order",
318    });
319    std::fs::write(
320        view_out.with_extension("manifest.json"),
321        serde_json::to_vec_pretty(&manifest).unwrap(),
322    )?;
323    println!(
324        "view built from keys: n={n} slot={slot_bytes}B view={} took={build_s:.1}s",
325        view_out.display()
326    );
327    Ok(build_s)
328}
329
330/// 流式文件版:从 rowid 文本文件(每行一个 u64)构建访问序视图,适合 keys 文件很大时使用。
331pub fn build_view_from_keys_file(
332    batch: &BadgeGather,
333    keys_path: &Path,
334    slot_bytes: u64,
335    view_out: &Path,
336    keys_out: Option<&Path>,
337) -> std::io::Result<f64> {
338    const CHUNK_ROWS: usize = 500_000 * HEAD_W as usize;
339    let mut reader = BufReader::new(File::open(keys_path)?);
340    let mut view = BufWriter::with_capacity(
341        64 << 20,
342        OpenOptions::new()
343            .write(true)
344            .create(true)
345            .truncate(true)
346            .open(view_out)?,
347    );
348    let mut keys_w = keys_out.map(|p| BufWriter::new(File::create(p).unwrap()));
349    let rec_len = (HEAD_W * ROW_BYTES) as usize;
350    let t0 = std::time::Instant::now();
351    let mut total = 0usize;
352    let mut chunk: Vec<u64> = Vec::with_capacity(CHUNK_ROWS);
353    let mut line = String::new();
354    loop {
355        line.clear();
356        let r = reader.read_line(&mut line)?;
357        if r == 0 {
358            break;
359        }
360        let t = line.trim();
361        if t.is_empty() {
362            continue;
363        }
364        chunk.push(
365            t.parse::<u64>()
366                .map_err(|e: std::num::ParseIntError| std::io::Error::other(e.to_string()))?,
367        );
368        if chunk.len() == CHUNK_ROWS {
369            let m = chunk.len() / HEAD_W as usize;
370            let mut out = vec![0u8; chunk.len() * ROW_BYTES as usize];
371            batch
372                .gather_pp(&chunk, &mut out, 8)
373                .map_err(|e| std::io::Error::other(e.to_string()))?;
374            let mut slot = vec![0u8; slot_bytes as usize];
375            for i in 0..m {
376                let rec = &out[i * rec_len..(i + 1) * rec_len];
377                slot[..rec.len()].copy_from_slice(rec);
378                view.write_all(&slot)?;
379            }
380            if let Some(w) = keys_w.as_mut() {
381                for &r in &chunk {
382                    writeln!(w, "{r}")?;
383                }
384            }
385            total += chunk.len();
386            chunk.clear();
387        }
388    }
389    if !chunk.is_empty() {
390        if !chunk.len().is_multiple_of(HEAD_W as usize) {
391            return Err(std::io::Error::new(
392                std::io::ErrorKind::InvalidInput,
393                format!(
394                    "keys file has {} rowids, not a multiple of heads {HEAD_W}",
395                    chunk.len() + total
396                ),
397            ));
398        }
399        let m = chunk.len() / HEAD_W as usize;
400        let mut out = vec![0u8; chunk.len() * ROW_BYTES as usize];
401        batch
402            .gather_pp(&chunk, &mut out, 8)
403            .map_err(|e| std::io::Error::other(e.to_string()))?;
404        let mut slot = vec![0u8; slot_bytes as usize];
405        for i in 0..m {
406            let rec = &out[i * rec_len..(i + 1) * rec_len];
407            slot[..rec.len()].copy_from_slice(rec);
408            view.write_all(&slot)?;
409        }
410        if let Some(w) = keys_w.as_mut() {
411            for &r in &chunk {
412                writeln!(w, "{r}")?;
413            }
414        }
415        total += chunk.len();
416    }
417    if total == 0 || !total.is_multiple_of(HEAD_W as usize) {
418        return Err(std::io::Error::new(
419            std::io::ErrorKind::InvalidInput,
420            format!("keys file total {total} rowids not multiple of heads"),
421        ));
422    }
423    view.flush()?;
424    let build_s = t0.elapsed().as_secs_f64();
425    let n = total / HEAD_W as usize;
426    let manifest = serde_json::json!({
427        "grans": n,
428        "heads": HEAD_W,
429        "slot_bytes": slot_bytes,
430        "record_bytes": RECORD_BYTES,
431        "build_seconds": build_s,
432        "build_mb_s": (n as f64 * slot_bytes as f64 / 1e6) / build_s.max(1e-9),
433        "rows": n as u64 * HEAD_W,
434        "source": format!("provided-keys-file:{} shards={}", total, batch.layout.shards),
435        "layout": "access-order",
436    });
437    std::fs::write(
438        view_out.with_extension("manifest.json"),
439        serde_json::to_vec_pretty(&manifest).unwrap(),
440    )?;
441    println!(
442        "view built from keys file: n={n} slot={slot_bytes}B view={} took={build_s:.1}s",
443        view_out.display()
444    );
445    Ok(build_s)
446}
447
448fn report(name: &str, rows: u64, dt: std::time::Duration) -> String {
449    let s = dt.as_secs_f64();
450    let rps = rows as f64 / s.max(1e-9);
451    let mbps = rows as f64 * ROW_BYTES as f64 / 1e6 / s.max(1e-9);
452    println!("{name}: rows={rows} time={s:.3}s  rows/s={rps:.0}  MB/s={mbps:.1}");
453    format!("{name},{rows},{:.0},{:.1}\n", rps, mbps)
454}
455
456/// 吞吐基准:A(scatter 对照, 有 keys 时)+ B(视图单记录读, 1t & 8t 两档)。
457/// 返回 (A_rows_per_s, B8t_rows_per_s, ampl)。CSV 打印保持探测格式。
458pub fn bench_view(
459    batch: &BadgeGather,
460    view_file: &Path,
461    keys: Option<&[u64]>,
462    sub_grams: usize,
463    threads: usize,
464    req_slot: u64,
465    order_mode: &str,
466) -> std::io::Result<(f64, f64, f64)> {
467    let (slot_bytes, grans) = if req_slot > 0 {
468        (req_slot, 0)
469    } else {
470        slot_of(view_file)
471    };
472    let keys = keys.unwrap_or(&[]);
473    let mut n_grams = if !keys.is_empty() {
474        keys.len() / HEAD_W as usize
475    } else {
476        grans as usize
477    };
478    if sub_grams > 0 {
479        n_grams = n_grams.min(sub_grams);
480    }
481    if n_grams == 0 {
482        return Err(std::io::Error::other("无 keys 且 manifest 缺 grans"));
483    }
484    let w = batch.layout.width as usize;
485
486    // A: 16 行 scatter(unique 页统计 + 8t 计时)
487    let out_a_len = keys.len() * w;
488    let mut out_a = vec![0u8; out_a_len];
489    if !keys.is_empty() {
490        let t0 = std::time::Instant::now();
491        batch
492            .gather_pp(keys, &mut out_a, 8)
493            .map_err(|e| std::io::Error::other(e.to_string()))?;
494        let dt_a = t0.elapsed();
495        let p = unique_pages(&keys[..n_grams * HEAD_W as usize], batch.layout);
496        println!("A unique 4KiB pages: {} (rows {})", p, keys.len());
497        let _ = dt_a;
498    }
499
500    // B: view 单记录读(固定 LCG 随机序)
501    let vf = File::open(view_file)?;
502    let meta = vf.metadata()?;
503    let grans_actual = (meta.len() / slot_bytes) as usize;
504    let n_grams = n_grams.min(grans_actual);
505    let mut order: Vec<u64>;
506    if order_mode == "seq" {
507        order = (0..n_grams as u64).collect();
508    } else {
509        let mut g_state: u64 = 0xCAFE_BEEF_0F1E_2D3C;
510        order = Vec::with_capacity(n_grams);
511        for _ in 0..n_grams {
512            g_state = g_state
513                .wrapping_mul(6364136223846793005)
514                .wrapping_add(1442695040888963407);
515            order.push(g_state % n_grams as u64);
516        }
517    }
518    let run_par = |threads: usize| -> std::time::Duration {
519        let t = std::time::Instant::now();
520        std::thread::scope(|sc| {
521            let chunk = n_grams.div_ceil(threads);
522            for th in 0..threads {
523                let lo = th * chunk;
524                if lo >= n_grams {
525                    break;
526                }
527                let hi = (lo + chunk).min(n_grams);
528                let slice = &order[lo..hi];
529                let fb = &vf;
530                sc.spawn(move || {
531                    let mut buf = vec![0u8; slot_bytes as usize];
532                    for &rec in slice {
533                        let _ = platform_read_exact_at(fb, &mut buf, rec * slot_bytes);
534                    }
535                });
536            }
537        });
538        t.elapsed()
539    };
540    let mut csv = String::new();
541    let dt_b = run_par(1);
542    csv.push_str(&report("B", n_grams as u64 * HEAD_W, dt_b));
543    let dt_b2 = run_par(threads);
544    csv.push_str(&report("B", n_grams as u64 * HEAD_W, dt_b2));
545    let a_rps = if !keys.is_empty() {
546        let t3 = std::time::Instant::now();
547        batch
548            .gather_pp(&keys[..n_grams * HEAD_W as usize], &mut out_a, 8)
549            .map_err(|e| std::io::Error::other(e.to_string()))?;
550        let dt_a2 = t3.elapsed();
551        let a = report("A", (n_grams * HEAD_W as usize) as u64, dt_a2);
552        csv.push_str(&a);
553        (n_grams * HEAD_W as usize) as f64 / dt_a2.as_secs_f64().max(1e-9)
554    } else {
555        0.0
556    };
557    csv.push_str(&format!(
558        "amplification,B,{}\n",
559        slot_bytes as f64 / RECORD_BYTES as f64
560    ));
561    let b_rps = (n_grams * HEAD_W as usize) as f64 / dt_b2.as_secs_f64().max(1e-9);
562    println!("{csv}");
563    Ok((a_rps, b_rps, slot_bytes as f64 / RECORD_BYTES as f64))
564}
565
566/// Fisher-Yates 随机访问序(xorshift64;固定 seed 可复现)。
567pub fn rand_order(n: usize, seed: u64) -> Vec<u64> {
568    let mut st = seed;
569    let mut v: Vec<u64> = (0..n as u64).collect();
570    for i in (1..n).rev() {
571        st ^= st << 13;
572        st ^= st >> 7;
573        st ^= st << 17;
574        let j = (st % (i as u64 + 1)) as usize;
575        v.swap(i, j);
576    }
577    v
578}
579
580/// 延迟分布:随机序单记录读 × N(每查独立计时)→ p50/p95/p99/max/mean(μs)。
581/// warm=全文件顺序预读;cold=Linux fadvise(DONTNEED)(真冷,需 Linux)。
582pub fn lat_view(
583    view_file: &Path,
584    threads: usize,
585    warm: bool,
586    cold: bool,
587    sub_grams: usize,
588    req_slot: u64,
589) -> std::io::Result<()> {
590    let (slot_bytes, manifest_n) = if req_slot > 0 {
591        (req_slot, 0usize)
592    } else {
593        let (s, g) = slot_of(view_file);
594        (s, g as usize)
595    };
596    let vf = File::open(view_file)?;
597    let meta = vf.metadata()?;
598    let n_all = (meta.len() / slot_bytes) as usize;
599    let n = if sub_grams > 0 {
600        n_all.min(sub_grams)
601    } else if manifest_n > 0 {
602        n_all.min(manifest_n)
603    } else {
604        n_all
605    };
606    if n == 0 {
607        return Err(std::io::Error::other("视图为空"));
608    }
609    if warm {
610        let mut buf = vec![0u8; 8 << 20];
611        let mut off = 0u64;
612        let f = vf.try_clone()?;
613        while off < meta.len() {
614            let want = (buf.len() as u64).min(meta.len() - off) as usize;
615            let rd = platform_read_at(&f, &mut buf[..want], off)?;
616            if rd == 0 {
617                break;
618            }
619            off += rd as u64;
620        }
621    }
622    let order = rand_order(n, 0xFEED_BEEF_0D0F_1E2C);
623    let view_ref = &vf;
624    #[cfg(not(target_os = "linux"))]
625    let _ = cold;
626    let per_thread = |tid: usize| -> Vec<u32> {
627        let mut buf = vec![0u8; slot_bytes as usize];
628        let mut out = Vec::new();
629        let stride = order.len() / threads;
630        let lo = tid * stride;
631        let hi = if tid + 1 == threads {
632            order.len()
633        } else {
634            lo + stride
635        };
636        for &rec in &order[lo..hi] {
637            let t0 = std::time::Instant::now();
638            #[cfg(target_os = "linux")]
639            if cold {
640                use std::os::fd::AsRawFd;
641                unsafe {
642                    libc::posix_fadvise(
643                        view_ref.as_raw_fd(),
644                        (rec * slot_bytes) as i64,
645                        slot_bytes as i64,
646                        libc::POSIX_FADV_DONTNEED,
647                    );
648                }
649            }
650            let _ = platform_read_exact_at(view_ref, &mut buf, rec * slot_bytes);
651            out.push(t0.elapsed().as_nanos().min(u32::MAX as u128) as u32);
652        }
653        out
654    };
655    let mut times: Vec<u32> = std::thread::scope(|sc| {
656        let mut h = Vec::new();
657        for tid in 0..threads {
658            h.push(sc.spawn(move || per_thread(tid)));
659        }
660        let mut all = Vec::with_capacity(n);
661        for x in h {
662            all.extend(x.join().unwrap());
663        }
664        all
665    });
666    times.sort_unstable();
667    let p = |q: f64| -> f64 {
668        let idx = ((times.len() - 1) as f64 * q).round() as usize;
669        times[idx] as f64
670    };
671    let mean = times.iter().map(|&x| x as f64).sum::<f64>() / times.len() as f64;
672    let (p50, p95, p99, mx) = (
673        p(0.50),
674        p(0.95),
675        p(0.99),
676        *times.last().unwrap_or(&0) as f64,
677    );
678    let us = |ns: f64| ns / 1000.0;
679    println!(
680        "lat: n={n} slot={slot_bytes} threads={threads} {} [μs] p50={:.2} p95={:.2} p99={:.2} max={:.2} mean={:.2}",
681        if warm { "warm" } else { "cold" },
682        us(p50), us(p95), us(p99), us(mx), us(mean)
683    );
684    println!(
685        "latency_us,p50,p95,p99,max,mean,{},{},{},{},{}\n",
686        us(p50),
687        us(p95),
688        us(p99),
689        us(mx),
690        us(mean)
691    );
692    Ok(())
693}
694
695/// 校验视图:用 keys(每 gram 连续 16 个 rowid)从源表重新 gather,和视图记录逐字节比对。
696/// 默认抽样 1000 个 gram;`sub_grams` 可显式指定检查数量。
697pub fn verify_view(
698    batch: &BadgeGather,
699    view_file: &Path,
700    keys: Option<&[u64]>,
701    sub_grams: usize,
702) -> std::io::Result<()> {
703    let reader = ViewReader::open(view_file)?;
704    let n = reader.len();
705    if n == 0 {
706        return Err(std::io::Error::other("视图为空"));
707    }
708    let keys = keys.unwrap_or(&[]);
709    if keys.is_empty() {
710        return Err(std::io::Error::other(
711            "verify 需要 keys 文件(每 gram 连续 16 个 rowid)",
712        ));
713    }
714    let key_grams = keys.len() / HEAD_W as usize;
715    if key_grams < n {
716        return Err(std::io::Error::new(
717            std::io::ErrorKind::InvalidInput,
718            format!("keys 文件只有 {key_grams} grams,不足以覆盖视图 {n} grams"),
719        ));
720    }
721    let total = n.min(key_grams);
722    let check = if sub_grams > 0 {
723        sub_grams.min(total)
724    } else {
725        total.min(1000)
726    };
727    let slot = reader.slot_bytes() as usize;
728    let mut view_buf = vec![0u8; slot];
729    let mut src_buf = vec![0u8; RECORD_BYTES as usize];
730    for gi in 0..check {
731        let start = gi * HEAD_W as usize;
732        let rowids = &keys[start..start + HEAD_W as usize];
733        batch
734            .gather_pp(rowids, &mut src_buf, 8)
735            .map_err(|e| std::io::Error::other(e.to_string()))?;
736        let got = reader.read_record(gi, &mut view_buf)?;
737        if got < RECORD_BYTES as usize || view_buf[..RECORD_BYTES as usize] != src_buf[..] {
738            return Err(std::io::Error::new(
739                std::io::ErrorKind::InvalidData,
740                format!("view record {gi} does not match source rows"),
741            ));
742        }
743    }
744    println!("view verified: checked {check}/{total} grams, slot={slot}B, all match");
745    Ok(())
746}
747
748#[cfg(test)]
749mod tests {
750    use super::*;
751    use std::io::Write;
752
753    #[test]
754    fn build_from_keys_and_view_reader_roundtrip() {
755        let dir = std::env::temp_dir().join("engramdb-view-reader-test");
756        let _ = std::fs::remove_dir_all(&dir);
757        std::fs::create_dir_all(&dir).unwrap();
758        let layout = Layout::new(1, 100, ROW_BYTES, 1); // 100 rows × 160B
759        let shard = dir.join("shard_000.bin");
760        let mut f = File::create(&shard).unwrap();
761        let mut row = vec![0u8; ROW_BYTES as usize];
762        for r in 0..100u64 {
763            row.fill((r % 251) as u8);
764            f.write_all(&row).unwrap();
765        }
766        drop(f);
767        let bg = BadgeGather::open(&dir, &layout).unwrap();
768        let keys: Vec<u64> = (0..32).collect(); // 2 grams × 16 heads
769        let view_path = dir.join("ordered.bin");
770        let keys_path = dir.join("ordered.keys.txt");
771        build_view_from_keys(&bg, &keys, RECORD_BYTES, &view_path, Some(&keys_path)).unwrap();
772
773        let vr = ViewReader::open(&view_path).unwrap();
774        assert_eq!(vr.len(), 2);
775        assert_eq!(vr.slot_bytes(), RECORD_BYTES);
776        let mut buf = vec![0u8; RECORD_BYTES as usize];
777        let got = vr.read_record(0, &mut buf).unwrap();
778        assert_eq!(got, RECORD_BYTES as usize);
779        for i in 0..16usize {
780            let r = i as u64;
781            let expect = (r % 251) as u8;
782            assert_eq!(buf[i * ROW_BYTES as usize], expect, "row {r} first byte");
783        }
784
785        let mut two = vec![0u8; 2 * RECORD_BYTES as usize];
786        vr.read_records(&[1, 0], &mut two).unwrap();
787        // record 1 begins with row 16 value 16
788        assert_eq!(two[0], 16);
789
790        verify_view(&bg, &view_path, Some(&keys), 0).unwrap();
791
792        let _ = std::fs::remove_dir_all(&dir);
793    }
794}