1use 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; fn 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
36pub 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
122pub 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
131pub 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
146pub struct ViewReader {
149 file: File,
150 slot_bytes: u64,
151 count: usize,
152}
153
154impl ViewReader {
155 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 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 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
224pub 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 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 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
256pub 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
330pub 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
456pub 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 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 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
566pub 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
580pub 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
695pub 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); 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(); 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 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}