1#![forbid(unsafe_code)]
34
35use std::cmp::Ordering;
36use std::collections::{HashMap, VecDeque};
37use std::fs::{File, OpenOptions};
38use std::io::{Read, Seek, SeekFrom};
39use std::mem::{size_of, size_of_val};
40use std::path::Path;
41use std::slice;
42use std::sync::atomic::{AtomicUsize, Ordering as Atomic};
43use std::sync::{Arc, Mutex, OnceLock};
44
45use rudb_common::bounds::{Bound, Op, scaled_as};
46use rudb_common::{Error, Field, LogicalType, PhysicalType, Result, Value};
47use rudb_encoding::{bitpack, chooser, integer, string};
48use rudb_storage::sieve::Sieve;
49use rudb_storage::{Probe, Range, Zone};
50use rudb_vector::string::StringColumn;
51use rudb_vector::validity::Validity;
52use rudb_vector::{Buffer, Chunk, Data, Packed, TextSource, Vector, search_below};
53
54mod zones;
55
56pub use zones::{Common, Stripes, distincts};
57
58const MAGIC: &[u8; 8] = b"RUDBNV10";
59const DIRECTORY: &[u8; 8] = b"RUDBDI10";
60const CATALOG: &[u8; 8] = b"RUDBCA10";
61const FORMAT: u32 = 22;
62const HEADER: u64 = 80;
63const SLOT_BYTES: usize = 28;
64const MAX_PAGE: usize = 256 * 1024 * 1024;
65const MAX_DIRECTORY: usize = 128 * 1024 * 1024;
66const FREQUENCIES: &[u8; 8] = b"RUDBFQ2\0";
67const FREQUENCY_CANDIDATES: usize = 32_768;
68const FREQUENCY_ENTRIES: usize = 512;
69const FREQUENCY_BUILD_RANK: usize = 10;
70const FREQUENCY_ORDINALS: usize = 65_536;
71const MAX_FREQUENCY_WORKERS: usize = 32;
78
79const MAX_ENCODE_WORKERS: usize = 32;
86
87const SIEVE_BUDGET: usize = 8 * 1024;
95
96const PART_BOUND_BYTES: usize = 24;
105
106fn io(error: std::io::Error) -> Error {
107 Error::io(error.to_string())
108}
109
110fn invalid(message: &str) -> Error {
111 Error::invalid_input(format!("invalid rudb native file: {message}"))
112}
113
114fn sum(counts: impl Iterator<Item = u64>) -> u64 {
116 counts.fold(0, u64::saturating_add)
117}
118
119fn span_bytes(spans: &[Span], at: usize) -> u64 {
121 spans.get(at).map_or(0, |span| u64::from(span.length))
122}
123
124fn page_bytes(pages: &[Option<Page>], at: usize) -> u64 {
126 pages.get(at).and_then(Option::as_ref).map_or(0, Page::bytes)
127}
128
129fn checksum(bytes: &[u8]) -> u64 {
139 const P1: u64 = 11_400_714_785_074_694_791;
140 const P2: u64 = 14_029_467_366_897_019_727;
141 const P3: u64 = 1_609_587_929_392_839_161;
142 const P4: u64 = 9_650_029_242_287_828_579;
143 const P5: u64 = 2_870_177_450_012_600_261;
144 let round = |state: u64, word: u64| {
145 state.wrapping_add(word.wrapping_mul(P2)).rotate_left(31).wrapping_mul(P1)
146 };
147 let merge = |state: u64, lane: u64| (state ^ round(0, lane)).wrapping_mul(P1).wrapping_add(P4);
148 let word = |chunk: &[u8]| u64::from_le_bytes(chunk.try_into().expect("eight checksum bytes"));
149
150 let mut blocks = bytes.chunks_exact(32);
153 let mut rest = blocks.remainder();
154 let mut hash = if bytes.len() >= 32 {
155 let mut one = P1.wrapping_add(P2);
156 let mut two = P2;
157 let mut three = 0;
158 let mut four = 0_u64.wrapping_sub(P1);
159 for block in blocks.by_ref() {
160 one = round(one, word(&block[..8]));
161 two = round(two, word(&block[8..16]));
162 three = round(three, word(&block[16..24]));
163 four = round(four, word(&block[24..]));
164 }
165 let combined = one
166 .rotate_left(1)
167 .wrapping_add(two.rotate_left(7))
168 .wrapping_add(three.rotate_left(12))
169 .wrapping_add(four.rotate_left(18));
170 merge(merge(merge(merge(combined, one), two), three), four)
171 } else {
172 P5
173 };
174 hash = hash.wrapping_add(bytes.len() as u64);
175 let mut words = rest.chunks_exact(8);
176 for chunk in words.by_ref() {
177 hash ^= round(0, word(chunk));
178 hash = hash.rotate_left(27).wrapping_mul(P1).wrapping_add(P4);
179 }
180 rest = words.remainder();
181 if rest.len() >= 4 {
182 let (head, tail) = rest.split_at(4);
183 let quarter = u32::from_le_bytes(head.try_into().expect("four checksum bytes"));
184 hash ^= u64::from(quarter).wrapping_mul(P1);
185 hash = hash.rotate_left(23).wrapping_mul(P2).wrapping_add(P3);
186 rest = tail;
187 }
188 for &byte in rest {
189 hash ^= u64::from(byte).wrapping_mul(P5);
190 hash = hash.rotate_left(11).wrapping_mul(P1);
191 }
192 hash ^= hash >> 33;
193 hash = hash.wrapping_mul(P2);
194 hash ^= hash >> 29;
195 hash = hash.wrapping_mul(P3);
196 hash ^ (hash >> 32)
197}
198
199#[derive(Debug, Clone, Copy)]
200struct Slot {
201 offset: u64,
202 length: u32,
203 generation: u64,
204 hash: u64,
205}
206
207impl Slot {
208 fn bytes(self) -> [u8; SLOT_BYTES] {
209 let mut result = [0; SLOT_BYTES];
210 result[..8].copy_from_slice(&self.offset.to_le_bytes());
211 result[8..12].copy_from_slice(&self.length.to_le_bytes());
212 result[12..20].copy_from_slice(&self.generation.to_le_bytes());
213 result[20..28].copy_from_slice(&self.hash.to_le_bytes());
214 result
215 }
216
217 fn read(bytes: &[u8]) -> Self {
218 Self {
219 offset: u64::from_le_bytes(bytes[..8].try_into().expect("eight bytes")),
220 length: u32::from_le_bytes(bytes[8..12].try_into().expect("four bytes")),
221 generation: u64::from_le_bytes(bytes[12..20].try_into().expect("eight bytes")),
222 hash: u64::from_le_bytes(bytes[20..28].try_into().expect("eight bytes")),
223 }
224 }
225}
226
227#[derive(Debug, Clone, Copy)]
228struct Page {
229 offset: u64,
230 length: u32,
231 hash: u64,
232}
233
234impl Page {
235 fn bytes(&self) -> u64 {
237 u64::from(self.length)
238 }
239}
240
241#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
242enum FrequencyValue {
243 Null,
244 Integer(i128),
245 Code(u32),
246}
247
248#[derive(Debug, Clone)]
249struct FrequencyEntry {
250 value: FrequencyValue,
251 count: u64,
252}
253
254#[derive(Debug, Clone)]
259struct FrequencySummary {
260 entries: Vec<FrequencyEntry>,
261 omitted_max: u64,
262 ordinals: Vec<u64>,
263}
264
265#[derive(Debug, Clone)]
270pub struct FrequencyPrefix {
271 pub entries: Vec<(Value, u64)>,
273 pub omitted_max: u64,
275}
276
277#[derive(Debug, Clone, PartialEq, Eq)]
279pub struct FrequencyOccurrences {
280 pub omitted_max: u64,
282 pub ordinals: Vec<u64>,
284}
285
286#[derive(Debug, Clone, Copy, Default)]
293struct Span {
294 offset: u64,
295 length: u32,
296}
297
298#[derive(Debug, Clone)]
300pub struct Stripe {
301 rows: usize,
302 parts: Vec<u32>,
305 index: Span,
309 pages: Vec<Span>,
310 memberships: Vec<Option<Page>>,
311 sieves: Vec<Option<Page>>,
314 part_ranges: Vec<Option<Page>>,
325 zone: Zone,
326}
327
328impl Stripe {
329 #[must_use]
331 pub fn rows(&self) -> usize {
332 self.rows
333 }
334
335 #[must_use]
337 pub fn parts(&self) -> usize {
338 self.parts.len()
339 }
340
341 #[must_use]
347 pub fn zone(&self) -> &Zone {
348 &self.zone
349 }
350}
351
352#[derive(Debug, Clone)]
354pub struct Table {
355 name: String,
356 fields: Vec<Field>,
357 stripes: Vec<Stripe>,
358 rows: usize,
359 dictionaries: Vec<Option<Page>>,
360 frequencies: Vec<Option<FrequencySummary>>,
361 distincts: Vec<Option<u64>>,
371}
372
373impl Table {
374 #[must_use]
376 pub fn name(&self) -> &str {
377 &self.name
378 }
379
380 #[must_use]
382 pub fn fields(&self) -> &[Field] {
383 &self.fields
384 }
385
386 #[must_use]
388 pub fn rows(&self) -> usize {
389 self.rows
390 }
391
392 #[must_use]
394 pub fn stripes(&self) -> &[Stripe] {
395 &self.stripes
396 }
397}
398
399#[derive(Debug, Clone)]
411struct Entry {
412 name: String,
413 fields: Vec<Field>,
414 rows: usize,
415 directory: Page,
417}
418
419#[derive(Debug, Clone)]
421pub struct ColumnLayout {
422 pub name: String,
424 pub kind: String,
426 pub pages: u64,
428 pub memberships: u64,
430 pub sieves: u64,
432 pub part_ranges: u64,
434 pub dictionary: u64,
436}
437
438impl ColumnLayout {
439 #[must_use]
441 pub fn total(&self) -> u64 {
442 self.pages
443 .saturating_add(self.memberships)
444 .saturating_add(self.sieves)
445 .saturating_add(self.part_ranges)
446 .saturating_add(self.dictionary)
447 }
448}
449
450#[derive(Debug, Clone)]
461pub struct Layout {
462 pub file: u64,
464 pub rows: usize,
466 pub stripes: usize,
468 pub parts: usize,
470 pub columns: Vec<ColumnLayout>,
472 pub indexes: u64,
475 pub directory: u64,
477 pub header: u64,
479}
480
481impl Layout {
482 #[must_use]
484 pub fn columns_total(&self) -> u64 {
485 self.columns.iter().map(ColumnLayout::total).fold(0, u64::saturating_add)
486 }
487
488 #[must_use]
494 pub fn unaccounted(&self) -> u64 {
495 self.file
496 .saturating_sub(self.columns_total())
497 .saturating_sub(self.indexes)
498 .saturating_sub(self.directory)
499 .saturating_sub(self.header)
500 }
501}
502
503#[derive(Debug)]
505struct GlobalDictionary {
506 primary: HashMap<u64, u32>,
507 collisions: HashMap<u64, Vec<u32>>,
508 offsets: Vec<u32>,
509 payload: Vec<u8>,
510 counts: Vec<u64>,
511 nulls: u64,
512}
513
514impl GlobalDictionary {
515 fn new() -> Self {
516 Self {
517 primary: HashMap::new(),
518 collisions: HashMap::new(),
519 offsets: vec![0],
520 payload: Vec::new(),
521 counts: Vec::new(),
522 nulls: 0,
523 }
524 }
525
526 fn bytes(&self, code: u32) -> Option<&[u8]> {
527 let start = *self.offsets.get(code as usize)? as usize;
528 let end = *self.offsets.get(code as usize + 1)? as usize;
529 self.payload.get(start..end)
530 }
531
532 fn code(&mut self, text: &str) -> Result<u32> {
533 let hash = checksum(text.as_bytes());
534 if let Some(&code) = self.primary.get(&hash) {
535 if self.bytes(code) == Some(text.as_bytes()) {
536 return Ok(code);
537 }
538 if let Some(codes) = self.collisions.get(&hash) {
539 if let Some(code) =
540 codes.iter().copied().find(|&code| self.bytes(code) == Some(text.as_bytes()))
541 {
542 return Ok(code);
543 }
544 }
545 let code = self.insert(text)?;
546 self.collisions.entry(hash).or_default().push(code);
547 return Ok(code);
548 }
549 let code = self.insert(text)?;
550 self.primary.insert(hash, code);
551 Ok(code)
552 }
553
554 fn insert(&mut self, text: &str) -> Result<u32> {
555 let code = u32::try_from(self.offsets.len() - 1)
556 .map_err(|_| invalid("global dictionary has too many values"))?;
557 self.payload.extend_from_slice(text.as_bytes());
558 self.offsets.push(
559 u32::try_from(self.payload.len())
560 .map_err(|_| invalid("global dictionary payload exceeds 4 GiB"))?,
561 );
562 self.counts.push(0);
563 Ok(code)
564 }
565
566 fn ranked(&self) -> Vec<(u64, u32)> {
586 let count = self.offsets.len() - 1;
587 let mut ranked = (0..count)
588 .map(|code| {
589 let code = code as u32;
590 (head(self.bytes(code).unwrap_or_default()), code)
591 })
592 .collect::<Vec<_>>();
593 ranked.sort_unstable_by(|left, right| {
594 left.0.cmp(&right.0).then_with(|| self.bytes(left.1).cmp(&self.bytes(right.1)))
595 });
596 ranked
597 }
598
599 fn observe(&mut self, code: u32, null: bool) -> Result<()> {
600 if null {
601 self.nulls = self.nulls.saturating_add(1);
602 return Ok(());
603 }
604 let count = self
605 .counts
606 .get_mut(code as usize)
607 .ok_or_else(|| invalid("global dictionary count code is out of range"))?;
608 *count = count.saturating_add(1);
609 Ok(())
610 }
611}
612
613#[derive(Debug)]
621pub struct Writer {
622 file: File,
623 at: u64,
631 table: Table,
632 generation: u64,
633 order: Vec<((u64, u64), (u64, u64))>,
636 next_order: u64,
637 dictionaries: Vec<Option<GlobalDictionary>>,
638 pending: Vec<PendingChunk>,
639 closed: Vec<Entry>,
641}
642
643#[derive(Debug)]
651struct PendingChunk {
652 order: (u64, u64),
653 chunk: Chunk,
654}
655
656#[derive(Debug)]
662struct ColumnStripe {
663 pages: Vec<Vec<u8>>,
664 codes: Vec<Option<Vec<u32>>>,
665 sieves: Vec<Option<Sieve>>,
666 ranges: Vec<Range>,
667}
668
669fn weight(ty: &LogicalType) -> usize {
677 match ty {
678 LogicalType::Varchar | LogicalType::Blob => 64,
679 LogicalType::BigInt
680 | LogicalType::UBigInt
681 | LogicalType::Timestamp
682 | LogicalType::Double
683 | LogicalType::Decimal { .. } => 8,
684 LogicalType::Integer | LogicalType::UInteger | LogicalType::Date | LogicalType::Float => 4,
685 LogicalType::SmallInt | LogicalType::USmallInt => 2,
686 _ => 1,
687 }
688}
689
690pub const STRIPE_PARTS: usize = 64;
697
698const DICTIONARY_DECIDE_ROWS: usize = 4_096;
706
707const DICTIONARY_DISTINCT_IN_TEN: usize = 9;
723
724const INDEX_ENTRY: usize = size_of::<u32>() + size_of::<u64>();
726
727fn index_section(parts: usize) -> Result<usize> {
729 parts
730 .checked_mul(INDEX_ENTRY)
731 .and_then(|bytes| bytes.checked_add(size_of::<u64>()))
732 .ok_or_else(|| invalid("index page length overflow"))
733}
734
735impl Writer {
736 pub fn open(
754 path: impl AsRef<Path>,
755 name: impl Into<String>,
756 fields: Vec<Field>,
757 ) -> Result<Self> {
758 for field in &fields {
759 type_tag(&field.ty)?;
760 }
761 let name = name.into();
762 let path = path.as_ref();
763 let (_, size, slot, bytes, _) = slot_bytes(path)?;
764 let closed = decode_catalog(&bytes, size)?;
765 if closed.iter().any(|held| held.name == name) {
766 return Err(invalid("two tables in one native file have the same name"));
767 }
768 let generation = slot
773 .generation
774 .checked_add(1)
775 .ok_or_else(|| invalid("native file generation overflow"))?;
776 let file = OpenOptions::new().write(true).read(true).open(path).map_err(io)?;
777 Ok(Self {
778 file,
779 at: size,
782 dictionaries: fields
783 .iter()
784 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
785 .collect(),
786 table: Table {
787 name,
788 dictionaries: vec![None; fields.len()],
789 distincts: vec![None; fields.len()],
790 fields,
791 stripes: Vec::new(),
792 rows: 0,
793 frequencies: Vec::new(),
794 },
795 generation,
796 order: Vec::new(),
797 next_order: 0,
798 pending: Vec::with_capacity(STRIPE_PARTS),
799 closed,
800 })
801 }
802
803 pub fn create(
809 path: impl AsRef<Path>,
810 name: impl Into<String>,
811 fields: Vec<Field>,
812 ) -> Result<Self> {
813 for field in &fields {
814 type_tag(&field.ty)?;
815 }
816 let file =
817 OpenOptions::new().write(true).read(true).create_new(true).open(path).map_err(io)?;
818 let mut header = [0; HEADER as usize];
819 header[..8].copy_from_slice(MAGIC);
820 header[8..12].copy_from_slice(&FORMAT.to_le_bytes());
821 write_at(&file, 0, &header)?;
822 Ok(Self {
823 file,
824 at: HEADER,
825 dictionaries: fields
826 .iter()
827 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
828 .collect(),
829 table: Table {
830 name: name.into(),
831 dictionaries: vec![None; fields.len()],
832 distincts: vec![None; fields.len()],
833 fields,
834 stripes: Vec::new(),
835 rows: 0,
836 frequencies: Vec::new(),
837 },
838 generation: 1,
839 order: Vec::new(),
840 next_order: 0,
841 pending: Vec::with_capacity(STRIPE_PARTS),
842 closed: Vec::new(),
843 })
844 }
845
846 pub fn next(mut self, name: impl Into<String>, fields: Vec<Field>) -> Result<Self> {
857 for field in &fields {
858 type_tag(&field.ty)?;
859 }
860 let name = name.into();
861 let entry = self.close()?;
862 if self.closed.iter().chain(std::iter::once(&entry)).any(|held| held.name == name) {
863 return Err(invalid("two tables in one native file have the same name"));
864 }
865 let Self { file, at, generation, mut closed, .. } = self;
866 closed.push(entry);
867 Ok(Self {
868 file,
869 at,
870 generation,
871 closed,
872 dictionaries: fields
873 .iter()
874 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
875 .collect(),
876 table: Table {
877 name,
878 dictionaries: vec![None; fields.len()],
879 distincts: vec![None; fields.len()],
880 fields,
881 stripes: Vec::new(),
882 rows: 0,
883 frequencies: Vec::new(),
884 },
885 order: Vec::new(),
886 next_order: 0,
887 pending: Vec::with_capacity(STRIPE_PARTS),
888 })
889 }
890
891 fn put(&mut self, bytes: &[u8]) -> Result<()> {
896 write_at(&self.file, self.at, bytes)?;
897 self.at = self
898 .at
899 .checked_add(bytes.len() as u64)
900 .ok_or_else(|| invalid("native file length overflow"))?;
901 Ok(())
902 }
903
904 pub fn append(&mut self, chunk: &Chunk) -> Result<()> {
910 let order = (self.next_order, 0);
911 self.next_order = self.next_order.saturating_add(1);
912 self.append_at(order, chunk)
913 }
914
915 pub fn append_at(&mut self, order: (u64, u64), chunk: &Chunk) -> Result<()> {
926 if chunk.is_empty() {
927 return Ok(());
928 }
929 self.admit(chunk)?;
930 if self.pending.last().is_some_and(|last| last.order > order) {
931 self.flush_pending()?;
932 }
933 self.pending.push(PendingChunk { order, chunk: chunk.clone() });
938 if self.pending.len() == STRIPE_PARTS {
939 self.flush_pending()?;
940 }
941 Ok(())
942 }
943
944 pub fn append_stripe(&mut self, parts: Vec<((u64, u64), Chunk)>) -> Result<()> {
960 if parts.len() > STRIPE_PARTS {
961 return Err(invalid("a stripe was handed more parts than it holds"));
962 }
963 self.flush_pending()?;
966 for (order, chunk) in parts {
967 if chunk.is_empty() {
968 continue;
969 }
970 self.admit(&chunk)?;
971 self.pending.push(PendingChunk { order, chunk });
972 }
973 self.flush_pending()
974 }
975
976 fn admit(&mut self, chunk: &Chunk) -> Result<()> {
978 if chunk.width() != self.table.fields.len() {
979 return Err(invalid("chunk width differs from table schema"));
980 }
981 for (index, field) in self.table.fields.iter().enumerate() {
982 if chunk.column(index)?.logical_type() != &field.ty {
983 return Err(invalid("chunk type differs from table schema"));
984 }
985 }
986 self.table.rows = self
987 .table
988 .rows
989 .checked_add(chunk.len())
990 .ok_or_else(|| invalid("row count overflow"))?;
991 Ok(())
992 }
993
994 fn encode_column(
1025 index: usize,
1026 held: &[PendingChunk],
1027 dictionary: &mut Option<GlobalDictionary>,
1028 ) -> Result<ColumnStripe> {
1029 let deciding = dictionary.as_ref().is_some_and(|held| held.offsets.len() == 1);
1032 let stripe = Self::encode_pages(index, held, dictionary.as_mut())?;
1033 if !deciding {
1034 return Ok(stripe);
1035 }
1036 let rows: usize = held.iter().map(|pending| pending.chunk.len()).sum();
1037 let distinct = dictionary.as_ref().map_or(0, |held| held.offsets.len() - 1);
1038 if rows < DICTIONARY_DECIDE_ROWS
1039 || distinct.saturating_mul(10) <= rows.saturating_mul(DICTIONARY_DISTINCT_IN_TEN)
1040 {
1041 return Ok(stripe);
1042 }
1043 *dictionary = None;
1044 Self::encode_pages(index, held, None)
1045 }
1046
1047 fn encode_pages(
1049 index: usize,
1050 held: &[PendingChunk],
1051 mut dictionary: Option<&mut GlobalDictionary>,
1052 ) -> Result<ColumnStripe> {
1053 let mut stripe = ColumnStripe {
1054 pages: Vec::with_capacity(held.len()),
1055 codes: Vec::with_capacity(held.len()),
1056 sieves: Vec::with_capacity(held.len()),
1057 ranges: Vec::with_capacity(held.len()),
1058 };
1059 for pending in held {
1060 let column = pending.chunk.column(index)?;
1061 let (bytes, unique) = encode(column, dictionary.as_deref_mut())?;
1062 if bytes.len() > MAX_PAGE {
1063 return Err(invalid("column page exceeds the configured bound"));
1064 }
1065 let range = Range::of(column);
1068 let sieve = match dictionary {
1081 Some(_) => None,
1082 None => Sieve::of(column, &range, SIEVE_BUDGET)
1083 .filter(|sieve| sieve.len() < bytes.len()),
1084 };
1085 stripe.pages.push(bytes);
1086 stripe.codes.push(unique);
1087 stripe.sieves.push(sieve);
1088 stripe.ranges.push(range);
1089 }
1090 Ok(stripe)
1091 }
1092
1093 fn encode_columns(&mut self, held: &[PendingChunk]) -> Result<Vec<ColumnStripe>> {
1102 let width = self.table.fields.len();
1103 let workers = std::thread::available_parallelism()
1104 .map_or(1, usize::from)
1105 .min(MAX_ENCODE_WORKERS)
1106 .min(width);
1107 if workers <= 1 || held.len() <= 1 {
1108 return self
1109 .dictionaries
1110 .iter_mut()
1111 .enumerate()
1112 .map(|(index, dictionary)| Self::encode_column(index, held, dictionary))
1113 .collect();
1114 }
1115 let mut jobs: Vec<(usize, Option<GlobalDictionary>)> =
1118 std::mem::take(&mut self.dictionaries).into_iter().enumerate().collect();
1119 jobs.sort_by_key(|(index, _)| weight(&self.table.fields[*index].ty));
1121 let queue = Mutex::new(jobs);
1122 let pieces = std::thread::scope(|scope| {
1123 (0..workers)
1124 .map(|_| {
1125 scope.spawn(|| {
1126 let mut mine = Vec::new();
1127 loop {
1128 let taken = queue
1129 .lock()
1130 .map_err(|_| Error::internal("a native encode worker panicked"))?
1131 .pop();
1132 let Some((index, mut dictionary)) = taken else { break };
1133 let encoded = Self::encode_column(index, held, &mut dictionary)?;
1134 mine.push((index, dictionary, encoded));
1135 }
1136 Ok(mine)
1137 })
1138 })
1139 .collect::<Vec<_>>()
1140 .into_iter()
1141 .map(|handle| {
1142 handle.join().map_err(|_| Error::internal("a native encode worker panicked"))?
1143 })
1144 .collect::<Result<Vec<_>>>()
1145 })?;
1146 let mut dictionaries: Vec<Option<GlobalDictionary>> = (0..width).map(|_| None).collect();
1147 let mut encoded: Vec<Option<ColumnStripe>> = (0..width).map(|_| None).collect();
1148 for piece in pieces {
1149 for (index, dictionary, stripe) in piece {
1150 dictionaries[index] = dictionary;
1151 encoded[index] = Some(stripe);
1152 }
1153 }
1154 self.dictionaries = dictionaries;
1155 encoded
1156 .into_iter()
1157 .map(|stripe| stripe.ok_or_else(|| Error::internal("a column was never encoded")))
1158 .collect()
1159 }
1160
1161 fn flush_pending(&mut self) -> Result<()> {
1163 if self.pending.is_empty() {
1164 return Ok(());
1165 }
1166 let width = self.table.fields.len();
1167 let mut held = std::mem::take(&mut self.pending);
1170 let parts = held.len();
1171 let encoded = self.encode_columns(&held)?;
1172 let mut pages = Vec::with_capacity(width);
1173 let mut memberships = vec![None; width];
1174 let mut ranges = Vec::with_capacity(width);
1175 let mut index = Vec::with_capacity(width.saturating_mul(index_section(parts)?));
1176 for stripe in &encoded {
1177 let offset = self.at;
1178 let section = index.len();
1179 let mut length = 0_usize;
1180 for bytes in &stripe.pages {
1181 write_at(&self.file, self.at + length as u64, bytes)?;
1182 put_u32(
1183 &mut index,
1184 u32::try_from(bytes.len()).map_err(|_| invalid("part length overflow"))?,
1185 );
1186 put_u64(&mut index, checksum(bytes));
1187 length = length
1188 .checked_add(bytes.len())
1189 .ok_or_else(|| invalid("column page length overflow"))?;
1190 }
1191 let hash = checksum(&index[section..]);
1192 put_u64(&mut index, hash);
1193 if length > MAX_PAGE {
1194 return Err(invalid("column page exceeds the configured bound"));
1195 }
1196 self.at = self
1197 .at
1198 .checked_add(length as u64)
1199 .ok_or_else(|| invalid("native file length overflow"))?;
1200 pages.push(Span {
1201 offset,
1202 length: u32::try_from(length).map_err(|_| invalid("page length overflow"))?,
1203 });
1204 ranges.push(merged_range(stripe.ranges.iter().cloned()));
1205 }
1206 for (membership, stripe) in memberships.iter_mut().zip(&encoded) {
1207 if stripe.codes.iter().all(Option::is_none) {
1208 continue;
1209 }
1210 let lists = stripe
1211 .codes
1212 .iter()
1213 .map(|codes| codes.clone().unwrap_or_default())
1214 .collect::<Vec<_>>();
1215 let bytes = encode_membership(&merged_codes(lists));
1216 let offset = self.at;
1217 self.put(&bytes)?;
1218 *membership = Some(Page {
1219 offset,
1220 length: u32::try_from(bytes.len())
1221 .map_err(|_| invalid("membership page length overflow"))?,
1222 hash: checksum(&bytes),
1223 });
1224 }
1225 let mut sieves = vec![None; width];
1226 for (page, stripe) in sieves.iter_mut().zip(&encoded) {
1227 if stripe.sieves.iter().all(Option::is_none) {
1228 continue;
1229 }
1230 let bytes = encode_sieves(stripe.sieves.iter())?;
1231 let offset = self.at;
1232 self.put(&bytes)?;
1233 *page = Some(Page {
1234 offset,
1235 length: u32::try_from(bytes.len())
1236 .map_err(|_| invalid("sieve page length overflow"))?,
1237 hash: checksum(&bytes),
1238 });
1239 }
1240 let mut part_ranges = vec![None; width];
1246 if parts > 1 {
1247 for ((page, stripe), span) in part_ranges.iter_mut().zip(&encoded).zip(&pages) {
1248 let bytes = encode_part_ranges(&stripe.ranges)?;
1249 if bytes.len() >= span.length as usize {
1250 continue;
1251 }
1252 let offset = self.at;
1253 self.put(&bytes)?;
1254 *page = Some(Page {
1255 offset,
1256 length: u32::try_from(bytes.len())
1257 .map_err(|_| invalid("part range page length overflow"))?,
1258 hash: checksum(&bytes),
1259 });
1260 }
1261 }
1262 let offset = self.at;
1263 self.put(&index)?;
1264 let index = Span {
1265 offset,
1266 length: u32::try_from(index.len())
1267 .map_err(|_| invalid("index page length overflow"))?,
1268 };
1269 let mut rows = 0_usize;
1270 let mut lengths = Vec::with_capacity(parts);
1271 let mut span = None;
1272 for pending in held.drain(..) {
1273 let part = pending.chunk.len();
1274 rows = rows.checked_add(part).ok_or_else(|| invalid("row count overflow"))?;
1275 lengths.push(u32::try_from(part).map_err(|_| invalid("part row count overflow"))?);
1276 span = Some(
1277 span.map_or((pending.order, pending.order), |(first, _)| (first, pending.order)),
1278 );
1279 }
1280 self.order.push(span.ok_or_else(|| invalid("a stripe was flushed with no parts"))?);
1281 self.table.stripes.push(Stripe {
1282 rows,
1283 parts: lengths,
1284 index,
1285 pages,
1286 memberships,
1287 sieves,
1288 part_ranges,
1289 zone: Zone::from_ranges(ranges),
1290 });
1291 self.pending = held;
1293 Ok(())
1294 }
1295
1296 fn numeric_frequency(&self, column: usize) -> Result<Option<FrequencySummary>> {
1300 let ty = &self.table.fields[column].ty;
1301 if !matches!(
1302 ty,
1303 LogicalType::TinyInt
1304 | LogicalType::SmallInt
1305 | LogicalType::Integer
1306 | LogicalType::BigInt
1307 | LogicalType::UTinyInt
1308 | LogicalType::USmallInt
1309 | LogicalType::UInteger
1310 | LogicalType::UBigInt
1311 | LogicalType::Date
1312 | LogicalType::Timestamp
1313 ) {
1314 return Ok(None);
1315 }
1316 let mut candidates: HashMap<FrequencyValue, u32> = HashMap::new();
1317 let mut decrements = 0_u64;
1318 self.visit_numeric(column, |_, value| {
1319 if let Some(count) = candidates.get_mut(&value) {
1320 *count = count.saturating_add(1);
1321 } else if candidates.len() < FREQUENCY_CANDIDATES {
1322 candidates.insert(value, 1);
1323 } else {
1324 candidates.retain(|_, count| {
1325 *count -= 1;
1326 *count != 0
1327 });
1328 decrements = decrements.saturating_add(1);
1329 }
1330 })?;
1331 let (exact, ordinals) = if decrements == 0 {
1332 (
1333 candidates
1334 .into_iter()
1335 .map(|(value, count)| (value, u64::from(count)))
1336 .collect::<HashMap<_, _>>(),
1337 Vec::new(),
1338 )
1339 } else {
1340 let mut lower = candidates.values().copied().collect::<Vec<_>>();
1341 lower.sort_unstable_by(|left, right| right.cmp(left));
1342 if lower.len() < FREQUENCY_BUILD_RANK
1343 || u64::from(lower[FREQUENCY_BUILD_RANK - 1]) <= decrements
1344 {
1345 return Ok(None);
1346 }
1347 let mut exact =
1348 candidates.into_keys().map(|value| (value, 0_u64)).collect::<HashMap<_, _>>();
1349 let mut ordinals = Vec::new();
1350 let mut exceeded = false;
1351 self.visit_numeric(column, |ordinal, value| {
1352 if let Some(count) = exact.get_mut(&value) {
1353 *count = count.saturating_add(1);
1354 if !exceeded {
1355 if ordinals.len() < FREQUENCY_ORDINALS {
1356 ordinals.push(ordinal);
1357 } else {
1358 ordinals.clear();
1359 exceeded = true;
1360 }
1361 }
1362 }
1363 })?;
1364 (exact, ordinals)
1365 };
1366 let mut entries = exact
1367 .into_iter()
1368 .map(|(value, count)| FrequencyEntry { value, count })
1369 .collect::<Vec<_>>();
1370 entries.sort_unstable_by(|left, right| {
1371 right.count.cmp(&left.count).then_with(|| frequency_order(left.value, right.value))
1372 });
1373 let omitted_max =
1374 entries.get(FREQUENCY_ENTRIES).map_or(decrements, |entry| decrements.max(entry.count));
1375 entries.truncate(FREQUENCY_ENTRIES);
1376 Ok(Some(FrequencySummary { entries, omitted_max, ordinals }))
1377 }
1378
1379 fn visit_numeric(
1380 &self,
1381 column: usize,
1382 mut visit: impl FnMut(u64, FrequencyValue),
1383 ) -> Result<()> {
1384 let ty = &self.table.fields[column].ty;
1385 let mut start = 0_u64;
1386 for stripe in &self.table.stripes {
1387 let spans = read_index(&self.file, stripe, column)?;
1388 let page = stripe.pages[column];
1389 let mut bytes = vec![0; page.length as usize];
1390 read_at(&self.file, page.offset, &mut bytes)?;
1391 for (span, &rows) in spans.iter().zip(&stripe.parts) {
1392 let part = part_bytes(&bytes, *span)?;
1393 if checksum(part) != span.hash {
1394 return Err(invalid("column page checksum differs while building frequencies"));
1395 }
1396 let rows = rows as usize;
1397 let vector = decode(ty, rows, part, None)?;
1398 for row in 0..rows {
1400 let value = if vector.is_null_at(row) {
1401 FrequencyValue::Null
1402 } else {
1403 let widened = match vector.signed_at(row) {
1407 Some(value) => Some(value),
1408 None => match vector.value_at(row) {
1409 Value::UTinyInt(value) => Some(i128::from(value)),
1410 Value::USmallInt(value) => Some(i128::from(value)),
1411 Value::UInteger(value) => Some(i128::from(value)),
1412 Value::UBigInt(value) => Some(i128::from(value)),
1413 _ => None,
1414 },
1415 };
1416 FrequencyValue::Integer(widened.ok_or_else(|| {
1417 invalid("numeric frequency page did not contain an integer value")
1418 })?)
1419 };
1420 visit(start.saturating_add(row as u64), value);
1421 }
1422 start = start.saturating_add(rows as u64);
1423 }
1424 }
1425 Ok(())
1426 }
1427
1428 fn numeric_frequencies(&self) -> Result<Vec<Option<FrequencySummary>>> {
1436 let mut columns = self
1437 .table
1438 .fields
1439 .iter()
1440 .enumerate()
1441 .filter_map(|(column, field)| {
1442 matches!(
1443 field.ty,
1444 LogicalType::TinyInt
1445 | LogicalType::SmallInt
1446 | LogicalType::Integer
1447 | LogicalType::BigInt
1448 | LogicalType::UTinyInt
1449 | LogicalType::USmallInt
1450 | LogicalType::UInteger
1451 | LogicalType::UBigInt
1452 | LogicalType::Date
1453 | LogicalType::Timestamp
1454 )
1455 .then_some(column)
1456 })
1457 .collect::<Vec<_>>();
1458 let workers = std::thread::available_parallelism()
1459 .map_or(1, usize::from)
1460 .min(MAX_FREQUENCY_WORKERS)
1461 .min(columns.len());
1462 if workers <= 1 {
1463 let mut frequencies = vec![None; self.table.fields.len()];
1464 for column in columns {
1465 frequencies[column] = self.numeric_frequency(column)?;
1466 }
1467 return Ok(frequencies);
1468 }
1469 columns.sort_by_key(|&column| weight(&self.table.fields[column].ty));
1472 let queue = Mutex::new(columns);
1473 let pieces = std::thread::scope(|scope| {
1474 (0..workers)
1475 .map(|_| {
1476 scope.spawn(|| {
1477 let mut mine = Vec::new();
1478 loop {
1479 let taken = queue
1480 .lock()
1481 .map_err(|_| Error::internal("a native frequency worker panicked"))?
1482 .pop();
1483 let Some(column) = taken else { break };
1484 mine.push((column, self.numeric_frequency(column)?));
1485 }
1486 Ok(mine)
1487 })
1488 })
1489 .collect::<Vec<_>>()
1490 .into_iter()
1491 .map(|handle| {
1492 handle
1493 .join()
1494 .map_err(|_| Error::internal("a native frequency worker panicked"))?
1495 })
1496 .collect::<Result<Vec<_>>>()
1497 })?;
1498 let mut frequencies = vec![None; self.table.fields.len()];
1499 for piece in pieces {
1500 for (column, summary) in piece {
1501 frequencies[column] = summary;
1502 }
1503 }
1504 Ok(frequencies)
1505 }
1506
1507 fn close(&mut self) -> Result<Entry> {
1518 self.flush_pending()?;
1519 let mut stripes = std::mem::take(&mut self.order)
1520 .into_iter()
1521 .zip(std::mem::take(&mut self.table.stripes))
1522 .collect::<Vec<_>>();
1523 stripes.sort_by_key(|(order, _)| order.0);
1524 let mut previous: Option<(u64, u64)> = None;
1525 for ((first, last), _) in &stripes {
1526 if previous.is_some_and(|previous| previous >= *first) {
1527 return Err(invalid("chunks did not arrive in source order"));
1528 }
1529 previous = Some(*last);
1530 }
1531 self.table.stripes = stripes.into_iter().map(|(_, stripe)| stripe).collect();
1532 self.table.frequencies = self.numeric_frequencies()?;
1533 let dictionaries = std::mem::take(&mut self.dictionaries);
1534 let orders = rankings(&dictionaries)?;
1535 for (index, (dictionary, order)) in dictionaries.into_iter().zip(orders).enumerate() {
1536 let Some(dictionary) = dictionary else { continue };
1537 self.table.distincts[index] =
1541 Some(dictionary.counts.iter().filter(|count| **count != 0).count() as u64);
1542 self.table.frequencies[index] = Some(code_frequency(&dictionary));
1543 let encoded = encode_global_dictionary(dictionary, &order)?;
1544 let offset = self.at;
1545 self.put(&encoded.index)?;
1546 self.put(&encoded.ranks)?;
1547 for block in &encoded.payload {
1548 self.put(block)?;
1549 }
1550 let payload_len =
1551 encoded.payload.iter().try_fold(0_usize, |len, block| len.checked_add(block.len()));
1552 let length = payload_len
1553 .and_then(|len| len.checked_add(encoded.index.len()))
1554 .and_then(|len| len.checked_add(encoded.ranks.len()))
1555 .ok_or_else(|| invalid("dictionary page length overflow"))?;
1556 self.table.dictionaries[index] = Some(Page {
1557 offset,
1558 length: u32::try_from(length)
1559 .map_err(|_| invalid("dictionary page length overflow"))?,
1560 hash: checksum(&encoded.index),
1561 });
1562 }
1563 let directory = encode_directory(&self.table)?;
1564 if directory.len() > MAX_DIRECTORY {
1565 return Err(invalid("directory exceeds the configured bound"));
1566 }
1567 let offset = self.at;
1568 self.put(&directory)?;
1569 Ok(Entry {
1570 name: self.table.name.clone(),
1571 fields: self.table.fields.clone(),
1572 rows: self.table.rows,
1573 directory: Page {
1574 offset,
1575 length: u32::try_from(directory.len())
1576 .map_err(|_| invalid("directory length overflow"))?,
1577 hash: checksum(&directory),
1578 },
1579 })
1580 }
1581
1582 pub fn finish(mut self) -> Result<Table> {
1592 let entry = self.close()?;
1593 let mut tables = std::mem::take(&mut self.closed);
1594 tables.push(entry);
1595 let catalog = encode_catalog(&tables)?;
1596 if catalog.len() > MAX_DIRECTORY {
1597 return Err(invalid("catalog exceeds the configured bound"));
1598 }
1599 let offset = self.at;
1600 self.put(&catalog)?;
1601 self.file.sync_all().map_err(io)?;
1605 let slot = Slot {
1606 offset,
1607 length: u32::try_from(catalog.len()).map_err(|_| invalid("catalog length overflow"))?,
1608 generation: self.generation,
1609 hash: checksum(&catalog),
1610 };
1611 write_at(&self.file, slot_offset(self.generation), &slot.bytes())?;
1616 self.file.sync_all().map_err(io)?;
1617 Ok(self.table)
1618 }
1619}
1620
1621#[derive(Debug, Clone)]
1623pub struct Reader {
1624 file: Arc<File>,
1625 table: Arc<Table>,
1626 dictionaries: Arc<Vec<OnceLock<Arc<Vector>>>>,
1627 loading: Arc<Vec<Mutex<()>>>,
1636 opened: Arc<AtomicUsize>,
1640 sieves: Arc<Vec<Vec<SieveSlot>>>,
1644 part_ranges: Arc<Vec<Vec<RangeSlot>>>,
1647 places: Arc<Vec<Place>>,
1649 cache: Arc<Vec<Mutex<Cached>>>,
1650 pages: Arc<AtomicUsize>,
1653 indexes: Arc<AtomicUsize>,
1656 kept: Arc<AtomicUsize>,
1659 size: u64,
1661 directory: u64,
1663 opening: Opening,
1665}
1666
1667#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1679pub struct Opening {
1680 pub reads: u32,
1683 pub bytes: u64,
1685}
1686
1687#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1689pub struct Reads {
1690 pub opening: Opening,
1692 pub pages: usize,
1694 pub indexes: usize,
1696 pub dictionaries: usize,
1699}
1700
1701#[derive(Debug, Clone, Copy)]
1703struct Place {
1704 stripe: u32,
1705 part: u32,
1706 rows: u32,
1707}
1708
1709#[derive(Debug, Clone, Copy)]
1711struct PartSpan {
1712 start: usize,
1713 length: usize,
1714 hash: u64,
1715}
1716
1717#[derive(Debug, Clone)]
1723struct CachedColumn {
1724 stripe: usize,
1725 index: Arc<Vec<PartSpan>>,
1726 page: Option<Arc<Vec<u8>>>,
1727}
1728
1729#[derive(Debug, Default)]
1749struct Cached {
1750 pages: Vec<Option<Arc<Vec<u8>>>>,
1751 order: VecDeque<usize>,
1752 loading: Vec<usize>,
1753 index: Vec<Option<Arc<Vec<PartSpan>>>>,
1754}
1755
1756const CACHED_STRIPES_PER_COLUMN: usize = 4;
1768
1769type SieveSlot = OnceLock<Arc<Vec<Option<Sieve>>>>;
1771
1772type RangeSlot = OnceLock<Arc<Vec<Range>>>;
1773
1774#[derive(Debug)]
1775struct NativeText {
1776 file: Arc<File>,
1777 values: usize,
1779 offsets: Vec<u8>,
1788 offset_bits: usize,
1791 ranks: usize,
1793 rank_at: u64,
1797 rank_ends: Vec<u64>,
1801 rank_hashes: Vec<u64>,
1802 rank_blocks: Vec<OnceLock<Result<Vec<u8>>>>,
1803 code_bits: usize,
1806 code_ranks: OnceLock<Option<Vec<u32>>>,
1813 payload: u64,
1814 ends: Vec<u64>,
1817 hashes: Vec<u64>,
1818 blocks: Vec<OnceLock<Result<Vec<u8>>>>,
1820 keep_budget: usize,
1823 payload_kept: AtomicUsize,
1831 searched: Mutex<HashMap<Vec<u8>, (usize, bool)>>,
1848}
1849
1850const TEXT_SEARCH_MEMO: usize = 64;
1855
1856const TEXT_PAYLOAD_VALUES: usize = 1024;
1872
1873const TEXT_KEEP_BUDGET: usize = 256 * 1024 * 1024;
1894
1895const TEXT_OFFSET_RUN: usize = 512;
1902
1903const DICTIONARY_HEADER: usize = 16;
1906
1907const TEXT_RANK_BLOCK: usize = 512;
1918
1919const RANK_BLOCK_HEADER: usize = size_of::<u64>() + 1;
1933
1934impl NativeText {
1935 fn payload_block(&self, block: usize) -> Result<Option<&[u8]>> {
1942 let Some(slot) = self.blocks.get(block) else { return Ok(None) };
1943 let bytes = slot.get_or_init(|| self.decode_block(block)).as_ref().map_err(Clone::clone)?;
1944 Ok(Some(bytes.as_slice()))
1945 }
1946
1947 fn decode_block(&self, block: usize) -> Result<Vec<u8>> {
1952 let start = if block == 0 { 0 } else { self.ends[block - 1] };
1953 let end = self.ends[block];
1954 let len = end
1955 .checked_sub(start)
1956 .ok_or_else(|| invalid("global dictionary block ends before it starts"))?;
1957 let mut stored = vec![
1958 0;
1959 usize::try_from(len).map_err(|_| invalid(
1960 "global dictionary block does not fit in memory"
1961 ))?
1962 ];
1963 read_at(&self.file, self.payload + start, &mut stored)?;
1964 if checksum(&stored) != self.hashes[block] {
1965 return Err(invalid("global dictionary payload checksum differs"));
1966 }
1967 let first = block * TEXT_PAYLOAD_VALUES;
1968 let last = (first + TEXT_PAYLOAD_VALUES).min(self.values);
1969 let want = self.end_within(last - 1)? as usize;
1970 let values = string::decode_flat(&stored)?;
1971 if values.len() != last - first {
1972 return Err(invalid("global dictionary block holds the wrong value count"));
1973 }
1974 let bytes = values.into_bytes();
1975 if bytes.len() != want {
1976 return Err(invalid("global dictionary block decodes to the wrong length"));
1977 }
1978 Ok(bytes)
1979 }
1980
1981 fn end_within(&self, index: usize) -> Result<u32> {
1983 let run = index / TEXT_OFFSET_RUN;
1984 let bytes = self
1985 .offsets
1986 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
1987 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1988 let end = bitpack::tail_at(bytes, self.offset_bits, index % TEXT_OFFSET_RUN)
1989 .map_err(|_| invalid("global dictionary offsets are short"))?;
1990 u32::try_from(end).map_err(|_| invalid("global dictionary offset is past the payload"))
1991 }
1992
1993 fn ends_within(&self, first: usize, last: usize) -> Result<Vec<u64>> {
2006 let mut ends = Vec::with_capacity(last.saturating_sub(first));
2007 let mut at = first;
2008 while at < last {
2009 let run = at / TEXT_OFFSET_RUN;
2010 let stop = ((run + 1) * TEXT_OFFSET_RUN).min(last);
2011 let held = self.values.saturating_sub(run * TEXT_OFFSET_RUN).min(TEXT_OFFSET_RUN);
2012 let bytes = self
2013 .offsets
2014 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
2015 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
2016 let run_ends = bitpack::unpack_tail(bytes, self.offset_bits, held)
2017 .map_err(|_| invalid("global dictionary offsets are short"))?;
2018 let within = run_ends
2019 .get(at % TEXT_OFFSET_RUN..stop - run * TEXT_OFFSET_RUN)
2020 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
2021 ends.extend_from_slice(within);
2022 at = stop;
2023 }
2024 Ok(ends)
2025 }
2026
2027 fn start_within(&self, index: usize) -> Result<u32> {
2030 if index % TEXT_PAYLOAD_VALUES == 0 { Ok(0) } else { self.end_within(index - 1) }
2031 }
2032
2033 fn span_within(&self, index: usize) -> Result<(u32, u32)> {
2041 let within = index % TEXT_OFFSET_RUN;
2042 let (start, end) = if within == 0 {
2043 (self.start_within(index)?, self.end_within(index)?)
2044 } else {
2045 let run = index / TEXT_OFFSET_RUN;
2046 let bytes = self
2047 .offsets
2048 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
2049 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
2050 let (start, end) = bitpack::tail_pair(bytes, self.offset_bits, within)
2051 .map_err(|_| invalid("global dictionary offsets are short"))?;
2052 let ends = u32::try_from(end)
2053 .map_err(|_| invalid("global dictionary offset is past the payload"))?;
2054 let starts = u32::try_from(start)
2055 .map_err(|_| invalid("global dictionary offset is past the payload"))?;
2056 (starts, ends)
2057 };
2058 if start > end {
2059 return Err(invalid("global dictionary value ends before it starts"));
2060 }
2061 Ok((start, end))
2062 }
2063
2064 fn rank_parts(&self, rank: usize) -> Result<(&[u8], usize)> {
2071 let slot = self
2072 .rank_blocks
2073 .get(rank / TEXT_RANK_BLOCK)
2074 .ok_or_else(|| invalid("global dictionary rank is past the order"))?;
2075 let block = slot
2076 .get_or_init(|| {
2077 let which = rank / TEXT_RANK_BLOCK;
2078 let start = if which == 0 { 0 } else { self.rank_ends[which - 1] };
2079 let end = self.rank_ends[which];
2080 let mut bytes = vec![0; (end - start) as usize];
2081 read_at(&self.file, self.rank_at + start, &mut bytes)?;
2082 if checksum(&bytes)
2083 != *self
2084 .rank_hashes
2085 .get(rank / TEXT_RANK_BLOCK)
2086 .ok_or_else(|| invalid("global dictionary rank block has no checksum"))?
2087 {
2088 return Err(invalid("global dictionary rank checksum differs"));
2089 }
2090 Ok(bytes)
2091 })
2092 .as_ref()
2093 .map_err(Clone::clone)?;
2094 Ok((block.as_slice(), rank % TEXT_RANK_BLOCK))
2095 }
2096
2097 fn head_at(&self, rank: usize) -> Result<u64> {
2099 let (block, within) = self.rank_parts(rank)?;
2100 let (base, width, packed) = rank_heads(block)?;
2101 let above = bitpack::tail_at(packed, width, within)
2102 .map_err(|_| invalid("global dictionary rank block is short of heads"))?;
2103 Ok(base.wrapping_add(above))
2104 }
2105
2106 fn rank_codes<'block>(&self, block: &'block [u8], count: usize) -> Result<&'block [u8]> {
2108 let (_, width, packed) = rank_heads(block)?;
2109 packed
2110 .get(bitpack::tail_len(count, width)..)
2111 .ok_or_else(|| invalid("global dictionary rank block is short of codes"))
2112 }
2113
2114 fn rank_block_len(&self, rank: usize) -> usize {
2116 let first = rank / TEXT_RANK_BLOCK * TEXT_RANK_BLOCK;
2117 TEXT_RANK_BLOCK.min(self.ranks - first)
2118 }
2119}
2120
2121fn rank_heads(block: &[u8]) -> Result<(u64, usize, &[u8])> {
2123 let header = block
2124 .get(..RANK_BLOCK_HEADER)
2125 .ok_or_else(|| invalid("global dictionary rank block is short"))?;
2126 let base = u64::from_le_bytes(header[..8].try_into().expect("eight bytes"));
2127 let width = header[8] as usize;
2128 if width > 64 {
2129 return Err(invalid("global dictionary rank block packs heads past a word"));
2130 }
2131 Ok((base, width, &block[RANK_BLOCK_HEADER..]))
2132}
2133
2134fn offset_width(offsets: &[u32]) -> usize {
2141 let values = offsets.len() - 1;
2142 let mut span = 0;
2143 for first in (0..values).step_by(TEXT_PAYLOAD_VALUES) {
2144 let last = (first + TEXT_PAYLOAD_VALUES).min(values);
2145 span = span.max(offsets[last] - offsets[first]);
2146 }
2147 (u32::BITS - span.leading_zeros()) as usize
2148}
2149
2150fn offset_bytes(values: usize, bits: usize) -> usize {
2153 let full = values / TEXT_OFFSET_RUN;
2154 let rest = values % TEXT_OFFSET_RUN;
2155 full * TEXT_OFFSET_RUN / 8 * bits + bitpack::tail_len(rest, bits)
2156}
2157
2158fn encode_offsets(offsets: &[u32], bits: usize, out: &mut Vec<u8>) -> Result<()> {
2160 let values = offsets.len() - 1;
2161 let mut run = Vec::with_capacity(TEXT_OFFSET_RUN);
2162 for first in (0..values).step_by(TEXT_OFFSET_RUN) {
2163 let last = (first + TEXT_OFFSET_RUN).min(values);
2164 let base = offsets[first / TEXT_PAYLOAD_VALUES * TEXT_PAYLOAD_VALUES];
2165 run.clear();
2166 run.extend((first..last).map(|value| u64::from(offsets[value + 1] - base)));
2167 bitpack::pack_tail(&run, bits, out)
2168 .map_err(|_| invalid("global dictionary offsets do not pack"))?;
2169 }
2170 Ok(())
2171}
2172
2173fn code_width(values: usize) -> usize {
2175 match u64::try_from(values).unwrap_or(u64::MAX) {
2176 0 | 1 => 0,
2177 last => (u64::BITS - (last - 1).leading_zeros()) as usize,
2178 }
2179}
2180
2181impl TextSource for NativeText {
2182 fn len(&self) -> usize {
2183 self.values
2184 }
2185
2186 fn bytes_at(&self, index: usize) -> Result<Option<&[u8]>> {
2187 if index >= self.values {
2188 return Ok(None);
2189 }
2190 let (start, end) = self.span_within(index)?;
2191 if start == end {
2192 return Ok(Some(&[]));
2193 }
2194 let block = index / TEXT_PAYLOAD_VALUES;
2197 let Some(bytes) = self.payload_block(block)? else { return Ok(None) };
2198 Ok(bytes.get(start as usize..end as usize))
2199 }
2200
2201 fn bytes_len_at(&self, index: usize) -> Result<Option<usize>> {
2202 if index >= self.values {
2203 return Ok(None);
2204 }
2205 let (start, end) = self.span_within(index)?;
2206 Ok(Some((end - start) as usize))
2207 }
2208
2209 fn sweep(
2222 &self,
2223 first: usize,
2224 limit: usize,
2225 body: &mut dyn FnMut(usize, &[u8]) -> Result<()>,
2226 ) -> Result<usize> {
2227 let limit = limit.min(self.values);
2228 if first >= limit {
2229 return Ok(first);
2230 }
2231 let block = first / TEXT_PAYLOAD_VALUES;
2232 let last = ((block + 1) * TEXT_PAYLOAD_VALUES).min(limit);
2233 let decoded;
2234 let bytes: &[u8] = match self.blocks.get(block).and_then(OnceLock::get) {
2235 Some(Ok(kept)) => kept,
2236 _ if self.payload_kept.load(Atomic::Relaxed) < self.keep_budget => {
2237 let kept = self
2238 .payload_block(block)?
2239 .ok_or_else(|| invalid("global dictionary block is past the payload"))?;
2240 self.payload_kept.fetch_add(kept.len(), Atomic::Relaxed);
2241 kept
2242 }
2243 _ => {
2244 decoded = self.decode_block(block)?;
2245 &decoded
2246 }
2247 };
2248 let ends = self.ends_within(first, last)?;
2249 if ends.len() != last - first {
2250 return Err(invalid("global dictionary offsets are short"));
2251 }
2252 let mut start = u64::from(self.start_within(first)?);
2253 for (index, &end) in (first..last).zip(&ends) {
2256 let value = usize::try_from(start)
2257 .ok()
2258 .zip(usize::try_from(end).ok())
2259 .and_then(|(from, to)| bytes.get(from..to))
2260 .ok_or_else(|| invalid("global dictionary value is past its block"))?;
2261 body(index, value)?;
2262 start = end;
2263 }
2264 Ok(last)
2265 }
2266
2267 fn ranks(&self) -> Option<usize> {
2268 (self.ranks > 0).then_some(self.ranks)
2269 }
2270
2271 fn below(&self, ranks: usize, wanted: &[u8]) -> Result<(usize, bool)> {
2279 let mut memo = self.searched.lock().map_err(|_| invalid("a poisoned dictionary search"))?;
2280 if let Some(&answer) = memo.get(wanted) {
2281 return Ok(answer);
2282 }
2283 let answer = search_below(self, ranks, wanted)?;
2284 if memo.len() >= TEXT_SEARCH_MEMO {
2285 memo.clear();
2286 }
2287 memo.insert(wanted.to_vec(), answer);
2288 Ok(answer)
2289 }
2290
2291 fn compare_rank(&self, rank: usize, wanted: &[u8]) -> Result<Ordering> {
2292 let settled = self.head_at(rank)?.cmp(&head(wanted));
2296 if settled != Ordering::Equal {
2297 return Ok(settled);
2298 }
2299 let code = self.code_at_rank(rank)?;
2300 let bytes = self
2301 .bytes_at(code as usize)?
2302 .ok_or_else(|| invalid("global dictionary order names a code it does not have"))?;
2303 Ok(bytes.cmp(wanted))
2304 }
2305
2306 fn code_at_rank(&self, rank: usize) -> Result<u32> {
2307 let (block, within) = self.rank_parts(rank)?;
2308 let codes = self.rank_codes(block, self.rank_block_len(rank))?;
2309 let code = bitpack::tail_at(codes, self.code_bits, within)
2310 .map_err(|_| invalid("global dictionary rank block is short of codes"))?;
2311 let code = u32::try_from(code)
2312 .map_err(|_| invalid("global dictionary order names a code it does not have"))?;
2313 if code as usize >= self.len() {
2314 return Err(invalid("global dictionary order names a code it does not have"));
2315 }
2316 Ok(code)
2317 }
2318
2319 fn code_ranks(&self) -> Option<&[u32]> {
2320 if self.ranks == 0 || self.ranks != self.len() {
2324 return None;
2325 }
2326 self.code_ranks
2327 .get_or_init(|| {
2328 let mut ranks = vec![u32::MAX; self.ranks];
2329 for first in (0..self.ranks).step_by(TEXT_RANK_BLOCK) {
2332 let (block, _) = self.rank_parts(first).ok()?;
2333 let count = self.rank_block_len(first);
2334 let codes = self.rank_codes(block, count).ok()?;
2335 for (within, code) in bitpack::unpack_tail(codes, self.code_bits, count)
2336 .ok()?
2337 .into_iter()
2338 .enumerate()
2339 {
2340 let code = usize::try_from(code).ok()?;
2341 *ranks.get_mut(code)? = u32::try_from(first + within).ok()?;
2342 }
2343 }
2344 if ranks.contains(&u32::MAX) {
2345 return None;
2346 }
2347 Some(ranks)
2348 })
2349 .as_deref()
2350 }
2351
2352 fn footprint(&self) -> usize {
2353 self.offsets.capacity()
2354 + self
2355 .code_ranks
2356 .get()
2357 .and_then(Option::as_ref)
2358 .map_or(0, |ranks| ranks.capacity() * size_of::<u32>())
2359 + self.rank_hashes.capacity() * size_of::<u64>()
2360 + self.rank_ends.capacity() * size_of::<u64>()
2361 + self.rank_blocks.capacity() * size_of::<OnceLock<Result<Vec<u8>>>>()
2362 + self
2363 .rank_blocks
2364 .iter()
2365 .filter_map(OnceLock::get)
2366 .filter_map(|result| result.as_ref().ok())
2367 .map(Vec::capacity)
2368 .sum::<usize>()
2369 + self.blocks.capacity() * size_of::<OnceLock<Result<Vec<u8>>>>()
2370 + self.hashes.capacity() * size_of::<u64>()
2371 + self.ends.capacity() * size_of::<u64>()
2372 + self
2373 .blocks
2374 .iter()
2375 .filter_map(OnceLock::get)
2376 .filter_map(|result| result.as_ref().ok())
2377 .map(Vec::capacity)
2378 .sum::<usize>()
2379 }
2380}
2381
2382fn places(table: &Table) -> Result<Vec<Place>> {
2384 let mut places = Vec::with_capacity(table.stripes.len().saturating_mul(STRIPE_PARTS));
2385 for (at, stripe) in table.stripes.iter().enumerate() {
2386 let index = u32::try_from(at).map_err(|_| invalid("too many stripes"))?;
2387 for (part, &rows) in stripe.parts.iter().enumerate() {
2388 places.push(Place {
2389 stripe: index,
2390 part: u32::try_from(part).map_err(|_| invalid("too many parts in a stripe"))?,
2391 rows,
2392 });
2393 }
2394 }
2395 Ok(places)
2396}
2397
2398fn read_index(file: &File, stripe: &Stripe, column: usize) -> Result<Vec<PartSpan>> {
2403 let parts = stripe.parts.len();
2404 let section = index_section(parts)?;
2405 let at = column.checked_mul(section).ok_or_else(|| invalid("index page offset overflow"))?;
2406 let end = at.checked_add(section).ok_or_else(|| invalid("index page offset overflow"))?;
2407 if end > stripe.index.length as usize {
2408 return Err(invalid("index page is shorter than its columns"));
2409 }
2410 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
2411 let mut bytes = vec![0; section];
2412 let offset = stripe
2413 .index
2414 .offset
2415 .checked_add(at as u64)
2416 .ok_or_else(|| invalid("index page offset overflow"))?;
2417 read_at(file, offset, &mut bytes)?;
2418 let entries = section - size_of::<u64>();
2419 let stored = u64::from_le_bytes(bytes[entries..].try_into().expect("eight bytes"));
2420 if checksum(&bytes[..entries]) != stored {
2421 return Err(invalid(&format!(
2424 "index page section checksum differs, column {column} of {parts} parts at {offset}, \
2425 wanted {stored:016x} and got {:016x}",
2426 checksum(&bytes[..entries]),
2427 )));
2428 }
2429 let mut spans = Vec::with_capacity(parts);
2430 let mut start = 0_usize;
2431 for part in 0..parts {
2432 let at = part * INDEX_ENTRY;
2433 let length = u32::from_le_bytes(bytes[at..at + 4].try_into().expect("four bytes")) as usize;
2434 let hash = u64::from_le_bytes(bytes[at + 4..at + 12].try_into().expect("eight bytes"));
2435 spans.push(PartSpan { start, length, hash });
2436 start = start.checked_add(length).ok_or_else(|| invalid("column page length overflow"))?;
2437 }
2438 if start != page.length as usize {
2439 return Err(invalid("column page length differs from its index"));
2440 }
2441 Ok(spans)
2442}
2443
2444fn part_bytes(page: &[u8], span: PartSpan) -> Result<&[u8]> {
2446 let end = span.start.checked_add(span.length).ok_or_else(|| invalid("part range overflow"))?;
2447 page.get(span.start..end).ok_or_else(|| invalid("part exceeds its column page"))
2448}
2449
2450fn remember(cached: &mut Cached, held: &CachedColumn, kept: usize) {
2455 if let Some(slot) = cached.index.get_mut(held.stripe) {
2456 if slot.is_none() {
2457 *slot = Some(Arc::clone(&held.index));
2458 }
2459 }
2460 let Some(page) = held.page.clone() else { return };
2461 let Some(slot) = cached.pages.get_mut(held.stripe) else { return };
2462 if slot.is_none() {
2463 cached.order.push_back(held.stripe);
2464 }
2465 *slot = Some(page);
2466 while cached.order.len() > kept.max(1) {
2467 let Some(oldest) = cached.order.pop_front() else { break };
2468 if let Some(slot) = cached.pages.get_mut(oldest) {
2469 *slot = None;
2470 }
2471 }
2472}
2473
2474#[derive(Debug, Clone)]
2483pub struct Catalog {
2484 file: Arc<File>,
2485 size: u64,
2486 entries: Arc<Vec<Entry>>,
2487 opening: Opening,
2488}
2489
2490impl Catalog {
2491 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
2497 let (file, size, _, bytes, opening) = slot_bytes(path)?;
2498 let entries = decode_catalog(&bytes, size)?;
2499 Ok(Self { file: Arc::new(file), size, entries: Arc::new(entries), opening })
2500 }
2501
2502 pub fn names(&self) -> impl ExactSizeIterator<Item = &str> {
2504 self.entries.iter().map(|entry| entry.name.as_str())
2505 }
2506
2507 #[must_use]
2509 pub fn len(&self) -> usize {
2510 self.entries.len()
2511 }
2512
2513 #[must_use]
2515 pub fn is_empty(&self) -> bool {
2516 self.entries.is_empty()
2517 }
2518
2519 pub fn table(&self, name: &str) -> Result<Reader> {
2525 let entry = self
2526 .entries
2527 .iter()
2528 .find(|entry| entry.name == name)
2529 .ok_or_else(|| invalid(&format!("the file holds no table called {name}")))?;
2530 let mut bytes = vec![0; entry.directory.length as usize];
2531 read_at(&self.file, entry.directory.offset, &mut bytes)?;
2532 if checksum(&bytes) != entry.directory.hash {
2533 return Err(invalid(&format!("the directory of table {name} does not checksum")));
2534 }
2535 let mut opening = self.opening;
2536 opening.reads += 1;
2537 opening.bytes += u64::from(entry.directory.length);
2538 Reader::build(
2539 Arc::clone(&self.file),
2540 self.size,
2541 decode_directory(&bytes, self.size)?,
2542 u64::from(entry.directory.length),
2543 opening,
2544 )
2545 }
2546}
2547
2548fn slot_offset(generation: u64) -> u64 {
2553 16 + (generation - 1) % 2 * SLOT_BYTES as u64
2554}
2555
2556fn slot_bytes(path: impl AsRef<Path>) -> Result<(File, u64, Slot, Vec<u8>, Opening)> {
2561 let mut file = File::open(path).map_err(io)?;
2562 let size = file.metadata().map_err(io)?.len();
2563 if size < HEADER {
2564 return Err(invalid("file is shorter than its header"));
2565 }
2566 let mut header = [0; HEADER as usize];
2567 file.read_exact(&mut header).map_err(io)?;
2568 let mut opening = Opening { reads: 1, bytes: HEADER };
2569 let version = u32::from_le_bytes([header[8], header[9], header[10], header[11]]);
2570 if &header[..8] != MAGIC {
2575 return Err(invalid("the header does not begin with a rudb native magic"));
2576 }
2577 if version != FORMAT {
2578 return Err(invalid(&format!(
2579 "the file is format {version} and this build reads format {FORMAT}, so it has to \
2580 be written again"
2581 )));
2582 }
2583 let mut selected = None;
2584 for start in [16, 16 + SLOT_BYTES] {
2585 let slot = Slot::read(&header[start..start + SLOT_BYTES]);
2586 if slot.generation == 0 || slot.length == 0 || slot.length as usize > MAX_DIRECTORY {
2587 continue;
2588 }
2589 let Some(end) = slot.offset.checked_add(u64::from(slot.length)) else { continue };
2590 if slot.offset < HEADER || end > size {
2591 continue;
2592 }
2593 let mut bytes = vec![0; slot.length as usize];
2594 file.seek(SeekFrom::Start(slot.offset)).map_err(io)?;
2595 file.read_exact(&mut bytes).map_err(io)?;
2596 opening.reads += 1;
2597 opening.bytes += u64::from(slot.length);
2598 if checksum(&bytes) == slot.hash
2599 && selected
2600 .as_ref()
2601 .is_none_or(|(old, _): &(Slot, Vec<u8>)| old.generation < slot.generation)
2602 {
2603 selected = Some((slot, bytes));
2604 }
2605 }
2606 let (slot, bytes) = selected.ok_or_else(|| invalid("no committed directory slot is valid"))?;
2607 Ok((file, size, slot, bytes, opening))
2608}
2609
2610impl Reader {
2611 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
2618 let catalog = Catalog::open(path)?;
2619 let mut names = catalog.names();
2620 let name = names.next().ok_or_else(|| invalid("the file holds no table"))?.to_string();
2621 if names.next().is_some() {
2622 return Err(invalid(
2623 "the file holds more than one table, so it has to be opened by name",
2624 ));
2625 }
2626 catalog.table(&name)
2627 }
2628
2629 fn build(
2631 file: Arc<File>,
2632 size: u64,
2633 table: Table,
2634 directory: u64,
2635 opening: Opening,
2636 ) -> Result<Self> {
2637 let places = places(&table)?;
2638 let dictionaries = (0..table.fields.len()).map(|_| OnceLock::new()).collect();
2639 let table_fields = table.fields.len();
2640 let stripes = table.stripes.len();
2641 let cache = (0..table.fields.len())
2642 .map(|_| {
2643 Mutex::new(Cached {
2644 pages: (0..stripes).map(|_| None).collect(),
2645 index: (0..stripes).map(|_| None).collect(),
2646 ..Cached::default()
2647 })
2648 })
2649 .collect::<Vec<_>>();
2650 let sieves: Vec<Vec<SieveSlot>> = (0..table.fields.len())
2651 .map(|_| table.stripes.iter().map(|_| OnceLock::new()).collect())
2652 .collect();
2653 let part_ranges: Vec<Vec<RangeSlot>> = (0..table.fields.len())
2654 .map(|_| table.stripes.iter().map(|_| OnceLock::new()).collect())
2655 .collect();
2656 Ok(Self {
2657 file,
2658 table: Arc::new(table),
2659 dictionaries: Arc::new(dictionaries),
2660 loading: Arc::new((0..table_fields).map(|_| Mutex::new(())).collect()),
2661 opened: Arc::new(AtomicUsize::new(0)),
2662 sieves: Arc::new(sieves),
2663 part_ranges: Arc::new(part_ranges),
2664 places: Arc::new(places),
2665 cache: Arc::new(cache),
2666 pages: Arc::new(AtomicUsize::new(0)),
2667 indexes: Arc::new(AtomicUsize::new(0)),
2668 kept: Arc::new(AtomicUsize::new(CACHED_STRIPES_PER_COLUMN)),
2669 size,
2670 directory,
2671 opening,
2672 })
2673 }
2674
2675 #[must_use]
2682 pub fn reads(&self) -> Reads {
2683 Reads {
2684 opening: self.opening,
2685 pages: self.pages.load(Atomic::Relaxed),
2686 indexes: self.indexes.load(Atomic::Relaxed),
2687 dictionaries: self.opened.load(Atomic::Relaxed),
2688 }
2689 }
2690
2691 #[must_use]
2696 pub fn layout(&self) -> Layout {
2697 let table = &self.table;
2698 let stripes = table.stripes.as_slice();
2699 let columns = table
2700 .fields
2701 .iter()
2702 .enumerate()
2703 .map(|(at, field)| ColumnLayout {
2704 name: field.name.clone(),
2705 kind: field.ty.to_string(),
2706 pages: sum(stripes.iter().map(|stripe| span_bytes(&stripe.pages, at))),
2707 memberships: sum(stripes.iter().map(|stripe| page_bytes(&stripe.memberships, at))),
2708 sieves: sum(stripes.iter().map(|stripe| page_bytes(&stripe.sieves, at))),
2709 part_ranges: sum(stripes.iter().map(|stripe| page_bytes(&stripe.part_ranges, at))),
2710 dictionary: page_bytes(&table.dictionaries, at),
2711 })
2712 .collect();
2713 Layout {
2714 file: self.size,
2715 rows: table.rows,
2716 stripes: stripes.len(),
2717 parts: self.places.len(),
2718 columns,
2719 indexes: sum(stripes.iter().map(|stripe| u64::from(stripe.index.length))),
2720 directory: self.directory,
2721 header: HEADER,
2722 }
2723 }
2724
2725 #[must_use]
2727 pub fn parts(&self) -> usize {
2728 self.places.len()
2729 }
2730
2731 #[must_use]
2738 pub fn stripe_parts(&self) -> Vec<std::ops::Range<usize>> {
2739 let mut runs = Vec::with_capacity(self.table.stripes.len());
2740 let mut start = 0;
2741 for stripe in &self.table.stripes {
2742 let end = start + stripe.parts.len();
2743 runs.push(start..end);
2744 start = end;
2745 }
2746 runs
2747 }
2748
2749 #[must_use]
2754 pub fn stripe_rows(&self, stripe: usize) -> usize {
2755 self.table.stripes.get(stripe).map_or(0, |held| held.rows)
2756 }
2757
2758 pub fn keep_stripes(&self, stripes: usize) {
2765 self.kept.fetch_max(stripes, Atomic::Relaxed);
2766 }
2767
2768 #[must_use]
2770 pub fn part_rows(&self, at: usize) -> usize {
2771 self.places.get(at).map_or(0, |place| place.rows as usize)
2772 }
2773
2774 #[must_use]
2776 pub fn table(&self) -> &Table {
2777 &self.table
2778 }
2779
2780 pub fn top_frequencies(&self, column: usize, top: usize) -> Result<Option<Vec<(Value, u64)>>> {
2789 let field = self
2790 .table
2791 .fields
2792 .get(column)
2793 .ok_or_else(|| invalid("frequency column index out of range"))?;
2794 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2795 return Ok(None);
2796 };
2797 if top == 0 || summary.entries.len() < top {
2798 return Ok(None);
2799 }
2800 let boundary = summary.entries[top - 1].count;
2801 if boundary <= summary.omitted_max {
2802 return Ok(None);
2803 }
2804 self.decode_frequencies(column, &field.ty, &summary.entries).map(Some)
2805 }
2806
2807 pub fn exact_frequencies(&self, column: usize) -> Result<Option<Vec<(Value, u64)>>> {
2827 let Some(prefix) = self.frequency_prefix(column)? else {
2828 return Ok(None);
2829 };
2830 Ok((prefix.omitted_max == 0).then_some(prefix.entries))
2831 }
2832
2833 pub fn frequency_prefix(&self, column: usize) -> Result<Option<FrequencyPrefix>> {
2856 let field = self
2857 .table
2858 .fields
2859 .get(column)
2860 .ok_or_else(|| invalid("frequency column index out of range"))?;
2861 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2862 return Ok(None);
2863 };
2864 let entries = self.decode_frequencies(column, &field.ty, &summary.entries)?;
2865 Ok(Some(FrequencyPrefix { entries, omitted_max: summary.omitted_max }))
2866 }
2867
2868 fn decode_frequencies(
2870 &self,
2871 column: usize,
2872 ty: &LogicalType,
2873 entries: &[FrequencyEntry],
2874 ) -> Result<Vec<(Value, u64)>> {
2875 let dictionary = if *ty == LogicalType::Varchar { self.dictionary(column)? } else { None };
2876 let mut out = Vec::with_capacity(entries.len());
2877 for entry in entries {
2878 let value = match entry.value {
2879 FrequencyValue::Null => Value::Null,
2880 FrequencyValue::Integer(value) => match *ty {
2881 LogicalType::TinyInt => Value::TinyInt(
2882 i8::try_from(value)
2883 .map_err(|_| invalid("frequency TINYINT is out of range"))?,
2884 ),
2885 LogicalType::UTinyInt => Value::UTinyInt(
2886 u8::try_from(value)
2887 .map_err(|_| invalid("frequency UTINYINT is out of range"))?,
2888 ),
2889 LogicalType::USmallInt => Value::USmallInt(
2890 u16::try_from(value)
2891 .map_err(|_| invalid("frequency USMALLINT is out of range"))?,
2892 ),
2893 LogicalType::UInteger => Value::UInteger(
2894 u32::try_from(value)
2895 .map_err(|_| invalid("frequency UINTEGER is out of range"))?,
2896 ),
2897 LogicalType::UBigInt => Value::UBigInt(
2898 u64::try_from(value)
2899 .map_err(|_| invalid("frequency UBIGINT is out of range"))?,
2900 ),
2901 LogicalType::SmallInt => Value::SmallInt(
2902 i16::try_from(value)
2903 .map_err(|_| invalid("frequency SMALLINT is out of range"))?,
2904 ),
2905 LogicalType::Integer => Value::Integer(
2906 i32::try_from(value)
2907 .map_err(|_| invalid("frequency INTEGER is out of range"))?,
2908 ),
2909 LogicalType::BigInt => Value::BigInt(
2910 i64::try_from(value)
2911 .map_err(|_| invalid("frequency BIGINT is out of range"))?,
2912 ),
2913 LogicalType::Date => Value::Date(
2914 i32::try_from(value)
2915 .map_err(|_| invalid("frequency DATE is out of range"))?,
2916 ),
2917 LogicalType::Timestamp => Value::Timestamp(
2918 i64::try_from(value)
2919 .map_err(|_| invalid("frequency TIMESTAMP is out of range"))?,
2920 ),
2921 _ => return Err(invalid("integer frequency belongs to another type")),
2922 },
2923 FrequencyValue::Code(code) => dictionary
2924 .as_ref()
2925 .ok_or_else(|| invalid("frequency code has no dictionary"))?
2926 .try_value_at(code as usize)?,
2927 };
2928 out.push((value, entry.count));
2929 }
2930 Ok(out)
2931 }
2932
2933 pub fn frequency_occurrences(&self, column: usize) -> Result<Option<FrequencyOccurrences>> {
2943 self.table
2944 .fields
2945 .get(column)
2946 .ok_or_else(|| invalid("frequency column index out of range"))?;
2947 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2948 return Ok(None);
2949 };
2950 if summary.ordinals.is_empty() {
2951 return Ok(None);
2952 }
2953 Ok(Some(FrequencyOccurrences {
2954 omitted_max: summary.omitted_max,
2955 ordinals: summary.ordinals.clone(),
2956 }))
2957 }
2958
2959 pub fn distinct_values(&self, column: usize) -> Result<Option<u64>> {
2983 self.table
2984 .distincts
2985 .get(column)
2986 .copied()
2987 .ok_or_else(|| invalid("distinct column index out of range"))
2988 }
2989
2990 pub fn null_count(&self, column: usize) -> Result<u64> {
3001 if column >= self.table.fields.len() {
3002 return Err(invalid("null count column index out of range"));
3003 }
3004 let mut nulls = 0_u64;
3005 for stripe in &self.table.stripes {
3006 let range = stripe
3007 .zone
3008 .column(column)
3009 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
3010 nulls = nulls
3011 .checked_add(range.nulls as u64)
3012 .ok_or_else(|| invalid("null count overflow"))?;
3013 }
3014 Ok(nulls)
3015 }
3016
3017 pub fn text_extremes(&self, column: usize) -> Result<Option<(Value, Value)>> {
3032 if self.null_count(column)? > 0 {
3033 return Ok(None);
3034 }
3035 let Some(dictionary) = self.dictionary(column)? else { return Ok(None) };
3036 let Some(ranks) = dictionary.ranks() else { return Ok(None) };
3037 if ranks == 0 {
3038 return Ok(None);
3039 }
3040 let low = text_at_rank(&dictionary, 0)?;
3041 let high = text_at_rank(&dictionary, ranks - 1)?;
3042 Ok(Some((low, high)))
3043 }
3044
3045 pub fn exact_extremes(&self, column: usize) -> Result<Option<(Bound, Bound)>> {
3068 if column >= self.table.fields.len() {
3069 return Err(invalid("extremes column index out of range"));
3070 }
3071 let mut low: Option<Bound> = None;
3072 let mut high: Option<Bound> = None;
3073 for stripe in &self.table.stripes {
3074 let range = stripe
3075 .zone
3076 .column(column)
3077 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
3078 if !range.exact {
3079 return Ok(None);
3080 }
3081 let (Some(small), Some(large)) = (range.low.as_ref(), range.high.as_ref()) else {
3086 if stripe.rows > range.nulls {
3087 return Ok(None);
3088 }
3089 continue;
3090 };
3091 low = Some(low.map_or_else(|| small.clone(), |held| held.smaller(small.clone())));
3092 high = Some(high.map_or_else(|| large.clone(), |held| held.larger(large.clone())));
3093 }
3094 Ok(low.zip(high))
3095 }
3096
3097 pub fn exact_sum(&self, column: usize) -> Result<Option<(i128, u64)>> {
3110 if column >= self.table.fields.len() {
3111 return Err(invalid("sum column index out of range"));
3112 }
3113 let mut total = 0_i128;
3114 let mut rows = 0_u64;
3115 for stripe in &self.table.stripes {
3116 let range = stripe
3117 .zone
3118 .column(column)
3119 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
3120 let Some(part) = range.sum else { return Ok(None) };
3121 let Some(sum) = total.checked_add(part) else { return Ok(None) };
3122 total = sum;
3123 rows = rows.saturating_add(stripe.rows as u64 - range.nulls as u64);
3124 }
3125 Ok(Some((total, rows)))
3126 }
3127
3128 fn dictionary(&self, column: usize) -> Result<Option<Arc<Vector>>> {
3137 let Some(page) = self.table.dictionaries[column] else { return Ok(None) };
3138 if let Some(dictionary) = self.dictionaries[column].get() {
3139 return Ok(Some(Arc::clone(dictionary)));
3140 }
3141 let _queued = self.loading[column].lock().map_err(|_| invalid("a poisoned dictionary"))?;
3142 if let Some(dictionary) = self.dictionaries[column].get() {
3143 return Ok(Some(Arc::clone(dictionary)));
3144 }
3145 self.opened.fetch_add(1, Atomic::Relaxed);
3146 let dictionary = Arc::new(open_global_dictionary(
3147 Arc::clone(&self.file),
3148 page,
3149 &self.table.fields[column].ty,
3150 TEXT_KEEP_BUDGET,
3151 )?);
3152 let _ = self.dictionaries[column].set(Arc::clone(&dictionary));
3153 Ok(Some(dictionary))
3154 }
3155
3156 pub fn read(&self, part: usize, columns: &[usize]) -> Result<Chunk> {
3165 self.read_impl(part, columns, true)
3166 }
3167
3168 pub fn read_sparse(&self, part: usize, columns: &[usize]) -> Result<Chunk> {
3178 self.read_impl(part, columns, false)
3179 }
3180
3181 pub fn skips_codes(&self, part: usize, column: usize, candidates: &[u32]) -> Result<bool> {
3188 if candidates.is_empty() {
3189 return Ok(true);
3190 }
3191 if candidates.windows(2).any(|pair| pair[0] >= pair[1]) {
3192 return Err(Error::internal("native code candidates are not sorted and unique"));
3193 }
3194 let stripe = self.stripe_of(part)?;
3195 let Some(page) = stripe.memberships.get(column).copied().flatten() else {
3196 return Ok(false);
3197 };
3198 let mut bytes = vec![0; page.length as usize];
3199 read_at(&self.file, page.offset, &mut bytes)?;
3200 if checksum(&bytes) != page.hash {
3201 return Err(invalid("membership page checksum differs"));
3202 }
3203 let codes = decode_membership(&bytes)?;
3204 let mut left = 0;
3205 let mut right = 0;
3206 while left < codes.len() && right < candidates.len() {
3207 match codes[left].cmp(&candidates[right]) {
3208 Ordering::Less => left += 1,
3209 Ordering::Greater => right += 1,
3210 Ordering::Equal => return Ok(false),
3211 }
3212 }
3213 Ok(true)
3214 }
3215
3216 fn stripe_of(&self, part: usize) -> Result<&Stripe> {
3217 let place = self.places.get(part).ok_or_else(|| invalid("part index out of range"))?;
3218 self.table
3219 .stripes
3220 .get(place.stripe as usize)
3221 .ok_or_else(|| invalid("stripe index out of range"))
3222 }
3223
3224 fn held(&self, at: usize, stripe: &Stripe, column: usize, whole: bool) -> Result<CachedColumn> {
3241 let cache = self.cache.get(column).ok_or_else(|| invalid("column index out of range"))?;
3242 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3243 let known = cached.index.get(at).and_then(Clone::clone);
3244 let page = cached.pages.get(at).and_then(Clone::clone);
3245 if let Some(index) = known.clone() {
3246 if !whole || page.is_some() {
3247 return Ok(CachedColumn { stripe: at, index, page });
3248 }
3249 }
3250 if cached.loading.contains(&at) {
3251 drop(cached);
3252 if let Some(index) = known {
3256 return Ok(CachedColumn { stripe: at, index, page: None });
3257 }
3258 let held = self.page_of(stripe, column, at, false, None)?;
3259 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3260 remember(&mut cached, &held, self.kept.load(Atomic::Relaxed));
3261 return Ok(held);
3262 }
3263 cached.loading.push(at);
3264 drop(cached);
3265
3266 let read = self.page_of(stripe, column, at, whole, known);
3267
3268 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3272 if let Some(position) = cached.loading.iter().position(|loading| *loading == at) {
3273 cached.loading.remove(position);
3274 }
3275 let held = read?;
3276 remember(&mut cached, &held, self.kept.load(Atomic::Relaxed));
3277 Ok(held)
3278 }
3279
3280 fn page_of(
3286 &self,
3287 stripe: &Stripe,
3288 column: usize,
3289 at: usize,
3290 whole: bool,
3291 known: Option<Arc<Vec<PartSpan>>>,
3292 ) -> Result<CachedColumn> {
3293 let index = match known {
3294 Some(index) => index,
3295 None => {
3296 self.indexes.fetch_add(1, Atomic::Relaxed);
3297 Arc::new(read_index(&self.file, stripe, column)?)
3298 }
3299 };
3300 let page = if whole {
3301 self.pages.fetch_add(1, Atomic::Relaxed);
3302 let span = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
3303 let mut bytes = vec![0; span.length as usize];
3304 read_at(&self.file, span.offset, &mut bytes)?;
3305 Some(Arc::new(bytes))
3306 } else {
3307 None
3308 };
3309 Ok(CachedColumn { stripe: at, index, page })
3310 }
3311
3312 fn read_impl(&self, at: usize, columns: &[usize], whole: bool) -> Result<Chunk> {
3313 let place = *self.places.get(at).ok_or_else(|| invalid("part index out of range"))?;
3314 let index = place.stripe as usize;
3315 let stripe =
3316 self.table.stripes.get(index).ok_or_else(|| invalid("stripe index out of range"))?;
3317 let rows = place.rows as usize;
3318 let mut picked = Vec::with_capacity(columns.len());
3319 for &column in columns {
3320 let field = self
3321 .table
3322 .fields
3323 .get(column)
3324 .ok_or_else(|| invalid("column index out of range"))?;
3325 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
3326 let held = self.held(index, stripe, column, whole)?;
3327 let span = *held
3328 .index
3329 .get(place.part as usize)
3330 .ok_or_else(|| invalid("part index out of range"))?;
3331 let owned;
3332 let bytes = match &held.page {
3333 Some(held) => part_bytes(held, span)?,
3334 None => {
3335 let offset = page
3336 .offset
3337 .checked_add(span.start as u64)
3338 .ok_or_else(|| invalid("part range overflow"))?;
3339 let mut bytes = vec![0; span.length];
3340 read_at(&self.file, offset, &mut bytes)?;
3341 owned = bytes;
3342 &owned
3343 }
3344 };
3345 if checksum(bytes) != span.hash {
3346 return Err(invalid(&format!(
3347 "column page checksum differs, column {column} part {} at {}+{} of {} bytes, \
3348 wanted {:016x} and got {:016x}",
3349 place.part,
3350 page.offset,
3351 span.start,
3352 span.length,
3353 span.hash,
3354 checksum(bytes),
3355 )));
3356 }
3357 let dictionary = self.dictionary(column)?;
3358 picked.push(decode(&field.ty, rows, bytes, dictionary)?.into_pages());
3364 }
3365 Chunk::with_rows(picked, rows)
3366 }
3367
3368 #[must_use]
3384 pub fn skips(&self, part: usize, probes: &[Probe]) -> bool {
3385 let Some(place) = self.places.get(part).copied() else { return false };
3386 let Some(stripe) = self.table.stripes.get(place.stripe as usize) else { return false };
3387 if stripe.zone.skips(probes) {
3388 return true;
3389 }
3390 probes.iter().any(|probe| self.outside(place, probe) || self.sifted(place, probe))
3391 }
3392
3393 fn outside(&self, place: Place, probe: &Probe) -> bool {
3399 match self.stripe_part_ranges(place.stripe as usize, probe.column) {
3400 Some(ranges) => ranges
3401 .get(place.part as usize)
3402 .is_some_and(|range| range.excludes(probe.op, &probe.value)),
3403 None => false,
3404 }
3405 }
3406
3407 fn stripe_part_ranges(&self, stripe: usize, column: usize) -> Option<&[Range]> {
3413 let slot = self.part_ranges.get(column)?.get(stripe)?;
3414 if let Some(held) = slot.get() {
3415 return Some(held);
3416 }
3417 let page = self.table.stripes.get(stripe)?.part_ranges.get(column).copied().flatten()?;
3418 let mut bytes = vec![0; page.length as usize];
3419 read_at(&self.file, page.offset, &mut bytes).ok()?;
3420 if checksum(&bytes) != page.hash {
3421 return None;
3422 }
3423 let ranges = Arc::new(decode_part_ranges(&bytes).ok()?);
3424 let _ = slot.set(ranges);
3425 slot.get().map(|held| held.as_slice())
3426 }
3427
3428 #[must_use]
3445 pub fn certain(&self, part: usize, probes: &[Probe]) -> bool {
3446 let Some(place) = self.places.get(part).copied() else { return false };
3447 let Some(stripe) = self.table.stripes.get(place.stripe as usize) else { return false };
3448 if stripe.zone.certain(probes) {
3449 return true;
3450 }
3451 probes
3452 .iter()
3453 .all(|probe| stripe.zone.certain(slice::from_ref(probe)) || self.inside(place, probe))
3454 }
3455
3456 fn inside(&self, place: Place, probe: &Probe) -> bool {
3462 match self.stripe_part_ranges(place.stripe as usize, probe.column) {
3463 Some(ranges) => ranges
3464 .get(place.part as usize)
3465 .is_some_and(|range| range.certain(probe.op, &probe.value)),
3466 None => false,
3467 }
3468 }
3469
3470 #[must_use]
3481 pub fn stripe_skips(&self, stripe: usize, probes: &[Probe]) -> bool {
3482 self.table.stripes.get(stripe).is_some_and(|held| held.zone.skips(probes))
3483 }
3484
3485 fn sifted(&self, place: Place, probe: &Probe) -> bool {
3491 if probe.op != Op::Equal {
3492 return false;
3493 }
3494 match self.stripe_sieves(place.stripe as usize, probe.column) {
3495 Some(sieves) => sieves
3496 .get(place.part as usize)
3497 .and_then(Option::as_ref)
3498 .is_some_and(|sieve| sieve.excludes(&probe.value)),
3499 None => false,
3500 }
3501 }
3502
3503 fn stripe_sieves(&self, stripe: usize, column: usize) -> Option<&[Option<Sieve>]> {
3510 let slot = self.sieves.get(column)?.get(stripe)?;
3511 if let Some(held) = slot.get() {
3512 return Some(held);
3513 }
3514 let page = self.table.stripes.get(stripe)?.sieves.get(column).copied().flatten()?;
3515 let mut bytes = vec![0; page.length as usize];
3516 read_at(&self.file, page.offset, &mut bytes).ok()?;
3517 if checksum(&bytes) != page.hash {
3518 return None;
3519 }
3520 let sieves = Arc::new(decode_sieves(&bytes).ok()?);
3521 let _ = slot.set(sieves);
3522 slot.get().map(|held| held.as_slice())
3523 }
3524}
3525
3526fn text_at_rank(dictionary: &Vector, rank: usize) -> Result<Value> {
3528 let code = dictionary.code_at_rank(rank)? as usize;
3529 let text = dictionary
3530 .try_text_at(code)?
3531 .ok_or_else(|| invalid("global dictionary order names a code it does not have"))?;
3532 Ok(Value::Varchar(text.into()))
3533}
3534
3535#[cfg(unix)]
3540fn write_at(file: &File, mut offset: u64, mut bytes: &[u8]) -> Result<()> {
3541 use std::os::unix::fs::FileExt;
3542 while !bytes.is_empty() {
3543 let written = file.write_at(bytes, offset).map_err(io)?;
3544 if written == 0 {
3545 return Err(invalid("a write to the native file wrote nothing"));
3546 }
3547 offset += written as u64;
3548 bytes = &bytes[written..];
3549 }
3550 Ok(())
3551}
3552
3553#[cfg(windows)]
3555fn write_at(file: &File, mut offset: u64, mut bytes: &[u8]) -> Result<()> {
3556 use std::os::windows::fs::FileExt;
3557 while !bytes.is_empty() {
3558 let written = file.seek_write(bytes, offset).map_err(io)?;
3559 if written == 0 {
3560 return Err(invalid("a write to the native file wrote nothing"));
3561 }
3562 offset += written as u64;
3563 bytes = &bytes[written..];
3564 }
3565 Ok(())
3566}
3567
3568#[cfg(not(any(unix, windows)))]
3570fn write_at(file: &File, offset: u64, bytes: &[u8]) -> Result<()> {
3571 use std::io::Write;
3572 let mut file = file.try_clone().map_err(io)?;
3573 file.seek(SeekFrom::Start(offset)).map_err(io)?;
3574 file.write_all(bytes).map_err(io)
3575}
3576
3577#[cfg(unix)]
3587fn read_at(file: &File, mut offset: u64, mut bytes: &mut [u8]) -> Result<()> {
3588 use std::os::unix::fs::FileExt;
3589 while !bytes.is_empty() {
3590 let read = file.read_at(bytes, offset).map_err(io)?;
3591 if read == 0 {
3592 return Err(invalid("column page ends before its declared length"));
3593 }
3594 offset += read as u64;
3595 bytes = &mut bytes[read..];
3596 }
3597 Ok(())
3598}
3599
3600#[cfg(windows)]
3606fn read_at(file: &File, mut offset: u64, mut bytes: &mut [u8]) -> Result<()> {
3607 use std::os::windows::fs::FileExt;
3608 while !bytes.is_empty() {
3609 let read = file.seek_read(bytes, offset).map_err(io)?;
3610 if read == 0 {
3611 return Err(invalid("column page ends before its declared length"));
3612 }
3613 offset += read as u64;
3614 bytes = &mut bytes[read..];
3615 }
3616 Ok(())
3617}
3618
3619#[cfg(not(any(unix, windows)))]
3624fn read_at(file: &File, offset: u64, bytes: &mut [u8]) -> Result<()> {
3625 let mut file = file.try_clone().map_err(io)?;
3626 file.seek(SeekFrom::Start(offset)).map_err(io)?;
3627 file.read_exact(bytes).map_err(io)
3628}
3629
3630fn type_tag(ty: &LogicalType) -> Result<u8> {
3631 match ty {
3632 LogicalType::SmallInt => Ok(1),
3633 LogicalType::Integer => Ok(2),
3634 LogicalType::BigInt => Ok(3),
3635 LogicalType::Varchar => Ok(4),
3636 LogicalType::Date => Ok(5),
3637 LogicalType::Timestamp => Ok(6),
3638 LogicalType::Boolean => Ok(7),
3639 LogicalType::TinyInt => Ok(8),
3640 LogicalType::UTinyInt => Ok(9),
3641 LogicalType::USmallInt => Ok(10),
3642 LogicalType::UInteger => Ok(11),
3643 LogicalType::UBigInt => Ok(12),
3644 LogicalType::Decimal { .. } => Ok(13),
3645 _ => Err(Error::not_implemented(format!("native storage for {ty}"))),
3646 }
3647}
3648
3649fn put_type(out: &mut Vec<u8>, ty: &LogicalType) -> Result<()> {
3655 out.push(type_tag(ty)?);
3656 if let LogicalType::Decimal { width, scale } = ty {
3657 out.push(*width);
3658 out.push(*scale);
3659 }
3660 Ok(())
3661}
3662
3663fn read_type(cur: &mut Cursor<'_>) -> Result<LogicalType> {
3665 let tag = cur.u8()?;
3666 if tag == 13 {
3667 let width = cur.u8()?;
3668 let scale = cur.u8()?;
3669 return LogicalType::decimal(width, scale)
3670 .map_err(|_| invalid("decimal column width and scale are not a decimal"));
3671 }
3672 tag_type(tag)
3673}
3674
3675fn tag_type(tag: u8) -> Result<LogicalType> {
3676 match tag {
3677 1 => Ok(LogicalType::SmallInt),
3678 2 => Ok(LogicalType::Integer),
3679 3 => Ok(LogicalType::BigInt),
3680 4 => Ok(LogicalType::Varchar),
3681 5 => Ok(LogicalType::Date),
3682 6 => Ok(LogicalType::Timestamp),
3683 7 => Ok(LogicalType::Boolean),
3684 8 => Ok(LogicalType::TinyInt),
3685 9 => Ok(LogicalType::UTinyInt),
3686 10 => Ok(LogicalType::USmallInt),
3687 11 => Ok(LogicalType::UInteger),
3688 12 => Ok(LogicalType::UBigInt),
3689 _ => Err(invalid("column type tag is unknown")),
3690 }
3691}
3692
3693fn put_u16(out: &mut Vec<u8>, value: u16) {
3694 out.extend_from_slice(&value.to_le_bytes());
3695}
3696fn put_u32(out: &mut Vec<u8>, value: u32) {
3697 out.extend_from_slice(&value.to_le_bytes());
3698}
3699fn put_u64(out: &mut Vec<u8>, value: u64) {
3700 out.extend_from_slice(&value.to_le_bytes());
3701}
3702fn put_var_u64(out: &mut Vec<u8>, mut value: u64) {
3703 while value >= 0x80 {
3704 out.push((value as u8 & 0x7f) | 0x80);
3705 value >>= 7;
3706 }
3707 out.push(value as u8);
3708}
3709
3710fn frequency_order(left: FrequencyValue, right: FrequencyValue) -> Ordering {
3711 match (left, right) {
3712 (FrequencyValue::Null, FrequencyValue::Null) => Ordering::Equal,
3713 (FrequencyValue::Null, _) => Ordering::Less,
3714 (_, FrequencyValue::Null) => Ordering::Greater,
3715 (FrequencyValue::Integer(left), FrequencyValue::Integer(right)) => left.cmp(&right),
3716 (FrequencyValue::Code(left), FrequencyValue::Code(right)) => left.cmp(&right),
3717 (FrequencyValue::Integer(_), FrequencyValue::Code(_)) => Ordering::Less,
3718 (FrequencyValue::Code(_), FrequencyValue::Integer(_)) => Ordering::Greater,
3719 }
3720}
3721
3722fn code_frequency(dictionary: &GlobalDictionary) -> FrequencySummary {
3723 let mut entries = dictionary
3724 .counts
3725 .iter()
3726 .enumerate()
3727 .filter(|(_, count)| **count != 0)
3728 .map(|(code, &count)| FrequencyEntry { value: FrequencyValue::Code(code as u32), count })
3729 .collect::<Vec<_>>();
3730 if dictionary.nulls != 0 {
3731 entries.push(FrequencyEntry { value: FrequencyValue::Null, count: dictionary.nulls });
3732 }
3733 entries.sort_unstable_by(|left, right| {
3734 right.count.cmp(&left.count).then_with(|| frequency_order(left.value, right.value))
3735 });
3736 let omitted_max = entries.get(FREQUENCY_ENTRIES).map_or(0, |entry| entry.count);
3737 entries.truncate(FREQUENCY_ENTRIES);
3738 FrequencySummary { entries, omitted_max, ordinals: Vec::new() }
3739}
3740
3741fn encode_directory(table: &Table) -> Result<Vec<u8>> {
3742 let mut out = DIRECTORY.to_vec();
3743 let name = table.name.as_bytes();
3744 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("table name too long"))?);
3745 out.extend_from_slice(name);
3746 put_u16(&mut out, u16::try_from(table.fields.len()).map_err(|_| invalid("too many columns"))?);
3747 for field in &table.fields {
3748 let name = field.name.as_bytes();
3749 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("column name too long"))?);
3750 out.extend_from_slice(name);
3751 put_type(&mut out, &field.ty)?;
3752 out.push(u8::from(field.not_null));
3753 }
3754 for dictionary in &table.dictionaries {
3755 match dictionary {
3756 None => out.push(0),
3757 Some(page) => {
3758 out.push(1);
3759 put_u64(&mut out, page.offset);
3760 put_u32(&mut out, page.length);
3761 put_u64(&mut out, page.hash);
3762 }
3763 }
3764 }
3765 for distinct in &table.distincts {
3766 match distinct {
3767 None => out.push(0),
3768 Some(count) => {
3769 out.push(1);
3770 put_u64(&mut out, *count);
3771 }
3772 }
3773 }
3774 put_u64(&mut out, u64::try_from(table.rows).map_err(|_| invalid("row count overflow"))?);
3775 put_u32(&mut out, u32::try_from(table.stripes.len()).map_err(|_| invalid("too many stripes"))?);
3776 for stripe in &table.stripes {
3777 put_u32(
3778 &mut out,
3779 u32::try_from(stripe.parts.len()).map_err(|_| invalid("too many parts in a stripe"))?,
3780 );
3781 for &rows in &stripe.parts {
3782 put_u32(&mut out, rows);
3783 }
3784 put_u64(&mut out, stripe.index.offset);
3785 put_u32(&mut out, stripe.index.length);
3786 for page in &stripe.pages {
3787 put_u64(&mut out, page.offset);
3788 put_u32(&mut out, page.length);
3789 }
3790 for ((field, dictionary), membership) in
3795 table.fields.iter().zip(&table.dictionaries).zip(&stripe.memberships)
3796 {
3797 if field.ty != LogicalType::Varchar || dictionary.is_none() {
3798 continue;
3799 }
3800 let page =
3801 membership.ok_or_else(|| invalid("string page has no code membership index"))?;
3802 put_u64(&mut out, page.offset);
3803 put_u32(&mut out, page.length);
3804 put_u64(&mut out, page.hash);
3805 }
3806 for sieve in &stripe.sieves {
3807 match sieve {
3808 None => out.push(0),
3809 Some(page) => {
3810 out.push(1);
3811 put_u64(&mut out, page.offset);
3812 put_u32(&mut out, page.length);
3813 put_u64(&mut out, page.hash);
3814 }
3815 }
3816 }
3817 for held in &stripe.part_ranges {
3818 match held {
3819 None => out.push(0),
3820 Some(page) => {
3821 out.push(1);
3822 put_u64(&mut out, page.offset);
3823 put_u32(&mut out, page.length);
3824 put_u64(&mut out, page.hash);
3825 }
3826 }
3827 }
3828 for range in stripe.zone.columns() {
3829 put_bound(&mut out, range.low.as_ref())?;
3830 put_bound(&mut out, range.high.as_ref())?;
3831 put_u32(
3832 &mut out,
3833 u32::try_from(range.nulls).map_err(|_| invalid("null count overflow"))?,
3834 );
3835 out.push(u8::from(range.exact));
3836 match range.sum {
3837 None => out.push(0),
3838 Some(total) => {
3839 out.push(1);
3840 out.extend_from_slice(&total.to_le_bytes());
3841 }
3842 }
3843 }
3844 }
3845 out.extend_from_slice(FREQUENCIES);
3846 put_u16(
3847 &mut out,
3848 u16::try_from(table.frequencies.len())
3849 .map_err(|_| invalid("too many frequency columns"))?,
3850 );
3851 for summary in &table.frequencies {
3852 let Some(summary) = summary else {
3853 out.push(0);
3854 continue;
3855 };
3856 out.push(1);
3857 put_u64(&mut out, summary.omitted_max);
3858 put_u32(
3859 &mut out,
3860 u32::try_from(summary.entries.len())
3861 .map_err(|_| invalid("too many frequency entries"))?,
3862 );
3863 for entry in &summary.entries {
3864 match entry.value {
3865 FrequencyValue::Null => out.push(0),
3866 FrequencyValue::Integer(value) => {
3867 out.push(1);
3868 out.extend_from_slice(&value.to_le_bytes());
3869 }
3870 FrequencyValue::Code(value) => {
3871 out.push(2);
3872 put_u32(&mut out, value);
3873 }
3874 }
3875 put_u64(&mut out, entry.count);
3876 }
3877 put_u32(
3878 &mut out,
3879 u32::try_from(summary.ordinals.len())
3880 .map_err(|_| invalid("too many frequency ordinals"))?,
3881 );
3882 let mut previous = 0_u64;
3883 for (at, &ordinal) in summary.ordinals.iter().enumerate() {
3884 let delta = if at == 0 {
3885 ordinal
3886 } else {
3887 ordinal
3888 .checked_sub(previous)
3889 .ok_or_else(|| invalid("frequency ordinals are not ordered"))?
3890 };
3891 if at != 0 && delta == 0 {
3892 return Err(invalid("frequency ordinals are not unique"));
3893 }
3894 put_var_u64(&mut out, delta);
3895 previous = ordinal;
3896 }
3897 }
3898 Ok(out)
3899}
3900
3901fn encode_catalog(entries: &[Entry]) -> Result<Vec<u8>> {
3907 let mut out = CATALOG.to_vec();
3908 put_u32(&mut out, u32::try_from(entries.len()).map_err(|_| invalid("too many tables"))?);
3909 for entry in entries {
3910 let name = entry.name.as_bytes();
3911 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("table name too long"))?);
3912 out.extend_from_slice(name);
3913 put_u64(&mut out, u64::try_from(entry.rows).map_err(|_| invalid("row count overflow"))?);
3914 put_u16(
3915 &mut out,
3916 u16::try_from(entry.fields.len()).map_err(|_| invalid("too many columns"))?,
3917 );
3918 for field in &entry.fields {
3919 let name = field.name.as_bytes();
3920 put_u16(
3921 &mut out,
3922 u16::try_from(name.len()).map_err(|_| invalid("column name too long"))?,
3923 );
3924 out.extend_from_slice(name);
3925 put_type(&mut out, &field.ty)?;
3926 out.push(u8::from(field.not_null));
3927 }
3928 put_u64(&mut out, entry.directory.offset);
3929 put_u32(&mut out, entry.directory.length);
3930 put_u64(&mut out, entry.directory.hash);
3931 }
3932 Ok(out)
3933}
3934
3935fn decode_catalog(bytes: &[u8], size: u64) -> Result<Vec<Entry>> {
3938 let mut cur = Cursor { bytes, at: 0 };
3939 if cur.take(8)? != CATALOG {
3940 return Err(invalid("catalog magic differs"));
3941 }
3942 let count = cur.u32()? as usize;
3943 let mut entries: Vec<Entry> = Vec::with_capacity(count.min(1024));
3944 for _ in 0..count {
3945 let name = cur.text()?;
3946 let rows = usize::try_from(cur.u64()?).map_err(|_| invalid("row count does not fit"))?;
3947 let width = cur.u16()? as usize;
3948 let mut fields = Vec::with_capacity(width);
3949 for _ in 0..width {
3950 let name = cur.text()?;
3951 let ty = read_type(&mut cur)?;
3952 let not_null = match cur.u8()? {
3953 0 => false,
3954 1 => true,
3955 _ => return Err(invalid("nullability flag differs")),
3956 };
3957 fields.push(Field { name, ty, not_null });
3958 }
3959 let directory = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3960 let end = directory
3961 .offset
3962 .checked_add(u64::from(directory.length))
3963 .ok_or_else(|| invalid("table directory offset overflow"))?;
3964 if directory.offset < HEADER
3965 || end > size
3966 || directory.length as usize > MAX_DIRECTORY
3967 || directory.length == 0
3968 {
3969 return Err(invalid("table directory range is outside the file"));
3970 }
3971 if entries.iter().any(|held| held.name == name) {
3972 return Err(invalid("two tables in the catalog have the same name"));
3973 }
3974 entries.push(Entry { name, fields, rows, directory });
3975 }
3976 Ok(entries)
3977}
3978
3979struct Cursor<'a> {
3980 bytes: &'a [u8],
3981 at: usize,
3982}
3983impl<'a> Cursor<'a> {
3984 fn take(&mut self, len: usize) -> Result<&'a [u8]> {
3985 let end = self.at.checked_add(len).ok_or_else(|| invalid("directory offset overflow"))?;
3986 let bytes =
3987 self.bytes.get(self.at..end).ok_or_else(|| invalid("directory is truncated"))?;
3988 self.at = end;
3989 Ok(bytes)
3990 }
3991 fn u8(&mut self) -> Result<u8> {
3992 Ok(self.take(1)?[0])
3993 }
3994 fn u16(&mut self) -> Result<u16> {
3995 Ok(u16::from_le_bytes(self.take(2)?.try_into().expect("two bytes")))
3996 }
3997 fn u32(&mut self) -> Result<u32> {
3998 Ok(u32::from_le_bytes(self.take(4)?.try_into().expect("four bytes")))
3999 }
4000 fn u64(&mut self) -> Result<u64> {
4001 Ok(u64::from_le_bytes(self.take(8)?.try_into().expect("eight bytes")))
4002 }
4003 fn var_u64(&mut self) -> Result<u64> {
4004 let mut value = 0_u64;
4005 for shift in (0..=63).step_by(7) {
4006 let byte = self.u8()?;
4007 let part = u64::from(byte & 0x7f);
4008 if shift == 63 && part > 1 {
4009 return Err(invalid("frequency ordinal varint overflows"));
4010 }
4011 value |= part << shift;
4012 if byte & 0x80 == 0 {
4013 return Ok(value);
4014 }
4015 }
4016 Err(invalid("frequency ordinal varint is too long"))
4017 }
4018 fn bound(&mut self) -> Result<Option<Bound>> {
4019 Ok(match self.u8()? {
4020 0 => None,
4021 1 => Some(Bound::Int(i128::from_le_bytes(
4022 self.take(16)?.try_into().expect("sixteen bytes"),
4023 ))),
4024 2 => Some(Bound::Real(f64::from_le_bytes(
4025 self.take(8)?.try_into().expect("eight bytes"),
4026 ))),
4027 3 => {
4028 let length = self.u32()? as usize;
4029 Some(Bound::Bytes(self.take(length)?.to_vec()))
4030 }
4031 4 => {
4032 let unscaled =
4033 i128::from_le_bytes(self.take(16)?.try_into().expect("sixteen bytes"));
4034 Some(Bound::Scaled { unscaled, scale: self.u8()? })
4035 }
4036 _ => return Err(invalid("bound tag differs")),
4037 })
4038 }
4039 fn text(&mut self) -> Result<String> {
4040 let len = self.u16()? as usize;
4041 String::from_utf8(self.take(len)?.to_vec()).map_err(|_| invalid("name is not UTF-8"))
4042 }
4043}
4044
4045fn decode_directory(bytes: &[u8], size: u64) -> Result<Table> {
4046 let mut cur = Cursor { bytes, at: 0 };
4047 if cur.take(8)? != DIRECTORY {
4048 return Err(invalid("directory magic differs"));
4049 }
4050 let name = cur.text()?;
4051 let width = cur.u16()? as usize;
4052 let mut fields = Vec::with_capacity(width);
4053 for _ in 0..width {
4054 let name = cur.text()?;
4055 let ty = read_type(&mut cur)?;
4056 let not_null = match cur.u8()? {
4057 0 => false,
4058 1 => true,
4059 _ => return Err(invalid("nullability flag differs")),
4060 };
4061 fields.push(Field { name, ty, not_null });
4062 }
4063 let mut dictionaries = Vec::with_capacity(width);
4064 for _ in 0..width {
4065 dictionaries.push(match cur.u8()? {
4066 0 => None,
4067 1 => {
4068 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4069 let end = page
4070 .offset
4071 .checked_add(u64::from(page.length))
4072 .ok_or_else(|| invalid("dictionary page offset overflow"))?;
4073 if page.offset < HEADER || end > size {
4078 return Err(invalid("dictionary page range is outside the file"));
4079 }
4080 Some(page)
4081 }
4082 _ => return Err(invalid("dictionary page tag differs")),
4083 });
4084 }
4085 let mut distincts = Vec::with_capacity(width);
4086 for _ in 0..width {
4087 distincts.push(match cur.u8()? {
4088 0 => None,
4089 1 => Some(cur.u64()?),
4090 _ => return Err(invalid("distinct count tag differs")),
4091 });
4092 }
4093 let rows = usize::try_from(cur.u64()?).map_err(|_| invalid("row count does not fit"))?;
4094 let count = cur.u32()? as usize;
4095 let mut stripes = Vec::with_capacity(count);
4096 let mut total = 0_usize;
4097 for _ in 0..count {
4098 let count = cur.u32()? as usize;
4099 if count == 0 || count > STRIPE_PARTS {
4100 return Err(invalid("stripe part count is outside its bound"));
4101 }
4102 let mut parts = Vec::with_capacity(count);
4103 let mut stripe_rows = 0_usize;
4104 for _ in 0..count {
4105 let rows = cur.u32()?;
4106 if rows == 0 {
4107 return Err(invalid("empty part"));
4108 }
4109 parts.push(rows);
4110 stripe_rows = stripe_rows
4111 .checked_add(rows as usize)
4112 .ok_or_else(|| invalid("stripe row count overflow"))?;
4113 }
4114 total =
4115 total.checked_add(stripe_rows).ok_or_else(|| invalid("stripe row count overflow"))?;
4116 let index = Span { offset: cur.u64()?, length: cur.u32()? };
4117 let section = index_section(count)?;
4118 let wanted = section
4119 .checked_mul(width)
4120 .and_then(|bytes| u32::try_from(bytes).ok())
4121 .ok_or_else(|| invalid("index page length overflow"))?;
4122 let end = index
4123 .offset
4124 .checked_add(u64::from(index.length))
4125 .ok_or_else(|| invalid("index page offset overflow"))?;
4126 if index.offset < HEADER || end > size || index.length != wanted {
4127 return Err(invalid("index page range is outside the file"));
4128 }
4129 let mut pages = Vec::with_capacity(width);
4130 for _ in 0..width {
4131 let offset = cur.u64()?;
4132 let length = cur.u32()?;
4133 let end = offset
4134 .checked_add(u64::from(length))
4135 .ok_or_else(|| invalid("page offset overflow"))?;
4136 if offset < HEADER || end > size || length as usize > MAX_PAGE {
4137 return Err(invalid("page range is outside the file"));
4138 }
4139 pages.push(Span { offset, length });
4140 }
4141 let mut memberships = vec![None; width];
4142 for (column, field) in fields.iter().enumerate() {
4143 if field.ty != LogicalType::Varchar || dictionaries[column].is_none() {
4144 continue;
4145 }
4146 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4147 let end = page
4148 .offset
4149 .checked_add(u64::from(page.length))
4150 .ok_or_else(|| invalid("membership page offset overflow"))?;
4151 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
4152 return Err(invalid("membership page range is outside the file"));
4153 }
4154 memberships[column] = Some(page);
4155 }
4156 let mut sieves = vec![None; width];
4157 for sieve in sieves.iter_mut().take(width) {
4158 match cur.u8()? {
4159 0 => continue,
4160 1 => {}
4161 _ => return Err(invalid("a sieve page has an unknown tag")),
4162 }
4163 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4164 let end = page
4165 .offset
4166 .checked_add(u64::from(page.length))
4167 .ok_or_else(|| invalid("sieve page offset overflow"))?;
4168 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
4169 return Err(invalid("sieve page range is outside the file"));
4170 }
4171 *sieve = Some(page);
4172 }
4173 let mut part_ranges = vec![None; width];
4174 for held in part_ranges.iter_mut().take(width) {
4175 match cur.u8()? {
4176 0 => continue,
4177 1 => {}
4178 _ => return Err(invalid("a part range page has an unknown tag")),
4179 }
4180 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4181 let end = page
4182 .offset
4183 .checked_add(u64::from(page.length))
4184 .ok_or_else(|| invalid("part range page offset overflow"))?;
4185 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
4186 return Err(invalid("part range page range is outside the file"));
4187 }
4188 *held = Some(page);
4189 }
4190 let mut ranges = Vec::with_capacity(width);
4191 for column in 0..width {
4192 let low = cur.bound()?;
4193 let high = cur.bound()?;
4194 let nulls = cur.u32()? as usize;
4195 if nulls > stripe_rows {
4196 return Err(invalid("null count exceeds stripe rows"));
4197 }
4198 let exact = cur.u8()? != 0;
4199 let sum = match cur.u8()? {
4200 0 => None,
4201 1 => Some(i128::from_le_bytes(
4202 cur.take(16)?.try_into().map_err(|_| invalid("a stripe sum is truncated"))?,
4203 )),
4204 _ => return Err(invalid("a stripe sum has an unknown tag")),
4205 };
4206 let ty = &fields.get(column).ok_or_else(|| invalid("a stripe range has no column"))?.ty;
4212 let low = low.map(|bound| scaled_as(bound, ty));
4213 let high = high.map(|bound| scaled_as(bound, ty));
4214 ranges.push(Range { low, high, nulls, exact, sum });
4215 }
4216 stripes.push(Stripe {
4217 rows: stripe_rows,
4218 parts,
4219 index,
4220 pages,
4221 memberships,
4222 sieves,
4223 part_ranges,
4224 zone: Zone::from_ranges(ranges),
4225 });
4226 }
4227 if total != rows {
4228 return Err(invalid("table row count differs from stripes"));
4229 }
4230 let frequencies = if cur.at == bytes.len() {
4231 vec![None; width]
4232 } else {
4233 if cur.take(8)? != FREQUENCIES {
4234 return Err(invalid("directory extension magic differs"));
4235 }
4236 if cur.u16()? as usize != width {
4237 return Err(invalid("frequency column count differs"));
4238 }
4239 let mut frequencies = Vec::with_capacity(width);
4240 for field in &fields {
4241 let summary = match cur.u8()? {
4242 0 => None,
4243 1 => {
4244 let omitted_max = cur.u64()?;
4245 let count = cur.u32()? as usize;
4246 if count > FREQUENCY_ENTRIES {
4247 return Err(invalid("frequency entry count exceeds its bound"));
4248 }
4249 let mut entries = Vec::with_capacity(count);
4250 for _ in 0..count {
4252 let value = match cur.u8()? {
4253 0 => FrequencyValue::Null,
4254 1 => FrequencyValue::Integer(i128::from_le_bytes(
4255 cur.take(16)?.try_into().expect("sixteen bytes"),
4256 )),
4257 2 => FrequencyValue::Code(cur.u32()?),
4258 _ => return Err(invalid("frequency value tag differs")),
4259 };
4260 let valid = matches!(
4261 (&field.ty, value),
4262 (_, FrequencyValue::Null)
4263 | (LogicalType::Varchar, FrequencyValue::Code(_))
4264 | (
4265 LogicalType::TinyInt
4266 | LogicalType::SmallInt
4267 | LogicalType::Integer
4268 | LogicalType::BigInt
4269 | LogicalType::UTinyInt
4270 | LogicalType::USmallInt
4271 | LogicalType::UInteger
4272 | LogicalType::UBigInt
4273 | LogicalType::Date
4274 | LogicalType::Timestamp,
4275 FrequencyValue::Integer(_),
4276 )
4277 );
4278 if !valid {
4279 return Err(invalid("frequency value does not match its column"));
4280 }
4281 let count = cur.u64()?;
4282 if count == 0 || count > rows as u64 {
4283 return Err(invalid("frequency count is outside the table"));
4284 }
4285 entries.push(FrequencyEntry { value, count });
4286 }
4287 if entries.windows(2).any(|pair| pair[0].count < pair[1].count) {
4288 return Err(invalid("frequency entries are not descending"));
4289 }
4290 let ordinals = {
4291 let ordinal_count = cur.u32()? as usize;
4292 if ordinal_count > FREQUENCY_ORDINALS || ordinal_count > rows {
4293 return Err(invalid("frequency ordinal count exceeds its bound"));
4294 }
4295 let mut ordinals = Vec::with_capacity(ordinal_count);
4296 let mut previous = 0_u64;
4297 for at in 0..ordinal_count {
4298 let delta = cur.var_u64()?;
4299 if at != 0 && delta == 0 {
4300 return Err(invalid("frequency ordinals are not increasing"));
4301 }
4302 let ordinal = if at == 0 {
4303 delta
4304 } else {
4305 previous
4306 .checked_add(delta)
4307 .ok_or_else(|| invalid("frequency ordinal overflows"))?
4308 };
4309 if ordinal >= rows as u64 {
4310 return Err(invalid("frequency ordinal is outside the table"));
4311 }
4312 ordinals.push(ordinal);
4313 previous = ordinal;
4314 }
4315 ordinals
4316 };
4317 Some(FrequencySummary { entries, omitted_max, ordinals })
4318 }
4319 _ => return Err(invalid("frequency summary tag differs")),
4320 };
4321 frequencies.push(summary);
4322 }
4323 frequencies
4324 };
4325 if cur.at != bytes.len() {
4326 return Err(invalid("directory has trailing bytes"));
4327 }
4328 Ok(Table { name, fields, stripes, rows, dictionaries, distincts, frequencies })
4329}
4330
4331fn put_bound(out: &mut Vec<u8>, bound: Option<&Bound>) -> Result<()> {
4332 match bound {
4333 None => out.push(0),
4334 Some(Bound::Int(value)) => {
4335 out.push(1);
4336 out.extend_from_slice(&value.to_le_bytes());
4337 }
4338 Some(Bound::Real(value)) => {
4339 out.push(2);
4340 out.extend_from_slice(&value.to_le_bytes());
4341 }
4342 Some(Bound::Bytes(value)) => {
4343 out.push(3);
4344 put_u32(out, u32::try_from(value.len()).map_err(|_| invalid("bound length overflow"))?);
4345 out.extend_from_slice(value);
4346 }
4347 Some(Bound::Scaled { unscaled, scale }) => {
4348 out.push(4);
4349 out.extend_from_slice(&unscaled.to_le_bytes());
4350 out.push(*scale);
4351 }
4352 }
4353 Ok(())
4354}
4355
4356#[derive(Debug)]
4373struct Codes;
4374
4375impl chooser::Chooser for Codes {
4376 fn name(&self) -> &'static str {
4377 "codes"
4378 }
4379
4380 fn narrow_strings(
4381 &self,
4382 _values: &[&[u8]],
4383 offered: &[string::Kind],
4384 _depth: u8,
4385 ) -> Vec<string::Kind> {
4386 offered.to_vec()
4389 }
4390
4391 fn narrow_integers(
4392 &self,
4393 _values: &[i64],
4394 offered: &[integer::Kind],
4395 depth: u8,
4396 ) -> Vec<integer::Kind> {
4397 let keep: &[integer::Kind] = if depth == 0 {
4398 &[integer::Kind::Constant, integer::Kind::Packed, integer::Kind::Rle]
4399 } else {
4400 &[integer::Kind::Constant, integer::Kind::Packed]
4401 };
4402 let narrowed: Vec<integer::Kind> =
4403 offered.iter().copied().filter(|kind| keep.contains(kind)).collect();
4404 if narrowed.is_empty() { offered.to_vec() } else { narrowed }
4407 }
4408}
4409
4410#[derive(Debug)]
4422struct Fixed;
4423
4424impl chooser::Chooser for Fixed {
4425 fn name(&self) -> &'static str {
4426 "fixed"
4427 }
4428
4429 fn narrow_strings(
4430 &self,
4431 _values: &[&[u8]],
4432 offered: &[string::Kind],
4433 _depth: u8,
4434 ) -> Vec<string::Kind> {
4435 offered.to_vec()
4436 }
4437
4438 fn narrow_integers(
4439 &self,
4440 _values: &[i64],
4441 offered: &[integer::Kind],
4442 depth: u8,
4443 ) -> Vec<integer::Kind> {
4444 let keep: &[integer::Kind] = if depth == 0 {
4445 &[
4446 integer::Kind::Constant,
4447 integer::Kind::Packed,
4448 integer::Kind::Delta,
4449 integer::Kind::Rle,
4450 integer::Kind::Sparse,
4451 integer::Kind::Strided,
4452 ]
4453 } else {
4454 &[integer::Kind::Constant, integer::Kind::Packed, integer::Kind::Delta]
4455 };
4456 let narrowed: Vec<integer::Kind> =
4457 offered.iter().copied().filter(|kind| keep.contains(kind)).collect();
4458 if narrowed.is_empty() { offered.to_vec() } else { narrowed }
4459 }
4460}
4461
4462fn widened(data: &Data) -> Option<Vec<i64>> {
4469 match data {
4470 Data::Int8(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4471 Data::UInt8(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4472 Data::Int16(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4473 Data::UInt16(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4474 Data::Int32(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4475 Data::UInt32(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4476 Data::Int64(values) => Some(values.to_vec()),
4477 _ => None,
4478 }
4479}
4480
4481trait Narrow: Copy {
4488 const BIASED: (u32, u64);
4493
4494 fn narrow(value: i64) -> Self;
4496}
4497
4498#[allow(clippy::cast_sign_loss, reason = "a residue is a bit pattern and not a number")]
4515fn residue<T: Narrow>(value: i64) -> u64 {
4516 let (bits, bias) = T::BIASED;
4517 (value as u64).wrapping_add(bias) >> bits
4518}
4519
4520macro_rules! narrows {
4525 ($($ty:ty => $bias:expr),* $(,)?) => {$(
4526 impl Narrow for $ty {
4527 const BIASED: (u32, u64) = (<$ty>::BITS, $bias);
4528
4529 #[allow(
4530 clippy::cast_possible_truncation,
4531 clippy::cast_sign_loss,
4532 reason = "the caller has checked the bits this truncates away"
4533 )]
4534 fn narrow(value: i64) -> Self {
4535 value as Self
4536 }
4537 }
4538 )*};
4539}
4540
4541narrows! {
4542 i8 => 1 << 7,
4543 u8 => 0,
4544 i16 => 1 << 15,
4545 u16 => 0,
4546 i32 => 1 << 31,
4547 u32 => 0,
4548}
4549
4550fn fit<T: Narrow>(values: &[i64]) -> Result<Vec<T>> {
4563 let mut spilled = 0u64;
4564 for value in values {
4565 spilled |= residue::<T>(*value);
4566 }
4567 if spilled != 0 {
4568 return Err(invalid("page value is not of its type"));
4569 }
4570 Ok(values.iter().map(|value| T::narrow(*value)).collect())
4571}
4572
4573fn narrowed(ty: &LogicalType, values: Vec<i64>) -> Result<Data> {
4578 Ok(match ty {
4579 LogicalType::TinyInt => Data::Int8(fit::<i8>(&values)?.into()),
4580 LogicalType::UTinyInt => Data::UInt8(fit::<u8>(&values)?.into()),
4581 LogicalType::SmallInt => Data::Int16(fit::<i16>(&values)?.into()),
4582 LogicalType::USmallInt => Data::UInt16(fit::<u16>(&values)?.into()),
4583 LogicalType::Integer | LogicalType::Date => Data::Int32(fit::<i32>(&values)?.into()),
4584 LogicalType::UInteger => Data::UInt32(fit::<u32>(&values)?.into()),
4585 LogicalType::BigInt | LogicalType::Timestamp => Data::Int64(values.into()),
4586 LogicalType::Decimal { .. } => match ty.physical() {
4589 PhysicalType::Int16 => Data::Int16(fit::<i16>(&values)?.into()),
4590 PhysicalType::Int32 => Data::Int32(fit::<i32>(&values)?.into()),
4591 PhysicalType::Int64 => Data::Int64(values.into()),
4592 _ => return Err(invalid("cascade codec belongs to a decimal that is not an integer")),
4593 },
4594 _ => return Err(invalid("cascade codec belongs to a page that is not integers")),
4595 })
4596}
4597
4598fn plain_width(ty: &LogicalType) -> Option<usize> {
4601 Some(match ty {
4602 LogicalType::TinyInt | LogicalType::UTinyInt => 1,
4603 LogicalType::SmallInt | LogicalType::USmallInt => 2,
4604 LogicalType::Integer | LogicalType::UInteger | LogicalType::Date => 4,
4605 LogicalType::BigInt | LogicalType::Timestamp => 8,
4606 LogicalType::Decimal { .. } => match ty.physical() {
4607 PhysicalType::Int16 => 2,
4608 PhysicalType::Int32 => 4,
4609 PhysicalType::Int64 => 8,
4610 _ => return None,
4613 },
4614 _ => return None,
4615 })
4616}
4617
4618fn cascaded(
4624 flat: &Vector,
4625 ty: &LogicalType,
4626 packed: Option<&Packed<'_>>,
4627) -> Result<Option<Vec<u8>>> {
4628 let (Some(width), Some(data)) = (plain_width(ty), flat.data()) else { return Ok(None) };
4629 let Some(values) = widened(data) else { return Ok(None) };
4630 let plain = values.len().saturating_mul(width);
4631 let best = match packed {
4632 Some(packed) => plain.min(21 + size_of_val(packed.words())),
4634 None => plain,
4635 };
4636 let out = integer::encode_with(&values, &Fixed)?;
4637 Ok((out.len() < best).then_some(out))
4638}
4639
4640fn text_compressed(flat: &Vector) -> Result<Option<Vec<u8>>> {
4678 let mut values: Vec<&[u8]> = Vec::with_capacity(flat.len());
4679 let mut payload = 0_usize;
4680 for row in 0..flat.len() {
4681 let text = flat.text_at(row).unwrap_or("").as_bytes();
4682 payload = payload.saturating_add(text.len());
4683 values.push(text);
4684 }
4685 let plain = (flat.len() + 1).saturating_mul(4).saturating_add(payload);
4687 let Some(out) = string::encode_only(string::Kind::Fsst, &values)? else {
4688 return Ok(None);
4689 };
4690 Ok((out.len() < plain).then_some(out))
4691}
4692
4693fn encoded_codes(codes: &[u32]) -> Result<Option<Vec<u8>>> {
4694 let wide: Vec<i64> = codes.iter().map(|code| i64::from(*code)).collect();
4695 let coded = integer::encode_with(&wide, &Codes)?;
4696 let plain = codes.len().saturating_mul(size_of::<u32>());
4697 Ok((coded.len() < plain).then_some(coded))
4698}
4699
4700fn encode(
4701 vector: &Vector,
4702 global: Option<&mut GlobalDictionary>,
4703) -> Result<(Vec<u8>, Option<Vec<u32>>)> {
4704 let ty = vector.logical_type();
4705 let flat = vector.flatten()?;
4707 let mut out = Vec::new();
4708 let mut global_codes = None;
4709 if let Some(global) = global {
4710 let mut codes = Vec::with_capacity(flat.len());
4711 for row in 0..flat.len() {
4712 let text = flat.text_at(row).unwrap_or("");
4713 let code = global.code(text)?;
4714 global.observe(code, flat.is_null_at(row))?;
4715 codes.push(code);
4716 }
4717 global_codes = Some(codes);
4718 }
4719 let membership = global_codes.as_deref().map(unique_codes);
4720 let dictionary = if global_codes.is_none() && ty == &LogicalType::Varchar {
4721 string_dictionary(&flat)?
4722 } else {
4723 None
4724 };
4725 let compressed_text =
4726 if global_codes.is_none() && dictionary.is_none() && ty == &LogicalType::Varchar {
4727 text_compressed(&flat)?
4728 } else {
4729 None
4730 };
4731 let packed_vector = if dictionary.is_none() && global_codes.is_none() {
4732 Some(flat.bit_packed()?)
4733 } else {
4734 None
4735 };
4736 let packed = packed_vector.as_ref().and_then(Vector::packed_parts);
4737 let coded = match global_codes.as_deref() {
4738 Some(codes) => encoded_codes(codes)?,
4739 None => None,
4740 };
4741 let cascade = if dictionary.is_none() && global_codes.is_none() {
4745 cascaded(&flat, ty, packed.as_ref())?
4746 } else {
4747 None
4748 };
4749 out.push(if coded.is_some() {
4750 4
4751 } else if cascade.is_some() {
4752 5
4753 } else if global_codes.is_some() {
4754 3
4755 } else if dictionary.is_some() {
4756 1
4757 } else if compressed_text.is_some() {
4758 6
4759 } else if packed.is_some() {
4760 2
4761 } else {
4762 0
4763 });
4764 let nulls = flat.validity();
4765 let flag = match nulls {
4766 Validity::AllValid => 0,
4767 Validity::AllInvalid => 1,
4768 Validity::Mask(_) => 2,
4769 };
4770 out.push(flag);
4771 if flag == 2 {
4772 for group in (0..vector.len()).step_by(8) {
4773 let mut bits = 0_u8;
4774 for bit in 0..8 {
4775 if group + bit < vector.len() && !flat.is_null_at(group + bit) {
4776 bits |= 1 << bit;
4777 }
4778 }
4779 out.push(bits);
4780 }
4781 }
4782 if let Some(coded) = coded {
4783 out.extend_from_slice(&coded);
4784 return Ok((out, membership));
4785 }
4786 if let Some(cascade) = cascade {
4787 out.extend_from_slice(&cascade);
4788 return Ok((out, membership));
4789 }
4790 if let Some(codes) = global_codes {
4791 for code in codes {
4792 put_u32(&mut out, code);
4793 }
4794 return Ok((out, membership));
4795 }
4796 if let Some(dictionary) = dictionary {
4797 out.extend_from_slice(&dictionary);
4798 return Ok((out, membership));
4799 }
4800 if let Some(compressed_text) = compressed_text {
4801 out.extend_from_slice(&compressed_text);
4802 return Ok((out, membership));
4803 }
4804 if let Some(packed) = packed {
4805 if packed.offset() != 0 {
4806 return Err(invalid("writer received a sliced packed vector"));
4807 }
4808 out.push(u8::try_from(packed.width()).map_err(|_| invalid("packed width overflow"))?);
4809 out.extend_from_slice(&packed.base().to_le_bytes());
4810 put_u32(
4811 &mut out,
4812 u32::try_from(packed.words().len()).map_err(|_| invalid("too many packed words"))?,
4813 );
4814 for word in packed.words() {
4815 put_u64(&mut out, *word);
4816 }
4817 return Ok((out, membership));
4818 }
4819 let data = flat.data().ok_or_else(|| invalid("scalar column did not flatten"))?;
4820 match (ty, data) {
4821 (LogicalType::TinyInt, Data::Int8(values)) => {
4822 for value in &**values {
4823 out.extend_from_slice(&value.to_le_bytes());
4824 }
4825 }
4826 (LogicalType::UTinyInt, Data::UInt8(values)) => {
4827 for value in &**values {
4828 out.extend_from_slice(&value.to_le_bytes());
4829 }
4830 }
4831 (LogicalType::SmallInt, Data::Int16(values)) => {
4832 for value in &**values {
4833 out.extend_from_slice(&value.to_le_bytes());
4834 }
4835 }
4836 (LogicalType::USmallInt, Data::UInt16(values)) => {
4837 for value in &**values {
4838 out.extend_from_slice(&value.to_le_bytes());
4839 }
4840 }
4841 (LogicalType::UInteger, Data::UInt32(values)) => {
4842 for value in &**values {
4843 out.extend_from_slice(&value.to_le_bytes());
4844 }
4845 }
4846 (LogicalType::UBigInt, Data::UInt64(values)) => {
4847 for value in &**values {
4848 out.extend_from_slice(&value.to_le_bytes());
4849 }
4850 }
4851 (LogicalType::Integer | LogicalType::Date, Data::Int32(values)) => {
4852 for value in &**values {
4853 out.extend_from_slice(&value.to_le_bytes());
4854 }
4855 }
4856 (LogicalType::BigInt | LogicalType::Timestamp, Data::Int64(values)) => {
4857 for value in &**values {
4858 out.extend_from_slice(&value.to_le_bytes());
4859 }
4860 }
4861 (LogicalType::Boolean, Data::Bool(values)) => {
4862 for value in &**values {
4863 out.push(u8::from(*value));
4864 }
4865 }
4866 (LogicalType::Decimal { .. }, Data::Int16(values)) => {
4869 for value in &**values {
4870 out.extend_from_slice(&value.to_le_bytes());
4871 }
4872 }
4873 (LogicalType::Decimal { .. }, Data::Int32(values)) => {
4874 for value in &**values {
4875 out.extend_from_slice(&value.to_le_bytes());
4876 }
4877 }
4878 (LogicalType::Decimal { .. }, Data::Int64(values)) => {
4879 for value in &**values {
4880 out.extend_from_slice(&value.to_le_bytes());
4881 }
4882 }
4883 (LogicalType::Decimal { .. }, Data::Int128(values)) => {
4884 for value in &**values {
4885 out.extend_from_slice(&value.to_le_bytes());
4886 }
4887 }
4888 (LogicalType::Varchar, Data::Varlen(values)) => {
4889 let mut bytes = Vec::new();
4890 put_u32(&mut out, 0);
4891 for row in 0..vector.len() {
4892 let value = values.bytes(row).ok_or_else(|| invalid("string view is invalid"))?;
4893 bytes.extend_from_slice(value);
4894 put_u32(
4895 &mut out,
4896 u32::try_from(bytes.len())
4897 .map_err(|_| invalid("string payload exceeds 4GiB"))?,
4898 );
4899 }
4900 out.extend_from_slice(&bytes);
4901 }
4902 _ => return Err(Error::not_implemented(format!("native page for {ty}"))),
4903 }
4904 Ok((out, membership))
4905}
4906
4907fn put_varint(out: &mut Vec<u8>, mut value: u32) {
4908 while value >= 0x80 {
4909 out.push((value as u8 & 0x7f) | 0x80);
4910 value >>= 7;
4911 }
4912 out.push(value as u8);
4913}
4914
4915fn unique_codes(codes: &[u32]) -> Vec<u32> {
4917 let mut unique = codes.to_vec();
4918 unique.sort_unstable();
4919 unique.dedup();
4920 unique
4921}
4922
4923fn merged_codes(lists: Vec<Vec<u32>>) -> Vec<u32> {
4929 let mut lists = lists;
4930 while lists.len() > 1 {
4931 let mut next = Vec::with_capacity(lists.len().div_ceil(2));
4932 for pair in lists.chunks(2) {
4933 match pair {
4934 [left, right] => next.push(merged_pair(left, right)),
4935 [only] => next.push(only.clone()),
4936 _ => {}
4937 }
4938 }
4939 lists = next;
4940 }
4941 lists.pop().unwrap_or_default()
4942}
4943
4944fn merged_pair(left: &[u32], right: &[u32]) -> Vec<u32> {
4945 let mut out = Vec::with_capacity(left.len().saturating_add(right.len()));
4946 let mut at = 0;
4947 let mut to = 0;
4948 while at < left.len() && to < right.len() {
4949 match left[at].cmp(&right[to]) {
4950 Ordering::Less => {
4951 out.push(left[at]);
4952 at += 1;
4953 }
4954 Ordering::Greater => {
4955 out.push(right[to]);
4956 to += 1;
4957 }
4958 Ordering::Equal => {
4959 out.push(left[at]);
4960 at += 1;
4961 to += 1;
4962 }
4963 }
4964 }
4965 out.extend_from_slice(&left[at..]);
4966 out.extend_from_slice(&right[to..]);
4967 out
4968}
4969
4970fn merged_range(ranges: impl Iterator<Item = Range>) -> Range {
4975 let mut merged = Range::default();
4976 let mut first = true;
4977 for range in ranges {
4978 merged.nulls = merged.nulls.saturating_add(range.nulls);
4979 merged.sum = match (merged.sum.take(), range.sum) {
4983 (Some(held), Some(next)) if !first => held.checked_add(next),
4984 (_, next) if first => next,
4985 _ => None,
4986 };
4987 merged.exact = if first { range.exact } else { merged.exact && range.exact };
4988 if first {
4989 merged.low = range.low;
4990 merged.high = range.high;
4991 first = false;
4992 continue;
4993 }
4994 merged.low = match (merged.low.take(), range.low) {
4995 (Some(held), Some(next)) => Some(held.smaller(next)),
4996 _ => None,
4997 };
4998 merged.high = match (merged.high.take(), range.high) {
4999 (Some(held), Some(next)) => Some(held.larger(next)),
5000 _ => None,
5001 };
5002 }
5003 merged
5004}
5005
5006fn shortened(bound: Option<Bound>, high: bool) -> Option<Bound> {
5019 match bound {
5020 Some(Bound::Bytes(mut value)) if value.len() > PART_BOUND_BYTES => {
5021 value.truncate(PART_BOUND_BYTES);
5022 if !high {
5023 return Some(Bound::Bytes(value));
5024 }
5025 while let Some(last) = value.pop() {
5026 if last < u8::MAX {
5027 value.push(last + 1);
5028 return Some(Bound::Bytes(value));
5029 }
5030 }
5031 None
5032 }
5033 other => other,
5034 }
5035}
5036
5037fn encode_part_ranges(ranges: &[Range]) -> Result<Vec<u8>> {
5045 let mut out = Vec::new();
5046 put_u32(
5047 &mut out,
5048 u32::try_from(ranges.len()).map_err(|_| invalid("too many parts in a stripe"))?,
5049 );
5050 for range in ranges {
5051 put_bound(&mut out, shortened(range.low.clone(), false).as_ref())?;
5052 put_bound(&mut out, shortened(range.high.clone(), true).as_ref())?;
5053 put_u32(&mut out, u32::try_from(range.nulls).map_err(|_| invalid("null count overflow"))?);
5054 }
5055 Ok(out)
5056}
5057
5058fn decode_part_ranges(bytes: &[u8]) -> Result<Vec<Range>> {
5060 let mut cur = Cursor { bytes, at: 0 };
5061 let parts = cur.u32()? as usize;
5062 let mut out = Vec::new();
5063 for _ in 0..parts {
5064 let low = cur.bound()?;
5065 let high = cur.bound()?;
5066 let nulls = cur.u32()? as usize;
5067 out.push(Range { low, high, nulls, exact: false, sum: None });
5068 }
5069 Ok(out)
5070}
5071
5072fn encode_sieves<'a>(sieves: impl Iterator<Item = &'a Option<Sieve>>) -> Result<Vec<u8>> {
5073 let held: Vec<&Option<Sieve>> = sieves.collect();
5074 let mut out = Vec::new();
5075 put_u32(
5076 &mut out,
5077 u32::try_from(held.len()).map_err(|_| invalid("too many parts in a stripe"))?,
5078 );
5079 for sieve in &held {
5080 let length = sieve.as_ref().map_or(0, Sieve::len);
5081 put_u32(&mut out, u32::try_from(length).map_err(|_| invalid("sieve length overflow"))?);
5082 }
5083 for sieve in held.into_iter().flatten() {
5085 out.extend_from_slice(&sieve.to_bytes());
5086 }
5087 Ok(out)
5088}
5089
5090fn decode_sieves(bytes: &[u8]) -> Result<Vec<Option<Sieve>>> {
5096 let parts = u32::from_le_bytes(
5097 bytes
5098 .get(..4)
5099 .ok_or_else(|| invalid("sieve page is truncated"))?
5100 .try_into()
5101 .map_err(|_| invalid("sieve page is truncated"))?,
5102 ) as usize;
5103 let mut lengths = Vec::with_capacity(parts);
5104 for part in 0..parts {
5105 let at = 4 + part * 4;
5106 let field = bytes.get(at..at + 4).ok_or_else(|| invalid("sieve page is truncated"))?;
5107 lengths.push(u32::from_le_bytes(
5108 field.try_into().map_err(|_| invalid("sieve page is truncated"))?,
5109 ) as usize);
5110 }
5111 let mut at = 4 + parts * 4;
5112 let mut out = Vec::with_capacity(parts);
5113 for length in lengths {
5114 if length == 0 {
5115 out.push(None);
5116 continue;
5117 }
5118 let end = at.checked_add(length).ok_or_else(|| invalid("sieve page is truncated"))?;
5119 let field = bytes.get(at..end).ok_or_else(|| invalid("sieve page is truncated"))?;
5120 out.push(Sieve::from_bytes(field));
5121 at = end;
5122 }
5123 if at != bytes.len() {
5124 return Err(invalid("sieve page has trailing bytes"));
5125 }
5126 Ok(out)
5127}
5128
5129fn encode_membership(unique: &[u32]) -> Vec<u8> {
5135 let mut out = Vec::with_capacity(unique.len().saturating_mul(2).saturating_add(5));
5136 put_varint(&mut out, u32::try_from(unique.len()).unwrap_or(u32::MAX));
5137 let mut previous = 0;
5138 for (at, &code) in unique.iter().enumerate() {
5139 put_varint(&mut out, if at == 0 { code } else { code - previous });
5140 previous = code;
5141 }
5142 out
5143}
5144
5145fn take_varint(bytes: &[u8], at: &mut usize) -> Result<u32> {
5146 let mut value = 0_u32;
5147 for shift in (0..35).step_by(7) {
5148 let byte = *bytes.get(*at).ok_or_else(|| invalid("membership varint is truncated"))?;
5149 *at += 1;
5150 let part = u32::from(byte & 0x7f);
5151 if shift == 28 && part > 0x0f {
5152 return Err(invalid("membership varint overflow"));
5153 }
5154 value = value
5155 .checked_add(
5156 part.checked_shl(shift).ok_or_else(|| invalid("membership varint overflow"))?,
5157 )
5158 .ok_or_else(|| invalid("membership varint overflow"))?;
5159 if byte & 0x80 == 0 {
5160 return Ok(value);
5161 }
5162 }
5163 Err(invalid("membership varint is too long"))
5164}
5165
5166fn decode_membership(bytes: &[u8]) -> Result<Vec<u32>> {
5167 let mut at = 0;
5168 let count = take_varint(bytes, &mut at)? as usize;
5169 let mut codes = Vec::with_capacity(count);
5170 let mut previous = 0_u32;
5171 for index in 0..count {
5172 let delta = take_varint(bytes, &mut at)?;
5173 let code = if index == 0 {
5174 delta
5175 } else {
5176 previous.checked_add(delta).ok_or_else(|| invalid("membership code overflow"))?
5177 };
5178 if index > 0 && code <= previous {
5179 return Err(invalid("membership codes are not increasing"));
5180 }
5181 codes.push(code);
5182 previous = code;
5183 }
5184 if at != bytes.len() {
5185 return Err(invalid("membership page has trailing bytes"));
5186 }
5187 Ok(codes)
5188}
5189
5190fn string_dictionary(vector: &Vector) -> Result<Option<Vec<u8>>> {
5191 let mut by_text = HashMap::new();
5192 let mut values = Vec::new();
5193 let mut codes = Vec::with_capacity(vector.len());
5194 let mut plain_bytes = 0_usize;
5195 for row in 0..vector.len() {
5196 let text = vector.text_at(row).unwrap_or("");
5197 plain_bytes = plain_bytes.saturating_add(text.len());
5198 let code = match by_text.get(text) {
5199 Some(&code) => code,
5200 None => {
5201 let code = u32::try_from(values.len())
5202 .map_err(|_| invalid("too many dictionary values"))?;
5203 by_text.insert(text, code);
5204 values.push(text);
5205 code
5206 }
5207 };
5208 codes.push(code);
5209 }
5210 let dictionary_bytes = values.iter().map(|value| value.len()).sum::<usize>();
5211 let encoded = 8_usize
5212 .saturating_add((values.len() + 1).saturating_mul(4))
5213 .saturating_add(dictionary_bytes)
5214 .saturating_add(codes.len().saturating_mul(4));
5215 let plain = (vector.len() + 1).saturating_mul(4).saturating_add(plain_bytes);
5216 if encoded >= plain {
5217 return Ok(None);
5218 }
5219 let mut out = Vec::with_capacity(encoded);
5220 put_u32(
5221 &mut out,
5222 u32::try_from(values.len()).map_err(|_| invalid("too many dictionary values"))?,
5223 );
5224 put_u32(
5225 &mut out,
5226 u32::try_from(dictionary_bytes).map_err(|_| invalid("dictionary payload exceeds 4GiB"))?,
5227 );
5228 let mut offset = 0_u32;
5229 put_u32(&mut out, offset);
5230 for value in &values {
5231 offset = offset
5232 .checked_add(
5233 u32::try_from(value.len()).map_err(|_| invalid("dictionary value is too long"))?,
5234 )
5235 .ok_or_else(|| invalid("dictionary payload exceeds 4GiB"))?;
5236 put_u32(&mut out, offset);
5237 }
5238 for value in values {
5239 out.extend_from_slice(value.as_bytes());
5240 }
5241 for code in codes {
5242 put_u32(&mut out, code);
5243 }
5244 Ok(Some(out))
5245}
5246
5247struct EncodedDictionary {
5248 index: Vec<u8>,
5249 ranks: Vec<u8>,
5250 payload: Vec<Vec<u8>>,
5253}
5254
5255fn head(bytes: &[u8]) -> u64 {
5257 let mut word = [0; 8];
5258 let take = bytes.len().min(8);
5259 word[..take].copy_from_slice(&bytes[..take]);
5260 u64::from_be_bytes(word)
5261}
5262
5263fn rankings(dictionaries: &[Option<GlobalDictionary>]) -> Result<Vec<Vec<(u64, u32)>>> {
5271 let present =
5272 dictionaries.iter().enumerate().filter(|(_, held)| held.is_some()).map(|(at, _)| at);
5273 let present = present.collect::<Vec<_>>();
5274 let mut orders = vec![Vec::new(); dictionaries.len()];
5275 let workers = std::thread::available_parallelism()
5276 .map_or(1, usize::from)
5277 .min(MAX_FREQUENCY_WORKERS)
5278 .min(present.len());
5279 if workers <= 1 {
5280 for at in present {
5281 if let Some(dictionary) = &dictionaries[at] {
5282 orders[at] = dictionary.ranked();
5283 }
5284 }
5285 return Ok(orders);
5286 }
5287 let width = present.len().div_ceil(workers);
5288 let pieces = std::thread::scope(|scope| {
5289 present
5290 .chunks(width)
5291 .map(|columns| {
5292 scope.spawn(|| {
5293 columns
5294 .iter()
5295 .filter_map(|&at| dictionaries[at].as_ref().map(|held| (at, held.ranked())))
5296 .collect::<Vec<_>>()
5297 })
5298 })
5299 .collect::<Vec<_>>()
5300 .into_iter()
5301 .map(|handle| {
5302 handle.join().map_err(|_| Error::internal("a dictionary sort worker panicked"))
5303 })
5304 .collect::<Result<Vec<_>>>()
5305 })?;
5306 for piece in pieces {
5307 for (at, order) in piece {
5308 orders[at] = order;
5309 }
5310 }
5311 Ok(orders)
5312}
5313
5314fn encode_global_dictionary(
5315 dictionary: GlobalDictionary,
5316 order: &[(u64, u32)],
5317) -> Result<EncodedDictionary> {
5318 let values = dictionary.offsets.len() - 1;
5319 if order.len() != values {
5320 return Err(invalid("global dictionary order does not cover its values"));
5321 }
5322 let blocks = values.div_ceil(TEXT_PAYLOAD_VALUES);
5323 let payload = encode_payload(&dictionary)?;
5324 if payload.len() != blocks {
5325 return Err(invalid("global dictionary payload is not the blocks it says it is"));
5326 }
5327 let (ranks, rank_ends) = encode_ranks(order, code_width(values))?;
5328 let rank_blocks = values.div_ceil(TEXT_RANK_BLOCK);
5329 let offset_bits = offset_width(&dictionary.offsets);
5330 let mut index = Vec::with_capacity(
5331 DICTIONARY_HEADER + offset_bytes(values, offset_bits) + (blocks + rank_blocks) * 16,
5332 );
5333 put_u32(
5334 &mut index,
5335 u32::try_from(values).map_err(|_| invalid("global dictionary has too many values"))?,
5336 );
5337 put_u32(&mut index, TEXT_PAYLOAD_VALUES as u32);
5338 put_u32(
5339 &mut index,
5340 u32::try_from(blocks).map_err(|_| invalid("global dictionary has too many blocks"))?,
5341 );
5342 put_u32(&mut index, offset_bits as u32);
5343 encode_offsets(&dictionary.offsets, offset_bits, &mut index)?;
5344 let mut at = 0_u64;
5348 for block in &payload {
5349 at = at
5350 .checked_add(block.len() as u64)
5351 .ok_or_else(|| invalid("global dictionary payload overflow"))?;
5352 put_u64(&mut index, at);
5353 }
5354 for block in &payload {
5355 put_u64(&mut index, checksum(block));
5356 }
5357 if rank_ends.len() != rank_blocks {
5360 return Err(invalid("global dictionary order is not the blocks it says it is"));
5361 }
5362 for end in &rank_ends {
5363 put_u64(&mut index, *end);
5364 }
5365 let mut at = 0_usize;
5366 for end in &rank_ends {
5367 let end = usize::try_from(*end).map_err(|_| invalid("global dictionary order overflow"))?;
5368 put_u64(&mut index, checksum(&ranks[at..end]));
5369 at = end;
5370 }
5371 Ok(EncodedDictionary { index, ranks, payload })
5372}
5373
5374const PAYLOAD_SAMPLE_BLOCKS: usize = 8;
5381
5382fn payload_shapes() -> Vec<chooser::Settled> {
5408 let integers = vec![integer::Kind::Packed];
5409 [
5410 vec![string::Kind::Front, string::Kind::Lz],
5411 vec![string::Kind::Lz, string::Kind::Fsst],
5412 vec![string::Kind::Lz, string::Kind::Plain],
5413 vec![string::Kind::Fsst],
5414 vec![string::Kind::Plain],
5415 ]
5416 .into_iter()
5417 .map(|strings| chooser::Settled::new(strings, integers.clone()))
5418 .collect()
5419}
5420
5421fn encode_payload(dictionary: &GlobalDictionary) -> Result<Vec<Vec<u8>>> {
5427 let values = dictionary.offsets.len() - 1;
5428 let blocks = values.div_ceil(TEXT_PAYLOAD_VALUES);
5429 let run = |block: usize| {
5430 let first = block * TEXT_PAYLOAD_VALUES;
5431 let last = (first + TEXT_PAYLOAD_VALUES).min(values);
5432 (first..last)
5433 .map(|value| {
5434 let from = dictionary.offsets[value] as usize;
5435 let to = dictionary.offsets[value + 1] as usize;
5436 &dictionary.payload[from..to]
5437 })
5438 .collect::<Vec<_>>()
5439 };
5440 let shape = (blocks > PAYLOAD_SAMPLE_BLOCKS).then(|| settle_shape(&run, blocks)).transpose()?;
5443 let one = |block: usize| match &shape {
5444 Some(shape) => string::encode_with(&run(block), shape),
5445 None => string::encode(&run(block)),
5446 };
5447 let workers = std::thread::available_parallelism()
5448 .map_or(1, usize::from)
5449 .min(MAX_FREQUENCY_WORKERS)
5450 .min(blocks);
5451 if workers <= 1 {
5452 return (0..blocks).map(one).collect();
5453 }
5454 let next = AtomicUsize::new(0);
5455 let pieces = std::thread::scope(|scope| {
5456 (0..workers)
5457 .map(|_| {
5458 scope.spawn(|| {
5459 let mut mine = Vec::new();
5460 loop {
5461 let block = next.fetch_add(1, Atomic::Relaxed);
5462 if block >= blocks {
5463 break;
5464 }
5465 mine.push((block, one(block)?));
5466 }
5467 Ok(mine)
5468 })
5469 })
5470 .collect::<Vec<_>>()
5471 .into_iter()
5472 .map(|handle| {
5473 handle.join().map_err(|_| Error::internal("a dictionary encode worker panicked"))?
5474 })
5475 .collect::<Result<Vec<_>>>()
5476 })?;
5477 let mut payload = vec![Vec::new(); blocks];
5478 for piece in pieces {
5479 for (block, bytes) in piece {
5480 payload[block] = bytes;
5481 }
5482 }
5483 Ok(payload)
5484}
5485
5486fn settle_shape<'a>(
5494 run: &dyn Fn(usize) -> Vec<&'a [u8]>,
5495 blocks: usize,
5496) -> Result<chooser::Settled> {
5497 let last = blocks - 1;
5498 let sample = (0..PAYLOAD_SAMPLE_BLOCKS)
5499 .map(|region| run(region * last / (PAYLOAD_SAMPLE_BLOCKS - 1)))
5500 .collect::<Vec<_>>();
5501 let mut best: Option<(chooser::Settled, usize)> = None;
5502 for shape in payload_shapes() {
5503 let mut size = 0;
5504 for block in &sample {
5505 size += string::encode_with(block, &shape)?.len();
5506 }
5507 if best.as_ref().is_none_or(|(_, smallest)| size < *smallest) {
5508 best = Some((shape, size));
5509 }
5510 }
5511 best.map(|(shape, _)| shape)
5512 .ok_or_else(|| invalid("no shape applies to a global dictionary payload"))
5513}
5514
5515fn encode_ranks(order: &[(u64, u32)], code_bits: usize) -> Result<(Vec<u8>, Vec<u64>)> {
5522 let mut out = Vec::with_capacity(order.len() * 4);
5523 let mut ends = Vec::with_capacity(order.len().div_ceil(TEXT_RANK_BLOCK));
5524 let mut heads = Vec::with_capacity(TEXT_RANK_BLOCK);
5525 let mut codes = Vec::with_capacity(TEXT_RANK_BLOCK);
5526 for block in order.chunks(TEXT_RANK_BLOCK) {
5527 let base = block.first().map_or(0, |&(head, _)| head);
5530 let span = block.last().map_or(0, |&(head, _)| head.wrapping_sub(base));
5531 let width = (u64::BITS - span.leading_zeros()) as usize;
5532 heads.clear();
5533 codes.clear();
5534 for &(head, code) in block {
5535 heads.push(head.wrapping_sub(base));
5536 codes.push(u64::from(code));
5537 }
5538 put_u64(&mut out, base);
5539 out.push(width as u8);
5540 bitpack::pack_tail(&heads, width, &mut out)
5541 .map_err(|_| invalid("global dictionary heads do not pack"))?;
5542 bitpack::pack_tail(&codes, code_bits, &mut out)
5543 .map_err(|_| invalid("global dictionary codes do not pack"))?;
5544 ends.push(out.len() as u64);
5545 }
5546 Ok((out, ends))
5547}
5548
5549fn open_global_dictionary(
5556 file: Arc<File>,
5557 page: Page,
5558 ty: &LogicalType,
5559 keep_budget: usize,
5560) -> Result<Vector> {
5561 if ty != &LogicalType::Varchar {
5562 return Err(invalid("global dictionary belongs to a non-string column"));
5563 }
5564 let mut header = [0; DICTIONARY_HEADER];
5565 read_at(&file, page.offset, &mut header)?;
5566 let count = u32::from_le_bytes(header[0..4].try_into().expect("four bytes")) as usize;
5567 let per_block = u32::from_le_bytes(header[4..8].try_into().expect("four bytes")) as usize;
5568 let blocks = u32::from_le_bytes(header[8..12].try_into().expect("four bytes")) as usize;
5569 let offset_bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
5570 if per_block != TEXT_PAYLOAD_VALUES {
5571 return Err(invalid("global dictionary block width differs"));
5572 }
5573 if blocks != count.div_ceil(TEXT_PAYLOAD_VALUES) {
5574 return Err(invalid("global dictionary block count differs from its value count"));
5575 }
5576 if offset_bits > u32::BITS as usize {
5577 return Err(invalid("global dictionary packs offsets past a payload"));
5578 }
5579 let offset_len = offset_bytes(count, offset_bits);
5580 let ranks = count;
5585 let rank_blocks = ranks.div_ceil(TEXT_RANK_BLOCK);
5586 let hash_len = blocks
5589 .checked_add(rank_blocks)
5590 .and_then(|words| words.checked_mul(16))
5591 .ok_or_else(|| invalid("global dictionary block count overflow"))?;
5592 let index_len = DICTIONARY_HEADER
5593 .checked_add(offset_len)
5594 .and_then(|len| len.checked_add(hash_len))
5595 .ok_or_else(|| invalid("global dictionary header overflow"))?;
5596 if index_len > page.length as usize {
5597 return Err(invalid("global dictionary offset index exceeds its page"));
5598 }
5599 let mut index = vec![0; index_len];
5600 index[..DICTIONARY_HEADER].copy_from_slice(&header);
5601 read_at(&file, page.offset + DICTIONARY_HEADER as u64, &mut index[DICTIONARY_HEADER..])?;
5602 if checksum(&index) != page.hash {
5603 return Err(invalid("global dictionary index checksum differs"));
5604 }
5605 let offsets = index[DICTIONARY_HEADER..DICTIONARY_HEADER + offset_len].to_vec();
5606 let mut words = index[DICTIONARY_HEADER + offset_len..]
5607 .chunks_exact(8)
5608 .map(|part| u64::from_le_bytes(part.try_into().expect("eight bytes")))
5609 .collect::<Vec<_>>();
5610 let mut hashes = words.split_off(blocks);
5611 let mut rank_ends = hashes.split_off(blocks);
5612 let rank_hashes = rank_ends.split_off(rank_blocks);
5613 let ends = words;
5614 if rank_ends.windows(2).any(|pair| pair[0] >= pair[1]) {
5617 return Err(invalid("global dictionary order blocks do not rise"));
5618 }
5619 let rank_len = usize::try_from(rank_ends.last().copied().unwrap_or_default())
5620 .map_err(|_| invalid("global dictionary rank overflow"))?;
5621 let body_len = index_len
5622 .checked_add(rank_len)
5623 .ok_or_else(|| invalid("global dictionary header overflow"))?;
5624 if body_len > page.length as usize {
5625 return Err(invalid("global dictionary order exceeds its page"));
5626 }
5627 let stored_len = page.length as usize - body_len;
5630 if ends.last().copied().unwrap_or_default() as usize != stored_len
5631 || ends.windows(2).any(|pair| pair[0] > pair[1])
5632 {
5633 return Err(invalid("global dictionary blocks do not bound the payload"));
5634 }
5635 Vector::external_text(
5636 LogicalType::Varchar,
5637 Arc::new(NativeText {
5638 file,
5639 values: count,
5640 offsets,
5641 offset_bits,
5642 ranks,
5643 rank_at: page.offset + index_len as u64,
5644 rank_ends,
5645 rank_hashes,
5646 rank_blocks: (0..rank_blocks).map(|_| OnceLock::new()).collect(),
5647 code_bits: code_width(count),
5648 code_ranks: OnceLock::new(),
5649 payload: page.offset + body_len as u64,
5650 ends,
5651 hashes,
5652 blocks: (0..blocks).map(|_| OnceLock::new()).collect(),
5653 keep_budget,
5654 payload_kept: AtomicUsize::new(0),
5655 searched: Mutex::new(HashMap::new()),
5656 }),
5657 )
5658}
5659
5660fn decode(
5661 ty: &LogicalType,
5662 rows: usize,
5663 bytes: &[u8],
5664 global: Option<Arc<Vector>>,
5665) -> Result<Vector> {
5666 let mut cur = Cursor { bytes, at: 0 };
5667 let codec = cur.u8()?;
5668 let flag = cur.u8()?;
5669 let validity = match flag {
5670 0 => Validity::AllValid,
5671 1 => Validity::AllInvalid,
5672 2 => {
5673 let mask = cur.take(rows.div_ceil(8))?;
5674 Validity::from_iter(rows, |row| mask[row / 8] >> (row % 8) & 1 == 1)
5675 }
5676 _ => return Err(invalid("page validity tag differs")),
5677 };
5678 if codec == 1 {
5679 if ty != &LogicalType::Varchar {
5680 return Err(invalid("dictionary codec belongs to a non-string page"));
5681 }
5682 let count = cur.u32()? as usize;
5683 let payload_len = cur.u32()? as usize;
5684 let offset_bytes = cur.take(
5685 (count + 1)
5686 .checked_mul(4)
5687 .ok_or_else(|| invalid("dictionary offset count overflow"))?,
5688 )?;
5689 let offsets = offset_bytes
5690 .chunks_exact(4)
5691 .map(|part| u32::from_le_bytes(part.try_into().expect("four bytes")))
5692 .collect::<Vec<_>>();
5693 let payload = cur.take(payload_len)?.to_vec();
5694 if offsets.first() != Some(&0)
5695 || offsets.last().copied().map(|last| last as usize) != Some(payload.len())
5696 || offsets.windows(2).any(|pair| pair[0] > pair[1])
5697 {
5698 return Err(invalid("dictionary offsets do not bound the payload"));
5699 }
5700 let mut strings = StringColumn::over(Buffer::from_vec(payload).into_page());
5703 for pair in offsets.windows(2) {
5704 strings.push_in_place(pair[0] as usize, (pair[1] - pair[0]) as usize)?;
5705 }
5706 let mut codes = Vec::with_capacity(rows);
5707 for _ in 0..rows {
5708 codes.push(cur.u32()?);
5709 }
5710 if codes.iter().any(|code| *code as usize >= count) {
5711 return Err(invalid("dictionary code is out of range"));
5712 }
5713 if cur.at != bytes.len() {
5714 return Err(invalid("dictionary page has trailing bytes"));
5715 }
5716 let dictionary = Vector::flat(LogicalType::Varchar, Data::Varlen(strings))?;
5717 return Ok(Vector::dictionary(codes, dictionary)?.with_validity(validity));
5718 }
5719 if codec == 3 || codec == 4 {
5720 let dictionary = global.ok_or_else(|| invalid("global code page has no dictionary"))?;
5721 let codes = if codec == 4 {
5722 let wide = integer::decode(&bytes[cur.at..])?;
5725 if wide.len() != rows {
5726 return Err(invalid("encoded code page holds the wrong number of rows"));
5727 }
5728 let mut codes = Vec::with_capacity(wide.len());
5735 let mut seen = 0_i64;
5736 for &code in &wide {
5737 seen |= code;
5738 codes.push(code as u32);
5739 }
5740 if seen < 0 || seen > i64::from(u32::MAX) {
5741 return Err(invalid("code is not a code"));
5742 }
5743 codes
5744 } else {
5745 let mut codes = Vec::with_capacity(rows);
5746 for _ in 0..rows {
5747 codes.push(cur.u32()?);
5748 }
5749 if cur.at != bytes.len() {
5750 return Err(invalid("global code page has trailing bytes"));
5751 }
5752 codes
5753 };
5754 let highest = codes.iter().copied().max();
5755 return Ok(Vector::stable_dictionary_validated(codes, dictionary, highest)?
5756 .with_validity(validity));
5757 }
5758 if codec == 6 {
5759 if ty != &LogicalType::Varchar {
5760 return Err(invalid("compressed text codec belongs to a non-string page"));
5761 }
5762 let (payload, ends) = string::decode_flat(&bytes[cur.at..])?.into_parts();
5766 if ends.len() != rows {
5767 return Err(invalid("compressed text page holds the wrong number of rows"));
5768 }
5769 let mut values = StringColumn::over(Buffer::from_vec(payload).into_page());
5772 let mut start = 0;
5773 for end in ends {
5774 let len = end
5775 .checked_sub(start)
5776 .ok_or_else(|| invalid("compressed text value ends before it starts"))?;
5777 values.push_in_place(start, len)?;
5778 start = end;
5779 }
5780 return Ok(Vector::flat(ty.clone(), Data::Varlen(values))?.with_validity(validity));
5781 }
5782 if codec == 5 {
5783 let values = integer::decode(&bytes[cur.at..])?;
5785 if values.len() != rows {
5786 return Err(invalid("cascade page holds the wrong number of rows"));
5787 }
5788 let data = narrowed(ty, values)?;
5789 return Ok(Vector::flat(ty.clone(), data)?.with_validity(validity));
5790 }
5791 if codec == 2 {
5792 let width = u32::from(cur.u8()?);
5793 let base = i128::from_le_bytes(cur.take(16)?.try_into().expect("sixteen bytes"));
5794 let count = cur.u32()? as usize;
5795 let mut words = Vec::with_capacity(count);
5796 for _ in 0..count {
5797 words.push(cur.u64()?);
5798 }
5799 if cur.at != bytes.len() {
5800 return Err(invalid("packed page has trailing bytes"));
5801 }
5802 return Ok(Vector::packed(ty.clone(), words, width, base, rows)?.with_validity(validity));
5803 }
5804 if codec != 0 {
5805 return Err(invalid("page codec is unknown"));
5806 }
5807 let data = match ty {
5808 LogicalType::TinyInt => {
5809 let values = cur.take(rows)?;
5810 Data::Int8(values.iter().map(|item| *item as i8).collect::<Vec<_>>().into())
5811 }
5812 LogicalType::UTinyInt => Data::UInt8(cur.take(rows)?.to_vec().into()),
5813 LogicalType::SmallInt => {
5814 let values =
5815 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5816 Data::Int16(
5817 values
5818 .chunks_exact(2)
5819 .map(|item| i16::from_le_bytes(item.try_into().expect("two bytes")))
5820 .collect::<Vec<_>>()
5821 .into(),
5822 )
5823 }
5824 LogicalType::USmallInt => {
5825 let values =
5826 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5827 Data::UInt16(
5828 values
5829 .chunks_exact(2)
5830 .map(|item| u16::from_le_bytes(item.try_into().expect("two bytes")))
5831 .collect::<Vec<_>>()
5832 .into(),
5833 )
5834 }
5835 LogicalType::UInteger => {
5836 let values =
5837 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5838 Data::UInt32(
5839 values
5840 .chunks_exact(4)
5841 .map(|item| u32::from_le_bytes(item.try_into().expect("four bytes")))
5842 .collect::<Vec<_>>()
5843 .into(),
5844 )
5845 }
5846 LogicalType::UBigInt => {
5847 let values =
5848 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5849 Data::UInt64(
5850 values
5851 .chunks_exact(8)
5852 .map(|item| u64::from_le_bytes(item.try_into().expect("eight bytes")))
5853 .collect::<Vec<_>>()
5854 .into(),
5855 )
5856 }
5857 LogicalType::Integer | LogicalType::Date => {
5858 let values =
5859 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5860 Data::Int32(
5861 values
5862 .chunks_exact(4)
5863 .map(|item| i32::from_le_bytes(item.try_into().expect("four bytes")))
5864 .collect::<Vec<_>>()
5865 .into(),
5866 )
5867 }
5868 LogicalType::BigInt | LogicalType::Timestamp => {
5869 let values =
5870 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5871 Data::Int64(
5872 values
5873 .chunks_exact(8)
5874 .map(|item| i64::from_le_bytes(item.try_into().expect("eight bytes")))
5875 .collect::<Vec<_>>()
5876 .into(),
5877 )
5878 }
5879 LogicalType::Boolean => {
5880 let values = cur.take(rows)?;
5881 if values.iter().any(|value| *value > 1) {
5882 return Err(invalid("boolean page has another value"));
5883 }
5884 Data::Bool(values.iter().map(|value| *value == 1).collect::<Vec<_>>().into())
5885 }
5886 LogicalType::Decimal { .. } => match ty.physical() {
5889 PhysicalType::Int16 => {
5890 let values =
5891 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5892 Data::Int16(
5893 values
5894 .chunks_exact(2)
5895 .map(|item| i16::from_le_bytes(item.try_into().expect("two bytes")))
5896 .collect::<Vec<_>>()
5897 .into(),
5898 )
5899 }
5900 PhysicalType::Int32 => {
5901 let values =
5902 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5903 Data::Int32(
5904 values
5905 .chunks_exact(4)
5906 .map(|item| i32::from_le_bytes(item.try_into().expect("four bytes")))
5907 .collect::<Vec<_>>()
5908 .into(),
5909 )
5910 }
5911 PhysicalType::Int64 => {
5912 let values =
5913 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5914 Data::Int64(
5915 values
5916 .chunks_exact(8)
5917 .map(|item| i64::from_le_bytes(item.try_into().expect("eight bytes")))
5918 .collect::<Vec<_>>()
5919 .into(),
5920 )
5921 }
5922 _ => {
5923 let values =
5924 cur.take(rows.checked_mul(16).ok_or_else(|| invalid("page size overflow"))?)?;
5925 Data::Int128(
5926 values
5927 .chunks_exact(16)
5928 .map(|item| i128::from_le_bytes(item.try_into().expect("sixteen bytes")))
5929 .collect::<Vec<_>>()
5930 .into(),
5931 )
5932 }
5933 },
5934 LogicalType::Varchar => {
5935 let offset_bytes = cur
5936 .take((rows + 1).checked_mul(4).ok_or_else(|| invalid("offset count overflow"))?)?;
5937 let offsets = offset_bytes
5938 .chunks_exact(4)
5939 .map(|part| u32::from_le_bytes(part.try_into().expect("four bytes")))
5940 .collect::<Vec<_>>();
5941 let payload = cur.take(bytes.len() - cur.at)?.to_vec();
5942 if offsets.first() != Some(&0)
5943 || offsets.last().copied().map(|last| last as usize) != Some(payload.len())
5944 || offsets.windows(2).any(|pair| pair[0] > pair[1])
5945 {
5946 return Err(invalid("string offsets do not bound the payload"));
5947 }
5948 let mut values = StringColumn::over(Buffer::from_vec(payload).into_page());
5952 for pair in offsets.windows(2) {
5953 values.push_in_place(pair[0] as usize, (pair[1] - pair[0]) as usize)?;
5954 }
5955 Data::Varlen(values)
5956 }
5957 _ => return Err(Error::not_implemented(format!("native page for {ty}"))),
5958 };
5959 if cur.at != bytes.len() {
5960 return Err(invalid("page has trailing bytes"));
5961 }
5962 Ok(Vector::flat(ty.clone(), data)?.with_validity(validity))
5963}
5964
5965#[cfg(test)]
5966mod tests {
5967 use std::fs;
5968 use std::io::{Seek, SeekFrom, Write};
5969 use std::path::PathBuf;
5970 use std::time::{SystemTime, UNIX_EPOCH};
5971
5972 use rudb_common::Stat;
5973 use rudb_common::Value;
5974 use rudb_common::bounds::{Frequencies, Op, Remainder, Zones};
5975 use rudb_common::stat::Provenance;
5976
5977 use super::*;
5978
5979 #[test]
5980 fn checksum_matches_fixed_vectors() {
5981 assert_eq!(checksum(b""), 0xef46_db37_51d8_e999);
5982 assert_eq!(checksum(b"a"), 0xd24e_c4f1_a98c_6e5b);
5983 assert_eq!(checksum(b"abc"), 0x44bc_2cf5_ad77_0999);
5984 }
5985
5986 fn path(label: &str) -> PathBuf {
5987 let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
5988 std::env::temp_dir().join(format!("rudb-native-{label}-{}-{stamp}.rdb", std::process::id()))
5989 }
5990
5991 #[test]
5993 fn a_read_at_an_offset_ignores_where_another_thread_left_the_cursor() {
5994 const SPANS: usize = 64;
5995 const SPAN: usize = 512;
5996 let path = path("positional");
5997 let content: Vec<u8> =
5998 (0..SPANS).flat_map(|span| std::iter::repeat_n(span as u8, SPAN)).collect();
5999 fs::write(&path, &content).expect("the file is written");
6000 let file = Arc::new(File::open(&path).expect("the file opens"));
6001 std::thread::scope(|scope| {
6002 for _ in 0..8 {
6003 let file = Arc::clone(&file);
6004 scope.spawn(move || {
6005 for _ in 0..64 {
6006 for span in 0..SPANS {
6007 let mut bytes = [0_u8; SPAN];
6008 read_at(&file, (span * SPAN) as u64, &mut bytes)
6009 .expect("the span reads");
6010 assert!(
6011 bytes.iter().all(|byte| *byte == span as u8),
6012 "span {span} came back as {}",
6013 bytes[0],
6014 );
6015 }
6016 }
6017 });
6018 }
6019 });
6020 let mut past = [0_u8; SPAN];
6021 let end = (SPANS * SPAN) as u64;
6022 let error = read_at(&file, end, &mut past).expect_err("a read past the end is refused");
6023 assert!(error.message().contains("ends before its declared length"), "{error}");
6024 drop(file);
6025 let _ = fs::remove_file(&path);
6026 }
6027
6028 #[test]
6034 fn a_writer_puts_a_page_where_it_said_it_did_wherever_the_cursor_has_got_to() {
6035 let path = path("cursor");
6036 let mut writer = Writer::create(
6037 &path,
6038 "items",
6039 vec![
6040 Field::required("id", LogicalType::Integer),
6041 Field::new("text", LogicalType::Varchar),
6042 ],
6043 )
6044 .expect("new file");
6045 writer.append(&sample()).expect("first part");
6046 writer.file.seek(SeekFrom::Start(0)).expect("the cursor goes back to the header");
6047 writer.append(&sample()).expect("second part");
6048 writer.file.seek(SeekFrom::Start(1)).expect("and somewhere useless again");
6049 writer.finish().expect("commit");
6050 let reader = Reader::open(&path).expect("reopen from disk");
6051 assert_eq!(reader.table().rows(), 6);
6052 let ids = reader.read(0, &[0]).expect("the integer page reads back");
6053 assert_eq!(ids.value_at(0, 0), Value::Integer(4));
6054 assert_eq!(ids.value_at(2, 0), Value::Integer(-2));
6055 let text = reader.read(1, &[1]).expect("the text page reads back");
6056 assert_eq!(text.value_at(1, 0), Value::Null);
6057 assert_eq!(text.value_at(2, 0), Value::Varchar("long text after a slash".into()));
6058 let end = reader.table().stripes().iter().flat_map(|stripe| {
6061 stripe
6062 .pages
6063 .iter()
6064 .map(|page| page.offset + u64::from(page.length))
6065 .chain(std::iter::once(stripe.index.offset + u64::from(stripe.index.length)))
6066 });
6067 let last = end.fold(HEADER, u64::max);
6068 let directory = fs::metadata(&path).expect("the file is there").len();
6069 assert!(last <= directory, "a page runs to {last} in a file of {directory} bytes");
6070 fs::remove_file(path).expect("remove scratch file");
6071 }
6072
6073 fn dictionary_index_len(header: &[u8; DICTIONARY_HEADER]) -> u64 {
6079 let count = u64::from(u32::from_le_bytes(header[0..4].try_into().expect("four bytes")));
6080 let blocks = u64::from(u32::from_le_bytes(header[8..12].try_into().expect("four bytes")));
6081 let bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
6082 let rank_blocks = count.div_ceil(TEXT_RANK_BLOCK as u64);
6083 DICTIONARY_HEADER as u64
6084 + offset_bytes(count as usize, bits) as u64
6085 + (blocks + rank_blocks) * 16
6086 }
6087
6088 fn last_rank_end(file: &File, offset: u64, header: &[u8; DICTIONARY_HEADER]) -> u64 {
6090 let count = u64::from(u32::from_le_bytes(header[0..4].try_into().expect("four bytes")));
6091 let blocks = u64::from(u32::from_le_bytes(header[8..12].try_into().expect("four bytes")));
6092 let bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
6093 let rank_blocks = count.div_ceil(TEXT_RANK_BLOCK as u64);
6094 let at = offset
6095 + DICTIONARY_HEADER as u64
6096 + offset_bytes(count as usize, bits) as u64
6097 + blocks * 16
6098 + (rank_blocks - 1) * 8;
6099 let mut end = [0; 8];
6100 read_at(file, at, &mut end).expect("the last rank block end");
6101 u64::from_le_bytes(end)
6102 }
6103
6104 fn sample() -> Chunk {
6105 Chunk::new(vec![
6106 Vector::from_values(
6107 LogicalType::Integer,
6108 &[Value::Integer(4), Value::Integer(9), Value::Integer(-2)],
6109 )
6110 .expect("integers"),
6111 Vector::from_values(
6112 LogicalType::Varchar,
6113 &[
6114 Value::Varchar("alpha".into()),
6115 Value::Null,
6116 Value::Varchar("long text after a slash".into()),
6117 ],
6118 )
6119 .expect("strings"),
6120 ])
6121 .expect("matching rows")
6122 }
6123
6124 fn sample_ids() -> Chunk {
6125 Chunk::new(vec![
6126 Vector::flat(LogicalType::Integer, Data::Int32(vec![7, 8, 9].into()))
6127 .expect("integers"),
6128 ])
6129 .expect("one column")
6130 }
6131
6132 #[test]
6133 fn the_planner_gets_the_null_count_off_the_same_directory_the_bounds_are_in() {
6134 let path = path("nulls_for_the_planner");
6137 let mut writer =
6138 Writer::create(&path, "items", vec![Field::new("a", LogicalType::Integer)])
6139 .expect("new file");
6140 let rows = Chunk::new(vec![
6141 Vector::from_values(
6142 LogicalType::Integer,
6143 &[
6144 Value::Integer(4),
6145 Value::Null,
6146 Value::Integer(9),
6147 Value::Null,
6148 Value::Integer(1),
6149 Value::Integer(2),
6150 ],
6151 )
6152 .expect("integers"),
6153 ])
6154 .expect("one column");
6155 writer.append(&rows).expect("the only part");
6156 writer.finish().expect("commit");
6157 let reader = Reader::open(&path).expect("reopen from disk");
6158 let stripes = Stripes::new(reader);
6159 let column = stripes.column("a").expect("the file has that column");
6160 assert_eq!(stripes.nulls(column), Stat::exact(2, Provenance::NullCount));
6161 assert_eq!(stripes.nulls(column + 1), Stat::Unknown);
6164 fs::remove_file(&path).expect("clean up");
6165 }
6166
6167 #[test]
6168 fn the_planner_gets_a_row_count_per_value_off_a_complete_synopsis() {
6169 let path = path("frequencies_for_the_planner");
6174 let mut writer =
6175 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6176 .expect("new file");
6177 let rows = Chunk::new(vec![
6178 Vector::from_values(
6179 LogicalType::Integer,
6180 &[
6181 Value::Integer(4),
6182 Value::Integer(4),
6183 Value::Integer(4),
6184 Value::Integer(9),
6185 Value::Integer(9),
6186 Value::Integer(1),
6187 ],
6188 )
6189 .expect("integers"),
6190 ])
6191 .expect("one column");
6192 writer.append(&rows).expect("the only part");
6193 writer.finish().expect("commit");
6194 let reader = Reader::open(&path).expect("reopen from disk");
6195 let common = Common::new(reader);
6196 assert_eq!(common.rows(), 6);
6197 let column = common.column("id").expect("the file has that column");
6198 assert_eq!(common.column("nothing"), None);
6199 assert_eq!(
6200 common.rows_with(column, &Bound::Int(4)),
6201 Stat::exact(3, Provenance::FrequencySynopsis)
6202 );
6203 assert_eq!(
6205 common.rows_with(column, &Bound::Int(7)),
6206 Stat::exact(0, Provenance::FrequencySynopsis)
6207 );
6208 assert_eq!(common.rows_with(column, &Bound::Bytes(b"four".to_vec())), Stat::Unknown);
6211 assert_eq!(common.remainder(column), None);
6214 fs::remove_file(&path).expect("clean up");
6215 }
6216
6217 #[test]
6218 fn the_planner_gets_an_exact_count_for_a_leading_value_of_an_incomplete_synopsis() {
6219 let path = path("frequency_prefix_for_the_planner");
6226 let mut writer =
6227 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6228 .expect("new file");
6229 let mut values = vec![Value::Integer(1); 10_000];
6230 for _ in 0..10 {
6231 values.extend((0..600).map(|tail| Value::Integer(1_000 + tail)));
6232 }
6233 for part in values.chunks(8_000) {
6236 let rows = Chunk::new(vec![
6237 Vector::from_values(LogicalType::Integer, part).expect("integers"),
6238 ])
6239 .expect("one column");
6240 writer.append(&rows).expect("a part");
6241 }
6242 writer.finish().expect("commit");
6243 let reader = Reader::open(&path).expect("reopen from disk");
6244 let prefix =
6245 reader.frequency_prefix(0).expect("a readable synopsis").expect("the column has one");
6246 assert_eq!(prefix.entries.len(), 512);
6249 assert_eq!(prefix.omitted_max, 10);
6250 let common = Common::new(reader);
6251 assert_eq!(common.rows(), 16_000);
6252 let column = common.column("id").expect("the file has that column");
6253 assert_eq!(
6254 common.rows_with(column, &Bound::Int(1)),
6255 Stat::exact(10_000, Provenance::FrequencySynopsis)
6256 );
6257 assert_eq!(
6259 common.rows_with(column, &Bound::Int(1_100)),
6260 Stat::exact(10, Provenance::FrequencySynopsis)
6261 );
6262 assert_eq!(common.rows_with(column, &Bound::Int(1_550)), Stat::Unknown);
6265 assert_eq!(common.rows_with(column, &Bound::Int(9_999)), Stat::Unknown);
6268 let remainder = common.remainder(column).expect("the list is a prefix");
6272 assert_eq!(remainder, Remainder { rows: 890, listed: 512, most: 10 });
6273 assert_eq!(remainder.rows / (601 - remainder.listed), 10);
6274 fs::remove_file(&path).expect("clean up");
6275 }
6276
6277 #[test]
6278 fn committed_file_reopens_and_reads_only_requested_columns() {
6279 let path = path("reopen");
6280 let mut writer = Writer::create(
6281 &path,
6282 "items",
6283 vec![
6284 Field::required("id", LogicalType::Integer),
6285 Field::new("text", LogicalType::Varchar),
6286 ],
6287 )
6288 .expect("new file");
6289 writer.append(&sample()).expect("first part");
6290 writer.append(&sample()).expect("second part");
6291 writer.finish().expect("commit");
6292 let reader = Reader::open(&path).expect("reopen from disk");
6293 assert_eq!(reader.table().rows(), 6);
6294 assert_eq!(reader.table().stripes().len(), 1);
6297 assert_eq!(reader.parts(), 2);
6298 assert_eq!(reader.part_rows(0), 3);
6299 assert_eq!(reader.part_rows(1), 3);
6300 let text = reader.read(1, &[1]).expect("only text page");
6301 assert_eq!(text.width(), 1);
6302 assert_eq!(text.value_at(1, 0), Value::Null);
6303 assert_eq!(text.value_at(2, 0), Value::Varchar("long text after a slash".into()));
6304 let sparse = reader.read_sparse(1, &[1]).expect("one part without its whole page");
6305 assert_eq!(sparse.width(), 1);
6306 assert_eq!(sparse.value_at(1, 0), Value::Null);
6307 assert_eq!(sparse.value_at(2, 0), Value::Varchar("long text after a slash".into()));
6308 assert!(!reader.skips_codes(0, 1, &[0]).expect("alpha is in the stripe"));
6309 assert!(!reader.skips_codes(0, 1, &[2]).expect("long text is in the stripe"));
6310 assert!(reader.skips_codes(0, 1, &[3]).expect("unknown code is absent"));
6311 let count = reader.read(0, &[]).expect("no page is needed for count");
6312 assert_eq!(count.len(), 3);
6313 assert!(reader.skips(0, &[Probe { column: 0, op: Op::Greater, value: Bound::Int(100) }]));
6314 assert!(!reader.skips(0, &[Probe { column: 0, op: Op::Greater, value: Bound::Int(0) }]));
6315 let integers = reader.top_frequencies(0, 1).expect("valid integer synopsis").expect("kept");
6316 assert_eq!(
6317 integers,
6318 vec![(Value::Integer(-2), 2), (Value::Integer(4), 2), (Value::Integer(9), 2),]
6319 );
6320 let strings = reader.top_frequencies(1, 1).expect("valid string synopsis").expect("kept");
6321 assert_eq!(strings.len(), 3);
6322 assert!(strings.contains(&(Value::Null, 2)));
6323 assert!(strings.contains(&(Value::Varchar("alpha".into()), 2)));
6324 assert!(strings.contains(&(Value::Varchar("long text after a slash".into()), 2)));
6325 fs::remove_file(path).expect("remove scratch file");
6326 }
6327
6328 #[test]
6336 fn runs_handed_over_out_of_order_still_read_back_in_source_order() {
6337 let path = path("interleaved-runs");
6338 let mut writer =
6339 Writer::create(&path, "interleaved", vec![Field::new("v", LogicalType::BigInt)])
6340 .expect("new file");
6341 for morsel in [2_u64, 0, 3, 1] {
6342 let parts = (0..4_u64)
6343 .map(|chunk| {
6344 let first = i64::try_from(morsel * 32 + chunk * 8).expect("small");
6345 let values =
6346 (0..8_i64).map(|row| Value::BigInt(first + row)).collect::<Vec<_>>();
6347 let column =
6348 Vector::from_values(LogicalType::BigInt, &values).expect("a column");
6349 ((morsel, chunk), Chunk::new(vec![column]).expect("one column"))
6350 })
6351 .collect::<Vec<_>>();
6352 writer.append_stripe(parts).expect("a stripe");
6353 }
6354 writer.finish().expect("commit");
6355
6356 let reader = Reader::open(&path).expect("valid directory");
6357 assert_eq!(reader.table().stripes().len(), 4, "a run is a stripe of its own");
6358 assert_eq!(reader.table().rows(), 128);
6359 for part in 0..16_usize {
6360 let read = reader.read(part, &[0]).expect("a part back");
6361 for row in 0..8_usize {
6362 let want = i64::try_from(part * 8 + row).expect("small");
6363 assert_eq!(read.value_at(row, 0), Value::BigInt(want), "part {part} row {row}");
6364 }
6365 }
6366 fs::remove_file(path).expect("remove scratch file");
6367 }
6368
6369 #[test]
6372 fn runs_that_overlap_each_other_are_refused_at_commit() {
6373 let path = path("overlapping-runs");
6374 let mut writer =
6375 Writer::create(&path, "overlapping", vec![Field::new("v", LogicalType::BigInt)])
6376 .expect("new file");
6377 let one = |order: (u64, u64)| {
6378 let column =
6379 Vector::from_values(LogicalType::BigInt, &[Value::BigInt(1)]).expect("a column");
6380 (order, Chunk::new(vec![column]).expect("one column"))
6381 };
6382 writer.append_stripe(vec![one((0, 0)), one((0, 2))]).expect("a stripe");
6385 writer.append_stripe(vec![one((0, 1))]).expect("a stripe");
6386 let error = writer.finish().expect_err("the runs overlap");
6387 assert!(error.message().contains("source order"), "{error}");
6388 fs::remove_file(path).expect("remove scratch file");
6389 }
6390
6391 #[test]
6394 fn a_run_longer_than_a_stripe_is_refused() {
6395 let path = path("overlong-run");
6396 let mut writer =
6397 Writer::create(&path, "overlong", vec![Field::new("v", LogicalType::BigInt)])
6398 .expect("new file");
6399 let parts = (0..=STRIPE_PARTS)
6400 .map(|at| {
6401 let column = Vector::from_values(LogicalType::BigInt, &[Value::BigInt(1)])
6402 .expect("a column");
6403 let chunk = Chunk::new(vec![column]).expect("one column");
6404 ((0, u64::try_from(at).expect("small")), chunk)
6405 })
6406 .collect::<Vec<_>>();
6407 let error = writer.append_stripe(parts).expect_err("one part too many");
6408 assert!(error.message().contains("more parts than it holds"), "{error}");
6409 fs::remove_file(path).expect("remove scratch file");
6410 }
6411
6412 #[test]
6418 fn parts_past_the_stripe_bound_start_a_new_stripe() {
6419 let path = path("stripe-bound");
6420 let mut writer = Writer::create(
6421 &path,
6422 "items",
6423 vec![
6424 Field::required("id", LogicalType::Integer),
6425 Field::new("text", LogicalType::Varchar),
6426 ],
6427 )
6428 .expect("new file");
6429 let parts = STRIPE_PARTS * 2 + 3;
6430 for part in 0..parts {
6431 let id = part as i32;
6432 let chunk = Chunk::new(vec![
6433 Vector::from_values(
6434 LogicalType::Integer,
6435 &[Value::Integer(id), Value::Integer(-id)],
6436 )
6437 .expect("integers"),
6438 Vector::from_values(
6439 LogicalType::Varchar,
6440 &[Value::Varchar(format!("value {part}")), Value::Null],
6441 )
6442 .expect("strings"),
6443 ])
6444 .expect("matching rows");
6445 writer.append(&chunk).expect("one part");
6446 }
6447 writer.finish().expect("commit");
6448
6449 let reader = Reader::open(&path).expect("reopen from disk");
6450 assert_eq!(reader.parts(), parts);
6451 assert_eq!(reader.table().rows(), parts * 2);
6452 assert_eq!(reader.table().stripes().len(), parts.div_ceil(STRIPE_PARTS));
6453 assert_eq!(reader.table().stripes()[0].parts(), STRIPE_PARTS);
6454 assert_eq!(reader.table().stripes()[0].rows(), STRIPE_PARTS * 2);
6455 assert_eq!(reader.table().stripes()[2].parts(), 3);
6456 for part in (0..parts).rev() {
6459 let dense = reader.read(part, &[0, 1]).expect("a whole page read");
6460 let sparse = reader.read_sparse(part, &[0, 1]).expect("one part read");
6461 for chunk in [&dense, &sparse] {
6462 assert_eq!(chunk.len(), 2, "part {part} has its own row count");
6463 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6464 assert_eq!(chunk.value_at(1, 0), Value::Integer(-(part as i32)));
6465 assert_eq!(chunk.value_at(0, 1), Value::Varchar(format!("value {part}")));
6466 assert_eq!(chunk.value_at(1, 1), Value::Null);
6467 }
6468 }
6469 let above = [Probe { column: 0, op: Op::Greater, value: Bound::Int(100) }];
6472 assert!(reader.skips(0, &above), "the first stripe stops at 63");
6473 assert!(!reader.skips(STRIPE_PARTS * 2, &above), "the third stripe reaches 130");
6474 fs::remove_file(path).expect("remove scratch file");
6475 }
6476
6477 fn scattered(n: i64) -> i64 {
6479 n.wrapping_mul(-7_046_029_254_386_353_131)
6480 }
6481
6482 #[test]
6488 fn a_part_is_skipped_when_its_sieve_does_not_hold_the_constant() {
6489 let path = path("sieve-skip");
6490 let mut writer =
6491 Writer::create(&path, "hits", vec![Field::required("id", LogicalType::BigInt)])
6492 .expect("new file");
6493 let parts = STRIPE_PARTS + 3;
6494 let per_part = 128;
6498 for part in 0..parts {
6499 let held: Vec<Value> = (0..per_part)
6500 .map(|row| Value::BigInt(scattered((part * per_part + row) as i64)))
6501 .collect();
6502 let chunk =
6503 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6504 .expect("one column");
6505 writer.append(&chunk).expect("one part");
6506 }
6507 writer.finish().expect("commit");
6508
6509 let reader = Reader::open(&path).expect("reopen from disk");
6510 let probe = |value: i64| Probe {
6511 column: 0,
6512 op: Op::Equal,
6513 value: Bound::Int(i128::from(scattered(value))),
6514 };
6515 for wanted in [0_i64, (per_part + 1) as i64, (parts * per_part - 1) as i64] {
6516 let tests = [probe(wanted)];
6517 let kept: Vec<usize> = (0..parts).filter(|&part| !reader.skips(part, &tests)).collect();
6518 let home = wanted as usize / per_part;
6519 assert!(kept.contains(&home), "the part holding {wanted} is read");
6520 assert!(kept.len() <= 2, "{wanted} keeps {kept:?}, which is more than one stray part");
6524 }
6525 let absent = [probe((parts * per_part) as i64 + 1)];
6526 let kept = (0..parts).filter(|&part| !reader.skips(part, &absent)).count();
6527 assert!(kept <= 1, "{kept} parts of {parts} kept a value no part holds");
6528 let tests = [probe(0)];
6531 assert!(
6532 reader.table().stripes().iter().all(|stripe| !stripe.zone.skips(&tests)),
6533 "the bounds rule out no stripe at all"
6534 );
6535 fs::remove_file(path).expect("remove scratch file");
6536 }
6537
6538 #[test]
6544 fn a_part_is_skipped_when_its_own_bounds_rule_out_a_comparison_the_stripe_keeps() {
6545 let path = path("part-range-skip");
6546 let mut writer =
6547 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
6548 .expect("new file");
6549 let parts = STRIPE_PARTS + 3;
6550 let per_part = 128;
6551 for part in 0..parts {
6552 let held: Vec<Value> = (0..per_part)
6556 .map(|row| {
6557 Value::BigInt((part * 1_000) as i64 + (scattered(row as i64).rem_euclid(900)))
6558 })
6559 .collect();
6560 let chunk =
6561 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6562 .expect("one column");
6563 writer.append(&chunk).expect("one part");
6564 }
6565 writer.finish().expect("commit");
6566
6567 let reader = Reader::open(&path).expect("reopen from disk");
6568 let under = [Probe { column: 0, op: Op::Less, value: Bound::Int(3_000) }];
6569 let kept: Vec<usize> = (0..parts).filter(|&part| !reader.skips(part, &under)).collect();
6570 assert_eq!(kept, vec![0, 1, 2], "only the three parts that start under three thousand");
6571 assert!(!reader.stripe_skips(0, &under), "the stripe reaches from zero and keeps itself");
6573 fs::remove_file(path).expect("remove scratch file");
6574 }
6575
6576 #[test]
6580 fn a_part_is_waved_through_when_its_own_bounds_pass_a_comparison_the_stripe_cannot() {
6581 let path = path("part-range-certain");
6582 let mut writer =
6583 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
6584 .expect("new file");
6585 let parts = STRIPE_PARTS + 3;
6586 let per_part = 128;
6587 for part in 0..parts {
6588 let held: Vec<Value> = (0..per_part)
6589 .map(|row| {
6590 Value::BigInt((part * 1_000) as i64 + (scattered(row as i64).rem_euclid(900)))
6591 })
6592 .collect();
6593 let chunk =
6594 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6595 .expect("one column");
6596 writer.append(&chunk).expect("one part");
6597 }
6598 writer.finish().expect("commit");
6599
6600 let reader = Reader::open(&path).expect("reopen from disk");
6601 let under = [Probe { column: 0, op: Op::Less, value: Bound::Int(3_000) }];
6602 let waved: Vec<usize> = (0..parts).filter(|&part| reader.certain(part, &under)).collect();
6603 assert_eq!(waved, vec![0, 1, 2], "the three parts that end under three thousand");
6604 assert!(!reader.stripe_skips(0, &under), "the stripe straddles the comparison");
6607 fs::remove_file(path).expect("remove scratch file");
6608 }
6609
6610 #[test]
6613 fn a_stripe_of_one_part_writes_no_range_page_and_a_stripe_of_many_does() {
6614 for (parts, wanted) in [(1_usize, false), (STRIPE_PARTS, true)] {
6615 let path = path("part-range-page");
6616 let mut writer =
6617 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
6618 .expect("new file");
6619 for part in 0..parts {
6620 let held: Vec<Value> = (0..128)
6621 .map(|row| {
6622 Value::BigInt((part * 1_000) as i64 + scattered(row as i64).rem_euclid(900))
6623 })
6624 .collect();
6625 let chunk = Chunk::new(vec![
6626 Vector::from_values(LogicalType::BigInt, &held).expect("numbers"),
6627 ])
6628 .expect("one column");
6629 writer.append(&chunk).expect("one part");
6630 }
6631 writer.finish().expect("commit");
6632 let reader = Reader::open(&path).expect("reopen from disk");
6633 let bytes = reader.layout().columns[0].part_ranges;
6634 assert_eq!(bytes > 0, wanted, "{parts} parts wrote {bytes} bytes of ranges");
6635 fs::remove_file(path).expect("remove scratch file");
6636 }
6637 }
6638
6639 #[test]
6642 fn a_string_end_that_is_cut_down_still_covers_the_value_it_came_from() {
6643 let long = vec![b'a'; PART_BOUND_BYTES * 2];
6644 let low = shortened(Some(Bound::Bytes(long.clone())), false).expect("a low end");
6645 let high = shortened(Some(Bound::Bytes(long.clone())), true).expect("a high end");
6646 let Bound::Bytes(low) = low else { panic!("a string stays a string") };
6647 let Bound::Bytes(high) = high else { panic!("a string stays a string") };
6648 assert!(low.len() <= PART_BOUND_BYTES && high.len() <= PART_BOUND_BYTES);
6649 assert!(low.as_slice() <= long.as_slice(), "the low end is at or under the value");
6650 assert!(high.as_slice() >= long.as_slice(), "the high end is at or over the value");
6651 }
6652
6653 #[test]
6656 fn a_string_end_with_no_room_to_step_up_gives_up_the_bound() {
6657 let long = vec![u8::MAX; PART_BOUND_BYTES * 2];
6658 assert_eq!(shortened(Some(Bound::Bytes(long.clone())), true), None);
6659 let low = shortened(Some(Bound::Bytes(long)), false).expect("a low end is still a prefix");
6660 assert_eq!(low, Bound::Bytes(vec![u8::MAX; PART_BOUND_BYTES]));
6661 }
6662
6663 #[test]
6673 fn a_sieve_larger_than_the_part_it_indexes_is_not_written() {
6674 let path = path("sieve-pays");
6675 let fields = vec![
6676 Field::required("spread", LogicalType::BigInt),
6677 Field::required("repeated", LogicalType::BigInt),
6678 ];
6679 let mut writer = Writer::create(&path, "hits", fields).expect("new file");
6680 let parts = 3;
6681 let per_part = 1024;
6682 for part in 0..parts {
6683 let base = (part * per_part) as i64;
6684 let spread: Vec<Value> =
6685 (0..per_part).map(|row| Value::BigInt(scattered(base + row as i64))).collect();
6686 let repeated: Vec<Value> =
6687 (0..per_part).map(|row| Value::BigInt(scattered((row / 256) as i64))).collect();
6688 let chunk = Chunk::new(vec![
6689 Vector::from_values(LogicalType::BigInt, &spread).expect("numbers"),
6690 Vector::from_values(LogicalType::BigInt, &repeated).expect("numbers"),
6691 ])
6692 .expect("two columns");
6693 writer.append(&chunk).expect("one part");
6694 }
6695 writer.finish().expect("commit");
6696
6697 let reader = Reader::open(&path).expect("reopen from disk");
6698 let layout = reader.layout();
6699 let spread = &layout.columns[0];
6700 let repeated = &layout.columns[1];
6701 assert!(spread.sieves > 0, "a column whose parts are worth a filter keeps one");
6702 assert_eq!(
6703 repeated.sieves, 0,
6704 "a column whose filter costs more than its parts keeps none"
6705 );
6706 for column in &layout.columns {
6709 assert!(
6710 column.sieves < column.pages,
6711 "{} spends {} on sieves over {} of data",
6712 column.name,
6713 column.sieves,
6714 column.pages
6715 );
6716 }
6717 let absent = [Probe {
6719 column: 0,
6720 op: Op::Equal,
6721 value: Bound::Int(i128::from(scattered((parts * per_part) as i64 + 1))),
6722 }];
6723 assert!((0..parts).all(|part| reader.skips(part, &absent)), "no part holds it");
6724 fs::remove_file(path).expect("remove scratch file");
6725 }
6726
6727 #[test]
6733 fn a_damaged_sieve_page_is_read_through_rather_than_refused() {
6734 let path = path("sieve-damaged");
6735 let mut writer =
6736 Writer::create(&path, "hits", vec![Field::required("id", LogicalType::BigInt)])
6737 .expect("new file");
6738 let rows = 128;
6739 let held: Vec<Value> = (0..rows).map(|row| Value::BigInt(scattered(row))).collect();
6740 let chunk =
6741 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6742 .expect("one column");
6743 writer.append(&chunk).expect("one part");
6744 writer.finish().expect("commit");
6745
6746 let page =
6747 Reader::open(&path).expect("reopen").table.stripes[0].sieves[0].expect("a sieve page");
6748 let mut file = OpenOptions::new().write(true).open(&path).expect("open the sieve page");
6749 file.seek(SeekFrom::Start(page.offset + u64::from(page.length) - 1)).expect("seek");
6750 file.write_all(&[0xff]).expect("damage one byte");
6751 drop(file);
6752
6753 let reader = Reader::open(&path).expect("reopen the damaged file");
6754 let absent =
6755 [Probe { column: 0, op: Op::Equal, value: Bound::Int(i128::from(scattered(99))) }];
6756 assert!(!reader.skips(0, &absent), "a sieve that cannot be read skips nothing");
6757 assert_eq!(
6758 reader.read(0, &[0]).expect("the rows are untouched").len(),
6759 usize::try_from(rows).expect("a small count")
6760 );
6761 fs::remove_file(path).expect("remove scratch file");
6762 }
6763
6764 #[test]
6775 fn workers_that_want_the_same_stripe_read_it_once() {
6776 let path = path("single-flight");
6777 let mut writer =
6778 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6779 .expect("new file");
6780 for part in 0..STRIPE_PARTS {
6781 let id = part as i32;
6782 let chunk = Chunk::new(vec![
6783 Vector::from_values(
6784 LogicalType::Integer,
6785 &[Value::Integer(id), Value::Integer(-id)],
6786 )
6787 .expect("integers"),
6788 ])
6789 .expect("matching rows");
6790 writer.append(&chunk).expect("one part");
6791 }
6792 writer.finish().expect("commit");
6793
6794 let reader = Reader::open(&path).expect("reopen from disk");
6795 assert_eq!(reader.table().stripes().len(), 1, "one stripe is the point of the test");
6796 let barrier = std::sync::Barrier::new(8);
6797 std::thread::scope(|scope| {
6798 for worker in 0..8 {
6799 let reader = &reader;
6800 let barrier = &barrier;
6801 scope.spawn(move || {
6802 barrier.wait();
6803 for part in (worker..STRIPE_PARTS).step_by(8) {
6804 let chunk = reader.read(part, &[0]).expect("a whole page read");
6805 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6806 assert_eq!(chunk.value_at(1, 0), Value::Integer(-(part as i32)));
6807 }
6808 });
6809 }
6810 });
6811 assert_eq!(reader.pages.load(Atomic::Relaxed), 1, "one stripe, one page read, whoever won");
6812 fs::remove_file(path).expect("remove scratch file");
6813 }
6814
6815 #[test]
6828 fn opening_costs_the_same_over_a_thousand_times_the_rows() {
6829 let opened = |label: &str, rows_per_part: i32| {
6830 let path = path(label);
6831 let mut writer =
6832 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6833 .expect("new file");
6834 for part in 0..STRIPE_PARTS * 3 {
6835 let values = (0..rows_per_part)
6839 .map(|row| {
6840 Value::Integer((part as i32 * rows_per_part + row).wrapping_mul(2_654_435))
6841 })
6842 .collect::<Vec<_>>();
6843 let chunk = Chunk::new(vec![
6844 Vector::from_values(LogicalType::Integer, &values).expect("integers"),
6845 ])
6846 .expect("matching rows");
6847 writer.append(&chunk).expect("one part");
6848 }
6849 writer.finish().expect("commit");
6850 let reader = Reader::open(&path).expect("reopen from disk");
6851 let size = fs::metadata(&path).expect("the file is there").len();
6852 let out = (reader.reads(), reader.table().stripes().len(), size);
6853 fs::remove_file(path).expect("remove scratch file");
6854 out
6855 };
6856
6857 let (thin, thin_stripes, thin_size) = opened("open-thin", 1);
6858 let (fat, fat_stripes, fat_size) = opened("open-fat", 1000);
6859 assert_eq!(
6860 thin_stripes, fat_stripes,
6861 "the same stripe count is what makes this a fair ask"
6862 );
6863 assert!(
6864 fat_size > thin_size * 50,
6865 "the fat file has to actually be larger, and it is {fat_size} against {thin_size}"
6866 );
6867
6868 assert_eq!(thin.opening.reads, fat.opening.reads, "the same reads either way");
6869 assert_eq!(thin.pages, 0, "opening read a page");
6870 assert_eq!(fat.pages, 0, "opening read a page");
6871 assert_eq!(thin.indexes, 0, "opening read an index");
6872 assert_eq!(fat.indexes, 0, "opening read an index");
6873 assert!(
6876 fat.opening.bytes < thin.opening.bytes * 2,
6877 "opening the thin file read {} bytes and the fat one read {}",
6878 thin.opening.bytes,
6879 fat.opening.bytes
6880 );
6881 }
6882
6883 #[test]
6891 fn two_opens_of_one_file_cost_the_same_and_the_second_is_not_cheaper() {
6892 let path = path("open-twice");
6893 let mut writer =
6894 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6895 .expect("new file");
6896 for part in 0..STRIPE_PARTS * 3 {
6897 let chunk = Chunk::new(vec![
6898 Vector::from_values(LogicalType::Integer, &[Value::Integer(part as i32)])
6899 .expect("integers"),
6900 ])
6901 .expect("matching rows");
6902 writer.append(&chunk).expect("one part");
6903 }
6904 writer.finish().expect("commit");
6905
6906 let first = Reader::open(&path).expect("open");
6907 for part in 0..first.parts() {
6910 first.read(part, &[0]).expect("a part");
6911 }
6912 assert!(first.reads().pages > 0, "the scan has to have read something");
6913 let second = Reader::open(&path).expect("open again");
6914
6915 assert_eq!(first.reads().opening, second.reads().opening);
6916 assert_eq!(
6917 second.reads().pages,
6918 0,
6919 "the second open read a page off the back of the first"
6920 );
6921 assert_eq!(second.reads().indexes, 0, "the second open read an index it inherited");
6922 fs::remove_file(path).expect("remove scratch file");
6923 }
6924
6925 #[test]
6933 fn an_index_is_read_once_per_stripe_however_often_the_page_is_evicted() {
6934 let path = path("index-cache");
6935 let mut writer =
6936 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6937 .expect("new file");
6938 let parts = STRIPE_PARTS * (CACHED_STRIPES_PER_COLUMN + 2);
6939 for part in 0..parts {
6940 let id = part as i32;
6941 let chunk = Chunk::new(vec![
6942 Vector::from_values(LogicalType::Integer, &[Value::Integer(id)]).expect("integers"),
6943 ])
6944 .expect("matching rows");
6945 writer.append(&chunk).expect("one part");
6946 }
6947 writer.finish().expect("commit");
6948
6949 let reader = Reader::open(&path).expect("reopen from disk");
6950 let stripes = reader.table().stripes().len();
6951 assert!(stripes > CACHED_STRIPES_PER_COLUMN, "the page cache has to be too small for this");
6952 for _ in 0..2 {
6954 for part in 0..parts {
6955 let chunk = reader.read(part, &[0]).expect("a part");
6956 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6957 }
6958 }
6959 assert_eq!(reader.indexes.load(Atomic::Relaxed), stripes, "one index read per stripe");
6960 assert!(
6961 reader.pages.load(Atomic::Relaxed) > stripes,
6962 "the pages are the ones that get read again, which is what makes the index count mean \
6963 something"
6964 );
6965 fs::remove_file(path).expect("remove scratch file");
6966 }
6967
6968 #[test]
6977 fn a_worker_per_stripe_reads_its_page_once_when_the_cache_was_told_to_expect_it() {
6978 let workers = CACHED_STRIPES_PER_COLUMN + 4;
6979 let path = path("stripe-per-worker");
6980 let mut writer =
6981 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6982 .expect("new file");
6983 for part in 0..STRIPE_PARTS * workers {
6984 let chunk = Chunk::new(vec![
6985 Vector::from_values(LogicalType::Integer, &[Value::Integer(part as i32)])
6986 .expect("integers"),
6987 ])
6988 .expect("matching rows");
6989 writer.append(&chunk).expect("one part");
6990 }
6991 writer.finish().expect("commit");
6992
6993 let read = |told: bool| {
6994 let reader = Reader::open(&path).expect("reopen from disk");
6995 assert_eq!(reader.table().stripes().len(), workers, "a stripe per worker");
6996 if told {
6997 reader.keep_stripes(workers);
6998 }
6999 let barrier = std::sync::Barrier::new(workers);
7000 std::thread::scope(|scope| {
7001 for (worker, run) in reader.stripe_parts().into_iter().enumerate() {
7002 let reader = &reader;
7003 let barrier = &barrier;
7004 scope.spawn(move || {
7005 for part in run {
7006 barrier.wait();
7007 let chunk = reader.read(part, &[0]).expect("a part of my own stripe");
7008 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
7009 }
7010 assert!(worker < workers);
7011 });
7012 }
7013 });
7014 reader.pages.load(Atomic::Relaxed)
7015 };
7016
7017 assert_eq!(read(true), workers, "one page read per stripe and no more");
7018 assert!(read(false) > workers, "a cache that small is read again on every part");
7019 fs::remove_file(path).expect("remove scratch file");
7020 }
7021
7022 #[test]
7027 fn a_damaged_index_page_is_an_error() {
7028 let path = path("damaged-index");
7029 let mut writer =
7030 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
7031 .expect("new file");
7032 writer.append(&sample_ids()).expect("first part");
7033 writer.append(&sample_ids()).expect("second part");
7034 writer.finish().expect("commit");
7035
7036 let reader = Reader::open(&path).expect("valid directory");
7037 let index = reader.table.stripes[0].index;
7038 let mut byte = [0; 1];
7039 read_at(&reader.file, index.offset, &mut byte).expect("the first part length");
7040 let mut file = OpenOptions::new().write(true).open(&path).expect("open index page");
7041 file.seek(SeekFrom::Start(index.offset)).expect("index start");
7042 file.write_all(&[!byte[0]]).expect("damage the first part length");
7043 let error = reader.read(1, &[0]).expect_err("a damaged index must not be used");
7044 assert!(error.message().contains("index page section checksum differs"), "{error}");
7045 fs::remove_file(path).expect("remove scratch file");
7046 }
7047
7048 #[test]
7055 fn every_integer_width_round_trips_through_a_page() {
7056 let path = path("integer-widths");
7057 let columns = [
7058 (LogicalType::TinyInt, vec![Value::TinyInt(i8::MIN), Value::TinyInt(i8::MAX)]),
7059 (LogicalType::UTinyInt, vec![Value::UTinyInt(0), Value::UTinyInt(u8::MAX)]),
7060 (LogicalType::SmallInt, vec![Value::SmallInt(i16::MIN), Value::SmallInt(i16::MAX)]),
7061 (LogicalType::USmallInt, vec![Value::USmallInt(0), Value::USmallInt(u16::MAX)]),
7062 (LogicalType::Integer, vec![Value::Integer(i32::MIN), Value::Integer(i32::MAX)]),
7063 (LogicalType::UInteger, vec![Value::UInteger(0), Value::UInteger(u32::MAX)]),
7064 (LogicalType::BigInt, vec![Value::BigInt(i64::MIN), Value::BigInt(i64::MAX)]),
7065 (LogicalType::UBigInt, vec![Value::UBigInt(0), Value::UBigInt(u64::MAX)]),
7066 ];
7067 let fields = columns
7068 .iter()
7069 .enumerate()
7070 .map(|(at, (ty, _))| Field::required(format!("c{at}"), ty.clone()))
7071 .collect::<Vec<_>>();
7072 let vectors = columns
7073 .iter()
7074 .map(|(ty, values)| Vector::from_values(ty.clone(), values).expect("a vector"))
7075 .collect::<Vec<_>>();
7076 let mut writer = Writer::create(&path, "widths", fields).expect("new file");
7077 writer.append(&Chunk::new(vectors).expect("matching rows")).expect("one stripe");
7078 writer.finish().expect("commit");
7079
7080 let reader = Reader::open(&path).expect("reopen from disk");
7081 let wanted = (0..columns.len()).collect::<Vec<_>>();
7082 let read = reader.read(0, &wanted).expect("every column");
7083 assert_eq!(read.len(), 2);
7084 for (at, (ty, values)) in columns.iter().enumerate() {
7086 assert_eq!(read.value_at(0, at), values[0], "the low end of {ty}");
7087 assert_eq!(read.value_at(1, at), values[1], "the high end of {ty}");
7088 }
7089 fs::remove_file(path).expect("remove scratch file");
7090 }
7091
7092 #[test]
7093 fn numeric_frequency_candidates_keep_bounded_row_ordinals() {
7094 let path = path("frequency-ordinals");
7095 let mut writer =
7096 Writer::create(&path, "items", vec![Field::required("id", LogicalType::BigInt)])
7097 .expect("new file");
7098 let mut values = Vec::new();
7099 for leader in 0..10_i64 {
7100 values.extend(std::iter::repeat_n(leader, 100));
7101 }
7102 values.extend(1_000_i64..41_000);
7103 for part in values.chunks(1_024) {
7104 let vector = Vector::flat(LogicalType::BigInt, Data::Int64(part.to_vec().into()))
7105 .expect("big integers");
7106 writer.append(&Chunk::new(vec![vector]).expect("one column")).expect("one stripe");
7107 }
7108 writer.finish().expect("commit");
7109
7110 let reader = Reader::open(&path).expect("reopen from disk");
7111 let occurrences =
7112 reader.frequency_occurrences(0).expect("valid metadata").expect("bounded ordinals");
7113 assert!(occurrences.omitted_max < 100);
7114 assert!(occurrences.ordinals.len() <= FREQUENCY_ORDINALS);
7115 assert!(occurrences.ordinals.windows(2).all(|pair| pair[0] < pair[1]));
7116 assert_eq!(&occurrences.ordinals[..1_000], &(0_u64..1_000).collect::<Vec<_>>());
7117 fs::remove_file(path).expect("remove scratch file");
7118 }
7119
7120 #[test]
7126 fn a_file_from_another_format_says_which_format_it_is() {
7127 let older = path("older-format");
7128 let mut writer =
7129 Writer::create(&older, "items", vec![Field::new("id", LogicalType::Integer)])
7130 .expect("new file");
7131 let chunk = Chunk::new(vec![
7132 Vector::flat(LogicalType::Integer, Data::Int32(vec![1, 2, 3].into()))
7133 .expect("integers"),
7134 ])
7135 .expect("chunk");
7136 writer.append(&chunk).expect("page written");
7137 writer.finish().expect("commit");
7138
7139 let mut file = OpenOptions::new().write(true).open(&older).expect("open for the header");
7140 file.seek(SeekFrom::Start(8)).expect("the version follows the magic");
7141 file.write_all(&(FORMAT - 1).to_le_bytes()).expect("write an older version");
7142 drop(file);
7143 let complaint = Reader::open(&older).expect_err("an older format is refused").to_string();
7144 assert!(complaint.contains(&format!("format {}", FORMAT - 1)), "{complaint}");
7145 assert!(complaint.contains(&format!("format {FORMAT}")), "{complaint}");
7146
7147 let mut file = OpenOptions::new().write(true).open(&older).expect("open for the header");
7148 file.seek(SeekFrom::Start(0)).expect("the magic is first");
7149 file.write_all(b"NOTRUDB!").expect("write another engine's magic");
7150 drop(file);
7151 let complaint = Reader::open(&older).expect_err("a foreign file is refused").to_string();
7152 assert!(complaint.contains("magic"), "{complaint}");
7153 assert!(!complaint.contains("format"), "a version has nothing to do with it: {complaint}");
7154 fs::remove_file(older).expect("remove scratch file");
7155 }
7156
7157 #[test]
7158 fn an_unfinished_or_damaged_file_does_not_answer_with_partial_rows() {
7159 let unfinished = path("unfinished");
7160 let mut writer =
7161 Writer::create(&unfinished, "items", vec![Field::new("id", LogicalType::Integer)])
7162 .expect("new file");
7163 let chunk = Chunk::new(vec![
7164 Vector::flat(LogicalType::Integer, Data::Int32(vec![1, 2, 3].into()))
7165 .expect("integers"),
7166 ])
7167 .expect("chunk");
7168 writer.append(&chunk).expect("page written");
7169 drop(writer);
7170 assert!(Reader::open(&unfinished).is_err(), "no directory was committed");
7171 fs::remove_file(unfinished).expect("remove scratch file");
7172
7173 let damaged = path("damaged");
7174 let mut writer =
7175 Writer::create(&damaged, "items", vec![Field::new("id", LogicalType::Integer)])
7176 .expect("new file");
7177 writer.append(&chunk).expect("page written");
7178 writer.finish().expect("commit");
7179 let reader = Reader::open(&damaged).expect("valid directory");
7180 let mut file =
7181 OpenOptions::new().write(true).open(&damaged).expect("open for a damaged page");
7182 file.seek(SeekFrom::Start(HEADER + 1)).expect("inside first page");
7183 file.write_all(&[255]).expect("damage one byte");
7184 assert!(reader.read(0, &[0]).is_err(), "page checksum rejects corruption");
7185 fs::remove_file(damaged).expect("remove scratch file");
7186 }
7187
7188 #[test]
7189 fn damaged_lazy_dictionary_payload_is_an_error() {
7190 let path = path("damaged-dictionary");
7191 let mut writer = Writer::create(
7192 &path,
7193 "items",
7194 vec![
7195 Field::required("id", LogicalType::Integer),
7196 Field::new("text", LogicalType::Varchar),
7197 ],
7198 )
7199 .expect("new file");
7200 writer.append(&sample()).expect("stripe written");
7201 writer.finish().expect("commit");
7202
7203 let reader = Reader::open(&path).expect("valid directory");
7204 let dictionary = reader.table.dictionaries[1].expect("string dictionary page");
7205 let mut header = [0; DICTIONARY_HEADER];
7208 read_at(&reader.file, dictionary.offset, &mut header).expect("dictionary header");
7209 let index_len = dictionary_index_len(&header);
7210 let rank_len = last_rank_end(&reader.file, dictionary.offset, &header);
7211 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
7212 file.seek(SeekFrom::Start(dictionary.offset + index_len + rank_len))
7213 .expect("inside dictionary payload");
7214 file.write_all(&[255]).expect("damage dictionary payload");
7215
7216 let chunk = reader.read(0, &[1]).expect("code page and dictionary index remain valid");
7217 let error =
7218 chunk.validate_external().expect_err("payload corruption must reach the caller");
7219 assert!(error.message().contains("payload checksum differs"), "{error}");
7220 fs::remove_file(path).expect("remove scratch file");
7221 }
7222
7223 #[test]
7233 fn a_column_of_all_different_values_is_written_without_a_dictionary() {
7234 let path = path("dictionary-decide");
7235 let rows = 20_000;
7236 let unique =
7238 |row: usize| format!("{row:09} a value that appears exactly once in the table");
7239 let repeated = |row: usize| unique(row / 40);
7241 let mut writer = Writer::create(
7242 &path,
7243 "items",
7244 vec![
7245 Field::required("unique", LogicalType::Varchar),
7246 Field::required("repeated", LogicalType::Varchar),
7247 ],
7248 )
7249 .expect("new file");
7250 for part in (0..rows).step_by(1_000) {
7251 let span = part..(part + 1_000).min(rows);
7252 let left = span.clone().map(|row| Value::Varchar(unique(row))).collect::<Vec<_>>();
7253 let right = span.map(|row| Value::Varchar(repeated(row))).collect::<Vec<_>>();
7254 writer
7255 .append(
7256 &Chunk::new(vec![
7257 Vector::from_values(LogicalType::Varchar, &left).expect("strings"),
7258 Vector::from_values(LogicalType::Varchar, &right).expect("strings"),
7259 ])
7260 .expect("two columns"),
7261 )
7262 .expect("a part");
7263 }
7264 writer.finish().expect("commit");
7265
7266 let reader = Reader::open(&path).expect("reopen from disk");
7267 assert!(
7268 reader.table.dictionaries[0].is_none(),
7269 "a column with no repeats has nothing to say twice"
7270 );
7271 assert!(
7272 reader.table.dictionaries[1].is_some(),
7273 "a column whose values come round again keeps its dictionary"
7274 );
7275 let mut first = 0;
7276 for part in 0..reader.parts() {
7277 let chunk = reader.read(part, &[0, 1]).expect("a part");
7278 for row in 0..chunk.len() {
7279 assert_eq!(chunk.value_at(row, 0), Value::Varchar(unique(first + row)));
7280 assert_eq!(chunk.value_at(row, 1), Value::Varchar(repeated(first + row)));
7281 }
7282 first += chunk.len();
7283 }
7284 assert_eq!(first, rows, "every row was read back");
7285 let raw = (0..rows).map(|row| unique(row).len()).sum::<usize>();
7286 let size = fs::metadata(&path).expect("the file is there").len() as usize;
7287 assert!(size < raw, "a column without a dictionary is still encoded: {size} against {raw}");
7288 fs::remove_file(path).expect("remove scratch file");
7289 }
7290
7291 #[test]
7304 fn a_dictionary_over_many_blocks_checks_every_block_of_it() {
7305 let path = path("dictionary-blocks");
7306 let value = |row: usize| {
7307 let row = row.saturating_sub(8_000);
7308 format!("{row:07} a value long enough to be worth a payload block")
7309 };
7310 let parts = 40;
7311 let per_part = 1000;
7312 let mut writer =
7313 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
7314 .expect("new file");
7315 for part in 0..parts {
7316 let values = (0..per_part)
7317 .map(|row| Value::Varchar(value(part * per_part + row)))
7318 .collect::<Vec<_>>();
7319 let chunk = Chunk::new(vec![
7320 Vector::from_values(LogicalType::Varchar, &values).expect("strings"),
7321 ])
7322 .expect("matching rows");
7323 writer.append(&chunk).expect("a part");
7324 }
7325 writer.finish().expect("commit");
7326
7327 let reader = Reader::open(&path).expect("reopen from disk");
7328 let dictionary = reader.table.dictionaries[0].expect("string dictionary page");
7329 assert!(
7330 parts * per_part > TEXT_PAYLOAD_VALUES * 4,
7331 "the dictionary has to be several blocks for this to be testing anything"
7332 );
7333 for part in [0, parts - 1] {
7334 let chunk = reader.read(part, &[0]).expect("a part");
7335 chunk.validate_external().expect("every payload block checks out");
7336 assert_eq!(chunk.value_at(0, 0), Value::Varchar(value(part * per_part)));
7337 }
7338
7339 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
7340 file.seek(SeekFrom::Start(dictionary.offset + u64::from(dictionary.length) - 4))
7341 .expect("the last bytes of the page are payload");
7342 file.write_all(&[255]).expect("damage the last payload block");
7343 let reader = Reader::open(&path).expect("the directory and the index are untouched");
7344 let chunk = reader.read(parts - 1, &[0]).expect("the code page remains valid");
7345 let error = chunk.validate_external().expect_err("the damage must reach the caller");
7346 assert!(error.message().contains("payload checksum differs"), "{error}");
7347 fs::remove_file(path).expect("remove scratch file");
7348 }
7349
7350 #[test]
7364 fn values_of_different_lengths_read_back_out_of_packed_offsets() {
7365 let path = path("dictionary-offsets");
7366 let value = |row: usize| {
7367 let row = row % 5_000;
7368 if row % 511 == 3 { String::new() } else { "x".repeat(row % 97) + &format!("{row:05}") }
7369 };
7370 let rows = 6_000;
7371 let mut writer =
7372 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
7373 .expect("new file");
7374 let values = (0..rows).map(|row| Value::Varchar(value(row))).collect::<Vec<_>>();
7375 for part in values.chunks(1_000) {
7376 let chunk =
7377 Chunk::new(vec![Vector::from_values(LogicalType::Varchar, part).expect("strings")])
7378 .expect("matching rows");
7379 writer.append(&chunk).expect("a part");
7380 }
7381 writer.finish().expect("commit");
7382
7383 let reader = Reader::open(&path).expect("reopen from disk");
7384 assert!(
7385 rows > TEXT_PAYLOAD_VALUES * 4,
7386 "the dictionary has to be several blocks for this to be testing anything"
7387 );
7388 for part in 0..rows / 1_000 {
7389 let chunk = reader.read(part, &[0]).expect("a part");
7390 for row in 0..1_000 {
7391 let row = part * 1_000 + row;
7392 assert_eq!(
7393 chunk.value_at(row % 1_000, 0),
7394 Value::Varchar(value(row)),
7395 "value {row}"
7396 );
7397 }
7398 }
7399 fs::remove_file(path).expect("remove scratch file");
7400 }
7401
7402 #[test]
7414 fn a_global_dictionary_is_opened_once_however_many_workers_ask_at_once() {
7415 let path = path("dictionary-once");
7416 let parts = 8;
7417 let per_part = 500;
7418 let value =
7419 |row: usize| format!("{row:07} a value long enough to be worth a payload block");
7420 let mut writer =
7421 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
7422 .expect("new file");
7423 for part in 0..parts {
7424 let values = (0..per_part)
7425 .map(|row| Value::Varchar(value(part * per_part + row)))
7426 .collect::<Vec<_>>();
7427 let chunk = Chunk::new(vec![
7428 Vector::from_values(LogicalType::Varchar, &values).expect("strings"),
7429 ])
7430 .expect("matching rows");
7431 writer.append(&chunk).expect("a part");
7432 }
7433 writer.finish().expect("commit");
7434
7435 let reader = Reader::open(&path).expect("reopen from disk");
7436 assert!(reader.table.dictionaries[0].is_some(), "the column has to have one to share");
7437 assert_eq!(reader.reads().dictionaries, 0, "opening the file does not open a dictionary");
7438
7439 let workers = 16;
7440 let gate = std::sync::Barrier::new(workers);
7441 std::thread::scope(|scope| {
7442 for worker in 0..workers {
7443 let reader = reader.clone();
7444 let gate = &gate;
7445 scope.spawn(move || {
7446 gate.wait();
7447 let chunk = reader.read(worker % parts, &[0]).expect("a part");
7448 assert_eq!(
7449 chunk.value_at(0, 0),
7450 Value::Varchar(value((worker % parts) * per_part))
7451 );
7452 });
7453 }
7454 });
7455
7456 assert_eq!(reader.reads().dictionaries, 1, "sixteen workers, one dictionary, one open");
7457 fs::remove_file(path).expect("remove scratch file");
7458 }
7459
7460 #[test]
7465 fn a_damaged_sorted_order_is_an_error() {
7466 let path = path("damaged-order");
7467 let mut writer = Writer::create(
7468 &path,
7469 "items",
7470 vec![
7471 Field::required("id", LogicalType::Integer),
7472 Field::new("text", LogicalType::Varchar),
7473 ],
7474 )
7475 .expect("new file");
7476 writer.append(&sample()).expect("stripe written");
7477 writer.finish().expect("commit");
7478
7479 let reader = Reader::open(&path).expect("valid directory");
7480 let page = reader.table.dictionaries[1].expect("string dictionary page");
7481 let mut header = [0; DICTIONARY_HEADER];
7482 read_at(&reader.file, page.offset, &mut header).expect("dictionary header");
7483 let index_len = dictionary_index_len(&header);
7484 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
7485 file.seek(SeekFrom::Start(page.offset + index_len)).expect("the first head");
7486 file.write_all(&[255]).expect("damage the order");
7487
7488 let dictionary = reader.dictionary(1).expect("read").expect("a string column has one");
7489 let error = dictionary.compare_rank(0, b"anything").expect_err("a damaged order is caught");
7490 assert!(error.message().contains("rank checksum differs"), "{error}");
7491 fs::remove_file(path).expect("remove scratch file");
7492 }
7493
7494 #[test]
7498 fn a_global_dictionary_carries_the_sorted_order_of_its_values() {
7499 let spellings = ["overlong1z", "b", "", "overlong1a", "overlong", "ab", "a", "overlong1"];
7502 let path = path("dictionary-order");
7503 let mut writer =
7504 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7505 .expect("new file");
7506 writer
7507 .append(
7508 &Chunk::new(vec![
7509 Vector::from_values(
7510 LogicalType::Varchar,
7511 &spellings.map(|text| Value::Varchar(text.into())),
7512 )
7513 .expect("strings"),
7514 ])
7515 .expect("one column"),
7516 )
7517 .expect("stripe written");
7518 writer.finish().expect("commit");
7519
7520 let reader = Reader::open(&path).expect("valid directory");
7521 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
7522 let count = dictionary.ranks().expect("a v10 file stores one");
7523 assert_eq!(count, spellings.len(), "every distinct value has a rank");
7524 let order = (0..count)
7525 .map(|rank| dictionary.code_at_rank(rank).expect("a code"))
7526 .collect::<Vec<_>>();
7527 let mut seen = order.clone();
7528 seen.sort_unstable();
7529 assert_eq!(seen, (0..spellings.len() as u32).collect::<Vec<_>>(), "a permutation of codes");
7530
7531 let ranked = order
7532 .iter()
7533 .map(|&code| {
7534 dictionary.try_bytes_at(code as usize).expect("read").expect("a value").to_vec()
7535 })
7536 .collect::<Vec<_>>();
7537 let mut expected = spellings.map(|text| text.as_bytes().to_vec()).to_vec();
7538 expected.sort();
7539 assert_eq!(ranked, expected, "rank order is value order");
7540
7541 for (rank, value) in expected.iter().enumerate() {
7544 assert_eq!(
7545 dictionary.compare_rank(rank, value).expect("compare"),
7546 Ordering::Equal,
7547 "rank {rank} is its own value"
7548 );
7549 if rank > 0 {
7550 assert_eq!(
7551 dictionary.compare_rank(rank - 1, value).expect("compare"),
7552 Ordering::Less,
7553 "rank {rank} follows the one before it"
7554 );
7555 }
7556 }
7557 fs::remove_file(path).expect("remove scratch file");
7558 }
7559
7560 #[test]
7568 fn a_dictionary_sweep_reads_every_value_and_keeps_it_under_the_budget() {
7569 let path = path("dictionary-sweep");
7570 let spellings = (0..2_500)
7573 .map(|index| Value::Varchar(format!("value {index:08} {}", "x".repeat(index % 40))))
7574 .collect::<Vec<_>>();
7575 let mut writer =
7576 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7577 .expect("new file");
7578 for part in spellings.chunks(1_024) {
7581 writer
7582 .append(
7583 &Chunk::new(vec![
7584 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
7585 ])
7586 .expect("one column"),
7587 )
7588 .expect("stripe written");
7589 }
7590 writer.finish().expect("commit");
7591
7592 let reader = Reader::open(&path).expect("valid directory");
7593 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
7594 assert_eq!(dictionary.len(), spellings.len(), "every value is distinct");
7595
7596 let resting = dictionary.footprint();
7597 let mut swept: Vec<Vec<u8>> = Vec::new();
7598 let mut at = 0;
7599 let mut calls = 0;
7600 while at < dictionary.len() {
7601 let stopped = dictionary
7602 .sweep_text(at, dictionary.len(), &mut |index: usize, text: &[u8]| {
7603 assert_eq!(index, swept.len(), "a sweep hands its values over in order");
7604 swept.push(text.to_vec());
7605 Ok(())
7606 })
7607 .expect("a sweep reads");
7608 assert!(stopped > at, "a sweep moves");
7609 at = stopped;
7610 calls += 1;
7611 }
7612 assert_eq!(calls, 3, "a sweep hands over one block at a time");
7613 let after = dictionary.footprint();
7614 assert!(after > resting, "a sweep under the budget keeps what it decoded");
7615
7616 let read = (0..dictionary.len())
7617 .map(|code| dictionary.try_bytes_at(code).expect("read").expect("a value").to_vec())
7618 .collect::<Vec<_>>();
7619 assert_eq!(swept, read, "a sweep answers what a point read answers");
7620 assert_eq!(dictionary.footprint(), after, "a point read of a kept block decodes nothing");
7621 fs::remove_file(path).expect("remove scratch file");
7622 }
7623
7624 #[test]
7635 fn a_sweep_over_a_block_with_a_short_second_run_reads_what_a_point_read_reads() {
7636 let path = path("dictionary-sweep-short-run");
7637 let spellings = (0..2_800)
7638 .map(|index| Value::Varchar(format!("value {index:08} {}", "x".repeat(index % 40))))
7639 .collect::<Vec<_>>();
7640 let mut writer =
7641 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7642 .expect("new file");
7643 for part in spellings.chunks(1_024) {
7644 writer
7645 .append(
7646 &Chunk::new(vec![
7647 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
7648 ])
7649 .expect("one column"),
7650 )
7651 .expect("stripe written");
7652 }
7653 writer.finish().expect("commit");
7654
7655 let reader = Reader::open(&path).expect("valid directory");
7656 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
7657 assert_eq!(dictionary.len(), spellings.len(), "every value is distinct");
7658 let last = dictionary.len() % TEXT_PAYLOAD_VALUES;
7659 assert!(last > TEXT_OFFSET_RUN, "the last block has to reach into a second run of offsets");
7660 assert!(last < TEXT_PAYLOAD_VALUES, "and that second run has to be short of a whole one");
7661
7662 let mut swept: Vec<Vec<u8>> = Vec::new();
7663 let mut at = 0;
7664 while at < dictionary.len() {
7665 let stopped = dictionary
7666 .sweep_text(at, dictionary.len(), &mut |index: usize, text: &[u8]| {
7667 assert_eq!(index, swept.len(), "a sweep hands its values over in order");
7668 swept.push(text.to_vec());
7669 Ok(())
7670 })
7671 .expect("a sweep reads");
7672 assert!(stopped > at, "a sweep moves");
7673 at = stopped;
7674 }
7675 let read = (0..dictionary.len())
7676 .map(|code| dictionary.try_bytes_at(code).expect("read").expect("a value").to_vec())
7677 .collect::<Vec<_>>();
7678 assert_eq!(swept, read, "a sweep answers what a point read answers");
7679 fs::remove_file(path).expect("remove scratch file");
7680 }
7681
7682 #[test]
7692 fn narrowing_a_page_takes_what_fits_and_refuses_what_does_not() {
7693 assert_eq!(fit::<i8>(&[]).expect("an empty page fits anything"), Vec::<i8>::new());
7694 assert_eq!(fit::<i8>(&[-128, 0, 127]).expect("the edges fit"), vec![-128_i8, 0, 127]);
7695 fit::<i8>(&[128]).expect_err("one past the top does not fit");
7696 fit::<i8>(&[-129]).expect_err("one past the bottom does not fit");
7697 assert_eq!(fit::<u8>(&[0, 255]).expect("the edges fit"), vec![0_u8, 255]);
7698 fit::<u8>(&[256]).expect_err("one past the top does not fit");
7699 fit::<u8>(&[-1]).expect_err("a negative does not fit an unsigned page");
7700 assert_eq!(
7701 fit::<i16>(&[-32_768, 0, 32_767]).expect("the edges fit"),
7702 vec![-32_768_i16, 0, 32_767]
7703 );
7704 fit::<i16>(&[32_768]).expect_err("one past the top does not fit");
7705 fit::<i16>(&[-32_769]).expect_err("one past the bottom does not fit");
7706 assert_eq!(fit::<u16>(&[0, 65_535]).expect("the edges fit"), vec![0_u16, 65_535]);
7707 fit::<u16>(&[65_536]).expect_err("one past the top does not fit");
7708 fit::<u16>(&[-1]).expect_err("a negative does not fit an unsigned page");
7709 assert_eq!(
7710 fit::<i32>(&[i64::from(i32::MIN), 0, i64::from(i32::MAX)]).expect("the edges fit"),
7711 vec![i32::MIN, 0, i32::MAX]
7712 );
7713 fit::<i32>(&[i64::from(i32::MAX) + 1]).expect_err("one past the top does not fit");
7714 fit::<i32>(&[i64::from(i32::MIN) - 1]).expect_err("one past the bottom does not fit");
7715 assert_eq!(
7716 fit::<u32>(&[0, 4_294_967_295]).expect("the edges fit"),
7717 vec![0_u32, 4_294_967_295]
7718 );
7719 fit::<u32>(&[4_294_967_296]).expect_err("one past the top does not fit");
7720 fit::<u32>(&[-1]).expect_err("a negative does not fit an unsigned page");
7721
7722 fit::<i8>(&[0, 1, 2, 128, 3]).expect_err("one bad value spoils the page");
7725 }
7726
7727 #[test]
7734 fn the_residue_agrees_with_a_checked_conversion_everywhere() {
7735 for value in -70_000_i64..70_000 {
7736 assert_eq!(fit::<i8>(&[value]).is_ok(), i8::try_from(value).is_ok(), "{value} as i8");
7737 assert_eq!(fit::<u8>(&[value]).is_ok(), u8::try_from(value).is_ok(), "{value} as u8");
7738 assert_eq!(fit::<i16>(&[value]).is_ok(), i16::try_from(value).is_ok(), "{value} i16");
7739 assert_eq!(fit::<u16>(&[value]).is_ok(), u16::try_from(value).is_ok(), "{value} u16");
7740 }
7741 let wide = [i64::MIN, i64::MIN + 1, i64::from(i32::MIN), 0, i64::from(u32::MAX), i64::MAX];
7742 for edge in wide {
7743 for step in -2_i64..=2 {
7744 let value = edge.saturating_add(step);
7745 assert_eq!(
7746 fit::<i32>(&[value]).is_ok(),
7747 i32::try_from(value).is_ok(),
7748 "{value} as i32"
7749 );
7750 assert_eq!(
7751 fit::<u32>(&[value]).is_ok(),
7752 u32::try_from(value).is_ok(),
7753 "{value} as u32"
7754 );
7755 }
7756 }
7757 }
7758
7759 #[test]
7767 fn a_dictionary_at_its_budget_sweeps_without_keeping() {
7768 let path = path("dictionary-budget");
7769 let spellings = (0..2_500)
7770 .map(|index| Value::Varchar(format!("value {index:08} {}", "y".repeat(index % 40))))
7771 .collect::<Vec<_>>();
7772 let mut writer =
7773 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7774 .expect("new file");
7775 for part in spellings.chunks(1_024) {
7776 writer
7777 .append(
7778 &Chunk::new(vec![
7779 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
7780 ])
7781 .expect("one column"),
7782 )
7783 .expect("stripe written");
7784 }
7785 writer.finish().expect("commit");
7786
7787 let reader = Reader::open(&path).expect("valid directory");
7788 let page = reader.table.dictionaries[0].expect("a string column has one");
7789 let file = Arc::clone(&reader.file);
7790 let starved = open_global_dictionary(file, page, &LogicalType::Varchar, 0)
7791 .expect("a dictionary opens whatever it may keep");
7792
7793 let resting = starved.footprint();
7794 let mut swept: Vec<Vec<u8>> = Vec::new();
7795 let mut at = 0;
7796 while at < starved.len() {
7797 at = starved
7798 .sweep_text(at, starved.len(), &mut |_index: usize, text: &[u8]| {
7799 swept.push(text.to_vec());
7800 Ok(())
7801 })
7802 .expect("a sweep reads");
7803 }
7804 assert_eq!(swept.len(), spellings.len(), "a starved sweep still reads every value");
7805 assert_eq!(starved.footprint(), resting, "and keeps no block it decoded");
7806
7807 let generous = reader.dictionary(0).expect("read").expect("a string column has one");
7808 let read = (0..generous.len())
7809 .map(|code| generous.try_bytes_at(code).expect("read").expect("a value").to_vec())
7810 .collect::<Vec<_>>();
7811 assert_eq!(swept, read, "a starved sweep answers what a point read answers");
7812 fs::remove_file(path).expect("remove scratch file");
7813 }
7814
7815 #[test]
7816 fn damaged_membership_cannot_skip_a_string_page() {
7817 let path = path("damaged-membership");
7818 let mut writer = Writer::create(
7819 &path,
7820 "items",
7821 vec![
7822 Field::required("id", LogicalType::Integer),
7823 Field::new("text", LogicalType::Varchar),
7824 ],
7825 )
7826 .expect("new file");
7827 writer.append(&sample()).expect("stripe written");
7828 writer.finish().expect("commit");
7829
7830 let reader = Reader::open(&path).expect("valid directory");
7831 let membership = reader.table.stripes[0].memberships[1].expect("string membership");
7832 let mut file = OpenOptions::new().write(true).open(&path).expect("open membership page");
7833 file.seek(SeekFrom::Start(membership.offset)).expect("membership start");
7834 file.write_all(&[255]).expect("damage membership");
7835 let error = reader.skips_codes(0, 1, &[3]).expect_err("corruption must not skip rows");
7836 assert!(error.message().contains("membership page checksum differs"), "{error}");
7837 fs::remove_file(path).expect("remove scratch file");
7838 }
7839
7840 #[test]
7841 fn membership_delta_stream_is_sorted_exact_and_bounded() {
7842 let unique = unique_codes(&[900, 4, 4, 72, 9, u32::MAX]);
7843 assert_eq!(unique, [4, 9, 72, 900, u32::MAX]);
7844 let encoded = encode_membership(&unique);
7845 assert_eq!(
7846 decode_membership(&encoded).expect("valid membership"),
7847 [4, 9, 72, 900, u32::MAX]
7848 );
7849 let merged = merged_codes(vec![vec![4, 900], vec![9, 900, u32::MAX], vec![72]]);
7852 assert_eq!(merged, [4, 9, 72, 900, u32::MAX]);
7853 assert_eq!(
7854 decode_membership(&encode_membership(&merged)).expect("valid membership"),
7855 unique
7856 );
7857 assert!(decode_membership(&[1, 0x80]).is_err(), "a truncated varint is invalid");
7858 assert!(
7859 decode_membership(&[1, 0xff, 0xff, 0xff, 0xff, 0x10]).is_err(),
7860 "a value past u32 is invalid"
7861 );
7862 }
7863
7864 #[test]
7865 fn a_global_dictionary_may_be_larger_than_one_column_page() {
7866 let dictionary = Page {
7867 offset: HEADER,
7868 length: u32::try_from(MAX_PAGE + 1).expect("the page bound fits on disk"),
7869 hash: 0,
7870 };
7871 let table = Table {
7872 name: "items".to_owned(),
7873 fields: vec![Field::new("text", LogicalType::Varchar)],
7874 stripes: Vec::new(),
7875 rows: 0,
7876 dictionaries: vec![Some(dictionary)],
7877 distincts: vec![None],
7878 frequencies: vec![None],
7879 };
7880 let directory = encode_directory(&table).expect("directory");
7881 let file_size = dictionary.offset + u64::from(dictionary.length) + 1;
7882
7883 let decoded = decode_directory(&directory, file_size).expect("large lazy dictionary");
7884 assert_eq!(decoded.dictionaries[0].expect("dictionary").length, dictionary.length);
7885 }
7886
7887 #[test]
7888 fn a_column_with_one_value_everywhere_costs_almost_nothing_a_row() {
7889 let path = path("constant-codes");
7890 let mut writer =
7891 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7892 .expect("new file");
7893 let empty = vec![Value::Varchar(String::new()); 1024];
7894 for _ in 0..4 {
7895 let column = Vector::from_values(LogicalType::Varchar, &empty).expect("strings");
7896 writer.append(&Chunk::new(vec![column]).expect("one column")).expect("a part");
7897 }
7898 writer.finish().expect("commit");
7899
7900 let reader = Reader::open(&path).expect("valid directory");
7901 let pages = reader.layout().columns.first().expect("one column").pages;
7902 assert!(pages < 256, "{pages} bytes of pages for 4,096 rows of one value");
7906 let read = reader.read(3, &[0]).expect("the last part back");
7907 assert_eq!(read.value_at(0, 0), Value::Varchar(String::new()));
7908 assert_eq!(read.value_at(1023, 0), Value::Varchar(String::new()));
7909 fs::remove_file(path).expect("remove scratch file");
7910 }
7911
7912 #[test]
7913 fn a_cascade_value_too_wide_for_its_column_is_refused_rather_than_cut() {
7914 let over = vec![i64::from(i32::MAX) + 1];
7917 let error = narrowed(&LogicalType::Integer, over).expect_err("a page that disagrees");
7918 assert!(format!("{error}").contains("not of its type"), "{error}");
7919 assert!(narrowed(&LogicalType::BigInt, vec![i64::MIN]).is_ok(), "bigint holds all of i64");
7920 assert!(narrowed(&LogicalType::Varchar, vec![0]).is_err(), "strings are not integers");
7921 }
7922
7923 #[test]
7924 fn a_code_stream_the_cascade_cannot_shrink_is_left_alone() {
7925 let mut state: u32 = 0x9e37_79b9;
7929 let spread: Vec<u32> = (0..1024)
7930 .map(|_| {
7931 state ^= state << 13;
7932 state ^= state >> 17;
7933 state ^= state << 5;
7934 state
7935 })
7936 .collect();
7937 assert_eq!(encoded_codes(&spread).expect("no failure"), None);
7938 let near: Vec<u32> = (0..1024).collect();
7939 let coded = encoded_codes(&near).expect("no failure").expect("counting up is packable");
7940 assert!(coded.len() < near.len() * 4, "{} bytes for a run of 1,024", coded.len());
7941 }
7942
7943 #[test]
7949 fn two_writes_of_the_same_rows_give_the_same_bytes() {
7950 fn written(path: &PathBuf) {
7951 let fields = (0..40)
7952 .map(|column| {
7953 let ty =
7954 if column % 4 == 0 { LogicalType::Varchar } else { LogicalType::BigInt };
7955 Field::new(format!("c{column}"), ty)
7956 })
7957 .collect::<Vec<_>>();
7958 let mut writer = Writer::create(path, "wide", fields).expect("new file");
7959 for part in 0..70_u64 {
7960 let columns = (0..40)
7961 .map(|column| {
7962 let values = (0..64_u64)
7963 .map(|row| {
7964 let seed = part.wrapping_mul(31).wrapping_add(row);
7965 if column % 4 == 0 {
7966 Value::Varchar(format!("v{}", seed % 17))
7967 } else {
7968 Value::BigInt(i64::try_from(seed % 97).expect("small"))
7969 }
7970 })
7971 .collect::<Vec<_>>();
7972 let ty = if column % 4 == 0 {
7973 LogicalType::Varchar
7974 } else {
7975 LogicalType::BigInt
7976 };
7977 Vector::from_values(ty, &values).expect("a column")
7978 })
7979 .collect::<Vec<_>>();
7980 writer.append(&Chunk::new(columns).expect("forty columns")).expect("a part");
7981 }
7982 writer.finish().expect("commit");
7983 }
7984
7985 let first = path("repeatable-one");
7986 let second = path("repeatable-two");
7987 written(&first);
7988 written(&second);
7989 let left = fs::read(&first).expect("the first file");
7990 let right = fs::read(&second).expect("the second file");
7991 assert_eq!(left.len(), right.len(), "two writes of the same rows differ in length");
7992 assert!(left == right, "two writes of the same rows differ in their bytes");
7993
7994 let reader = Reader::open(&first).expect("valid directory");
7997 assert_eq!(reader.table().rows(), 70 * 64);
7998 let read = reader.read(0, &[0, 1]).expect("the first part back");
7999 assert_eq!(read.value_at(0, 0), Value::Varchar("v0".to_owned()));
8000 assert_eq!(read.value_at(0, 1), Value::BigInt(0));
8001 fs::remove_file(first).expect("remove scratch file");
8002 fs::remove_file(second).expect("remove scratch file");
8003 }
8004
8005 fn three_tables(path: &PathBuf) {
8007 let writer = Writer::create(
8008 path,
8009 "region",
8010 vec![
8011 Field::new("r_key", LogicalType::Integer),
8012 Field::new("r_name", LogicalType::Varchar),
8013 ],
8014 )
8015 .expect("new file");
8016 let mut writer = writer;
8017 writer
8018 .append(
8019 &Chunk::new(vec![
8020 Vector::from_values(
8021 LogicalType::Integer,
8022 &[Value::Integer(0), Value::Integer(1)],
8023 )
8024 .expect("keys"),
8025 Vector::from_values(
8026 LogicalType::Varchar,
8027 &[Value::Varchar("AFRICA".to_owned()), Value::Varchar("ASIA".to_owned())],
8028 )
8029 .expect("names"),
8030 ])
8031 .expect("two columns"),
8032 )
8033 .expect("a part");
8034 let mut writer = writer
8035 .next("empty", vec![Field::new("nothing", LogicalType::BigInt)])
8036 .expect("a second table");
8037 writer
8038 .append(
8039 &Chunk::new(vec![
8040 Vector::from_values(LogicalType::BigInt, &[Value::BigInt(7)]).expect("a row"),
8041 ])
8042 .expect("one column"),
8043 )
8044 .expect("a part");
8045 let mut writer =
8046 writer.next("wide", vec![Field::new("n", LogicalType::BigInt)]).expect("a third table");
8047 for part in 0..70_i64 {
8048 let values = (0..64).map(|row| Value::BigInt(part * 64 + row)).collect::<Vec<_>>();
8049 writer
8050 .append(
8051 &Chunk::new(vec![
8052 Vector::from_values(LogicalType::BigInt, &values).expect("a column"),
8053 ])
8054 .expect("one column"),
8055 )
8056 .expect("a part");
8057 }
8058 writer.finish().expect("commit");
8059 }
8060
8061 #[test]
8062 fn three_tables_in_one_file_read_back_by_name() {
8063 let file = path("three-tables");
8064 three_tables(&file);
8065 let catalog = Catalog::open(&file).expect("a committed catalog");
8066 assert_eq!(catalog.names().collect::<Vec<_>>(), ["region", "empty", "wide"]);
8067
8068 let region = catalog.table("region").expect("the first table");
8069 assert_eq!(region.table().rows(), 2);
8070 assert_eq!(
8071 region.read(0, &[1]).expect("names").value_at(1, 0),
8072 Value::Varchar("ASIA".to_owned())
8073 );
8074
8075 let wide = catalog.table("wide").expect("the third table");
8076 assert_eq!(wide.table().rows(), 70 * 64);
8077 assert_eq!(wide.read(0, &[0]).expect("the first part").value_at(0, 0), Value::BigInt(0));
8078
8079 let empty = catalog.table("empty").expect("the second table");
8082 assert_eq!(empty.table().rows(), 1);
8083 assert_eq!(empty.read(0, &[0]).expect("the row").value_at(0, 0), Value::BigInt(7));
8084
8085 fs::remove_file(file).expect("remove scratch file");
8086 }
8087
8088 #[test]
8089 fn a_name_the_file_does_not_hold_is_an_error_rather_than_the_first_table() {
8090 let file = path("three-tables-missing");
8091 three_tables(&file);
8092 let catalog = Catalog::open(&file).expect("a committed catalog");
8093 let error = catalog.table("nation").expect_err("no such table");
8094 assert!(error.message().contains("nation"), "{}", error.message());
8095 fs::remove_file(file).expect("remove scratch file");
8096 }
8097
8098 #[test]
8099 fn a_file_of_three_tables_will_not_open_as_one() {
8100 let file = path("three-tables-unnamed");
8101 three_tables(&file);
8102 let error = Reader::open(&file).expect_err("more than one table");
8103 assert!(error.message().contains("more than one table"), "{}", error.message());
8104 fs::remove_file(file).expect("remove scratch file");
8105 }
8106
8107 #[test]
8109 fn decimals_of_every_storage_width_round_trip() {
8110 let file = path("decimals");
8111 let widths = [(4_u8, 2_u8), (9, 2), (18, 4), (38, 6)];
8112 let fields = widths
8113 .iter()
8114 .enumerate()
8115 .map(|(index, (width, scale))| {
8116 Field::new(
8117 format!("d{index}"),
8118 LogicalType::decimal(*width, *scale).expect("a decimal type"),
8119 )
8120 })
8121 .collect::<Vec<_>>();
8122 let mut writer = Writer::create(&file, "money", fields).expect("new file");
8123 let rows: [i128; 3] = [-1234, 0, 999];
8124 let columns = widths
8125 .iter()
8126 .map(|(width, scale)| {
8127 let values = rows
8128 .iter()
8129 .map(|unscaled| Value::Decimal {
8130 unscaled: *unscaled,
8131 width: *width,
8132 scale: *scale,
8133 })
8134 .collect::<Vec<_>>();
8135 Vector::from_values(
8136 LogicalType::decimal(*width, *scale).expect("a decimal type"),
8137 &values,
8138 )
8139 .expect("a decimal column")
8140 })
8141 .collect::<Vec<_>>();
8142 writer.append(&Chunk::new(columns).expect("four columns")).expect("a part");
8143 writer.finish().expect("commit");
8144
8145 let reader = Reader::open(&file).expect("a committed file");
8146 for (index, (width, scale)) in widths.iter().enumerate() {
8147 assert_eq!(
8148 reader.table().fields()[index].ty,
8149 LogicalType::decimal(*width, *scale).expect("a decimal type"),
8150 "column {index} came back as another type"
8151 );
8152 let column = reader.read(0, &[index]).expect("the column");
8153 for (row, unscaled) in rows.iter().enumerate() {
8154 assert_eq!(
8155 column.value_at(row, 0),
8156 Value::Decimal { unscaled: *unscaled, width: *width, scale: *scale },
8157 "column {index} row {row}"
8158 );
8159 }
8160 }
8161 fs::remove_file(file).expect("remove scratch file");
8162 }
8163
8164 #[test]
8165 fn two_tables_of_one_name_are_refused_before_anything_is_committed() {
8166 let file = path("two-of-a-name");
8167 let writer = Writer::create(&file, "t", vec![Field::new("a", LogicalType::BigInt)])
8168 .expect("new file");
8169 let error = writer
8170 .next("t", vec![Field::new("a", LogicalType::BigInt)])
8171 .expect_err("the same name twice");
8172 assert!(error.message().contains("same name"), "{}", error.message());
8173 fs::remove_file(file).expect("remove scratch file");
8174 }
8175
8176 #[test]
8177 fn opening_the_catalog_reads_no_table_directory() {
8178 let file = path("catalog-only");
8179 three_tables(&file);
8180 let catalog = Catalog::open(&file).expect("a committed catalog");
8181 assert_eq!(catalog.opening.reads, 2, "opening the catalog read more than the slot");
8184 assert_eq!(catalog.names().len(), 3);
8185 fs::remove_file(file).expect("remove scratch file");
8186 }
8187
8188 #[test]
8199 fn the_checksum_answers_what_it_has_always_answered() {
8200 let bytes: Vec<u8> =
8201 (0..1000_u32).map(|at| (at.wrapping_mul(31).wrapping_add(7) % 251) as u8).collect();
8202 for (length, expected) in [
8203 (0, 0xef46_db37_51d8_e999),
8204 (1, 0xa96c_7f0c_e858_bbb7),
8205 (3, 0x56e6_9576_32a4_87f9),
8206 (4, 0xc60d_15b1_e3ff_8f04),
8207 (5, 0x8088_1585_8624_dd4e),
8208 (7, 0xafbe_fc3d_6c6f_9a8e),
8209 (8, 0x3da5_c7aa_2696_83e0),
8210 (9, 0x465e_c429_b13c_3892),
8211 (15, 0xdee8_9d8a_065a_6233),
8212 (16, 0x1330_489a_7767_9c80),
8213 (31, 0x3391_303d_485e_846e),
8214 (32, 0x40b7_aff7_5d45_bbc8),
8215 (33, 0x4997_cae4_951c_17a5),
8216 (39, 0x5807_28fd_5c14_5739),
8217 (40, 0xf95c_f6f5_c08a_3d3b),
8218 (63, 0x2944_b4da_fc69_b206),
8219 (64, 0xbb76_f6ef_19bd_5a1b),
8220 (65, 0x814e_0c65_4a9f_d640),
8221 (127, 0x00de_aab1_31cf_f89b),
8222 (1000, 0x9e33_00c1_cde3_c58d),
8223 ] {
8224 assert_eq!(checksum(&bytes[..length]), expected, "the checksum of {length} bytes");
8225 }
8226 assert_eq!(checksum(b"the quick brown fox jumps over the lazy dog"), 0xed71_4233_c5a9_a792);
8227 }
8228}