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, followed by the covered code
4//! for every original row. Duplicate `(order, covered)` rows are retained. Pages end between
5//! order values so query workers can count exact pairs independently.
6
7use std::collections::{HashMap, HashSet};
8use std::path::Path;
9
10use rudb_common::Result;
11
12use crate::projection::{eligible, id, integers};
13use crate::{Catalog, Reader, attach, invalid, section};
14
15const MAGIC: &[u8; 8] = b"RUDBRP1\0";
16const PAGE_BYTES: usize = 1 << 19;
17const FIXED_HEADER: usize = 8 + 2 + 2 + 8 + 2 + 4 + 1;
18const PAGE_HEADER: usize = 4 + 4 + 8 + 8;
19const RUN_HEADER: usize = 8 + 4;
20
21/// Attach a row-preserving run projection to an existing native table.
22///
23/// The builder stores one covered-value code for every source row. Consecutive equal order
24/// values share their eight-byte order value but retain their original multiplicities. The
25/// caller must include this explicit build in any indexed-load measurement. An append makes
26/// the section stale until it is rebuilt.
27///
28/// # Errors
29///
30/// If the table or columns are missing, nullable, unsupported, or the file cannot be updated.
31/// This first page format also rejects one order-value run too large to fit in a single page.
32pub fn build_run_projection(
33    path: impl AsRef<Path>,
34    table: &str,
35    order_column: &str,
36    covered_column: &str,
37) -> Result<()> {
38    let path = path.as_ref();
39    let catalog = Catalog::open(path)?;
40    let reader = catalog.table(table)?;
41    let fields = reader.table().fields();
42    let order = fields
43        .iter()
44        .position(|field| field.name.eq_ignore_ascii_case(order_column))
45        .ok_or_else(|| invalid("run projection order column is missing"))?;
46    let covered = fields
47        .iter()
48        .position(|field| field.name.eq_ignore_ascii_case(covered_column))
49        .ok_or_else(|| invalid("run projection covered column is missing"))?;
50    if order == covered
51        || !fields[order].not_null
52        || !fields[covered].not_null
53        || !eligible(&fields[order].ty, &fields[covered].ty)
54    {
55        return Err(invalid("run projection needs two supported, non-null integer columns"));
56    }
57    let mut rows = Vec::<(i64, i32)>::with_capacity(reader.table().rows());
58    let mut values = HashSet::<i32>::new();
59    let (mut users, mut groups) = (Vec::new(), Vec::new());
60    for part in 0..reader.parts() {
61        let chunk = reader.read(part, &[order, covered])?;
62        integers(chunk.column(0)?, &mut users)?;
63        integers(chunk.column(1)?, &mut groups)?;
64        for (&order_value, &group) in users.iter().zip(&groups) {
65            let covered_value = i32::try_from(group)
66                .map_err(|_| invalid("projection covered value exceeds INTEGER"))?;
67            rows.push((order_value, covered_value));
68            values.insert(covered_value);
69        }
70    }
71    if rows.len() != reader.table().rows() {
72        return Err(invalid("run projection row count differs from its table"));
73    }
74    let mut dictionary = values.into_iter().collect::<Vec<_>>();
75    dictionary.sort_unstable();
76    let dictionary_len = u16::try_from(dictionary.len())
77        .map_err(|_| invalid("run projection covered column exceeds 65535 values"))?;
78    let code_bytes = if dictionary.len() <= 256 { 1 } else { 2 };
79    let codes = dictionary
80        .iter()
81        .enumerate()
82        .map(|(at, &value)| (value, at as u16))
83        .collect::<HashMap<_, _>>();
84    rows.sort_unstable_by_key(|&(order_value, _)| order_value);
85    let order_index =
86        u16::try_from(order).map_err(|_| invalid("run projection column index overflow"))?;
87    let covered_index =
88        u16::try_from(covered).map_err(|_| invalid("run projection column index overflow"))?;
89    let header = FIXED_HEADER + dictionary.len() * 4;
90    if header + PAGE_HEADER >= PAGE_BYTES {
91        return Err(invalid("run projection dictionary does not fit in its first page"));
92    }
93    let mut bytes = vec![0_u8; PAGE_BYTES];
94    bytes[..8].copy_from_slice(MAGIC);
95    bytes[8..10].copy_from_slice(&order_index.to_le_bytes());
96    bytes[10..12].copy_from_slice(&covered_index.to_le_bytes());
97    bytes[12..20].copy_from_slice(&(rows.len() as u64).to_le_bytes());
98    bytes[20..22].copy_from_slice(&dictionary_len.to_le_bytes());
99    bytes[26] = code_bytes as u8;
100    for (at, value) in dictionary.iter().enumerate() {
101        let offset = FIXED_HEADER + at * 4;
102        bytes[offset..offset + 4].copy_from_slice(&value.to_le_bytes());
103    }
104    let mut page = 0_usize;
105    let mut cursor = header + PAGE_HEADER;
106    let mut page_rows = 0_u32;
107    let mut page_first = None;
108    let mut page_last = None;
109    let mut at = 0;
110    while at < rows.len() {
111        let user = rows[at].0;
112        let mut end = at + 1;
113        while end < rows.len() && rows[end].0 == user {
114            end += 1;
115        }
116        let count = u32::try_from(end - at)
117            .map_err(|_| invalid("run projection user run exceeds its page count"))?;
118        let run_bytes = RUN_HEADER + (end - at) * code_bytes;
119        if run_bytes > PAGE_BYTES - PAGE_HEADER {
120            return Err(invalid("run projection user run exceeds a page"));
121        }
122        if cursor + run_bytes > (page + 1) * PAGE_BYTES {
123            finish_page(&mut bytes, page, header, cursor, page_rows, page_first, page_last)?;
124            page += 1;
125            bytes.resize((page + 1) * PAGE_BYTES, 0);
126            cursor = page * PAGE_BYTES + PAGE_HEADER;
127            page_rows = 0;
128            page_first = None;
129        }
130        if cursor + run_bytes > (page + 1) * PAGE_BYTES {
131            return Err(invalid("run projection user run exceeds the first page"));
132        }
133        page_first.get_or_insert(user);
134        page_last = Some(user);
135        bytes[cursor..cursor + 8].copy_from_slice(&user.to_le_bytes());
136        bytes[cursor + 8..cursor + 12].copy_from_slice(&count.to_le_bytes());
137        cursor += RUN_HEADER;
138        for &(_, group) in &rows[at..end] {
139            let code = codes[&group];
140            if code_bytes == 1 {
141                bytes[cursor] = code as u8;
142                cursor += 1;
143            } else {
144                bytes[cursor..cursor + 2].copy_from_slice(&code.to_le_bytes());
145                cursor += 2;
146            }
147        }
148        page_rows = page_rows
149            .checked_add(count)
150            .ok_or_else(|| invalid("run projection page row count overflow"))?;
151        at = end;
152    }
153    finish_page(&mut bytes, page, header, cursor, page_rows, page_first, page_last)?;
154    let pages = u32::try_from(page + 1).map_err(|_| invalid("too many run projection pages"))?;
155    bytes[22..26].copy_from_slice(&pages.to_le_bytes());
156    drop(reader);
157    drop(catalog);
158    attach(
159        path,
160        table,
161        &[section::Attachment {
162            kind: *section::RUN_PROJECTION,
163            id: id(order, covered)?,
164            flags: 0,
165            header_bytes: header as u32,
166            bytes: &bytes,
167        }],
168    )?;
169    Ok(())
170}
171
172fn finish_page(
173    bytes: &mut [u8],
174    page: usize,
175    header: usize,
176    cursor: usize,
177    rows: u32,
178    first: Option<i64>,
179    last: Option<i64>,
180) -> Result<()> {
181    let prefix = page * PAGE_BYTES + if page == 0 { header } else { 0 };
182    let used = u32::try_from(cursor - prefix - PAGE_HEADER)
183        .map_err(|_| invalid("run projection page length overflow"))?;
184    bytes[prefix..prefix + 4].copy_from_slice(&used.to_le_bytes());
185    bytes[prefix + 4..prefix + 8].copy_from_slice(&rows.to_le_bytes());
186    bytes[prefix + 8..prefix + 16].copy_from_slice(&first.unwrap_or(0).to_le_bytes());
187    bytes[prefix + 16..prefix + 24].copy_from_slice(&last.unwrap_or(0).to_le_bytes());
188    Ok(())
189}
190
191#[derive(Debug)]
192struct Scan {
193    counts: Vec<u64>,
194    rows: u64,
195    first: Option<i64>,
196    last: Option<i64>,
197}
198
199impl Reader {
200    /// Count exact distinct order values by covered value from a current run projection.
201    ///
202    /// Returns `None` if no current matching section exists. Every covered code is read when
203    /// this query runs; the file contains no saved distinct pair or grouped count.
204    ///
205    /// # Errors
206    ///
207    /// If a matching section or its source directory is damaged.
208    ///
209    /// # Panics
210    ///
211    /// Fixed-width header decoding assumes the lengths checked immediately before it.
212    pub fn grouped_distinct_run_projection(
213        &self,
214        order: usize,
215        covered: usize,
216        limit: usize,
217    ) -> Result<Option<Vec<(i32, u64)>>> {
218        let wanted = id(order, covered)?;
219        let Some(section) = self.table().sections().iter().find(|section| {
220            section.kind == *section::RUN_PROJECTION
221                && section.id == wanted
222                && section.usable(self.table().generation())
223        }) else {
224            return Ok(None);
225        };
226        let extents = self.extents(section)?;
227        let first_extent = extents.first().ok_or_else(|| invalid("run projection has no page"))?;
228        let first_page = self.extent(first_extent)?;
229        if first_page.len() != PAGE_BYTES || &first_page[..8] != MAGIC {
230            return Err(invalid("run projection header differs"));
231        }
232        let stored_order = u16::from_le_bytes(first_page[8..10].try_into().unwrap());
233        let stored_covered = u16::from_le_bytes(first_page[10..12].try_into().unwrap());
234        if usize::from(stored_order) != order || usize::from(stored_covered) != covered {
235            return Err(invalid("run projection columns differ from its section"));
236        }
237        let rows = u64::from_le_bytes(first_page[12..20].try_into().unwrap());
238        if rows != self.table().rows() as u64 {
239            return Err(invalid("run projection row count differs from its table"));
240        }
241        let size = u16::from_le_bytes(first_page[20..22].try_into().unwrap()) as usize;
242        let pages = u32::from_le_bytes(first_page[22..26].try_into().unwrap()) as usize;
243        let code_bytes = usize::from(first_page[26]);
244        if !matches!(code_bytes, 1 | 2) || (code_bytes == 1 && size > 256) {
245            return Err(invalid("run projection code width differs from its dictionary"));
246        }
247        if pages != extents.len() {
248            return Err(invalid("run projection page count differs from its extents"));
249        }
250        let header = FIXED_HEADER + size * 4;
251        if header + PAGE_HEADER > PAGE_BYTES || section.header_bytes as usize != header {
252            return Err(invalid("run projection dictionary exceeds its first page"));
253        }
254        let dictionary = first_page[FIXED_HEADER..header]
255            .chunks_exact(4)
256            .map(|bytes| i32::from_le_bytes(bytes.try_into().unwrap()))
257            .collect::<Vec<_>>();
258        if dictionary.windows(2).any(|pair| pair[0] >= pair[1]) {
259            return Err(invalid("run projection dictionary is not sorted and unique"));
260        }
261        let base = first_extent.offset;
262        for (at, extent) in extents.iter().enumerate() {
263            let offset = at as u64 * PAGE_BYTES as u64;
264            if extent.first != offset
265                || extent.offset != base + offset
266                || extent.length as usize != PAGE_BYTES
267            {
268                return Err(invalid("run projection pages are not contiguous"));
269            }
270        }
271        drop(first_page);
272        let workers = std::thread::available_parallelism().map_or(1, usize::from).min(8).min(pages);
273        let scans = std::thread::scope(|scope| -> Result<Vec<Scan>> {
274            let mut handles = Vec::with_capacity(workers);
275            for worker in 0..workers {
276                let begin = pages * worker / workers;
277                let end = pages * (worker + 1) / workers;
278                let extent_slice = &extents[begin..end];
279                handles.push(scope.spawn(move || {
280                    scan_pages(self, extent_slice, begin, header, size, code_bytes)
281                }));
282            }
283            handles
284                .into_iter()
285                .map(|handle| {
286                    handle.join().map_err(|_| invalid("run projection worker panicked"))?
287                })
288                .collect()
289        })?;
290        let mut totals = vec![0_u64; size];
291        let mut total_rows = 0_u64;
292        let mut previous_last = None;
293        for scan in scans {
294            if let Some(first) = scan.first {
295                if previous_last.is_some_and(|previous| first <= previous) {
296                    return Err(invalid("run projection page order differs"));
297                }
298                previous_last = scan.last;
299            }
300            total_rows = total_rows
301                .checked_add(scan.rows)
302                .ok_or_else(|| invalid("run projection row count overflow"))?;
303            for (total, value) in totals.iter_mut().zip(scan.counts) {
304                *total = total
305                    .checked_add(value)
306                    .ok_or_else(|| invalid("run projection count overflow"))?;
307            }
308        }
309        if total_rows != rows {
310            return Err(invalid("run projection decoded row count differs"));
311        }
312        let mut ranked =
313            dictionary.into_iter().zip(totals).filter(|(_, count)| *count != 0).collect::<Vec<_>>();
314        ranked.sort_unstable_by(|a, b| b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0)));
315        ranked.truncate(limit);
316        Ok(Some(ranked))
317    }
318}
319
320fn scan_pages(
321    reader: &Reader,
322    extents: &[section::Extent],
323    first_index: usize,
324    header: usize,
325    dictionary: usize,
326    code_bytes: usize,
327) -> Result<Scan> {
328    let mut marks = vec![0_u32; dictionary];
329    let mut counts = vec![0_u64; dictionary];
330    let mut epoch = 0_u32;
331    let mut rows = 0_u64;
332    let mut first = None;
333    let mut last = None;
334    let mut bytes = Vec::with_capacity(PAGE_BYTES);
335    for (relative, extent) in extents.iter().enumerate() {
336        reader.extent_into(extent, &mut bytes)?;
337        let prefix = if first_index + relative == 0 { header } else { 0 };
338        if bytes.len() != PAGE_BYTES || prefix + PAGE_HEADER > PAGE_BYTES {
339            return Err(invalid("run projection page length differs"));
340        }
341        let used = u32::from_le_bytes(bytes[prefix..prefix + 4].try_into().unwrap()) as usize;
342        let stored_rows =
343            u32::from_le_bytes(bytes[prefix + 4..prefix + 8].try_into().unwrap()) as u64;
344        let stored_first = i64::from_le_bytes(bytes[prefix + 8..prefix + 16].try_into().unwrap());
345        let stored_last = i64::from_le_bytes(bytes[prefix + 16..prefix + 24].try_into().unwrap());
346        let mut at = prefix + PAGE_HEADER;
347        let end = at
348            .checked_add(used)
349            .filter(|&end| end <= PAGE_BYTES)
350            .ok_or_else(|| invalid("run projection page data exceeds its extent"))?;
351        let mut page_rows = 0_u64;
352        let mut page_first = None;
353        let mut page_last = None;
354        while at < end {
355            if end - at < RUN_HEADER {
356                return Err(invalid("run projection page has a partial run header"));
357            }
358            let user = i64::from_le_bytes(bytes[at..at + 8].try_into().unwrap());
359            let length = u32::from_le_bytes(bytes[at + 8..at + 12].try_into().unwrap()) as usize;
360            at += RUN_HEADER;
361            if length == 0 || length > (end - at) / code_bytes {
362                return Err(invalid("run projection run length exceeds its page"));
363            }
364            if last.is_some_and(|previous| user <= previous) {
365                return Err(invalid("run projection order is not increasing"));
366            }
367            first.get_or_insert(user);
368            page_first.get_or_insert(user);
369            page_last = Some(user);
370            last = Some(user);
371            if code_bytes == 1 {
372                epoch = epoch.wrapping_add(1);
373                if epoch == 0 {
374                    marks.fill(0);
375                    epoch = 1;
376                }
377                for &code in &bytes[at..at + length] {
378                    let code = usize::from(code);
379                    let mark = marks
380                        .get_mut(code)
381                        .ok_or_else(|| invalid("run projection code is outside its dictionary"))?;
382                    if *mark != epoch {
383                        *mark = epoch;
384                        counts[code] += 1;
385                    }
386                }
387            } else {
388                if length <= 2 {
389                    let first_code =
390                        u16::from_le_bytes(bytes[at..at + 2].try_into().unwrap()) as usize;
391                    let first_count = counts
392                        .get_mut(first_code)
393                        .ok_or_else(|| invalid("run projection code is outside its dictionary"))?;
394                    *first_count += 1;
395                    if length == 2 {
396                        let second_code =
397                            u16::from_le_bytes(bytes[at + 2..at + 4].try_into().unwrap()) as usize;
398                        if second_code != first_code {
399                            let second_count = counts.get_mut(second_code).ok_or_else(|| {
400                                invalid("run projection code is outside its dictionary")
401                            })?;
402                            *second_count += 1;
403                        }
404                    }
405                } else {
406                    epoch = epoch.wrapping_add(1);
407                    if epoch == 0 {
408                        marks.fill(0);
409                        epoch = 1;
410                    }
411                    for code_bytes in bytes[at..at + length * 2].chunks_exact(2) {
412                        let code = u16::from_le_bytes(code_bytes.try_into().unwrap()) as usize;
413                        let mark = marks.get_mut(code).ok_or_else(|| {
414                            invalid("run projection code is outside its dictionary")
415                        })?;
416                        if *mark != epoch {
417                            *mark = epoch;
418                            counts[code] += 1;
419                        }
420                    }
421                }
422            }
423            at += length * code_bytes;
424            page_rows += length as u64;
425        }
426        if page_rows != stored_rows
427            || page_first.unwrap_or(0) != stored_first
428            || page_last.unwrap_or(0) != stored_last
429        {
430            return Err(invalid("run projection page directory differs from its rows"));
431        }
432        rows = rows
433            .checked_add(page_rows)
434            .ok_or_else(|| invalid("run projection row count overflow"))?;
435    }
436    Ok(Scan { counts, rows, first, last })
437}
438
439#[cfg(test)]
440mod tests {
441    use std::sync::atomic::{AtomicUsize, Ordering};
442
443    use rudb_common::{Field, LogicalType, Value};
444    use rudb_vector::{Chunk, Vector};
445
446    use crate::{Catalog, Writer};
447
448    use super::build_run_projection;
449
450    static NEXT: AtomicUsize = AtomicUsize::new(0);
451
452    #[test]
453    fn run_projection_keeps_duplicate_rows_but_counts_distinct_pairs() {
454        let at = NEXT.fetch_add(1, Ordering::Relaxed);
455        let path = std::env::temp_dir()
456            .join(format!("rudb-run-projection-{}-{at}.rdb", std::process::id()));
457        let fields = vec![
458            Field::required("user", LogicalType::BigInt),
459            Field::required("region", LogicalType::Integer),
460        ];
461        let mut writer = Writer::create(&path, "events", fields).expect("create native file");
462        for (users, regions) in [
463            (vec![9_i64, 2, 9, 1], vec![7_i32, 1, 7, 2]),
464            (vec![2_i64, 2, 5, 9, 5, 8, 8], vec![2_i32, 2, 1, 1, 2, 1, 1]),
465        ] {
466            let users = users.into_iter().map(Value::BigInt).collect::<Vec<_>>();
467            let regions = regions.into_iter().map(Value::Integer).collect::<Vec<_>>();
468            let chunk = Chunk::new(vec![
469                Vector::from_values(LogicalType::BigInt, &users).expect("users"),
470                Vector::from_values(LogicalType::Integer, &regions).expect("regions"),
471            ])
472            .expect("matching columns");
473            writer.append(&chunk).expect("append rows");
474        }
475        writer.finish().expect("commit native file");
476        build_run_projection(&path, "events", "user", "region").expect("build run projection");
477        let reader = Catalog::open(&path).expect("catalog").table("events").expect("table");
478        assert_eq!(
479            reader.grouped_distinct_run_projection(0, 1, 10).expect("valid projection"),
480            Some(vec![(1, 4), (2, 3), (7, 1)]),
481        );
482        std::fs::remove_file(path).expect("remove scratch file");
483    }
484
485    #[test]
486    fn run_projection_uses_wide_codes_for_large_dictionaries() {
487        let at = NEXT.fetch_add(1, Ordering::Relaxed);
488        let path = std::env::temp_dir()
489            .join(format!("rudb-run-projection-wide-{}-{at}.rdb", std::process::id()));
490        let fields = vec![
491            Field::required("user", LogicalType::BigInt),
492            Field::required("region", LogicalType::Integer),
493        ];
494        let users = vec![Value::BigInt(7); 300];
495        let regions = (0..300).map(Value::Integer).collect::<Vec<_>>();
496        let chunk = Chunk::new(vec![
497            Vector::from_values(LogicalType::BigInt, &users).expect("users"),
498            Vector::from_values(LogicalType::Integer, &regions).expect("regions"),
499        ])
500        .expect("matching columns");
501        let mut writer = Writer::create(&path, "events", fields).expect("create native file");
502        writer.append(&chunk).expect("append rows");
503        writer.finish().expect("commit native file");
504        build_run_projection(&path, "events", "user", "region").expect("build run projection");
505        let reader = Catalog::open(&path).expect("catalog").table("events").expect("table");
506        let expected = (0..300).map(|region| (region, 1)).collect::<Vec<_>>();
507        assert_eq!(
508            reader.grouped_distinct_run_projection(0, 1, 300).expect("valid projection"),
509            Some(expected),
510        );
511        std::fs::remove_file(path).expect("remove scratch file");
512    }
513}