1use std::collections::{HashMap, HashSet};
9use std::path::Path;
10
11use rudb_common::Result;
12
13use crate::projection::{eligible, id, integers};
14use crate::{Catalog, Reader, attach, invalid, section};
15
16const MAGIC: &[u8; 8] = b"RUDBRP1\0";
17const LEGACY_PAGE_BYTES: usize = 1 << 19;
18pub(crate) const RLE_PAGE_BYTES: usize = 1 << 16;
19pub(crate) const RLE_PAGES: u32 = 1;
20const FIXED_HEADER: usize = 8 + 2 + 2 + 8 + 2 + 4 + 1;
21const PAGE_HEADER: usize = 4 + 4 + 8 + 8;
22const RUN_HEADER: usize = 8 + 4;
23
24pub fn build_run_projection(
37 path: impl AsRef<Path>,
38 table: &str,
39 order_column: &str,
40 covered_column: &str,
41) -> Result<()> {
42 let path = path.as_ref();
43 let catalog = Catalog::open(path)?;
44 let reader = catalog.table(table)?;
45 let fields = reader.table().fields();
46 let order = fields
47 .iter()
48 .position(|field| field.name.eq_ignore_ascii_case(order_column))
49 .ok_or_else(|| invalid("run projection order column is missing"))?;
50 let covered = fields
51 .iter()
52 .position(|field| field.name.eq_ignore_ascii_case(covered_column))
53 .ok_or_else(|| invalid("run projection covered column is missing"))?;
54 if order == covered
55 || !fields[order].not_null
56 || !fields[covered].not_null
57 || !eligible(&fields[order].ty, &fields[covered].ty)
58 {
59 return Err(invalid("run projection needs two supported, non-null integer columns"));
60 }
61 let mut rows = Vec::<(i64, i32)>::with_capacity(reader.table().rows());
62 let mut values = HashSet::<i32>::new();
63 let (mut users, mut groups) = (Vec::new(), Vec::new());
64 for part in 0..reader.parts() {
65 let chunk = reader.read(part, &[order, covered])?;
66 integers(chunk.column(0)?, &mut users)?;
67 integers(chunk.column(1)?, &mut groups)?;
68 for (&order_value, &group) in users.iter().zip(&groups) {
69 let covered_value = i32::try_from(group)
70 .map_err(|_| invalid("projection covered value exceeds INTEGER"))?;
71 rows.push((order_value, covered_value));
72 values.insert(covered_value);
73 }
74 }
75 if rows.len() != reader.table().rows() {
76 return Err(invalid("run projection row count differs from its table"));
77 }
78 let mut dictionary = values.into_iter().collect::<Vec<_>>();
79 dictionary.sort_unstable();
80 let dictionary_len = u16::try_from(dictionary.len())
81 .map_err(|_| invalid("run projection covered column exceeds 65535 values"))?;
82 let code_bytes = if dictionary.len() <= 256 { 1 } else { 2 };
83 let codes = dictionary
84 .iter()
85 .enumerate()
86 .map(|(at, &value)| (value, at as u16))
87 .collect::<HashMap<_, _>>();
88 rows.sort_unstable_by_key(|&(order_value, _)| order_value);
89 let order_index =
90 u16::try_from(order).map_err(|_| invalid("run projection column index overflow"))?;
91 let covered_index =
92 u16::try_from(covered).map_err(|_| invalid("run projection column index overflow"))?;
93 let header = FIXED_HEADER + dictionary.len() * 4;
94 if header + PAGE_HEADER >= LEGACY_PAGE_BYTES {
95 return Err(invalid("run projection dictionary does not fit in its first page"));
96 }
97 let longest_run =
98 rows.chunk_by(|left, right| left.0 == right.0).map(|run| run.len()).max().unwrap_or(0);
99 let raw_run_bytes =
100 longest_run.checked_mul(code_bytes).and_then(|bytes| bytes.checked_add(RUN_HEADER + 1));
101 let rle_fits = header + PAGE_HEADER < RLE_PAGE_BYTES
102 && raw_run_bytes.is_some_and(|bytes| bytes <= RLE_PAGE_BYTES - PAGE_HEADER - header);
103 let flags = if rle_fits { RLE_PAGES } else { 0 };
104 let page_bytes = if rle_fits { RLE_PAGE_BYTES } else { LEGACY_PAGE_BYTES };
105 let mut bytes = vec![0_u8; page_bytes];
106 bytes[..8].copy_from_slice(MAGIC);
107 bytes[8..10].copy_from_slice(&order_index.to_le_bytes());
108 bytes[10..12].copy_from_slice(&covered_index.to_le_bytes());
109 bytes[12..20].copy_from_slice(&(rows.len() as u64).to_le_bytes());
110 bytes[20..22].copy_from_slice(&dictionary_len.to_le_bytes());
111 bytes[26] = code_bytes as u8;
112 for (at, value) in dictionary.iter().enumerate() {
113 let offset = FIXED_HEADER + at * 4;
114 bytes[offset..offset + 4].copy_from_slice(&value.to_le_bytes());
115 }
116 let mut page = 0_usize;
117 let mut cursor = header + PAGE_HEADER;
118 let mut page_rows = 0_u32;
119 let mut page_first = None;
120 let mut page_last = None;
121 let mut at = 0;
122 while at < rows.len() {
123 let user = rows[at].0;
124 let mut end = at + 1;
125 while end < rows.len() && rows[end].0 == user {
126 end += 1;
127 }
128 let count = u32::try_from(end - at)
129 .map_err(|_| invalid("run projection user run exceeds its page count"))?;
130 let uniform = flags == RLE_PAGES
131 && rows[at..end].iter().all(|&(_, covered_value)| covered_value == rows[at].1);
132 let payload_bytes = if uniform { code_bytes } else { (end - at) * code_bytes };
133 let run_bytes = RUN_HEADER + usize::from(flags == RLE_PAGES) + payload_bytes;
134 if run_bytes > page_bytes - PAGE_HEADER {
135 return Err(invalid("run projection user run exceeds a page"));
136 }
137 if cursor + run_bytes > (page + 1) * page_bytes {
138 finish_page(&mut bytes, page, header, cursor, page_rows, page_first, page_last)?;
139 page += 1;
140 bytes.resize((page + 1) * page_bytes, 0);
141 cursor = page * page_bytes + PAGE_HEADER;
142 page_rows = 0;
143 page_first = None;
144 }
145 if cursor + run_bytes > (page + 1) * page_bytes {
146 return Err(invalid("run projection user run exceeds the first page"));
147 }
148 page_first.get_or_insert(user);
149 page_last = Some(user);
150 bytes[cursor..cursor + 8].copy_from_slice(&user.to_le_bytes());
151 bytes[cursor + 8..cursor + 12].copy_from_slice(&count.to_le_bytes());
152 cursor += RUN_HEADER;
153 if flags == RLE_PAGES {
154 bytes[cursor] = u8::from(uniform);
155 cursor += 1;
156 }
157 let stored = if uniform { &rows[at..at + 1] } else { &rows[at..end] };
158 for &(_, covered_value) in stored {
159 let code = codes[&covered_value];
160 if code_bytes == 1 {
161 bytes[cursor] = code as u8;
162 cursor += 1;
163 } else {
164 bytes[cursor..cursor + 2].copy_from_slice(&code.to_le_bytes());
165 cursor += 2;
166 }
167 }
168 page_rows = page_rows
169 .checked_add(count)
170 .ok_or_else(|| invalid("run projection page row count overflow"))?;
171 at = end;
172 }
173 finish_page(&mut bytes, page, header, cursor, page_rows, page_first, page_last)?;
174 let pages = u32::try_from(page + 1).map_err(|_| invalid("too many run projection pages"))?;
175 bytes[22..26].copy_from_slice(&pages.to_le_bytes());
176 drop(reader);
177 drop(catalog);
178 attach(
179 path,
180 table,
181 &[section::Attachment {
182 kind: *section::RUN_PROJECTION,
183 id: id(order, covered)?,
184 flags,
185 header_bytes: header as u32,
186 bytes: &bytes,
187 }],
188 )?;
189 Ok(())
190}
191
192fn finish_page(
193 bytes: &mut [u8],
194 page: usize,
195 header: usize,
196 cursor: usize,
197 rows: u32,
198 first: Option<i64>,
199 last: Option<i64>,
200) -> Result<()> {
201 let page_bytes = bytes.len() / (page + 1);
202 let prefix = page * page_bytes + if page == 0 { header } else { 0 };
203 let used = u32::try_from(cursor - prefix - PAGE_HEADER)
204 .map_err(|_| invalid("run projection page length overflow"))?;
205 bytes[prefix..prefix + 4].copy_from_slice(&used.to_le_bytes());
206 bytes[prefix + 4..prefix + 8].copy_from_slice(&rows.to_le_bytes());
207 bytes[prefix + 8..prefix + 16].copy_from_slice(&first.unwrap_or(0).to_le_bytes());
208 bytes[prefix + 16..prefix + 24].copy_from_slice(&last.unwrap_or(0).to_le_bytes());
209 Ok(())
210}
211
212#[derive(Debug)]
213pub struct RunProjectionPart {
214 counts: Vec<u64>,
215 rows: u64,
216 first: Option<i64>,
217 last: Option<i64>,
218}
219
220#[derive(Debug)]
221pub struct RunProjectionScan<'a> {
222 reader: &'a Reader,
223 extents: Vec<section::Extent>,
224 first_page: Vec<u8>,
225 dictionary: Vec<i32>,
226 rows: u64,
227 header: usize,
228 code_bytes: usize,
229 rle: bool,
230}
231
232#[derive(Clone, Copy)]
233struct ScanLayout {
234 header: usize,
235 dictionary: usize,
236 code_bytes: usize,
237 rle: bool,
238}
239
240impl RunProjectionScan<'_> {
241 #[must_use]
242 pub fn pages(&self) -> usize {
243 self.extents.len()
244 }
245
246 pub fn partition(&self, part: usize, parts: usize) -> Result<RunProjectionPart> {
248 if parts == 0 || part >= parts || parts > self.pages() {
249 return Err(invalid("run projection partition is outside its pages"));
250 }
251 let begin = self.pages() * part / parts;
252 let end = self.pages() * (part + 1) / parts;
253 scan_pages(
254 self.reader,
255 &self.extents[begin..end],
256 begin,
257 ScanLayout {
258 header: self.header,
259 dictionary: self.dictionary.len(),
260 code_bytes: self.code_bytes,
261 rle: self.rle,
262 },
263 (begin == 0).then(|| self.first_page.clone()),
264 )
265 }
266
267 pub fn finish(
269 &self,
270 scans: impl IntoIterator<Item = RunProjectionPart>,
271 limit: usize,
272 ) -> Result<Vec<(i32, u64)>> {
273 let mut totals = vec![0_u64; self.dictionary.len()];
274 let mut total_rows = 0_u64;
275 let mut previous_last = None;
276 for scan in scans {
277 if let Some(first) = scan.first {
278 if previous_last.is_some_and(|previous| first <= previous) {
279 return Err(invalid("run projection page order differs"));
280 }
281 previous_last = scan.last;
282 }
283 total_rows = total_rows
284 .checked_add(scan.rows)
285 .ok_or_else(|| invalid("run projection row count overflow"))?;
286 for (total, value) in totals.iter_mut().zip(scan.counts) {
287 *total = total
288 .checked_add(value)
289 .ok_or_else(|| invalid("run projection count overflow"))?;
290 }
291 }
292 if total_rows != self.rows {
293 return Err(invalid("run projection decoded row count differs"));
294 }
295 let mut ranked = self
296 .dictionary
297 .iter()
298 .copied()
299 .zip(totals)
300 .filter(|(_, count)| *count != 0)
301 .collect::<Vec<_>>();
302 ranked.sort_unstable_by(|a, b| b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0)));
303 ranked.truncate(limit);
304 Ok(ranked)
305 }
306}
307
308impl Reader {
309 pub fn has_run_projection(&self, order: usize, covered: usize) -> Result<bool> {
313 let wanted = id(order, covered)?;
314 Ok(self.table().sections().iter().any(|section| {
315 section.kind == *section::RUN_PROJECTION
316 && section.id == wanted
317 && section.usable(self.table().generation())
318 }))
319 }
320
321 pub fn grouped_distinct_run_projection(
334 &self,
335 order: usize,
336 covered: usize,
337 limit: usize,
338 ) -> Result<Option<Vec<(i32, u64)>>> {
339 let workers = std::thread::available_parallelism().map_or(1, usize::from);
340 self.grouped_distinct_run_projection_with_workers(order, covered, limit, workers)
341 }
342
343 pub fn grouped_distinct_run_projection_with_workers(
354 &self,
355 order: usize,
356 covered: usize,
357 limit: usize,
358 workers: usize,
359 ) -> Result<Option<Vec<(i32, u64)>>> {
360 let Some(scan) = self.run_projection_scan(order, covered)? else {
361 return Ok(None);
362 };
363 let workers = workers.clamp(1, 8).min(scan.pages());
364 let parts = if workers == 1 {
365 let mut scan = scan;
366 let first_page = std::mem::take(&mut scan.first_page);
367 let part = scan_pages(
368 scan.reader,
369 &scan.extents,
370 0,
371 ScanLayout {
372 header: scan.header,
373 dictionary: scan.dictionary.len(),
374 code_bytes: scan.code_bytes,
375 rle: scan.rle,
376 },
377 Some(first_page),
378 )?;
379 return Ok(Some(scan.finish([part], limit)?));
380 } else {
381 std::thread::scope(|scope| -> Result<Vec<RunProjectionPart>> {
382 let handles = (0..workers)
383 .map(|part| {
384 let scan = &scan;
385 scope.spawn(move || scan.partition(part, workers))
386 })
387 .collect::<Vec<_>>();
388 handles
389 .into_iter()
390 .map(|handle| {
391 handle.join().map_err(|_| invalid("run projection worker panicked"))?
392 })
393 .collect()
394 })?
395 };
396 Ok(Some(scan.finish(parts, limit)?))
397 }
398
399 pub fn run_projection_scan(
409 &self,
410 order: usize,
411 covered: usize,
412 ) -> Result<Option<RunProjectionScan<'_>>> {
413 let wanted = id(order, covered)?;
414 let Some(section) = self.table().sections().iter().find(|section| {
415 section.kind == *section::RUN_PROJECTION
416 && section.id == wanted
417 && section.usable(self.table().generation())
418 }) else {
419 return Ok(None);
420 };
421 let page_bytes = match section.flags {
422 0 => LEGACY_PAGE_BYTES,
423 RLE_PAGES => RLE_PAGE_BYTES,
424 _ => return Err(invalid("run projection has unknown page layout flags")),
425 };
426 let extents = self.extents(section)?;
427 let first_extent = extents.first().ok_or_else(|| invalid("run projection has no page"))?;
428 let first_page = self.extent(first_extent)?;
429 if first_page.len() != page_bytes || &first_page[..8] != MAGIC {
430 return Err(invalid("run projection header differs"));
431 }
432 let stored_order = u16::from_le_bytes(first_page[8..10].try_into().unwrap());
433 let stored_covered = u16::from_le_bytes(first_page[10..12].try_into().unwrap());
434 if usize::from(stored_order) != order || usize::from(stored_covered) != covered {
435 return Err(invalid("run projection columns differ from its section"));
436 }
437 let rows = u64::from_le_bytes(first_page[12..20].try_into().unwrap());
438 if rows != self.table().rows() as u64 {
439 return Err(invalid("run projection row count differs from its table"));
440 }
441 let size = u16::from_le_bytes(first_page[20..22].try_into().unwrap()) as usize;
442 let pages = u32::from_le_bytes(first_page[22..26].try_into().unwrap()) as usize;
443 let code_bytes = usize::from(first_page[26]);
444 if !matches!(code_bytes, 1 | 2) || (code_bytes == 1 && size > 256) {
445 return Err(invalid("run projection code width differs from its dictionary"));
446 }
447 if pages != extents.len() {
448 return Err(invalid("run projection page count differs from its extents"));
449 }
450 let header = FIXED_HEADER + size * 4;
451 if header + PAGE_HEADER > page_bytes || section.header_bytes as usize != header {
452 return Err(invalid("run projection dictionary exceeds its first page"));
453 }
454 let dictionary = first_page[FIXED_HEADER..header]
455 .chunks_exact(4)
456 .map(|bytes| i32::from_le_bytes(bytes.try_into().unwrap()))
457 .collect::<Vec<_>>();
458 if dictionary.windows(2).any(|pair| pair[0] >= pair[1]) {
459 return Err(invalid("run projection dictionary is not sorted and unique"));
460 }
461 let base = first_extent.offset;
462 for (at, extent) in extents.iter().enumerate() {
463 let offset = at as u64 * page_bytes as u64;
464 if extent.first != offset
465 || extent.offset != base + offset
466 || extent.length as usize != page_bytes
467 {
468 return Err(invalid("run projection pages are not contiguous"));
469 }
470 }
471 Ok(Some(RunProjectionScan {
472 reader: self,
473 extents,
474 first_page,
475 dictionary,
476 rows,
477 header,
478 code_bytes,
479 rle: section.flags == RLE_PAGES,
480 }))
481 }
482}
483
484fn scan_pages(
485 reader: &Reader,
486 extents: &[section::Extent],
487 first_index: usize,
488 layout: ScanLayout,
489 initial: Option<Vec<u8>>,
490) -> Result<RunProjectionPart> {
491 let ScanLayout { header, dictionary, code_bytes, rle } = layout;
492 let page_bytes =
493 extents.first().ok_or_else(|| invalid("run projection has no page"))?.length as usize;
494 let mut marks = vec![0_u32; dictionary];
495 let mut counts = vec![0_u64; dictionary];
496 let mut epoch = 0_u32;
497 let mut rows = 0_u64;
498 let mut first = None;
499 let mut last = None;
500 let reused_first = initial.is_some();
501 let mut bytes = initial.unwrap_or_else(|| Vec::with_capacity(page_bytes));
502 for (relative, extent) in extents.iter().enumerate() {
503 if !reused_first || relative != 0 {
504 reader.extent_into(extent, &mut bytes)?;
505 }
506 let prefix = if first_index + relative == 0 { header } else { 0 };
507 if bytes.len() != page_bytes || prefix + PAGE_HEADER > page_bytes {
508 return Err(invalid("run projection page length differs"));
509 }
510 let used = u32::from_le_bytes(bytes[prefix..prefix + 4].try_into().unwrap()) as usize;
511 let stored_rows =
512 u32::from_le_bytes(bytes[prefix + 4..prefix + 8].try_into().unwrap()) as u64;
513 let stored_first = i64::from_le_bytes(bytes[prefix + 8..prefix + 16].try_into().unwrap());
514 let stored_last = i64::from_le_bytes(bytes[prefix + 16..prefix + 24].try_into().unwrap());
515 let mut at = prefix + PAGE_HEADER;
516 let end = at
517 .checked_add(used)
518 .filter(|&end| end <= page_bytes)
519 .ok_or_else(|| invalid("run projection page data exceeds its extent"))?;
520 let mut page_rows = 0_u64;
521 let mut page_first = None;
522 let mut page_last = None;
523 while at < end {
524 if end - at < RUN_HEADER {
525 return Err(invalid("run projection page has a partial run header"));
526 }
527 let user = i64::from_le_bytes(bytes[at..at + 8].try_into().unwrap());
528 let length = u32::from_le_bytes(bytes[at + 8..at + 12].try_into().unwrap()) as usize;
529 at += RUN_HEADER;
530 let mode = if rle {
531 if at == end {
532 return Err(invalid("run projection run is missing its mode"));
533 }
534 let mode = bytes[at];
535 at += 1;
536 mode
537 } else {
538 0
539 };
540 let payload_bytes = match mode {
541 0 => length * code_bytes,
542 1 => code_bytes,
543 _ => return Err(invalid("run projection run has an unknown mode")),
544 };
545 if length == 0 || payload_bytes > end - at {
546 return Err(invalid("run projection run length exceeds its page"));
547 }
548 if last.is_some_and(|previous| user <= previous) {
549 return Err(invalid("run projection order is not increasing"));
550 }
551 first.get_or_insert(user);
552 page_first.get_or_insert(user);
553 page_last = Some(user);
554 last = Some(user);
555 if mode == 1 {
556 let code = if code_bytes == 1 {
557 usize::from(bytes[at])
558 } else {
559 u16::from_le_bytes(bytes[at..at + 2].try_into().unwrap()) as usize
560 };
561 let count = counts
562 .get_mut(code)
563 .ok_or_else(|| invalid("run projection code is outside its dictionary"))?;
564 *count += 1;
565 } else if code_bytes == 1 {
566 epoch = epoch.wrapping_add(1);
567 if epoch == 0 {
568 marks.fill(0);
569 epoch = 1;
570 }
571 for &code in &bytes[at..at + length] {
572 let code = usize::from(code);
573 let mark = marks
574 .get_mut(code)
575 .ok_or_else(|| invalid("run projection code is outside its dictionary"))?;
576 if *mark != epoch {
577 *mark = epoch;
578 counts[code] += 1;
579 }
580 }
581 } else {
582 if length <= 2 {
583 let first_code =
584 u16::from_le_bytes(bytes[at..at + 2].try_into().unwrap()) as usize;
585 let first_count = counts
586 .get_mut(first_code)
587 .ok_or_else(|| invalid("run projection code is outside its dictionary"))?;
588 *first_count += 1;
589 if length == 2 {
590 let second_code =
591 u16::from_le_bytes(bytes[at + 2..at + 4].try_into().unwrap()) as usize;
592 if second_code != first_code {
593 let second_count = counts.get_mut(second_code).ok_or_else(|| {
594 invalid("run projection code is outside its dictionary")
595 })?;
596 *second_count += 1;
597 }
598 }
599 } else {
600 epoch = epoch.wrapping_add(1);
601 if epoch == 0 {
602 marks.fill(0);
603 epoch = 1;
604 }
605 for code_bytes in bytes[at..at + length * 2].chunks_exact(2) {
606 let code = u16::from_le_bytes(code_bytes.try_into().unwrap()) as usize;
607 let mark = marks.get_mut(code).ok_or_else(|| {
608 invalid("run projection code is outside its dictionary")
609 })?;
610 if *mark != epoch {
611 *mark = epoch;
612 counts[code] += 1;
613 }
614 }
615 }
616 }
617 at += payload_bytes;
618 page_rows += length as u64;
619 }
620 if page_rows != stored_rows
621 || page_first.unwrap_or(0) != stored_first
622 || page_last.unwrap_or(0) != stored_last
623 {
624 return Err(invalid("run projection page directory differs from its rows"));
625 }
626 rows = rows
627 .checked_add(page_rows)
628 .ok_or_else(|| invalid("run projection row count overflow"))?;
629 }
630 Ok(RunProjectionPart { counts, rows, first, last })
631}
632
633#[cfg(test)]
634mod tests {
635 use std::sync::atomic::{AtomicUsize, Ordering};
636
637 use rudb_common::{Field, LogicalType, Value};
638 use rudb_vector::{Chunk, Vector};
639
640 use crate::{Catalog, Writer};
641
642 use super::{RLE_PAGES, build_run_projection};
643
644 static NEXT: AtomicUsize = AtomicUsize::new(0);
645
646 #[test]
647 fn run_projection_keeps_duplicate_rows_but_counts_distinct_pairs() {
648 let at = NEXT.fetch_add(1, Ordering::Relaxed);
649 let path = std::env::temp_dir()
650 .join(format!("rudb-run-projection-{}-{at}.rdb", std::process::id()));
651 let fields = vec![
652 Field::required("user", LogicalType::BigInt),
653 Field::required("region", LogicalType::Integer),
654 ];
655 let mut writer = Writer::create(&path, "events", fields).expect("create native file");
656 for (users, regions) in [
657 (vec![9_i64, 2, 9, 1], vec![7_i32, 1, 7, 2]),
658 (vec![2_i64, 2, 5, 9, 5, 8, 8], vec![2_i32, 2, 1, 1, 2, 1, 1]),
659 ] {
660 let users = users.into_iter().map(Value::BigInt).collect::<Vec<_>>();
661 let regions = regions.into_iter().map(Value::Integer).collect::<Vec<_>>();
662 let chunk = Chunk::new(vec![
663 Vector::from_values(LogicalType::BigInt, &users).expect("users"),
664 Vector::from_values(LogicalType::Integer, ®ions).expect("regions"),
665 ])
666 .expect("matching columns");
667 writer.append(&chunk).expect("append rows");
668 }
669 writer.finish().expect("commit native file");
670 build_run_projection(&path, "events", "user", "region").expect("build run projection");
671 let reader = Catalog::open(&path).expect("catalog").table("events").expect("table");
672 assert!(reader.table().sections().iter().any(|section| {
673 section.kind == *crate::section::RUN_PROJECTION && section.flags == RLE_PAGES
674 }));
675 assert_eq!(
676 reader.grouped_distinct_run_projection(0, 1, 10).expect("valid projection"),
677 Some(vec![(1, 4), (2, 3), (7, 1)]),
678 );
679 let scan = reader.run_projection_scan(0, 1).expect("open projection").expect("projection");
680 let mut reconstructed = Vec::new();
681 for (page, extent) in scan.extents.iter().enumerate() {
682 let bytes = reader.extent(extent).expect("read verified page");
683 let prefix = if page == 0 { scan.header } else { 0 };
684 let used = u32::from_le_bytes(bytes[prefix..prefix + 4].try_into().unwrap()) as usize;
685 let mut cursor = prefix + super::PAGE_HEADER;
686 let end = cursor + used;
687 while cursor < end {
688 let user = i64::from_le_bytes(bytes[cursor..cursor + 8].try_into().unwrap());
689 let length =
690 u32::from_le_bytes(bytes[cursor + 8..cursor + 12].try_into().unwrap()) as usize;
691 cursor += super::RUN_HEADER;
692 let mode = if scan.rle {
693 let mode = bytes[cursor];
694 cursor += 1;
695 mode
696 } else {
697 0
698 };
699 let codes = if mode == 1 { 1 } else { length };
700 let mut covered = Vec::with_capacity(codes);
701 for _ in 0..codes {
702 let code = if scan.code_bytes == 1 {
703 let code = usize::from(bytes[cursor]);
704 cursor += 1;
705 code
706 } else {
707 let code = u16::from_le_bytes(bytes[cursor..cursor + 2].try_into().unwrap())
708 as usize;
709 cursor += 2;
710 code
711 };
712 covered.push(scan.dictionary[code]);
713 }
714 if mode == 1 {
715 reconstructed.extend(std::iter::repeat_n((user, covered[0]), length));
716 } else {
717 reconstructed.extend(covered.into_iter().map(|value| (user, value)));
718 }
719 }
720 }
721 reconstructed.sort_unstable();
722 assert_eq!(
723 reconstructed,
724 vec![
725 (1, 2),
726 (2, 1),
727 (2, 2),
728 (2, 2),
729 (5, 1),
730 (5, 2),
731 (8, 1),
732 (8, 1),
733 (9, 1),
734 (9, 7),
735 (9, 7),
736 ]
737 );
738 std::fs::remove_file(path).expect("remove scratch file");
739 }
740
741 #[test]
742 fn run_projection_uses_wide_codes_for_large_dictionaries() {
743 let at = NEXT.fetch_add(1, Ordering::Relaxed);
744 let path = std::env::temp_dir()
745 .join(format!("rudb-run-projection-wide-{}-{at}.rdb", std::process::id()));
746 let fields = vec![
747 Field::required("user", LogicalType::BigInt),
748 Field::required("region", LogicalType::Integer),
749 ];
750 let users = vec![Value::BigInt(7); 300];
751 let regions = (0..300).map(Value::Integer).collect::<Vec<_>>();
752 let chunk = Chunk::new(vec![
753 Vector::from_values(LogicalType::BigInt, &users).expect("users"),
754 Vector::from_values(LogicalType::Integer, ®ions).expect("regions"),
755 ])
756 .expect("matching columns");
757 let mut writer = Writer::create(&path, "events", fields).expect("create native file");
758 writer.append(&chunk).expect("append rows");
759 writer.finish().expect("commit native file");
760 build_run_projection(&path, "events", "user", "region").expect("build run projection");
761 let reader = Catalog::open(&path).expect("catalog").table("events").expect("table");
762 let expected = (0..300).map(|region| (region, 1)).collect::<Vec<_>>();
763 assert_eq!(
764 reader.grouped_distinct_run_projection(0, 1, 300).expect("valid projection"),
765 Some(expected),
766 );
767 std::fs::remove_file(path).expect("remove scratch file");
768 }
769
770 #[test]
771 fn run_projection_page_partitions_merge_exact_distinct_counts() {
772 let at = NEXT.fetch_add(1, Ordering::Relaxed);
773 let path = std::env::temp_dir()
774 .join(format!("rudb-run-projection-pages-{}-{at}.rdb", std::process::id()));
775 let fields = vec![
776 Field::required("user", LogicalType::BigInt),
777 Field::required("region", LogicalType::Integer),
778 ];
779 let mut writer = Writer::create(&path, "events", fields).expect("create native file");
780 for batch in 0..150 {
781 let users = (0..1000)
782 .map(|row| Value::BigInt(i64::from(batch * 500 + row / 2)))
783 .collect::<Vec<_>>();
784 let regions = (0..1000).map(|row| Value::Integer(row % 2 + 1)).collect::<Vec<_>>();
785 writer
786 .append(
787 &Chunk::new(vec![
788 Vector::from_values(LogicalType::BigInt, &users).expect("users"),
789 Vector::from_values(LogicalType::Integer, ®ions).expect("regions"),
790 ])
791 .expect("matching columns"),
792 )
793 .expect("append rows");
794 }
795 writer.finish().expect("commit native file");
796 build_run_projection(&path, "events", "user", "region").expect("build run projection");
797 let reader = Catalog::open(&path).expect("catalog").table("events").expect("table");
798 let scan = reader.run_projection_scan(0, 1).expect("open projection").expect("projection");
799 assert!(scan.pages() > 1);
800 let left = scan.partition(0, 2).expect("left pages");
801 let right = scan.partition(1, 2).expect("right pages");
802 assert_eq!(
803 scan.finish([left, right], usize::MAX).expect("merge pages"),
804 vec![(1, 75_000), (2, 75_000)]
805 );
806 std::fs::remove_file(path).expect("remove scratch file");
807 }
808
809 #[test]
810 fn long_order_runs_keep_the_legacy_page_width() {
811 let at = NEXT.fetch_add(1, Ordering::Relaxed);
812 let path = std::env::temp_dir()
813 .join(format!("rudb-run-projection-long-{}-{at}.rdb", std::process::id()));
814 let fields = vec![
815 Field::required("user", LogicalType::BigInt),
816 Field::required("region", LogicalType::Integer),
817 ];
818 let users = vec![Value::BigInt(7); 1000];
819 let regions = vec![Value::Integer(1); 1000];
820 let chunk = Chunk::new(vec![
821 Vector::from_values(LogicalType::BigInt, &users).expect("users"),
822 Vector::from_values(LogicalType::Integer, ®ions).expect("regions"),
823 ])
824 .expect("matching columns");
825 let mut writer = Writer::create(&path, "events", fields).expect("create native file");
826 for _ in 0..150 {
827 writer.append(&chunk).expect("append rows");
828 }
829 writer.finish().expect("commit native file");
830 build_run_projection(&path, "events", "user", "region").expect("build run projection");
831 let reader = Catalog::open(&path).expect("catalog").table("events").expect("table");
832 assert!(reader.table().sections().iter().any(|section| {
833 section.kind == *crate::section::RUN_PROJECTION && section.flags == 0
834 }));
835 assert_eq!(
836 reader.grouped_distinct_run_projection(0, 1, 10).expect("valid projection"),
837 Some(vec![(1, 1)])
838 );
839 std::fs::remove_file(path).expect("remove scratch file");
840 }
841}