1use 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
21pub 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 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, ®ions).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, ®ions).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}