Skip to main content

rudb_native/
run_projection.rs

1//! A row-preserving run codec for a sorted, covering integer projection.
2//!
3//! A page stores each sorted order value once with its run length. Raw runs keep every covered
4//! code; uniform runs keep one code and the row count already in the header. Both layouts retain
5//! duplicate `(order, covered)` rows. Pages end between order values so query workers can count
6//! exact pairs independently.
7
8use 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
24/// Attach a row-preserving run projection to an existing native table.
25///
26/// The builder retains every source row's covered value and multiplicity. Consecutive equal
27/// order values share their eight-byte order value, and uniform covered runs use lossless
28/// run-length encoding. The caller must include this explicit build in indexed-load measurements.
29/// An append makes
30/// the section stale until it is rebuilt.
31///
32/// # Errors
33///
34/// If the table or columns are missing, nullable, unsupported, or the file cannot be updated.
35/// This first page format also rejects one order-value run too large to fit in a single page.
36pub 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    /// Scan one disjoint range of whole pages on an engine worker.
247    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    /// Merge page ranges in their stored order and verify the complete row count.
268    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    /// Whether a current row-preserving projection covers these columns.
310    ///
311    /// The payload is validated when it is read, not during this directory lookup.
312    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    /// Count exact distinct order values by covered value from a current run projection.
322    ///
323    /// Returns `None` if no current matching section exists. Stored codes and run lengths are
324    /// decoded when this query runs; the file contains no saved distinct pair or grouped count.
325    ///
326    /// # Errors
327    ///
328    /// If a matching section or its source directory is damaged.
329    ///
330    /// # Panics
331    ///
332    /// Fixed-width header decoding assumes the lengths checked immediately before it.
333    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    /// Evaluate a grouped distinct count from source rows with a bounded number of workers.
344    /// This is the entry point for the SQL operator, whose parallelism is set by the engine.
345    ///
346    /// # Errors
347    ///
348    /// If the projection directory or payload is damaged.
349    ///
350    /// # Panics
351    ///
352    /// Fixed-width header decoding assumes the lengths checked immediately before it.
353    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    /// Open and validate a row-preserving projection for independently scheduled page scans.
400    ///
401    /// # Errors
402    ///
403    /// If the projection directory or payload is damaged.
404    ///
405    /// # Panics
406    ///
407    /// Fixed-width header decoding assumes the lengths checked immediately before it.
408    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, &regions).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, &regions).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, &regions).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, &regions).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}