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 {
130 const P1: u64 = 11_400_714_785_074_694_791;
131 const P2: u64 = 14_029_467_366_897_019_727;
132 const P3: u64 = 1_609_587_929_392_839_161;
133 const P4: u64 = 9_650_029_242_287_828_579;
134 const P5: u64 = 2_870_177_450_012_600_261;
135 let round = |state: u64, word: u64| {
136 state.wrapping_add(word.wrapping_mul(P2)).rotate_left(31).wrapping_mul(P1)
137 };
138 let merge = |state: u64, lane: u64| (state ^ round(0, lane)).wrapping_mul(P1).wrapping_add(P4);
139 let word =
140 |at: usize| u64::from_le_bytes(bytes[at..at + 8].try_into().expect("eight checksum bytes"));
141
142 let mut at = 0;
143 let mut hash = if bytes.len() >= 32 {
144 let mut one = P1.wrapping_add(P2);
145 let mut two = P2;
146 let mut three = 0;
147 let mut four = 0_u64.wrapping_sub(P1);
148 while at + 32 <= bytes.len() {
149 one = round(one, word(at));
150 two = round(two, word(at + 8));
151 three = round(three, word(at + 16));
152 four = round(four, word(at + 24));
153 at += 32;
154 }
155 let combined = one
156 .rotate_left(1)
157 .wrapping_add(two.rotate_left(7))
158 .wrapping_add(three.rotate_left(12))
159 .wrapping_add(four.rotate_left(18));
160 merge(merge(merge(merge(combined, one), two), three), four)
161 } else {
162 P5
163 };
164 hash = hash.wrapping_add(bytes.len() as u64);
165 while at + 8 <= bytes.len() {
166 hash ^= round(0, word(at));
167 hash = hash.rotate_left(27).wrapping_mul(P1).wrapping_add(P4);
168 at += 8;
169 }
170 if at + 4 <= bytes.len() {
171 let tail = u32::from_le_bytes(bytes[at..at + 4].try_into().expect("four checksum bytes"));
172 hash ^= u64::from(tail).wrapping_mul(P1);
173 hash = hash.rotate_left(23).wrapping_mul(P2).wrapping_add(P3);
174 at += 4;
175 }
176 while at < bytes.len() {
177 hash ^= u64::from(bytes[at]).wrapping_mul(P5);
178 hash = hash.rotate_left(11).wrapping_mul(P1);
179 at += 1;
180 }
181 hash ^= hash >> 33;
182 hash = hash.wrapping_mul(P2);
183 hash ^= hash >> 29;
184 hash = hash.wrapping_mul(P3);
185 hash ^ (hash >> 32)
186}
187
188#[derive(Debug, Clone, Copy)]
189struct Slot {
190 offset: u64,
191 length: u32,
192 generation: u64,
193 hash: u64,
194}
195
196impl Slot {
197 fn bytes(self) -> [u8; SLOT_BYTES] {
198 let mut result = [0; SLOT_BYTES];
199 result[..8].copy_from_slice(&self.offset.to_le_bytes());
200 result[8..12].copy_from_slice(&self.length.to_le_bytes());
201 result[12..20].copy_from_slice(&self.generation.to_le_bytes());
202 result[20..28].copy_from_slice(&self.hash.to_le_bytes());
203 result
204 }
205
206 fn read(bytes: &[u8]) -> Self {
207 Self {
208 offset: u64::from_le_bytes(bytes[..8].try_into().expect("eight bytes")),
209 length: u32::from_le_bytes(bytes[8..12].try_into().expect("four bytes")),
210 generation: u64::from_le_bytes(bytes[12..20].try_into().expect("eight bytes")),
211 hash: u64::from_le_bytes(bytes[20..28].try_into().expect("eight bytes")),
212 }
213 }
214}
215
216#[derive(Debug, Clone, Copy)]
217struct Page {
218 offset: u64,
219 length: u32,
220 hash: u64,
221}
222
223impl Page {
224 fn bytes(&self) -> u64 {
226 u64::from(self.length)
227 }
228}
229
230#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
231enum FrequencyValue {
232 Null,
233 Integer(i128),
234 Code(u32),
235}
236
237#[derive(Debug, Clone)]
238struct FrequencyEntry {
239 value: FrequencyValue,
240 count: u64,
241}
242
243#[derive(Debug, Clone)]
248struct FrequencySummary {
249 entries: Vec<FrequencyEntry>,
250 omitted_max: u64,
251 ordinals: Vec<u64>,
252}
253
254#[derive(Debug, Clone, PartialEq, Eq)]
256pub struct FrequencyOccurrences {
257 pub omitted_max: u64,
259 pub ordinals: Vec<u64>,
261}
262
263#[derive(Debug, Clone, Copy, Default)]
270struct Span {
271 offset: u64,
272 length: u32,
273}
274
275#[derive(Debug, Clone)]
277pub struct Stripe {
278 rows: usize,
279 parts: Vec<u32>,
282 index: Span,
286 pages: Vec<Span>,
287 memberships: Vec<Option<Page>>,
288 sieves: Vec<Option<Page>>,
291 part_ranges: Vec<Option<Page>>,
302 zone: Zone,
303}
304
305impl Stripe {
306 #[must_use]
308 pub fn rows(&self) -> usize {
309 self.rows
310 }
311
312 #[must_use]
314 pub fn parts(&self) -> usize {
315 self.parts.len()
316 }
317
318 #[must_use]
324 pub fn zone(&self) -> &Zone {
325 &self.zone
326 }
327}
328
329#[derive(Debug, Clone)]
331pub struct Table {
332 name: String,
333 fields: Vec<Field>,
334 stripes: Vec<Stripe>,
335 rows: usize,
336 dictionaries: Vec<Option<Page>>,
337 frequencies: Vec<Option<FrequencySummary>>,
338 distincts: Vec<Option<u64>>,
348}
349
350impl Table {
351 #[must_use]
353 pub fn name(&self) -> &str {
354 &self.name
355 }
356
357 #[must_use]
359 pub fn fields(&self) -> &[Field] {
360 &self.fields
361 }
362
363 #[must_use]
365 pub fn rows(&self) -> usize {
366 self.rows
367 }
368
369 #[must_use]
371 pub fn stripes(&self) -> &[Stripe] {
372 &self.stripes
373 }
374}
375
376#[derive(Debug, Clone)]
388struct Entry {
389 name: String,
390 fields: Vec<Field>,
391 rows: usize,
392 directory: Page,
394}
395
396#[derive(Debug, Clone)]
398pub struct ColumnLayout {
399 pub name: String,
401 pub kind: String,
403 pub pages: u64,
405 pub memberships: u64,
407 pub sieves: u64,
409 pub part_ranges: u64,
411 pub dictionary: u64,
413}
414
415impl ColumnLayout {
416 #[must_use]
418 pub fn total(&self) -> u64 {
419 self.pages
420 .saturating_add(self.memberships)
421 .saturating_add(self.sieves)
422 .saturating_add(self.part_ranges)
423 .saturating_add(self.dictionary)
424 }
425}
426
427#[derive(Debug, Clone)]
438pub struct Layout {
439 pub file: u64,
441 pub rows: usize,
443 pub stripes: usize,
445 pub parts: usize,
447 pub columns: Vec<ColumnLayout>,
449 pub indexes: u64,
452 pub directory: u64,
454 pub header: u64,
456}
457
458impl Layout {
459 #[must_use]
461 pub fn columns_total(&self) -> u64 {
462 self.columns.iter().map(ColumnLayout::total).fold(0, u64::saturating_add)
463 }
464
465 #[must_use]
471 pub fn unaccounted(&self) -> u64 {
472 self.file
473 .saturating_sub(self.columns_total())
474 .saturating_sub(self.indexes)
475 .saturating_sub(self.directory)
476 .saturating_sub(self.header)
477 }
478}
479
480#[derive(Debug)]
482struct GlobalDictionary {
483 primary: HashMap<u64, u32>,
484 collisions: HashMap<u64, Vec<u32>>,
485 offsets: Vec<u32>,
486 payload: Vec<u8>,
487 counts: Vec<u64>,
488 nulls: u64,
489}
490
491impl GlobalDictionary {
492 fn new() -> Self {
493 Self {
494 primary: HashMap::new(),
495 collisions: HashMap::new(),
496 offsets: vec![0],
497 payload: Vec::new(),
498 counts: Vec::new(),
499 nulls: 0,
500 }
501 }
502
503 fn bytes(&self, code: u32) -> Option<&[u8]> {
504 let start = *self.offsets.get(code as usize)? as usize;
505 let end = *self.offsets.get(code as usize + 1)? as usize;
506 self.payload.get(start..end)
507 }
508
509 fn code(&mut self, text: &str) -> Result<u32> {
510 let hash = checksum(text.as_bytes());
511 if let Some(&code) = self.primary.get(&hash) {
512 if self.bytes(code) == Some(text.as_bytes()) {
513 return Ok(code);
514 }
515 if let Some(codes) = self.collisions.get(&hash) {
516 if let Some(code) =
517 codes.iter().copied().find(|&code| self.bytes(code) == Some(text.as_bytes()))
518 {
519 return Ok(code);
520 }
521 }
522 let code = self.insert(text)?;
523 self.collisions.entry(hash).or_default().push(code);
524 return Ok(code);
525 }
526 let code = self.insert(text)?;
527 self.primary.insert(hash, code);
528 Ok(code)
529 }
530
531 fn insert(&mut self, text: &str) -> Result<u32> {
532 let code = u32::try_from(self.offsets.len() - 1)
533 .map_err(|_| invalid("global dictionary has too many values"))?;
534 self.payload.extend_from_slice(text.as_bytes());
535 self.offsets.push(
536 u32::try_from(self.payload.len())
537 .map_err(|_| invalid("global dictionary payload exceeds 4 GiB"))?,
538 );
539 self.counts.push(0);
540 Ok(code)
541 }
542
543 fn ranked(&self) -> Vec<(u64, u32)> {
563 let count = self.offsets.len() - 1;
564 let mut ranked = (0..count)
565 .map(|code| {
566 let code = code as u32;
567 (head(self.bytes(code).unwrap_or_default()), code)
568 })
569 .collect::<Vec<_>>();
570 ranked.sort_unstable_by(|left, right| {
571 left.0.cmp(&right.0).then_with(|| self.bytes(left.1).cmp(&self.bytes(right.1)))
572 });
573 ranked
574 }
575
576 fn observe(&mut self, code: u32, null: bool) -> Result<()> {
577 if null {
578 self.nulls = self.nulls.saturating_add(1);
579 return Ok(());
580 }
581 let count = self
582 .counts
583 .get_mut(code as usize)
584 .ok_or_else(|| invalid("global dictionary count code is out of range"))?;
585 *count = count.saturating_add(1);
586 Ok(())
587 }
588}
589
590#[derive(Debug)]
598pub struct Writer {
599 file: File,
600 at: u64,
608 table: Table,
609 generation: u64,
610 order: Vec<((u64, u64), (u64, u64))>,
613 next_order: u64,
614 dictionaries: Vec<Option<GlobalDictionary>>,
615 pending: Vec<PendingChunk>,
616 closed: Vec<Entry>,
618}
619
620#[derive(Debug)]
628struct PendingChunk {
629 order: (u64, u64),
630 chunk: Chunk,
631}
632
633#[derive(Debug)]
639struct ColumnStripe {
640 pages: Vec<Vec<u8>>,
641 codes: Vec<Option<Vec<u32>>>,
642 sieves: Vec<Option<Sieve>>,
643 ranges: Vec<Range>,
644}
645
646fn weight(ty: &LogicalType) -> usize {
654 match ty {
655 LogicalType::Varchar | LogicalType::Blob => 64,
656 LogicalType::BigInt
657 | LogicalType::UBigInt
658 | LogicalType::Timestamp
659 | LogicalType::Double
660 | LogicalType::Decimal { .. } => 8,
661 LogicalType::Integer | LogicalType::UInteger | LogicalType::Date | LogicalType::Float => 4,
662 LogicalType::SmallInt | LogicalType::USmallInt => 2,
663 _ => 1,
664 }
665}
666
667pub const STRIPE_PARTS: usize = 64;
674
675const DICTIONARY_DECIDE_ROWS: usize = 4_096;
683
684const DICTIONARY_DISTINCT_IN_TEN: usize = 9;
700
701const INDEX_ENTRY: usize = size_of::<u32>() + size_of::<u64>();
703
704fn index_section(parts: usize) -> Result<usize> {
706 parts
707 .checked_mul(INDEX_ENTRY)
708 .and_then(|bytes| bytes.checked_add(size_of::<u64>()))
709 .ok_or_else(|| invalid("index page length overflow"))
710}
711
712impl Writer {
713 pub fn open(
731 path: impl AsRef<Path>,
732 name: impl Into<String>,
733 fields: Vec<Field>,
734 ) -> Result<Self> {
735 for field in &fields {
736 type_tag(&field.ty)?;
737 }
738 let name = name.into();
739 let path = path.as_ref();
740 let (_, size, slot, bytes, _) = slot_bytes(path)?;
741 let closed = decode_catalog(&bytes, size)?;
742 if closed.iter().any(|held| held.name == name) {
743 return Err(invalid("two tables in one native file have the same name"));
744 }
745 let generation = slot
750 .generation
751 .checked_add(1)
752 .ok_or_else(|| invalid("native file generation overflow"))?;
753 let file = OpenOptions::new().write(true).read(true).open(path).map_err(io)?;
754 Ok(Self {
755 file,
756 at: size,
759 dictionaries: fields
760 .iter()
761 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
762 .collect(),
763 table: Table {
764 name,
765 dictionaries: vec![None; fields.len()],
766 distincts: vec![None; fields.len()],
767 fields,
768 stripes: Vec::new(),
769 rows: 0,
770 frequencies: Vec::new(),
771 },
772 generation,
773 order: Vec::new(),
774 next_order: 0,
775 pending: Vec::with_capacity(STRIPE_PARTS),
776 closed,
777 })
778 }
779
780 pub fn create(
786 path: impl AsRef<Path>,
787 name: impl Into<String>,
788 fields: Vec<Field>,
789 ) -> Result<Self> {
790 for field in &fields {
791 type_tag(&field.ty)?;
792 }
793 let file =
794 OpenOptions::new().write(true).read(true).create_new(true).open(path).map_err(io)?;
795 let mut header = [0; HEADER as usize];
796 header[..8].copy_from_slice(MAGIC);
797 header[8..12].copy_from_slice(&FORMAT.to_le_bytes());
798 write_at(&file, 0, &header)?;
799 Ok(Self {
800 file,
801 at: HEADER,
802 dictionaries: fields
803 .iter()
804 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
805 .collect(),
806 table: Table {
807 name: name.into(),
808 dictionaries: vec![None; fields.len()],
809 distincts: vec![None; fields.len()],
810 fields,
811 stripes: Vec::new(),
812 rows: 0,
813 frequencies: Vec::new(),
814 },
815 generation: 1,
816 order: Vec::new(),
817 next_order: 0,
818 pending: Vec::with_capacity(STRIPE_PARTS),
819 closed: Vec::new(),
820 })
821 }
822
823 pub fn next(mut self, name: impl Into<String>, fields: Vec<Field>) -> Result<Self> {
834 for field in &fields {
835 type_tag(&field.ty)?;
836 }
837 let name = name.into();
838 let entry = self.close()?;
839 if self.closed.iter().chain(std::iter::once(&entry)).any(|held| held.name == name) {
840 return Err(invalid("two tables in one native file have the same name"));
841 }
842 let Self { file, at, generation, mut closed, .. } = self;
843 closed.push(entry);
844 Ok(Self {
845 file,
846 at,
847 generation,
848 closed,
849 dictionaries: fields
850 .iter()
851 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
852 .collect(),
853 table: Table {
854 name,
855 dictionaries: vec![None; fields.len()],
856 distincts: vec![None; fields.len()],
857 fields,
858 stripes: Vec::new(),
859 rows: 0,
860 frequencies: Vec::new(),
861 },
862 order: Vec::new(),
863 next_order: 0,
864 pending: Vec::with_capacity(STRIPE_PARTS),
865 })
866 }
867
868 fn put(&mut self, bytes: &[u8]) -> Result<()> {
873 write_at(&self.file, self.at, bytes)?;
874 self.at = self
875 .at
876 .checked_add(bytes.len() as u64)
877 .ok_or_else(|| invalid("native file length overflow"))?;
878 Ok(())
879 }
880
881 pub fn append(&mut self, chunk: &Chunk) -> Result<()> {
887 let order = (self.next_order, 0);
888 self.next_order = self.next_order.saturating_add(1);
889 self.append_at(order, chunk)
890 }
891
892 pub fn append_at(&mut self, order: (u64, u64), chunk: &Chunk) -> Result<()> {
903 if chunk.is_empty() {
904 return Ok(());
905 }
906 self.admit(chunk)?;
907 if self.pending.last().is_some_and(|last| last.order > order) {
908 self.flush_pending()?;
909 }
910 self.pending.push(PendingChunk { order, chunk: chunk.clone() });
915 if self.pending.len() == STRIPE_PARTS {
916 self.flush_pending()?;
917 }
918 Ok(())
919 }
920
921 pub fn append_stripe(&mut self, parts: Vec<((u64, u64), Chunk)>) -> Result<()> {
937 if parts.len() > STRIPE_PARTS {
938 return Err(invalid("a stripe was handed more parts than it holds"));
939 }
940 self.flush_pending()?;
943 for (order, chunk) in parts {
944 if chunk.is_empty() {
945 continue;
946 }
947 self.admit(&chunk)?;
948 self.pending.push(PendingChunk { order, chunk });
949 }
950 self.flush_pending()
951 }
952
953 fn admit(&mut self, chunk: &Chunk) -> Result<()> {
955 if chunk.width() != self.table.fields.len() {
956 return Err(invalid("chunk width differs from table schema"));
957 }
958 for (index, field) in self.table.fields.iter().enumerate() {
959 if chunk.column(index)?.logical_type() != &field.ty {
960 return Err(invalid("chunk type differs from table schema"));
961 }
962 }
963 self.table.rows = self
964 .table
965 .rows
966 .checked_add(chunk.len())
967 .ok_or_else(|| invalid("row count overflow"))?;
968 Ok(())
969 }
970
971 fn encode_column(
1002 index: usize,
1003 held: &[PendingChunk],
1004 dictionary: &mut Option<GlobalDictionary>,
1005 ) -> Result<ColumnStripe> {
1006 let deciding = dictionary.as_ref().is_some_and(|held| held.offsets.len() == 1);
1009 let stripe = Self::encode_pages(index, held, dictionary.as_mut())?;
1010 if !deciding {
1011 return Ok(stripe);
1012 }
1013 let rows: usize = held.iter().map(|pending| pending.chunk.len()).sum();
1014 let distinct = dictionary.as_ref().map_or(0, |held| held.offsets.len() - 1);
1015 if rows < DICTIONARY_DECIDE_ROWS
1016 || distinct.saturating_mul(10) <= rows.saturating_mul(DICTIONARY_DISTINCT_IN_TEN)
1017 {
1018 return Ok(stripe);
1019 }
1020 *dictionary = None;
1021 Self::encode_pages(index, held, None)
1022 }
1023
1024 fn encode_pages(
1026 index: usize,
1027 held: &[PendingChunk],
1028 mut dictionary: Option<&mut GlobalDictionary>,
1029 ) -> Result<ColumnStripe> {
1030 let mut stripe = ColumnStripe {
1031 pages: Vec::with_capacity(held.len()),
1032 codes: Vec::with_capacity(held.len()),
1033 sieves: Vec::with_capacity(held.len()),
1034 ranges: Vec::with_capacity(held.len()),
1035 };
1036 for pending in held {
1037 let column = pending.chunk.column(index)?;
1038 let (bytes, unique) = encode(column, dictionary.as_deref_mut())?;
1039 if bytes.len() > MAX_PAGE {
1040 return Err(invalid("column page exceeds the configured bound"));
1041 }
1042 let range = Range::of(column);
1045 let sieve = match dictionary {
1058 Some(_) => None,
1059 None => Sieve::of(column, &range, SIEVE_BUDGET)
1060 .filter(|sieve| sieve.len() < bytes.len()),
1061 };
1062 stripe.pages.push(bytes);
1063 stripe.codes.push(unique);
1064 stripe.sieves.push(sieve);
1065 stripe.ranges.push(range);
1066 }
1067 Ok(stripe)
1068 }
1069
1070 fn encode_columns(&mut self, held: &[PendingChunk]) -> Result<Vec<ColumnStripe>> {
1079 let width = self.table.fields.len();
1080 let workers = std::thread::available_parallelism()
1081 .map_or(1, usize::from)
1082 .min(MAX_ENCODE_WORKERS)
1083 .min(width);
1084 if workers <= 1 || held.len() <= 1 {
1085 return self
1086 .dictionaries
1087 .iter_mut()
1088 .enumerate()
1089 .map(|(index, dictionary)| Self::encode_column(index, held, dictionary))
1090 .collect();
1091 }
1092 let mut jobs: Vec<(usize, Option<GlobalDictionary>)> =
1095 std::mem::take(&mut self.dictionaries).into_iter().enumerate().collect();
1096 jobs.sort_by_key(|(index, _)| weight(&self.table.fields[*index].ty));
1098 let queue = Mutex::new(jobs);
1099 let pieces = std::thread::scope(|scope| {
1100 (0..workers)
1101 .map(|_| {
1102 scope.spawn(|| {
1103 let mut mine = Vec::new();
1104 loop {
1105 let taken = queue
1106 .lock()
1107 .map_err(|_| Error::internal("a native encode worker panicked"))?
1108 .pop();
1109 let Some((index, mut dictionary)) = taken else { break };
1110 let encoded = Self::encode_column(index, held, &mut dictionary)?;
1111 mine.push((index, dictionary, encoded));
1112 }
1113 Ok(mine)
1114 })
1115 })
1116 .collect::<Vec<_>>()
1117 .into_iter()
1118 .map(|handle| {
1119 handle.join().map_err(|_| Error::internal("a native encode worker panicked"))?
1120 })
1121 .collect::<Result<Vec<_>>>()
1122 })?;
1123 let mut dictionaries: Vec<Option<GlobalDictionary>> = (0..width).map(|_| None).collect();
1124 let mut encoded: Vec<Option<ColumnStripe>> = (0..width).map(|_| None).collect();
1125 for piece in pieces {
1126 for (index, dictionary, stripe) in piece {
1127 dictionaries[index] = dictionary;
1128 encoded[index] = Some(stripe);
1129 }
1130 }
1131 self.dictionaries = dictionaries;
1132 encoded
1133 .into_iter()
1134 .map(|stripe| stripe.ok_or_else(|| Error::internal("a column was never encoded")))
1135 .collect()
1136 }
1137
1138 fn flush_pending(&mut self) -> Result<()> {
1140 if self.pending.is_empty() {
1141 return Ok(());
1142 }
1143 let width = self.table.fields.len();
1144 let mut held = std::mem::take(&mut self.pending);
1147 let parts = held.len();
1148 let encoded = self.encode_columns(&held)?;
1149 let mut pages = Vec::with_capacity(width);
1150 let mut memberships = vec![None; width];
1151 let mut ranges = Vec::with_capacity(width);
1152 let mut index = Vec::with_capacity(width.saturating_mul(index_section(parts)?));
1153 for stripe in &encoded {
1154 let offset = self.at;
1155 let section = index.len();
1156 let mut length = 0_usize;
1157 for bytes in &stripe.pages {
1158 write_at(&self.file, self.at + length as u64, bytes)?;
1159 put_u32(
1160 &mut index,
1161 u32::try_from(bytes.len()).map_err(|_| invalid("part length overflow"))?,
1162 );
1163 put_u64(&mut index, checksum(bytes));
1164 length = length
1165 .checked_add(bytes.len())
1166 .ok_or_else(|| invalid("column page length overflow"))?;
1167 }
1168 let hash = checksum(&index[section..]);
1169 put_u64(&mut index, hash);
1170 if length > MAX_PAGE {
1171 return Err(invalid("column page exceeds the configured bound"));
1172 }
1173 self.at = self
1174 .at
1175 .checked_add(length as u64)
1176 .ok_or_else(|| invalid("native file length overflow"))?;
1177 pages.push(Span {
1178 offset,
1179 length: u32::try_from(length).map_err(|_| invalid("page length overflow"))?,
1180 });
1181 ranges.push(merged_range(stripe.ranges.iter().cloned()));
1182 }
1183 for (membership, stripe) in memberships.iter_mut().zip(&encoded) {
1184 if stripe.codes.iter().all(Option::is_none) {
1185 continue;
1186 }
1187 let lists = stripe
1188 .codes
1189 .iter()
1190 .map(|codes| codes.clone().unwrap_or_default())
1191 .collect::<Vec<_>>();
1192 let bytes = encode_membership(&merged_codes(lists));
1193 let offset = self.at;
1194 self.put(&bytes)?;
1195 *membership = Some(Page {
1196 offset,
1197 length: u32::try_from(bytes.len())
1198 .map_err(|_| invalid("membership page length overflow"))?,
1199 hash: checksum(&bytes),
1200 });
1201 }
1202 let mut sieves = vec![None; width];
1203 for (page, stripe) in sieves.iter_mut().zip(&encoded) {
1204 if stripe.sieves.iter().all(Option::is_none) {
1205 continue;
1206 }
1207 let bytes = encode_sieves(stripe.sieves.iter())?;
1208 let offset = self.at;
1209 self.put(&bytes)?;
1210 *page = Some(Page {
1211 offset,
1212 length: u32::try_from(bytes.len())
1213 .map_err(|_| invalid("sieve page length overflow"))?,
1214 hash: checksum(&bytes),
1215 });
1216 }
1217 let mut part_ranges = vec![None; width];
1223 if parts > 1 {
1224 for ((page, stripe), span) in part_ranges.iter_mut().zip(&encoded).zip(&pages) {
1225 let bytes = encode_part_ranges(&stripe.ranges)?;
1226 if bytes.len() >= span.length as usize {
1227 continue;
1228 }
1229 let offset = self.at;
1230 self.put(&bytes)?;
1231 *page = Some(Page {
1232 offset,
1233 length: u32::try_from(bytes.len())
1234 .map_err(|_| invalid("part range page length overflow"))?,
1235 hash: checksum(&bytes),
1236 });
1237 }
1238 }
1239 let offset = self.at;
1240 self.put(&index)?;
1241 let index = Span {
1242 offset,
1243 length: u32::try_from(index.len())
1244 .map_err(|_| invalid("index page length overflow"))?,
1245 };
1246 let mut rows = 0_usize;
1247 let mut lengths = Vec::with_capacity(parts);
1248 let mut span = None;
1249 for pending in held.drain(..) {
1250 let part = pending.chunk.len();
1251 rows = rows.checked_add(part).ok_or_else(|| invalid("row count overflow"))?;
1252 lengths.push(u32::try_from(part).map_err(|_| invalid("part row count overflow"))?);
1253 span = Some(
1254 span.map_or((pending.order, pending.order), |(first, _)| (first, pending.order)),
1255 );
1256 }
1257 self.order.push(span.ok_or_else(|| invalid("a stripe was flushed with no parts"))?);
1258 self.table.stripes.push(Stripe {
1259 rows,
1260 parts: lengths,
1261 index,
1262 pages,
1263 memberships,
1264 sieves,
1265 part_ranges,
1266 zone: Zone::from_ranges(ranges),
1267 });
1268 self.pending = held;
1270 Ok(())
1271 }
1272
1273 fn numeric_frequency(&self, column: usize) -> Result<Option<FrequencySummary>> {
1277 let ty = &self.table.fields[column].ty;
1278 if !matches!(
1279 ty,
1280 LogicalType::TinyInt
1281 | LogicalType::SmallInt
1282 | LogicalType::Integer
1283 | LogicalType::BigInt
1284 | LogicalType::UTinyInt
1285 | LogicalType::USmallInt
1286 | LogicalType::UInteger
1287 | LogicalType::UBigInt
1288 | LogicalType::Date
1289 | LogicalType::Timestamp
1290 ) {
1291 return Ok(None);
1292 }
1293 let mut candidates: HashMap<FrequencyValue, u32> = HashMap::new();
1294 let mut decrements = 0_u64;
1295 self.visit_numeric(column, |_, value| {
1296 if let Some(count) = candidates.get_mut(&value) {
1297 *count = count.saturating_add(1);
1298 } else if candidates.len() < FREQUENCY_CANDIDATES {
1299 candidates.insert(value, 1);
1300 } else {
1301 candidates.retain(|_, count| {
1302 *count -= 1;
1303 *count != 0
1304 });
1305 decrements = decrements.saturating_add(1);
1306 }
1307 })?;
1308 let (exact, ordinals) = if decrements == 0 {
1309 (
1310 candidates
1311 .into_iter()
1312 .map(|(value, count)| (value, u64::from(count)))
1313 .collect::<HashMap<_, _>>(),
1314 Vec::new(),
1315 )
1316 } else {
1317 let mut lower = candidates.values().copied().collect::<Vec<_>>();
1318 lower.sort_unstable_by(|left, right| right.cmp(left));
1319 if lower.len() < FREQUENCY_BUILD_RANK
1320 || u64::from(lower[FREQUENCY_BUILD_RANK - 1]) <= decrements
1321 {
1322 return Ok(None);
1323 }
1324 let mut exact =
1325 candidates.into_keys().map(|value| (value, 0_u64)).collect::<HashMap<_, _>>();
1326 let mut ordinals = Vec::new();
1327 let mut exceeded = false;
1328 self.visit_numeric(column, |ordinal, value| {
1329 if let Some(count) = exact.get_mut(&value) {
1330 *count = count.saturating_add(1);
1331 if !exceeded {
1332 if ordinals.len() < FREQUENCY_ORDINALS {
1333 ordinals.push(ordinal);
1334 } else {
1335 ordinals.clear();
1336 exceeded = true;
1337 }
1338 }
1339 }
1340 })?;
1341 (exact, ordinals)
1342 };
1343 let mut entries = exact
1344 .into_iter()
1345 .map(|(value, count)| FrequencyEntry { value, count })
1346 .collect::<Vec<_>>();
1347 entries.sort_unstable_by(|left, right| {
1348 right.count.cmp(&left.count).then_with(|| frequency_order(left.value, right.value))
1349 });
1350 let omitted_max =
1351 entries.get(FREQUENCY_ENTRIES).map_or(decrements, |entry| decrements.max(entry.count));
1352 entries.truncate(FREQUENCY_ENTRIES);
1353 Ok(Some(FrequencySummary { entries, omitted_max, ordinals }))
1354 }
1355
1356 fn visit_numeric(
1357 &self,
1358 column: usize,
1359 mut visit: impl FnMut(u64, FrequencyValue),
1360 ) -> Result<()> {
1361 let ty = &self.table.fields[column].ty;
1362 let mut start = 0_u64;
1363 for stripe in &self.table.stripes {
1364 let spans = read_index(&self.file, stripe, column)?;
1365 let page = stripe.pages[column];
1366 let mut bytes = vec![0; page.length as usize];
1367 read_at(&self.file, page.offset, &mut bytes)?;
1368 for (span, &rows) in spans.iter().zip(&stripe.parts) {
1369 let part = part_bytes(&bytes, *span)?;
1370 if checksum(part) != span.hash {
1371 return Err(invalid("column page checksum differs while building frequencies"));
1372 }
1373 let rows = rows as usize;
1374 let vector = decode(ty, rows, part, None)?;
1375 for row in 0..rows {
1377 let value = if vector.is_null_at(row) {
1378 FrequencyValue::Null
1379 } else {
1380 let widened = match vector.signed_at(row) {
1384 Some(value) => Some(value),
1385 None => match vector.value_at(row) {
1386 Value::UTinyInt(value) => Some(i128::from(value)),
1387 Value::USmallInt(value) => Some(i128::from(value)),
1388 Value::UInteger(value) => Some(i128::from(value)),
1389 Value::UBigInt(value) => Some(i128::from(value)),
1390 _ => None,
1391 },
1392 };
1393 FrequencyValue::Integer(widened.ok_or_else(|| {
1394 invalid("numeric frequency page did not contain an integer value")
1395 })?)
1396 };
1397 visit(start.saturating_add(row as u64), value);
1398 }
1399 start = start.saturating_add(rows as u64);
1400 }
1401 }
1402 Ok(())
1403 }
1404
1405 fn numeric_frequencies(&self) -> Result<Vec<Option<FrequencySummary>>> {
1413 let mut columns = self
1414 .table
1415 .fields
1416 .iter()
1417 .enumerate()
1418 .filter_map(|(column, field)| {
1419 matches!(
1420 field.ty,
1421 LogicalType::TinyInt
1422 | LogicalType::SmallInt
1423 | LogicalType::Integer
1424 | LogicalType::BigInt
1425 | LogicalType::UTinyInt
1426 | LogicalType::USmallInt
1427 | LogicalType::UInteger
1428 | LogicalType::UBigInt
1429 | LogicalType::Date
1430 | LogicalType::Timestamp
1431 )
1432 .then_some(column)
1433 })
1434 .collect::<Vec<_>>();
1435 let workers = std::thread::available_parallelism()
1436 .map_or(1, usize::from)
1437 .min(MAX_FREQUENCY_WORKERS)
1438 .min(columns.len());
1439 if workers <= 1 {
1440 let mut frequencies = vec![None; self.table.fields.len()];
1441 for column in columns {
1442 frequencies[column] = self.numeric_frequency(column)?;
1443 }
1444 return Ok(frequencies);
1445 }
1446 columns.sort_by_key(|&column| weight(&self.table.fields[column].ty));
1449 let queue = Mutex::new(columns);
1450 let pieces = std::thread::scope(|scope| {
1451 (0..workers)
1452 .map(|_| {
1453 scope.spawn(|| {
1454 let mut mine = Vec::new();
1455 loop {
1456 let taken = queue
1457 .lock()
1458 .map_err(|_| Error::internal("a native frequency worker panicked"))?
1459 .pop();
1460 let Some(column) = taken else { break };
1461 mine.push((column, self.numeric_frequency(column)?));
1462 }
1463 Ok(mine)
1464 })
1465 })
1466 .collect::<Vec<_>>()
1467 .into_iter()
1468 .map(|handle| {
1469 handle
1470 .join()
1471 .map_err(|_| Error::internal("a native frequency worker panicked"))?
1472 })
1473 .collect::<Result<Vec<_>>>()
1474 })?;
1475 let mut frequencies = vec![None; self.table.fields.len()];
1476 for piece in pieces {
1477 for (column, summary) in piece {
1478 frequencies[column] = summary;
1479 }
1480 }
1481 Ok(frequencies)
1482 }
1483
1484 fn close(&mut self) -> Result<Entry> {
1495 self.flush_pending()?;
1496 let mut stripes = std::mem::take(&mut self.order)
1497 .into_iter()
1498 .zip(std::mem::take(&mut self.table.stripes))
1499 .collect::<Vec<_>>();
1500 stripes.sort_by_key(|(order, _)| order.0);
1501 let mut previous: Option<(u64, u64)> = None;
1502 for ((first, last), _) in &stripes {
1503 if previous.is_some_and(|previous| previous >= *first) {
1504 return Err(invalid("chunks did not arrive in source order"));
1505 }
1506 previous = Some(*last);
1507 }
1508 self.table.stripes = stripes.into_iter().map(|(_, stripe)| stripe).collect();
1509 self.table.frequencies = self.numeric_frequencies()?;
1510 let dictionaries = std::mem::take(&mut self.dictionaries);
1511 let orders = rankings(&dictionaries)?;
1512 for (index, (dictionary, order)) in dictionaries.into_iter().zip(orders).enumerate() {
1513 let Some(dictionary) = dictionary else { continue };
1514 self.table.distincts[index] =
1518 Some(dictionary.counts.iter().filter(|count| **count != 0).count() as u64);
1519 self.table.frequencies[index] = Some(code_frequency(&dictionary));
1520 let encoded = encode_global_dictionary(dictionary, &order)?;
1521 let offset = self.at;
1522 self.put(&encoded.index)?;
1523 self.put(&encoded.ranks)?;
1524 for block in &encoded.payload {
1525 self.put(block)?;
1526 }
1527 let payload_len =
1528 encoded.payload.iter().try_fold(0_usize, |len, block| len.checked_add(block.len()));
1529 let length = payload_len
1530 .and_then(|len| len.checked_add(encoded.index.len()))
1531 .and_then(|len| len.checked_add(encoded.ranks.len()))
1532 .ok_or_else(|| invalid("dictionary page length overflow"))?;
1533 self.table.dictionaries[index] = Some(Page {
1534 offset,
1535 length: u32::try_from(length)
1536 .map_err(|_| invalid("dictionary page length overflow"))?,
1537 hash: checksum(&encoded.index),
1538 });
1539 }
1540 let directory = encode_directory(&self.table)?;
1541 if directory.len() > MAX_DIRECTORY {
1542 return Err(invalid("directory exceeds the configured bound"));
1543 }
1544 let offset = self.at;
1545 self.put(&directory)?;
1546 Ok(Entry {
1547 name: self.table.name.clone(),
1548 fields: self.table.fields.clone(),
1549 rows: self.table.rows,
1550 directory: Page {
1551 offset,
1552 length: u32::try_from(directory.len())
1553 .map_err(|_| invalid("directory length overflow"))?,
1554 hash: checksum(&directory),
1555 },
1556 })
1557 }
1558
1559 pub fn finish(mut self) -> Result<Table> {
1569 let entry = self.close()?;
1570 let mut tables = std::mem::take(&mut self.closed);
1571 tables.push(entry);
1572 let catalog = encode_catalog(&tables)?;
1573 if catalog.len() > MAX_DIRECTORY {
1574 return Err(invalid("catalog exceeds the configured bound"));
1575 }
1576 let offset = self.at;
1577 self.put(&catalog)?;
1578 self.file.sync_all().map_err(io)?;
1582 let slot = Slot {
1583 offset,
1584 length: u32::try_from(catalog.len()).map_err(|_| invalid("catalog length overflow"))?,
1585 generation: self.generation,
1586 hash: checksum(&catalog),
1587 };
1588 write_at(&self.file, slot_offset(self.generation), &slot.bytes())?;
1593 self.file.sync_all().map_err(io)?;
1594 Ok(self.table)
1595 }
1596}
1597
1598#[derive(Debug, Clone)]
1600pub struct Reader {
1601 file: Arc<File>,
1602 table: Arc<Table>,
1603 dictionaries: Arc<Vec<OnceLock<Arc<Vector>>>>,
1604 loading: Arc<Vec<Mutex<()>>>,
1613 opened: Arc<AtomicUsize>,
1617 sieves: Arc<Vec<Vec<SieveSlot>>>,
1621 part_ranges: Arc<Vec<Vec<RangeSlot>>>,
1624 places: Arc<Vec<Place>>,
1626 cache: Arc<Vec<Mutex<Cached>>>,
1627 pages: Arc<AtomicUsize>,
1630 indexes: Arc<AtomicUsize>,
1633 kept: Arc<AtomicUsize>,
1636 size: u64,
1638 directory: u64,
1640 opening: Opening,
1642}
1643
1644#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1656pub struct Opening {
1657 pub reads: u32,
1660 pub bytes: u64,
1662}
1663
1664#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1666pub struct Reads {
1667 pub opening: Opening,
1669 pub pages: usize,
1671 pub indexes: usize,
1673 pub dictionaries: usize,
1676}
1677
1678#[derive(Debug, Clone, Copy)]
1680struct Place {
1681 stripe: u32,
1682 part: u32,
1683 rows: u32,
1684}
1685
1686#[derive(Debug, Clone, Copy)]
1688struct PartSpan {
1689 start: usize,
1690 length: usize,
1691 hash: u64,
1692}
1693
1694#[derive(Debug, Clone)]
1700struct CachedColumn {
1701 stripe: usize,
1702 index: Arc<Vec<PartSpan>>,
1703 page: Option<Arc<Vec<u8>>>,
1704}
1705
1706#[derive(Debug, Default)]
1726struct Cached {
1727 pages: Vec<Option<Arc<Vec<u8>>>>,
1728 order: VecDeque<usize>,
1729 loading: Vec<usize>,
1730 index: Vec<Option<Arc<Vec<PartSpan>>>>,
1731}
1732
1733const CACHED_STRIPES_PER_COLUMN: usize = 4;
1745
1746type SieveSlot = OnceLock<Arc<Vec<Option<Sieve>>>>;
1748
1749type RangeSlot = OnceLock<Arc<Vec<Range>>>;
1750
1751#[derive(Debug)]
1752struct NativeText {
1753 file: Arc<File>,
1754 values: usize,
1756 offsets: Vec<u8>,
1765 offset_bits: usize,
1768 ranks: usize,
1770 rank_at: u64,
1774 rank_ends: Vec<u64>,
1778 rank_hashes: Vec<u64>,
1779 rank_blocks: Vec<OnceLock<Result<Vec<u8>>>>,
1780 code_bits: usize,
1783 code_ranks: OnceLock<Option<Vec<u32>>>,
1790 payload: u64,
1791 ends: Vec<u64>,
1794 hashes: Vec<u64>,
1795 blocks: Vec<OnceLock<Result<Vec<u8>>>>,
1797 keep_budget: usize,
1800 payload_kept: AtomicUsize,
1808 searched: Mutex<HashMap<Vec<u8>, (usize, bool)>>,
1825}
1826
1827const TEXT_SEARCH_MEMO: usize = 64;
1832
1833const TEXT_PAYLOAD_VALUES: usize = 1024;
1849
1850const TEXT_KEEP_BUDGET: usize = 256 * 1024 * 1024;
1871
1872const TEXT_OFFSET_RUN: usize = 512;
1879
1880const DICTIONARY_HEADER: usize = 16;
1883
1884const TEXT_RANK_BLOCK: usize = 512;
1895
1896const RANK_BLOCK_HEADER: usize = size_of::<u64>() + 1;
1910
1911impl NativeText {
1912 fn payload_block(&self, block: usize) -> Result<Option<&[u8]>> {
1919 let Some(slot) = self.blocks.get(block) else { return Ok(None) };
1920 let bytes = slot.get_or_init(|| self.decode_block(block)).as_ref().map_err(Clone::clone)?;
1921 Ok(Some(bytes.as_slice()))
1922 }
1923
1924 fn decode_block(&self, block: usize) -> Result<Vec<u8>> {
1929 let start = if block == 0 { 0 } else { self.ends[block - 1] };
1930 let end = self.ends[block];
1931 let len = end
1932 .checked_sub(start)
1933 .ok_or_else(|| invalid("global dictionary block ends before it starts"))?;
1934 let mut stored = vec![
1935 0;
1936 usize::try_from(len).map_err(|_| invalid(
1937 "global dictionary block does not fit in memory"
1938 ))?
1939 ];
1940 read_at(&self.file, self.payload + start, &mut stored)?;
1941 if checksum(&stored) != self.hashes[block] {
1942 return Err(invalid("global dictionary payload checksum differs"));
1943 }
1944 let first = block * TEXT_PAYLOAD_VALUES;
1945 let last = (first + TEXT_PAYLOAD_VALUES).min(self.values);
1946 let want = self.end_within(last - 1)? as usize;
1947 let values = string::decode_flat(&stored)?;
1948 if values.len() != last - first {
1949 return Err(invalid("global dictionary block holds the wrong value count"));
1950 }
1951 let bytes = values.into_bytes();
1952 if bytes.len() != want {
1953 return Err(invalid("global dictionary block decodes to the wrong length"));
1954 }
1955 Ok(bytes)
1956 }
1957
1958 fn end_within(&self, index: usize) -> Result<u32> {
1960 let run = index / TEXT_OFFSET_RUN;
1961 let bytes = self
1962 .offsets
1963 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
1964 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1965 let end = bitpack::tail_at(bytes, self.offset_bits, index % TEXT_OFFSET_RUN)
1966 .map_err(|_| invalid("global dictionary offsets are short"))?;
1967 u32::try_from(end).map_err(|_| invalid("global dictionary offset is past the payload"))
1968 }
1969
1970 fn ends_within(&self, first: usize, last: usize) -> Result<Vec<u64>> {
1983 let mut ends = Vec::with_capacity(last.saturating_sub(first));
1984 let mut at = first;
1985 while at < last {
1986 let run = at / TEXT_OFFSET_RUN;
1987 let stop = ((run + 1) * TEXT_OFFSET_RUN).min(last);
1988 let held = self.values.saturating_sub(run * TEXT_OFFSET_RUN).min(TEXT_OFFSET_RUN);
1989 let bytes = self
1990 .offsets
1991 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
1992 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1993 let run_ends = bitpack::unpack_tail(bytes, self.offset_bits, held)
1994 .map_err(|_| invalid("global dictionary offsets are short"))?;
1995 let within = run_ends
1996 .get(at % TEXT_OFFSET_RUN..stop - run * TEXT_OFFSET_RUN)
1997 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1998 ends.extend_from_slice(within);
1999 at = stop;
2000 }
2001 Ok(ends)
2002 }
2003
2004 fn start_within(&self, index: usize) -> Result<u32> {
2007 if index % TEXT_PAYLOAD_VALUES == 0 { Ok(0) } else { self.end_within(index - 1) }
2008 }
2009
2010 fn span_within(&self, index: usize) -> Result<(u32, u32)> {
2018 let within = index % TEXT_OFFSET_RUN;
2019 let (start, end) = if within == 0 {
2020 (self.start_within(index)?, self.end_within(index)?)
2021 } else {
2022 let run = index / TEXT_OFFSET_RUN;
2023 let bytes = self
2024 .offsets
2025 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
2026 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
2027 let (start, end) = bitpack::tail_pair(bytes, self.offset_bits, within)
2028 .map_err(|_| invalid("global dictionary offsets are short"))?;
2029 let ends = u32::try_from(end)
2030 .map_err(|_| invalid("global dictionary offset is past the payload"))?;
2031 let starts = u32::try_from(start)
2032 .map_err(|_| invalid("global dictionary offset is past the payload"))?;
2033 (starts, ends)
2034 };
2035 if start > end {
2036 return Err(invalid("global dictionary value ends before it starts"));
2037 }
2038 Ok((start, end))
2039 }
2040
2041 fn rank_parts(&self, rank: usize) -> Result<(&[u8], usize)> {
2048 let slot = self
2049 .rank_blocks
2050 .get(rank / TEXT_RANK_BLOCK)
2051 .ok_or_else(|| invalid("global dictionary rank is past the order"))?;
2052 let block = slot
2053 .get_or_init(|| {
2054 let which = rank / TEXT_RANK_BLOCK;
2055 let start = if which == 0 { 0 } else { self.rank_ends[which - 1] };
2056 let end = self.rank_ends[which];
2057 let mut bytes = vec![0; (end - start) as usize];
2058 read_at(&self.file, self.rank_at + start, &mut bytes)?;
2059 if checksum(&bytes)
2060 != *self
2061 .rank_hashes
2062 .get(rank / TEXT_RANK_BLOCK)
2063 .ok_or_else(|| invalid("global dictionary rank block has no checksum"))?
2064 {
2065 return Err(invalid("global dictionary rank checksum differs"));
2066 }
2067 Ok(bytes)
2068 })
2069 .as_ref()
2070 .map_err(Clone::clone)?;
2071 Ok((block.as_slice(), rank % TEXT_RANK_BLOCK))
2072 }
2073
2074 fn head_at(&self, rank: usize) -> Result<u64> {
2076 let (block, within) = self.rank_parts(rank)?;
2077 let (base, width, packed) = rank_heads(block)?;
2078 let above = bitpack::tail_at(packed, width, within)
2079 .map_err(|_| invalid("global dictionary rank block is short of heads"))?;
2080 Ok(base.wrapping_add(above))
2081 }
2082
2083 fn rank_codes<'block>(&self, block: &'block [u8], count: usize) -> Result<&'block [u8]> {
2085 let (_, width, packed) = rank_heads(block)?;
2086 packed
2087 .get(bitpack::tail_len(count, width)..)
2088 .ok_or_else(|| invalid("global dictionary rank block is short of codes"))
2089 }
2090
2091 fn rank_block_len(&self, rank: usize) -> usize {
2093 let first = rank / TEXT_RANK_BLOCK * TEXT_RANK_BLOCK;
2094 TEXT_RANK_BLOCK.min(self.ranks - first)
2095 }
2096}
2097
2098fn rank_heads(block: &[u8]) -> Result<(u64, usize, &[u8])> {
2100 let header = block
2101 .get(..RANK_BLOCK_HEADER)
2102 .ok_or_else(|| invalid("global dictionary rank block is short"))?;
2103 let base = u64::from_le_bytes(header[..8].try_into().expect("eight bytes"));
2104 let width = header[8] as usize;
2105 if width > 64 {
2106 return Err(invalid("global dictionary rank block packs heads past a word"));
2107 }
2108 Ok((base, width, &block[RANK_BLOCK_HEADER..]))
2109}
2110
2111fn offset_width(offsets: &[u32]) -> usize {
2118 let values = offsets.len() - 1;
2119 let mut span = 0;
2120 for first in (0..values).step_by(TEXT_PAYLOAD_VALUES) {
2121 let last = (first + TEXT_PAYLOAD_VALUES).min(values);
2122 span = span.max(offsets[last] - offsets[first]);
2123 }
2124 (u32::BITS - span.leading_zeros()) as usize
2125}
2126
2127fn offset_bytes(values: usize, bits: usize) -> usize {
2130 let full = values / TEXT_OFFSET_RUN;
2131 let rest = values % TEXT_OFFSET_RUN;
2132 full * TEXT_OFFSET_RUN / 8 * bits + bitpack::tail_len(rest, bits)
2133}
2134
2135fn encode_offsets(offsets: &[u32], bits: usize, out: &mut Vec<u8>) -> Result<()> {
2137 let values = offsets.len() - 1;
2138 let mut run = Vec::with_capacity(TEXT_OFFSET_RUN);
2139 for first in (0..values).step_by(TEXT_OFFSET_RUN) {
2140 let last = (first + TEXT_OFFSET_RUN).min(values);
2141 let base = offsets[first / TEXT_PAYLOAD_VALUES * TEXT_PAYLOAD_VALUES];
2142 run.clear();
2143 run.extend((first..last).map(|value| u64::from(offsets[value + 1] - base)));
2144 bitpack::pack_tail(&run, bits, out)
2145 .map_err(|_| invalid("global dictionary offsets do not pack"))?;
2146 }
2147 Ok(())
2148}
2149
2150fn code_width(values: usize) -> usize {
2152 match u64::try_from(values).unwrap_or(u64::MAX) {
2153 0 | 1 => 0,
2154 last => (u64::BITS - (last - 1).leading_zeros()) as usize,
2155 }
2156}
2157
2158impl TextSource for NativeText {
2159 fn len(&self) -> usize {
2160 self.values
2161 }
2162
2163 fn bytes_at(&self, index: usize) -> Result<Option<&[u8]>> {
2164 if index >= self.values {
2165 return Ok(None);
2166 }
2167 let (start, end) = self.span_within(index)?;
2168 if start == end {
2169 return Ok(Some(&[]));
2170 }
2171 let block = index / TEXT_PAYLOAD_VALUES;
2174 let Some(bytes) = self.payload_block(block)? else { return Ok(None) };
2175 Ok(bytes.get(start as usize..end as usize))
2176 }
2177
2178 fn bytes_len_at(&self, index: usize) -> Result<Option<usize>> {
2179 if index >= self.values {
2180 return Ok(None);
2181 }
2182 let (start, end) = self.span_within(index)?;
2183 Ok(Some((end - start) as usize))
2184 }
2185
2186 fn sweep(
2199 &self,
2200 first: usize,
2201 limit: usize,
2202 body: &mut dyn FnMut(usize, &[u8]) -> Result<()>,
2203 ) -> Result<usize> {
2204 let limit = limit.min(self.values);
2205 if first >= limit {
2206 return Ok(first);
2207 }
2208 let block = first / TEXT_PAYLOAD_VALUES;
2209 let last = ((block + 1) * TEXT_PAYLOAD_VALUES).min(limit);
2210 let decoded;
2211 let bytes: &[u8] = match self.blocks.get(block).and_then(OnceLock::get) {
2212 Some(Ok(kept)) => kept,
2213 _ if self.payload_kept.load(Atomic::Relaxed) < self.keep_budget => {
2214 let kept = self
2215 .payload_block(block)?
2216 .ok_or_else(|| invalid("global dictionary block is past the payload"))?;
2217 self.payload_kept.fetch_add(kept.len(), Atomic::Relaxed);
2218 kept
2219 }
2220 _ => {
2221 decoded = self.decode_block(block)?;
2222 &decoded
2223 }
2224 };
2225 let ends = self.ends_within(first, last)?;
2226 if ends.len() != last - first {
2227 return Err(invalid("global dictionary offsets are short"));
2228 }
2229 let mut start = u64::from(self.start_within(first)?);
2230 for (index, &end) in (first..last).zip(&ends) {
2233 let value = usize::try_from(start)
2234 .ok()
2235 .zip(usize::try_from(end).ok())
2236 .and_then(|(from, to)| bytes.get(from..to))
2237 .ok_or_else(|| invalid("global dictionary value is past its block"))?;
2238 body(index, value)?;
2239 start = end;
2240 }
2241 Ok(last)
2242 }
2243
2244 fn ranks(&self) -> Option<usize> {
2245 (self.ranks > 0).then_some(self.ranks)
2246 }
2247
2248 fn below(&self, ranks: usize, wanted: &[u8]) -> Result<(usize, bool)> {
2256 let mut memo = self.searched.lock().map_err(|_| invalid("a poisoned dictionary search"))?;
2257 if let Some(&answer) = memo.get(wanted) {
2258 return Ok(answer);
2259 }
2260 let answer = search_below(self, ranks, wanted)?;
2261 if memo.len() >= TEXT_SEARCH_MEMO {
2262 memo.clear();
2263 }
2264 memo.insert(wanted.to_vec(), answer);
2265 Ok(answer)
2266 }
2267
2268 fn compare_rank(&self, rank: usize, wanted: &[u8]) -> Result<Ordering> {
2269 let settled = self.head_at(rank)?.cmp(&head(wanted));
2273 if settled != Ordering::Equal {
2274 return Ok(settled);
2275 }
2276 let code = self.code_at_rank(rank)?;
2277 let bytes = self
2278 .bytes_at(code as usize)?
2279 .ok_or_else(|| invalid("global dictionary order names a code it does not have"))?;
2280 Ok(bytes.cmp(wanted))
2281 }
2282
2283 fn code_at_rank(&self, rank: usize) -> Result<u32> {
2284 let (block, within) = self.rank_parts(rank)?;
2285 let codes = self.rank_codes(block, self.rank_block_len(rank))?;
2286 let code = bitpack::tail_at(codes, self.code_bits, within)
2287 .map_err(|_| invalid("global dictionary rank block is short of codes"))?;
2288 let code = u32::try_from(code)
2289 .map_err(|_| invalid("global dictionary order names a code it does not have"))?;
2290 if code as usize >= self.len() {
2291 return Err(invalid("global dictionary order names a code it does not have"));
2292 }
2293 Ok(code)
2294 }
2295
2296 fn code_ranks(&self) -> Option<&[u32]> {
2297 if self.ranks == 0 || self.ranks != self.len() {
2301 return None;
2302 }
2303 self.code_ranks
2304 .get_or_init(|| {
2305 let mut ranks = vec![u32::MAX; self.ranks];
2306 for first in (0..self.ranks).step_by(TEXT_RANK_BLOCK) {
2309 let (block, _) = self.rank_parts(first).ok()?;
2310 let count = self.rank_block_len(first);
2311 let codes = self.rank_codes(block, count).ok()?;
2312 for (within, code) in bitpack::unpack_tail(codes, self.code_bits, count)
2313 .ok()?
2314 .into_iter()
2315 .enumerate()
2316 {
2317 let code = usize::try_from(code).ok()?;
2318 *ranks.get_mut(code)? = u32::try_from(first + within).ok()?;
2319 }
2320 }
2321 if ranks.contains(&u32::MAX) {
2322 return None;
2323 }
2324 Some(ranks)
2325 })
2326 .as_deref()
2327 }
2328
2329 fn footprint(&self) -> usize {
2330 self.offsets.capacity()
2331 + self
2332 .code_ranks
2333 .get()
2334 .and_then(Option::as_ref)
2335 .map_or(0, |ranks| ranks.capacity() * size_of::<u32>())
2336 + self.rank_hashes.capacity() * size_of::<u64>()
2337 + self.rank_ends.capacity() * size_of::<u64>()
2338 + self.rank_blocks.capacity() * size_of::<OnceLock<Result<Vec<u8>>>>()
2339 + self
2340 .rank_blocks
2341 .iter()
2342 .filter_map(OnceLock::get)
2343 .filter_map(|result| result.as_ref().ok())
2344 .map(Vec::capacity)
2345 .sum::<usize>()
2346 + self.blocks.capacity() * size_of::<OnceLock<Result<Vec<u8>>>>()
2347 + self.hashes.capacity() * size_of::<u64>()
2348 + self.ends.capacity() * size_of::<u64>()
2349 + self
2350 .blocks
2351 .iter()
2352 .filter_map(OnceLock::get)
2353 .filter_map(|result| result.as_ref().ok())
2354 .map(Vec::capacity)
2355 .sum::<usize>()
2356 }
2357}
2358
2359fn places(table: &Table) -> Result<Vec<Place>> {
2361 let mut places = Vec::with_capacity(table.stripes.len().saturating_mul(STRIPE_PARTS));
2362 for (at, stripe) in table.stripes.iter().enumerate() {
2363 let index = u32::try_from(at).map_err(|_| invalid("too many stripes"))?;
2364 for (part, &rows) in stripe.parts.iter().enumerate() {
2365 places.push(Place {
2366 stripe: index,
2367 part: u32::try_from(part).map_err(|_| invalid("too many parts in a stripe"))?,
2368 rows,
2369 });
2370 }
2371 }
2372 Ok(places)
2373}
2374
2375fn read_index(file: &File, stripe: &Stripe, column: usize) -> Result<Vec<PartSpan>> {
2380 let parts = stripe.parts.len();
2381 let section = index_section(parts)?;
2382 let at = column.checked_mul(section).ok_or_else(|| invalid("index page offset overflow"))?;
2383 let end = at.checked_add(section).ok_or_else(|| invalid("index page offset overflow"))?;
2384 if end > stripe.index.length as usize {
2385 return Err(invalid("index page is shorter than its columns"));
2386 }
2387 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
2388 let mut bytes = vec![0; section];
2389 let offset = stripe
2390 .index
2391 .offset
2392 .checked_add(at as u64)
2393 .ok_or_else(|| invalid("index page offset overflow"))?;
2394 read_at(file, offset, &mut bytes)?;
2395 let entries = section - size_of::<u64>();
2396 let stored = u64::from_le_bytes(bytes[entries..].try_into().expect("eight bytes"));
2397 if checksum(&bytes[..entries]) != stored {
2398 return Err(invalid(&format!(
2401 "index page section checksum differs, column {column} of {parts} parts at {offset}, \
2402 wanted {stored:016x} and got {:016x}",
2403 checksum(&bytes[..entries]),
2404 )));
2405 }
2406 let mut spans = Vec::with_capacity(parts);
2407 let mut start = 0_usize;
2408 for part in 0..parts {
2409 let at = part * INDEX_ENTRY;
2410 let length = u32::from_le_bytes(bytes[at..at + 4].try_into().expect("four bytes")) as usize;
2411 let hash = u64::from_le_bytes(bytes[at + 4..at + 12].try_into().expect("eight bytes"));
2412 spans.push(PartSpan { start, length, hash });
2413 start = start.checked_add(length).ok_or_else(|| invalid("column page length overflow"))?;
2414 }
2415 if start != page.length as usize {
2416 return Err(invalid("column page length differs from its index"));
2417 }
2418 Ok(spans)
2419}
2420
2421fn part_bytes(page: &[u8], span: PartSpan) -> Result<&[u8]> {
2423 let end = span.start.checked_add(span.length).ok_or_else(|| invalid("part range overflow"))?;
2424 page.get(span.start..end).ok_or_else(|| invalid("part exceeds its column page"))
2425}
2426
2427fn remember(cached: &mut Cached, held: &CachedColumn, kept: usize) {
2432 if let Some(slot) = cached.index.get_mut(held.stripe) {
2433 if slot.is_none() {
2434 *slot = Some(Arc::clone(&held.index));
2435 }
2436 }
2437 let Some(page) = held.page.clone() else { return };
2438 let Some(slot) = cached.pages.get_mut(held.stripe) else { return };
2439 if slot.is_none() {
2440 cached.order.push_back(held.stripe);
2441 }
2442 *slot = Some(page);
2443 while cached.order.len() > kept.max(1) {
2444 let Some(oldest) = cached.order.pop_front() else { break };
2445 if let Some(slot) = cached.pages.get_mut(oldest) {
2446 *slot = None;
2447 }
2448 }
2449}
2450
2451#[derive(Debug, Clone)]
2460pub struct Catalog {
2461 file: Arc<File>,
2462 size: u64,
2463 entries: Arc<Vec<Entry>>,
2464 opening: Opening,
2465}
2466
2467impl Catalog {
2468 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
2474 let (file, size, _, bytes, opening) = slot_bytes(path)?;
2475 let entries = decode_catalog(&bytes, size)?;
2476 Ok(Self { file: Arc::new(file), size, entries: Arc::new(entries), opening })
2477 }
2478
2479 pub fn names(&self) -> impl ExactSizeIterator<Item = &str> {
2481 self.entries.iter().map(|entry| entry.name.as_str())
2482 }
2483
2484 #[must_use]
2486 pub fn len(&self) -> usize {
2487 self.entries.len()
2488 }
2489
2490 #[must_use]
2492 pub fn is_empty(&self) -> bool {
2493 self.entries.is_empty()
2494 }
2495
2496 pub fn table(&self, name: &str) -> Result<Reader> {
2502 let entry = self
2503 .entries
2504 .iter()
2505 .find(|entry| entry.name == name)
2506 .ok_or_else(|| invalid(&format!("the file holds no table called {name}")))?;
2507 let mut bytes = vec![0; entry.directory.length as usize];
2508 read_at(&self.file, entry.directory.offset, &mut bytes)?;
2509 if checksum(&bytes) != entry.directory.hash {
2510 return Err(invalid(&format!("the directory of table {name} does not checksum")));
2511 }
2512 let mut opening = self.opening;
2513 opening.reads += 1;
2514 opening.bytes += u64::from(entry.directory.length);
2515 Reader::build(
2516 Arc::clone(&self.file),
2517 self.size,
2518 decode_directory(&bytes, self.size)?,
2519 u64::from(entry.directory.length),
2520 opening,
2521 )
2522 }
2523}
2524
2525fn slot_offset(generation: u64) -> u64 {
2530 16 + (generation - 1) % 2 * SLOT_BYTES as u64
2531}
2532
2533fn slot_bytes(path: impl AsRef<Path>) -> Result<(File, u64, Slot, Vec<u8>, Opening)> {
2538 let mut file = File::open(path).map_err(io)?;
2539 let size = file.metadata().map_err(io)?.len();
2540 if size < HEADER {
2541 return Err(invalid("file is shorter than its header"));
2542 }
2543 let mut header = [0; HEADER as usize];
2544 file.read_exact(&mut header).map_err(io)?;
2545 let mut opening = Opening { reads: 1, bytes: HEADER };
2546 let version = u32::from_le_bytes([header[8], header[9], header[10], header[11]]);
2547 if &header[..8] != MAGIC {
2552 return Err(invalid("the header does not begin with a rudb native magic"));
2553 }
2554 if version != FORMAT {
2555 return Err(invalid(&format!(
2556 "the file is format {version} and this build reads format {FORMAT}, so it has to \
2557 be written again"
2558 )));
2559 }
2560 let mut selected = None;
2561 for start in [16, 16 + SLOT_BYTES] {
2562 let slot = Slot::read(&header[start..start + SLOT_BYTES]);
2563 if slot.generation == 0 || slot.length == 0 || slot.length as usize > MAX_DIRECTORY {
2564 continue;
2565 }
2566 let Some(end) = slot.offset.checked_add(u64::from(slot.length)) else { continue };
2567 if slot.offset < HEADER || end > size {
2568 continue;
2569 }
2570 let mut bytes = vec![0; slot.length as usize];
2571 file.seek(SeekFrom::Start(slot.offset)).map_err(io)?;
2572 file.read_exact(&mut bytes).map_err(io)?;
2573 opening.reads += 1;
2574 opening.bytes += u64::from(slot.length);
2575 if checksum(&bytes) == slot.hash
2576 && selected
2577 .as_ref()
2578 .is_none_or(|(old, _): &(Slot, Vec<u8>)| old.generation < slot.generation)
2579 {
2580 selected = Some((slot, bytes));
2581 }
2582 }
2583 let (slot, bytes) = selected.ok_or_else(|| invalid("no committed directory slot is valid"))?;
2584 Ok((file, size, slot, bytes, opening))
2585}
2586
2587impl Reader {
2588 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
2595 let catalog = Catalog::open(path)?;
2596 let mut names = catalog.names();
2597 let name = names.next().ok_or_else(|| invalid("the file holds no table"))?.to_string();
2598 if names.next().is_some() {
2599 return Err(invalid(
2600 "the file holds more than one table, so it has to be opened by name",
2601 ));
2602 }
2603 catalog.table(&name)
2604 }
2605
2606 fn build(
2608 file: Arc<File>,
2609 size: u64,
2610 table: Table,
2611 directory: u64,
2612 opening: Opening,
2613 ) -> Result<Self> {
2614 let places = places(&table)?;
2615 let dictionaries = (0..table.fields.len()).map(|_| OnceLock::new()).collect();
2616 let table_fields = table.fields.len();
2617 let stripes = table.stripes.len();
2618 let cache = (0..table.fields.len())
2619 .map(|_| {
2620 Mutex::new(Cached {
2621 pages: (0..stripes).map(|_| None).collect(),
2622 index: (0..stripes).map(|_| None).collect(),
2623 ..Cached::default()
2624 })
2625 })
2626 .collect::<Vec<_>>();
2627 let sieves: Vec<Vec<SieveSlot>> = (0..table.fields.len())
2628 .map(|_| table.stripes.iter().map(|_| OnceLock::new()).collect())
2629 .collect();
2630 let part_ranges: Vec<Vec<RangeSlot>> = (0..table.fields.len())
2631 .map(|_| table.stripes.iter().map(|_| OnceLock::new()).collect())
2632 .collect();
2633 Ok(Self {
2634 file,
2635 table: Arc::new(table),
2636 dictionaries: Arc::new(dictionaries),
2637 loading: Arc::new((0..table_fields).map(|_| Mutex::new(())).collect()),
2638 opened: Arc::new(AtomicUsize::new(0)),
2639 sieves: Arc::new(sieves),
2640 part_ranges: Arc::new(part_ranges),
2641 places: Arc::new(places),
2642 cache: Arc::new(cache),
2643 pages: Arc::new(AtomicUsize::new(0)),
2644 indexes: Arc::new(AtomicUsize::new(0)),
2645 kept: Arc::new(AtomicUsize::new(CACHED_STRIPES_PER_COLUMN)),
2646 size,
2647 directory,
2648 opening,
2649 })
2650 }
2651
2652 #[must_use]
2659 pub fn reads(&self) -> Reads {
2660 Reads {
2661 opening: self.opening,
2662 pages: self.pages.load(Atomic::Relaxed),
2663 indexes: self.indexes.load(Atomic::Relaxed),
2664 dictionaries: self.opened.load(Atomic::Relaxed),
2665 }
2666 }
2667
2668 #[must_use]
2673 pub fn layout(&self) -> Layout {
2674 let table = &self.table;
2675 let stripes = table.stripes.as_slice();
2676 let columns = table
2677 .fields
2678 .iter()
2679 .enumerate()
2680 .map(|(at, field)| ColumnLayout {
2681 name: field.name.clone(),
2682 kind: field.ty.to_string(),
2683 pages: sum(stripes.iter().map(|stripe| span_bytes(&stripe.pages, at))),
2684 memberships: sum(stripes.iter().map(|stripe| page_bytes(&stripe.memberships, at))),
2685 sieves: sum(stripes.iter().map(|stripe| page_bytes(&stripe.sieves, at))),
2686 part_ranges: sum(stripes.iter().map(|stripe| page_bytes(&stripe.part_ranges, at))),
2687 dictionary: page_bytes(&table.dictionaries, at),
2688 })
2689 .collect();
2690 Layout {
2691 file: self.size,
2692 rows: table.rows,
2693 stripes: stripes.len(),
2694 parts: self.places.len(),
2695 columns,
2696 indexes: sum(stripes.iter().map(|stripe| u64::from(stripe.index.length))),
2697 directory: self.directory,
2698 header: HEADER,
2699 }
2700 }
2701
2702 #[must_use]
2704 pub fn parts(&self) -> usize {
2705 self.places.len()
2706 }
2707
2708 #[must_use]
2715 pub fn stripe_parts(&self) -> Vec<std::ops::Range<usize>> {
2716 let mut runs = Vec::with_capacity(self.table.stripes.len());
2717 let mut start = 0;
2718 for stripe in &self.table.stripes {
2719 let end = start + stripe.parts.len();
2720 runs.push(start..end);
2721 start = end;
2722 }
2723 runs
2724 }
2725
2726 #[must_use]
2731 pub fn stripe_rows(&self, stripe: usize) -> usize {
2732 self.table.stripes.get(stripe).map_or(0, |held| held.rows)
2733 }
2734
2735 pub fn keep_stripes(&self, stripes: usize) {
2742 self.kept.fetch_max(stripes, Atomic::Relaxed);
2743 }
2744
2745 #[must_use]
2747 pub fn part_rows(&self, at: usize) -> usize {
2748 self.places.get(at).map_or(0, |place| place.rows as usize)
2749 }
2750
2751 #[must_use]
2753 pub fn table(&self) -> &Table {
2754 &self.table
2755 }
2756
2757 pub fn top_frequencies(&self, column: usize, top: usize) -> Result<Option<Vec<(Value, u64)>>> {
2766 let field = self
2767 .table
2768 .fields
2769 .get(column)
2770 .ok_or_else(|| invalid("frequency column index out of range"))?;
2771 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2772 return Ok(None);
2773 };
2774 if top == 0 || summary.entries.len() < top {
2775 return Ok(None);
2776 }
2777 let boundary = summary.entries[top - 1].count;
2778 if boundary <= summary.omitted_max {
2779 return Ok(None);
2780 }
2781 self.decode_frequencies(column, &field.ty, &summary.entries).map(Some)
2782 }
2783
2784 pub fn exact_frequencies(&self, column: usize) -> Result<Option<Vec<(Value, u64)>>> {
2804 let field = self
2805 .table
2806 .fields
2807 .get(column)
2808 .ok_or_else(|| invalid("frequency column index out of range"))?;
2809 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2810 return Ok(None);
2811 };
2812 if summary.omitted_max > 0 {
2813 return Ok(None);
2814 }
2815 self.decode_frequencies(column, &field.ty, &summary.entries).map(Some)
2816 }
2817
2818 fn decode_frequencies(
2820 &self,
2821 column: usize,
2822 ty: &LogicalType,
2823 entries: &[FrequencyEntry],
2824 ) -> Result<Vec<(Value, u64)>> {
2825 let dictionary = if *ty == LogicalType::Varchar { self.dictionary(column)? } else { None };
2826 let mut out = Vec::with_capacity(entries.len());
2827 for entry in entries {
2828 let value = match entry.value {
2829 FrequencyValue::Null => Value::Null,
2830 FrequencyValue::Integer(value) => match *ty {
2831 LogicalType::TinyInt => Value::TinyInt(
2832 i8::try_from(value)
2833 .map_err(|_| invalid("frequency TINYINT is out of range"))?,
2834 ),
2835 LogicalType::UTinyInt => Value::UTinyInt(
2836 u8::try_from(value)
2837 .map_err(|_| invalid("frequency UTINYINT is out of range"))?,
2838 ),
2839 LogicalType::USmallInt => Value::USmallInt(
2840 u16::try_from(value)
2841 .map_err(|_| invalid("frequency USMALLINT is out of range"))?,
2842 ),
2843 LogicalType::UInteger => Value::UInteger(
2844 u32::try_from(value)
2845 .map_err(|_| invalid("frequency UINTEGER is out of range"))?,
2846 ),
2847 LogicalType::UBigInt => Value::UBigInt(
2848 u64::try_from(value)
2849 .map_err(|_| invalid("frequency UBIGINT is out of range"))?,
2850 ),
2851 LogicalType::SmallInt => Value::SmallInt(
2852 i16::try_from(value)
2853 .map_err(|_| invalid("frequency SMALLINT is out of range"))?,
2854 ),
2855 LogicalType::Integer => Value::Integer(
2856 i32::try_from(value)
2857 .map_err(|_| invalid("frequency INTEGER is out of range"))?,
2858 ),
2859 LogicalType::BigInt => Value::BigInt(
2860 i64::try_from(value)
2861 .map_err(|_| invalid("frequency BIGINT is out of range"))?,
2862 ),
2863 LogicalType::Date => Value::Date(
2864 i32::try_from(value)
2865 .map_err(|_| invalid("frequency DATE is out of range"))?,
2866 ),
2867 LogicalType::Timestamp => Value::Timestamp(
2868 i64::try_from(value)
2869 .map_err(|_| invalid("frequency TIMESTAMP is out of range"))?,
2870 ),
2871 _ => return Err(invalid("integer frequency belongs to another type")),
2872 },
2873 FrequencyValue::Code(code) => dictionary
2874 .as_ref()
2875 .ok_or_else(|| invalid("frequency code has no dictionary"))?
2876 .try_value_at(code as usize)?,
2877 };
2878 out.push((value, entry.count));
2879 }
2880 Ok(out)
2881 }
2882
2883 pub fn frequency_occurrences(&self, column: usize) -> Result<Option<FrequencyOccurrences>> {
2893 self.table
2894 .fields
2895 .get(column)
2896 .ok_or_else(|| invalid("frequency column index out of range"))?;
2897 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2898 return Ok(None);
2899 };
2900 if summary.ordinals.is_empty() {
2901 return Ok(None);
2902 }
2903 Ok(Some(FrequencyOccurrences {
2904 omitted_max: summary.omitted_max,
2905 ordinals: summary.ordinals.clone(),
2906 }))
2907 }
2908
2909 pub fn distinct_values(&self, column: usize) -> Result<Option<u64>> {
2933 self.table
2934 .distincts
2935 .get(column)
2936 .copied()
2937 .ok_or_else(|| invalid("distinct column index out of range"))
2938 }
2939
2940 pub fn null_count(&self, column: usize) -> Result<u64> {
2951 if column >= self.table.fields.len() {
2952 return Err(invalid("null count column index out of range"));
2953 }
2954 let mut nulls = 0_u64;
2955 for stripe in &self.table.stripes {
2956 let range = stripe
2957 .zone
2958 .column(column)
2959 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
2960 nulls = nulls
2961 .checked_add(range.nulls as u64)
2962 .ok_or_else(|| invalid("null count overflow"))?;
2963 }
2964 Ok(nulls)
2965 }
2966
2967 pub fn text_extremes(&self, column: usize) -> Result<Option<(Value, Value)>> {
2982 if self.null_count(column)? > 0 {
2983 return Ok(None);
2984 }
2985 let Some(dictionary) = self.dictionary(column)? else { return Ok(None) };
2986 let Some(ranks) = dictionary.ranks() else { return Ok(None) };
2987 if ranks == 0 {
2988 return Ok(None);
2989 }
2990 let low = text_at_rank(&dictionary, 0)?;
2991 let high = text_at_rank(&dictionary, ranks - 1)?;
2992 Ok(Some((low, high)))
2993 }
2994
2995 pub fn exact_extremes(&self, column: usize) -> Result<Option<(Bound, Bound)>> {
3018 if column >= self.table.fields.len() {
3019 return Err(invalid("extremes column index out of range"));
3020 }
3021 let mut low: Option<Bound> = None;
3022 let mut high: Option<Bound> = None;
3023 for stripe in &self.table.stripes {
3024 let range = stripe
3025 .zone
3026 .column(column)
3027 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
3028 if !range.exact {
3029 return Ok(None);
3030 }
3031 let (Some(small), Some(large)) = (range.low.as_ref(), range.high.as_ref()) else {
3036 if stripe.rows > range.nulls {
3037 return Ok(None);
3038 }
3039 continue;
3040 };
3041 low = Some(low.map_or_else(|| small.clone(), |held| held.smaller(small.clone())));
3042 high = Some(high.map_or_else(|| large.clone(), |held| held.larger(large.clone())));
3043 }
3044 Ok(low.zip(high))
3045 }
3046
3047 pub fn exact_sum(&self, column: usize) -> Result<Option<(i128, u64)>> {
3060 if column >= self.table.fields.len() {
3061 return Err(invalid("sum column index out of range"));
3062 }
3063 let mut total = 0_i128;
3064 let mut rows = 0_u64;
3065 for stripe in &self.table.stripes {
3066 let range = stripe
3067 .zone
3068 .column(column)
3069 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
3070 let Some(part) = range.sum else { return Ok(None) };
3071 let Some(sum) = total.checked_add(part) else { return Ok(None) };
3072 total = sum;
3073 rows = rows.saturating_add(stripe.rows as u64 - range.nulls as u64);
3074 }
3075 Ok(Some((total, rows)))
3076 }
3077
3078 fn dictionary(&self, column: usize) -> Result<Option<Arc<Vector>>> {
3087 let Some(page) = self.table.dictionaries[column] else { return Ok(None) };
3088 if let Some(dictionary) = self.dictionaries[column].get() {
3089 return Ok(Some(Arc::clone(dictionary)));
3090 }
3091 let _queued = self.loading[column].lock().map_err(|_| invalid("a poisoned dictionary"))?;
3092 if let Some(dictionary) = self.dictionaries[column].get() {
3093 return Ok(Some(Arc::clone(dictionary)));
3094 }
3095 self.opened.fetch_add(1, Atomic::Relaxed);
3096 let dictionary = Arc::new(open_global_dictionary(
3097 Arc::clone(&self.file),
3098 page,
3099 &self.table.fields[column].ty,
3100 TEXT_KEEP_BUDGET,
3101 )?);
3102 let _ = self.dictionaries[column].set(Arc::clone(&dictionary));
3103 Ok(Some(dictionary))
3104 }
3105
3106 pub fn read(&self, part: usize, columns: &[usize]) -> Result<Chunk> {
3115 self.read_impl(part, columns, true)
3116 }
3117
3118 pub fn read_sparse(&self, part: usize, columns: &[usize]) -> Result<Chunk> {
3128 self.read_impl(part, columns, false)
3129 }
3130
3131 pub fn skips_codes(&self, part: usize, column: usize, candidates: &[u32]) -> Result<bool> {
3138 if candidates.is_empty() {
3139 return Ok(true);
3140 }
3141 if candidates.windows(2).any(|pair| pair[0] >= pair[1]) {
3142 return Err(Error::internal("native code candidates are not sorted and unique"));
3143 }
3144 let stripe = self.stripe_of(part)?;
3145 let Some(page) = stripe.memberships.get(column).copied().flatten() else {
3146 return Ok(false);
3147 };
3148 let mut bytes = vec![0; page.length as usize];
3149 read_at(&self.file, page.offset, &mut bytes)?;
3150 if checksum(&bytes) != page.hash {
3151 return Err(invalid("membership page checksum differs"));
3152 }
3153 let codes = decode_membership(&bytes)?;
3154 let mut left = 0;
3155 let mut right = 0;
3156 while left < codes.len() && right < candidates.len() {
3157 match codes[left].cmp(&candidates[right]) {
3158 Ordering::Less => left += 1,
3159 Ordering::Greater => right += 1,
3160 Ordering::Equal => return Ok(false),
3161 }
3162 }
3163 Ok(true)
3164 }
3165
3166 fn stripe_of(&self, part: usize) -> Result<&Stripe> {
3167 let place = self.places.get(part).ok_or_else(|| invalid("part index out of range"))?;
3168 self.table
3169 .stripes
3170 .get(place.stripe as usize)
3171 .ok_or_else(|| invalid("stripe index out of range"))
3172 }
3173
3174 fn held(&self, at: usize, stripe: &Stripe, column: usize, whole: bool) -> Result<CachedColumn> {
3191 let cache = self.cache.get(column).ok_or_else(|| invalid("column index out of range"))?;
3192 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3193 let known = cached.index.get(at).and_then(Clone::clone);
3194 let page = cached.pages.get(at).and_then(Clone::clone);
3195 if let Some(index) = known.clone() {
3196 if !whole || page.is_some() {
3197 return Ok(CachedColumn { stripe: at, index, page });
3198 }
3199 }
3200 if cached.loading.contains(&at) {
3201 drop(cached);
3202 if let Some(index) = known {
3206 return Ok(CachedColumn { stripe: at, index, page: None });
3207 }
3208 let held = self.page_of(stripe, column, at, false, None)?;
3209 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3210 remember(&mut cached, &held, self.kept.load(Atomic::Relaxed));
3211 return Ok(held);
3212 }
3213 cached.loading.push(at);
3214 drop(cached);
3215
3216 let read = self.page_of(stripe, column, at, whole, known);
3217
3218 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3222 if let Some(position) = cached.loading.iter().position(|loading| *loading == at) {
3223 cached.loading.remove(position);
3224 }
3225 let held = read?;
3226 remember(&mut cached, &held, self.kept.load(Atomic::Relaxed));
3227 Ok(held)
3228 }
3229
3230 fn page_of(
3236 &self,
3237 stripe: &Stripe,
3238 column: usize,
3239 at: usize,
3240 whole: bool,
3241 known: Option<Arc<Vec<PartSpan>>>,
3242 ) -> Result<CachedColumn> {
3243 let index = match known {
3244 Some(index) => index,
3245 None => {
3246 self.indexes.fetch_add(1, Atomic::Relaxed);
3247 Arc::new(read_index(&self.file, stripe, column)?)
3248 }
3249 };
3250 let page = if whole {
3251 self.pages.fetch_add(1, Atomic::Relaxed);
3252 let span = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
3253 let mut bytes = vec![0; span.length as usize];
3254 read_at(&self.file, span.offset, &mut bytes)?;
3255 Some(Arc::new(bytes))
3256 } else {
3257 None
3258 };
3259 Ok(CachedColumn { stripe: at, index, page })
3260 }
3261
3262 fn read_impl(&self, at: usize, columns: &[usize], whole: bool) -> Result<Chunk> {
3263 let place = *self.places.get(at).ok_or_else(|| invalid("part index out of range"))?;
3264 let index = place.stripe as usize;
3265 let stripe =
3266 self.table.stripes.get(index).ok_or_else(|| invalid("stripe index out of range"))?;
3267 let rows = place.rows as usize;
3268 let mut picked = Vec::with_capacity(columns.len());
3269 for &column in columns {
3270 let field = self
3271 .table
3272 .fields
3273 .get(column)
3274 .ok_or_else(|| invalid("column index out of range"))?;
3275 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
3276 let held = self.held(index, stripe, column, whole)?;
3277 let span = *held
3278 .index
3279 .get(place.part as usize)
3280 .ok_or_else(|| invalid("part index out of range"))?;
3281 let owned;
3282 let bytes = match &held.page {
3283 Some(held) => part_bytes(held, span)?,
3284 None => {
3285 let offset = page
3286 .offset
3287 .checked_add(span.start as u64)
3288 .ok_or_else(|| invalid("part range overflow"))?;
3289 let mut bytes = vec![0; span.length];
3290 read_at(&self.file, offset, &mut bytes)?;
3291 owned = bytes;
3292 &owned
3293 }
3294 };
3295 if checksum(bytes) != span.hash {
3296 return Err(invalid(&format!(
3297 "column page checksum differs, column {column} part {} at {}+{} of {} bytes, \
3298 wanted {:016x} and got {:016x}",
3299 place.part,
3300 page.offset,
3301 span.start,
3302 span.length,
3303 span.hash,
3304 checksum(bytes),
3305 )));
3306 }
3307 let dictionary = self.dictionary(column)?;
3308 picked.push(decode(&field.ty, rows, bytes, dictionary)?.into_pages());
3314 }
3315 Chunk::with_rows(picked, rows)
3316 }
3317
3318 #[must_use]
3334 pub fn skips(&self, part: usize, probes: &[Probe]) -> bool {
3335 let Some(place) = self.places.get(part).copied() else { return false };
3336 let Some(stripe) = self.table.stripes.get(place.stripe as usize) else { return false };
3337 if stripe.zone.skips(probes) {
3338 return true;
3339 }
3340 probes.iter().any(|probe| self.outside(place, probe) || self.sifted(place, probe))
3341 }
3342
3343 fn outside(&self, place: Place, probe: &Probe) -> bool {
3349 match self.stripe_part_ranges(place.stripe as usize, probe.column) {
3350 Some(ranges) => ranges
3351 .get(place.part as usize)
3352 .is_some_and(|range| range.excludes(probe.op, &probe.value)),
3353 None => false,
3354 }
3355 }
3356
3357 fn stripe_part_ranges(&self, stripe: usize, column: usize) -> Option<&[Range]> {
3363 let slot = self.part_ranges.get(column)?.get(stripe)?;
3364 if let Some(held) = slot.get() {
3365 return Some(held);
3366 }
3367 let page = self.table.stripes.get(stripe)?.part_ranges.get(column).copied().flatten()?;
3368 let mut bytes = vec![0; page.length as usize];
3369 read_at(&self.file, page.offset, &mut bytes).ok()?;
3370 if checksum(&bytes) != page.hash {
3371 return None;
3372 }
3373 let ranges = Arc::new(decode_part_ranges(&bytes).ok()?);
3374 let _ = slot.set(ranges);
3375 slot.get().map(|held| held.as_slice())
3376 }
3377
3378 #[must_use]
3395 pub fn certain(&self, part: usize, probes: &[Probe]) -> bool {
3396 let Some(place) = self.places.get(part).copied() else { return false };
3397 let Some(stripe) = self.table.stripes.get(place.stripe as usize) else { return false };
3398 if stripe.zone.certain(probes) {
3399 return true;
3400 }
3401 probes
3402 .iter()
3403 .all(|probe| stripe.zone.certain(slice::from_ref(probe)) || self.inside(place, probe))
3404 }
3405
3406 fn inside(&self, place: Place, probe: &Probe) -> bool {
3412 match self.stripe_part_ranges(place.stripe as usize, probe.column) {
3413 Some(ranges) => ranges
3414 .get(place.part as usize)
3415 .is_some_and(|range| range.certain(probe.op, &probe.value)),
3416 None => false,
3417 }
3418 }
3419
3420 #[must_use]
3431 pub fn stripe_skips(&self, stripe: usize, probes: &[Probe]) -> bool {
3432 self.table.stripes.get(stripe).is_some_and(|held| held.zone.skips(probes))
3433 }
3434
3435 fn sifted(&self, place: Place, probe: &Probe) -> bool {
3441 if probe.op != Op::Equal {
3442 return false;
3443 }
3444 match self.stripe_sieves(place.stripe as usize, probe.column) {
3445 Some(sieves) => sieves
3446 .get(place.part as usize)
3447 .and_then(Option::as_ref)
3448 .is_some_and(|sieve| sieve.excludes(&probe.value)),
3449 None => false,
3450 }
3451 }
3452
3453 fn stripe_sieves(&self, stripe: usize, column: usize) -> Option<&[Option<Sieve>]> {
3460 let slot = self.sieves.get(column)?.get(stripe)?;
3461 if let Some(held) = slot.get() {
3462 return Some(held);
3463 }
3464 let page = self.table.stripes.get(stripe)?.sieves.get(column).copied().flatten()?;
3465 let mut bytes = vec![0; page.length as usize];
3466 read_at(&self.file, page.offset, &mut bytes).ok()?;
3467 if checksum(&bytes) != page.hash {
3468 return None;
3469 }
3470 let sieves = Arc::new(decode_sieves(&bytes).ok()?);
3471 let _ = slot.set(sieves);
3472 slot.get().map(|held| held.as_slice())
3473 }
3474}
3475
3476fn text_at_rank(dictionary: &Vector, rank: usize) -> Result<Value> {
3478 let code = dictionary.code_at_rank(rank)? as usize;
3479 let text = dictionary
3480 .try_text_at(code)?
3481 .ok_or_else(|| invalid("global dictionary order names a code it does not have"))?;
3482 Ok(Value::Varchar(text.into()))
3483}
3484
3485#[cfg(unix)]
3490fn write_at(file: &File, mut offset: u64, mut bytes: &[u8]) -> Result<()> {
3491 use std::os::unix::fs::FileExt;
3492 while !bytes.is_empty() {
3493 let written = file.write_at(bytes, offset).map_err(io)?;
3494 if written == 0 {
3495 return Err(invalid("a write to the native file wrote nothing"));
3496 }
3497 offset += written as u64;
3498 bytes = &bytes[written..];
3499 }
3500 Ok(())
3501}
3502
3503#[cfg(windows)]
3505fn write_at(file: &File, mut offset: u64, mut bytes: &[u8]) -> Result<()> {
3506 use std::os::windows::fs::FileExt;
3507 while !bytes.is_empty() {
3508 let written = file.seek_write(bytes, offset).map_err(io)?;
3509 if written == 0 {
3510 return Err(invalid("a write to the native file wrote nothing"));
3511 }
3512 offset += written as u64;
3513 bytes = &bytes[written..];
3514 }
3515 Ok(())
3516}
3517
3518#[cfg(not(any(unix, windows)))]
3520fn write_at(file: &File, offset: u64, bytes: &[u8]) -> Result<()> {
3521 use std::io::Write;
3522 let mut file = file.try_clone().map_err(io)?;
3523 file.seek(SeekFrom::Start(offset)).map_err(io)?;
3524 file.write_all(bytes).map_err(io)
3525}
3526
3527#[cfg(unix)]
3537fn read_at(file: &File, mut offset: u64, mut bytes: &mut [u8]) -> Result<()> {
3538 use std::os::unix::fs::FileExt;
3539 while !bytes.is_empty() {
3540 let read = file.read_at(bytes, offset).map_err(io)?;
3541 if read == 0 {
3542 return Err(invalid("column page ends before its declared length"));
3543 }
3544 offset += read as u64;
3545 bytes = &mut bytes[read..];
3546 }
3547 Ok(())
3548}
3549
3550#[cfg(windows)]
3556fn read_at(file: &File, mut offset: u64, mut bytes: &mut [u8]) -> Result<()> {
3557 use std::os::windows::fs::FileExt;
3558 while !bytes.is_empty() {
3559 let read = file.seek_read(bytes, offset).map_err(io)?;
3560 if read == 0 {
3561 return Err(invalid("column page ends before its declared length"));
3562 }
3563 offset += read as u64;
3564 bytes = &mut bytes[read..];
3565 }
3566 Ok(())
3567}
3568
3569#[cfg(not(any(unix, windows)))]
3574fn read_at(file: &File, offset: u64, bytes: &mut [u8]) -> Result<()> {
3575 let mut file = file.try_clone().map_err(io)?;
3576 file.seek(SeekFrom::Start(offset)).map_err(io)?;
3577 file.read_exact(bytes).map_err(io)
3578}
3579
3580fn type_tag(ty: &LogicalType) -> Result<u8> {
3581 match ty {
3582 LogicalType::SmallInt => Ok(1),
3583 LogicalType::Integer => Ok(2),
3584 LogicalType::BigInt => Ok(3),
3585 LogicalType::Varchar => Ok(4),
3586 LogicalType::Date => Ok(5),
3587 LogicalType::Timestamp => Ok(6),
3588 LogicalType::Boolean => Ok(7),
3589 LogicalType::TinyInt => Ok(8),
3590 LogicalType::UTinyInt => Ok(9),
3591 LogicalType::USmallInt => Ok(10),
3592 LogicalType::UInteger => Ok(11),
3593 LogicalType::UBigInt => Ok(12),
3594 LogicalType::Decimal { .. } => Ok(13),
3595 _ => Err(Error::not_implemented(format!("native storage for {ty}"))),
3596 }
3597}
3598
3599fn put_type(out: &mut Vec<u8>, ty: &LogicalType) -> Result<()> {
3605 out.push(type_tag(ty)?);
3606 if let LogicalType::Decimal { width, scale } = ty {
3607 out.push(*width);
3608 out.push(*scale);
3609 }
3610 Ok(())
3611}
3612
3613fn read_type(cur: &mut Cursor<'_>) -> Result<LogicalType> {
3615 let tag = cur.u8()?;
3616 if tag == 13 {
3617 let width = cur.u8()?;
3618 let scale = cur.u8()?;
3619 return LogicalType::decimal(width, scale)
3620 .map_err(|_| invalid("decimal column width and scale are not a decimal"));
3621 }
3622 tag_type(tag)
3623}
3624
3625fn tag_type(tag: u8) -> Result<LogicalType> {
3626 match tag {
3627 1 => Ok(LogicalType::SmallInt),
3628 2 => Ok(LogicalType::Integer),
3629 3 => Ok(LogicalType::BigInt),
3630 4 => Ok(LogicalType::Varchar),
3631 5 => Ok(LogicalType::Date),
3632 6 => Ok(LogicalType::Timestamp),
3633 7 => Ok(LogicalType::Boolean),
3634 8 => Ok(LogicalType::TinyInt),
3635 9 => Ok(LogicalType::UTinyInt),
3636 10 => Ok(LogicalType::USmallInt),
3637 11 => Ok(LogicalType::UInteger),
3638 12 => Ok(LogicalType::UBigInt),
3639 _ => Err(invalid("column type tag is unknown")),
3640 }
3641}
3642
3643fn put_u16(out: &mut Vec<u8>, value: u16) {
3644 out.extend_from_slice(&value.to_le_bytes());
3645}
3646fn put_u32(out: &mut Vec<u8>, value: u32) {
3647 out.extend_from_slice(&value.to_le_bytes());
3648}
3649fn put_u64(out: &mut Vec<u8>, value: u64) {
3650 out.extend_from_slice(&value.to_le_bytes());
3651}
3652fn put_var_u64(out: &mut Vec<u8>, mut value: u64) {
3653 while value >= 0x80 {
3654 out.push((value as u8 & 0x7f) | 0x80);
3655 value >>= 7;
3656 }
3657 out.push(value as u8);
3658}
3659
3660fn frequency_order(left: FrequencyValue, right: FrequencyValue) -> Ordering {
3661 match (left, right) {
3662 (FrequencyValue::Null, FrequencyValue::Null) => Ordering::Equal,
3663 (FrequencyValue::Null, _) => Ordering::Less,
3664 (_, FrequencyValue::Null) => Ordering::Greater,
3665 (FrequencyValue::Integer(left), FrequencyValue::Integer(right)) => left.cmp(&right),
3666 (FrequencyValue::Code(left), FrequencyValue::Code(right)) => left.cmp(&right),
3667 (FrequencyValue::Integer(_), FrequencyValue::Code(_)) => Ordering::Less,
3668 (FrequencyValue::Code(_), FrequencyValue::Integer(_)) => Ordering::Greater,
3669 }
3670}
3671
3672fn code_frequency(dictionary: &GlobalDictionary) -> FrequencySummary {
3673 let mut entries = dictionary
3674 .counts
3675 .iter()
3676 .enumerate()
3677 .filter(|(_, count)| **count != 0)
3678 .map(|(code, &count)| FrequencyEntry { value: FrequencyValue::Code(code as u32), count })
3679 .collect::<Vec<_>>();
3680 if dictionary.nulls != 0 {
3681 entries.push(FrequencyEntry { value: FrequencyValue::Null, count: dictionary.nulls });
3682 }
3683 entries.sort_unstable_by(|left, right| {
3684 right.count.cmp(&left.count).then_with(|| frequency_order(left.value, right.value))
3685 });
3686 let omitted_max = entries.get(FREQUENCY_ENTRIES).map_or(0, |entry| entry.count);
3687 entries.truncate(FREQUENCY_ENTRIES);
3688 FrequencySummary { entries, omitted_max, ordinals: Vec::new() }
3689}
3690
3691fn encode_directory(table: &Table) -> Result<Vec<u8>> {
3692 let mut out = DIRECTORY.to_vec();
3693 let name = table.name.as_bytes();
3694 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("table name too long"))?);
3695 out.extend_from_slice(name);
3696 put_u16(&mut out, u16::try_from(table.fields.len()).map_err(|_| invalid("too many columns"))?);
3697 for field in &table.fields {
3698 let name = field.name.as_bytes();
3699 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("column name too long"))?);
3700 out.extend_from_slice(name);
3701 put_type(&mut out, &field.ty)?;
3702 out.push(u8::from(field.not_null));
3703 }
3704 for dictionary in &table.dictionaries {
3705 match dictionary {
3706 None => out.push(0),
3707 Some(page) => {
3708 out.push(1);
3709 put_u64(&mut out, page.offset);
3710 put_u32(&mut out, page.length);
3711 put_u64(&mut out, page.hash);
3712 }
3713 }
3714 }
3715 for distinct in &table.distincts {
3716 match distinct {
3717 None => out.push(0),
3718 Some(count) => {
3719 out.push(1);
3720 put_u64(&mut out, *count);
3721 }
3722 }
3723 }
3724 put_u64(&mut out, u64::try_from(table.rows).map_err(|_| invalid("row count overflow"))?);
3725 put_u32(&mut out, u32::try_from(table.stripes.len()).map_err(|_| invalid("too many stripes"))?);
3726 for stripe in &table.stripes {
3727 put_u32(
3728 &mut out,
3729 u32::try_from(stripe.parts.len()).map_err(|_| invalid("too many parts in a stripe"))?,
3730 );
3731 for &rows in &stripe.parts {
3732 put_u32(&mut out, rows);
3733 }
3734 put_u64(&mut out, stripe.index.offset);
3735 put_u32(&mut out, stripe.index.length);
3736 for page in &stripe.pages {
3737 put_u64(&mut out, page.offset);
3738 put_u32(&mut out, page.length);
3739 }
3740 for ((field, dictionary), membership) in
3745 table.fields.iter().zip(&table.dictionaries).zip(&stripe.memberships)
3746 {
3747 if field.ty != LogicalType::Varchar || dictionary.is_none() {
3748 continue;
3749 }
3750 let page =
3751 membership.ok_or_else(|| invalid("string page has no code membership index"))?;
3752 put_u64(&mut out, page.offset);
3753 put_u32(&mut out, page.length);
3754 put_u64(&mut out, page.hash);
3755 }
3756 for sieve in &stripe.sieves {
3757 match sieve {
3758 None => out.push(0),
3759 Some(page) => {
3760 out.push(1);
3761 put_u64(&mut out, page.offset);
3762 put_u32(&mut out, page.length);
3763 put_u64(&mut out, page.hash);
3764 }
3765 }
3766 }
3767 for held in &stripe.part_ranges {
3768 match held {
3769 None => out.push(0),
3770 Some(page) => {
3771 out.push(1);
3772 put_u64(&mut out, page.offset);
3773 put_u32(&mut out, page.length);
3774 put_u64(&mut out, page.hash);
3775 }
3776 }
3777 }
3778 for range in stripe.zone.columns() {
3779 put_bound(&mut out, range.low.as_ref())?;
3780 put_bound(&mut out, range.high.as_ref())?;
3781 put_u32(
3782 &mut out,
3783 u32::try_from(range.nulls).map_err(|_| invalid("null count overflow"))?,
3784 );
3785 out.push(u8::from(range.exact));
3786 match range.sum {
3787 None => out.push(0),
3788 Some(total) => {
3789 out.push(1);
3790 out.extend_from_slice(&total.to_le_bytes());
3791 }
3792 }
3793 }
3794 }
3795 out.extend_from_slice(FREQUENCIES);
3796 put_u16(
3797 &mut out,
3798 u16::try_from(table.frequencies.len())
3799 .map_err(|_| invalid("too many frequency columns"))?,
3800 );
3801 for summary in &table.frequencies {
3802 let Some(summary) = summary else {
3803 out.push(0);
3804 continue;
3805 };
3806 out.push(1);
3807 put_u64(&mut out, summary.omitted_max);
3808 put_u32(
3809 &mut out,
3810 u32::try_from(summary.entries.len())
3811 .map_err(|_| invalid("too many frequency entries"))?,
3812 );
3813 for entry in &summary.entries {
3814 match entry.value {
3815 FrequencyValue::Null => out.push(0),
3816 FrequencyValue::Integer(value) => {
3817 out.push(1);
3818 out.extend_from_slice(&value.to_le_bytes());
3819 }
3820 FrequencyValue::Code(value) => {
3821 out.push(2);
3822 put_u32(&mut out, value);
3823 }
3824 }
3825 put_u64(&mut out, entry.count);
3826 }
3827 put_u32(
3828 &mut out,
3829 u32::try_from(summary.ordinals.len())
3830 .map_err(|_| invalid("too many frequency ordinals"))?,
3831 );
3832 let mut previous = 0_u64;
3833 for (at, &ordinal) in summary.ordinals.iter().enumerate() {
3834 let delta = if at == 0 {
3835 ordinal
3836 } else {
3837 ordinal
3838 .checked_sub(previous)
3839 .ok_or_else(|| invalid("frequency ordinals are not ordered"))?
3840 };
3841 if at != 0 && delta == 0 {
3842 return Err(invalid("frequency ordinals are not unique"));
3843 }
3844 put_var_u64(&mut out, delta);
3845 previous = ordinal;
3846 }
3847 }
3848 Ok(out)
3849}
3850
3851fn encode_catalog(entries: &[Entry]) -> Result<Vec<u8>> {
3857 let mut out = CATALOG.to_vec();
3858 put_u32(&mut out, u32::try_from(entries.len()).map_err(|_| invalid("too many tables"))?);
3859 for entry in entries {
3860 let name = entry.name.as_bytes();
3861 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("table name too long"))?);
3862 out.extend_from_slice(name);
3863 put_u64(&mut out, u64::try_from(entry.rows).map_err(|_| invalid("row count overflow"))?);
3864 put_u16(
3865 &mut out,
3866 u16::try_from(entry.fields.len()).map_err(|_| invalid("too many columns"))?,
3867 );
3868 for field in &entry.fields {
3869 let name = field.name.as_bytes();
3870 put_u16(
3871 &mut out,
3872 u16::try_from(name.len()).map_err(|_| invalid("column name too long"))?,
3873 );
3874 out.extend_from_slice(name);
3875 put_type(&mut out, &field.ty)?;
3876 out.push(u8::from(field.not_null));
3877 }
3878 put_u64(&mut out, entry.directory.offset);
3879 put_u32(&mut out, entry.directory.length);
3880 put_u64(&mut out, entry.directory.hash);
3881 }
3882 Ok(out)
3883}
3884
3885fn decode_catalog(bytes: &[u8], size: u64) -> Result<Vec<Entry>> {
3888 let mut cur = Cursor { bytes, at: 0 };
3889 if cur.take(8)? != CATALOG {
3890 return Err(invalid("catalog magic differs"));
3891 }
3892 let count = cur.u32()? as usize;
3893 let mut entries: Vec<Entry> = Vec::with_capacity(count.min(1024));
3894 for _ in 0..count {
3895 let name = cur.text()?;
3896 let rows = usize::try_from(cur.u64()?).map_err(|_| invalid("row count does not fit"))?;
3897 let width = cur.u16()? as usize;
3898 let mut fields = Vec::with_capacity(width);
3899 for _ in 0..width {
3900 let name = cur.text()?;
3901 let ty = read_type(&mut cur)?;
3902 let not_null = match cur.u8()? {
3903 0 => false,
3904 1 => true,
3905 _ => return Err(invalid("nullability flag differs")),
3906 };
3907 fields.push(Field { name, ty, not_null });
3908 }
3909 let directory = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3910 let end = directory
3911 .offset
3912 .checked_add(u64::from(directory.length))
3913 .ok_or_else(|| invalid("table directory offset overflow"))?;
3914 if directory.offset < HEADER
3915 || end > size
3916 || directory.length as usize > MAX_DIRECTORY
3917 || directory.length == 0
3918 {
3919 return Err(invalid("table directory range is outside the file"));
3920 }
3921 if entries.iter().any(|held| held.name == name) {
3922 return Err(invalid("two tables in the catalog have the same name"));
3923 }
3924 entries.push(Entry { name, fields, rows, directory });
3925 }
3926 Ok(entries)
3927}
3928
3929struct Cursor<'a> {
3930 bytes: &'a [u8],
3931 at: usize,
3932}
3933impl<'a> Cursor<'a> {
3934 fn take(&mut self, len: usize) -> Result<&'a [u8]> {
3935 let end = self.at.checked_add(len).ok_or_else(|| invalid("directory offset overflow"))?;
3936 let bytes =
3937 self.bytes.get(self.at..end).ok_or_else(|| invalid("directory is truncated"))?;
3938 self.at = end;
3939 Ok(bytes)
3940 }
3941 fn u8(&mut self) -> Result<u8> {
3942 Ok(self.take(1)?[0])
3943 }
3944 fn u16(&mut self) -> Result<u16> {
3945 Ok(u16::from_le_bytes(self.take(2)?.try_into().expect("two bytes")))
3946 }
3947 fn u32(&mut self) -> Result<u32> {
3948 Ok(u32::from_le_bytes(self.take(4)?.try_into().expect("four bytes")))
3949 }
3950 fn u64(&mut self) -> Result<u64> {
3951 Ok(u64::from_le_bytes(self.take(8)?.try_into().expect("eight bytes")))
3952 }
3953 fn var_u64(&mut self) -> Result<u64> {
3954 let mut value = 0_u64;
3955 for shift in (0..=63).step_by(7) {
3956 let byte = self.u8()?;
3957 let part = u64::from(byte & 0x7f);
3958 if shift == 63 && part > 1 {
3959 return Err(invalid("frequency ordinal varint overflows"));
3960 }
3961 value |= part << shift;
3962 if byte & 0x80 == 0 {
3963 return Ok(value);
3964 }
3965 }
3966 Err(invalid("frequency ordinal varint is too long"))
3967 }
3968 fn bound(&mut self) -> Result<Option<Bound>> {
3969 Ok(match self.u8()? {
3970 0 => None,
3971 1 => Some(Bound::Int(i128::from_le_bytes(
3972 self.take(16)?.try_into().expect("sixteen bytes"),
3973 ))),
3974 2 => Some(Bound::Real(f64::from_le_bytes(
3975 self.take(8)?.try_into().expect("eight bytes"),
3976 ))),
3977 3 => {
3978 let length = self.u32()? as usize;
3979 Some(Bound::Bytes(self.take(length)?.to_vec()))
3980 }
3981 4 => {
3982 let unscaled =
3983 i128::from_le_bytes(self.take(16)?.try_into().expect("sixteen bytes"));
3984 Some(Bound::Scaled { unscaled, scale: self.u8()? })
3985 }
3986 _ => return Err(invalid("bound tag differs")),
3987 })
3988 }
3989 fn text(&mut self) -> Result<String> {
3990 let len = self.u16()? as usize;
3991 String::from_utf8(self.take(len)?.to_vec()).map_err(|_| invalid("name is not UTF-8"))
3992 }
3993}
3994
3995fn decode_directory(bytes: &[u8], size: u64) -> Result<Table> {
3996 let mut cur = Cursor { bytes, at: 0 };
3997 if cur.take(8)? != DIRECTORY {
3998 return Err(invalid("directory magic differs"));
3999 }
4000 let name = cur.text()?;
4001 let width = cur.u16()? as usize;
4002 let mut fields = Vec::with_capacity(width);
4003 for _ in 0..width {
4004 let name = cur.text()?;
4005 let ty = read_type(&mut cur)?;
4006 let not_null = match cur.u8()? {
4007 0 => false,
4008 1 => true,
4009 _ => return Err(invalid("nullability flag differs")),
4010 };
4011 fields.push(Field { name, ty, not_null });
4012 }
4013 let mut dictionaries = Vec::with_capacity(width);
4014 for _ in 0..width {
4015 dictionaries.push(match cur.u8()? {
4016 0 => None,
4017 1 => {
4018 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4019 let end = page
4020 .offset
4021 .checked_add(u64::from(page.length))
4022 .ok_or_else(|| invalid("dictionary page offset overflow"))?;
4023 if page.offset < HEADER || end > size {
4028 return Err(invalid("dictionary page range is outside the file"));
4029 }
4030 Some(page)
4031 }
4032 _ => return Err(invalid("dictionary page tag differs")),
4033 });
4034 }
4035 let mut distincts = Vec::with_capacity(width);
4036 for _ in 0..width {
4037 distincts.push(match cur.u8()? {
4038 0 => None,
4039 1 => Some(cur.u64()?),
4040 _ => return Err(invalid("distinct count tag differs")),
4041 });
4042 }
4043 let rows = usize::try_from(cur.u64()?).map_err(|_| invalid("row count does not fit"))?;
4044 let count = cur.u32()? as usize;
4045 let mut stripes = Vec::with_capacity(count);
4046 let mut total = 0_usize;
4047 for _ in 0..count {
4048 let count = cur.u32()? as usize;
4049 if count == 0 || count > STRIPE_PARTS {
4050 return Err(invalid("stripe part count is outside its bound"));
4051 }
4052 let mut parts = Vec::with_capacity(count);
4053 let mut stripe_rows = 0_usize;
4054 for _ in 0..count {
4055 let rows = cur.u32()?;
4056 if rows == 0 {
4057 return Err(invalid("empty part"));
4058 }
4059 parts.push(rows);
4060 stripe_rows = stripe_rows
4061 .checked_add(rows as usize)
4062 .ok_or_else(|| invalid("stripe row count overflow"))?;
4063 }
4064 total =
4065 total.checked_add(stripe_rows).ok_or_else(|| invalid("stripe row count overflow"))?;
4066 let index = Span { offset: cur.u64()?, length: cur.u32()? };
4067 let section = index_section(count)?;
4068 let wanted = section
4069 .checked_mul(width)
4070 .and_then(|bytes| u32::try_from(bytes).ok())
4071 .ok_or_else(|| invalid("index page length overflow"))?;
4072 let end = index
4073 .offset
4074 .checked_add(u64::from(index.length))
4075 .ok_or_else(|| invalid("index page offset overflow"))?;
4076 if index.offset < HEADER || end > size || index.length != wanted {
4077 return Err(invalid("index page range is outside the file"));
4078 }
4079 let mut pages = Vec::with_capacity(width);
4080 for _ in 0..width {
4081 let offset = cur.u64()?;
4082 let length = cur.u32()?;
4083 let end = offset
4084 .checked_add(u64::from(length))
4085 .ok_or_else(|| invalid("page offset overflow"))?;
4086 if offset < HEADER || end > size || length as usize > MAX_PAGE {
4087 return Err(invalid("page range is outside the file"));
4088 }
4089 pages.push(Span { offset, length });
4090 }
4091 let mut memberships = vec![None; width];
4092 for (column, field) in fields.iter().enumerate() {
4093 if field.ty != LogicalType::Varchar || dictionaries[column].is_none() {
4094 continue;
4095 }
4096 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4097 let end = page
4098 .offset
4099 .checked_add(u64::from(page.length))
4100 .ok_or_else(|| invalid("membership page offset overflow"))?;
4101 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
4102 return Err(invalid("membership page range is outside the file"));
4103 }
4104 memberships[column] = Some(page);
4105 }
4106 let mut sieves = vec![None; width];
4107 for sieve in sieves.iter_mut().take(width) {
4108 match cur.u8()? {
4109 0 => continue,
4110 1 => {}
4111 _ => return Err(invalid("a sieve page has an unknown tag")),
4112 }
4113 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4114 let end = page
4115 .offset
4116 .checked_add(u64::from(page.length))
4117 .ok_or_else(|| invalid("sieve page offset overflow"))?;
4118 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
4119 return Err(invalid("sieve page range is outside the file"));
4120 }
4121 *sieve = Some(page);
4122 }
4123 let mut part_ranges = vec![None; width];
4124 for held in part_ranges.iter_mut().take(width) {
4125 match cur.u8()? {
4126 0 => continue,
4127 1 => {}
4128 _ => return Err(invalid("a part range page has an unknown tag")),
4129 }
4130 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4131 let end = page
4132 .offset
4133 .checked_add(u64::from(page.length))
4134 .ok_or_else(|| invalid("part range page offset overflow"))?;
4135 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
4136 return Err(invalid("part range page range is outside the file"));
4137 }
4138 *held = Some(page);
4139 }
4140 let mut ranges = Vec::with_capacity(width);
4141 for column in 0..width {
4142 let low = cur.bound()?;
4143 let high = cur.bound()?;
4144 let nulls = cur.u32()? as usize;
4145 if nulls > stripe_rows {
4146 return Err(invalid("null count exceeds stripe rows"));
4147 }
4148 let exact = cur.u8()? != 0;
4149 let sum = match cur.u8()? {
4150 0 => None,
4151 1 => Some(i128::from_le_bytes(
4152 cur.take(16)?.try_into().map_err(|_| invalid("a stripe sum is truncated"))?,
4153 )),
4154 _ => return Err(invalid("a stripe sum has an unknown tag")),
4155 };
4156 let ty = &fields.get(column).ok_or_else(|| invalid("a stripe range has no column"))?.ty;
4162 let low = low.map(|bound| scaled_as(bound, ty));
4163 let high = high.map(|bound| scaled_as(bound, ty));
4164 ranges.push(Range { low, high, nulls, exact, sum });
4165 }
4166 stripes.push(Stripe {
4167 rows: stripe_rows,
4168 parts,
4169 index,
4170 pages,
4171 memberships,
4172 sieves,
4173 part_ranges,
4174 zone: Zone::from_ranges(ranges),
4175 });
4176 }
4177 if total != rows {
4178 return Err(invalid("table row count differs from stripes"));
4179 }
4180 let frequencies = if cur.at == bytes.len() {
4181 vec![None; width]
4182 } else {
4183 if cur.take(8)? != FREQUENCIES {
4184 return Err(invalid("directory extension magic differs"));
4185 }
4186 if cur.u16()? as usize != width {
4187 return Err(invalid("frequency column count differs"));
4188 }
4189 let mut frequencies = Vec::with_capacity(width);
4190 for field in &fields {
4191 let summary = match cur.u8()? {
4192 0 => None,
4193 1 => {
4194 let omitted_max = cur.u64()?;
4195 let count = cur.u32()? as usize;
4196 if count > FREQUENCY_ENTRIES {
4197 return Err(invalid("frequency entry count exceeds its bound"));
4198 }
4199 let mut entries = Vec::with_capacity(count);
4200 for _ in 0..count {
4202 let value = match cur.u8()? {
4203 0 => FrequencyValue::Null,
4204 1 => FrequencyValue::Integer(i128::from_le_bytes(
4205 cur.take(16)?.try_into().expect("sixteen bytes"),
4206 )),
4207 2 => FrequencyValue::Code(cur.u32()?),
4208 _ => return Err(invalid("frequency value tag differs")),
4209 };
4210 let valid = matches!(
4211 (&field.ty, value),
4212 (_, FrequencyValue::Null)
4213 | (LogicalType::Varchar, FrequencyValue::Code(_))
4214 | (
4215 LogicalType::TinyInt
4216 | LogicalType::SmallInt
4217 | LogicalType::Integer
4218 | LogicalType::BigInt
4219 | LogicalType::UTinyInt
4220 | LogicalType::USmallInt
4221 | LogicalType::UInteger
4222 | LogicalType::UBigInt
4223 | LogicalType::Date
4224 | LogicalType::Timestamp,
4225 FrequencyValue::Integer(_),
4226 )
4227 );
4228 if !valid {
4229 return Err(invalid("frequency value does not match its column"));
4230 }
4231 let count = cur.u64()?;
4232 if count == 0 || count > rows as u64 {
4233 return Err(invalid("frequency count is outside the table"));
4234 }
4235 entries.push(FrequencyEntry { value, count });
4236 }
4237 if entries.windows(2).any(|pair| pair[0].count < pair[1].count) {
4238 return Err(invalid("frequency entries are not descending"));
4239 }
4240 let ordinals = {
4241 let ordinal_count = cur.u32()? as usize;
4242 if ordinal_count > FREQUENCY_ORDINALS || ordinal_count > rows {
4243 return Err(invalid("frequency ordinal count exceeds its bound"));
4244 }
4245 let mut ordinals = Vec::with_capacity(ordinal_count);
4246 let mut previous = 0_u64;
4247 for at in 0..ordinal_count {
4248 let delta = cur.var_u64()?;
4249 if at != 0 && delta == 0 {
4250 return Err(invalid("frequency ordinals are not increasing"));
4251 }
4252 let ordinal = if at == 0 {
4253 delta
4254 } else {
4255 previous
4256 .checked_add(delta)
4257 .ok_or_else(|| invalid("frequency ordinal overflows"))?
4258 };
4259 if ordinal >= rows as u64 {
4260 return Err(invalid("frequency ordinal is outside the table"));
4261 }
4262 ordinals.push(ordinal);
4263 previous = ordinal;
4264 }
4265 ordinals
4266 };
4267 Some(FrequencySummary { entries, omitted_max, ordinals })
4268 }
4269 _ => return Err(invalid("frequency summary tag differs")),
4270 };
4271 frequencies.push(summary);
4272 }
4273 frequencies
4274 };
4275 if cur.at != bytes.len() {
4276 return Err(invalid("directory has trailing bytes"));
4277 }
4278 Ok(Table { name, fields, stripes, rows, dictionaries, distincts, frequencies })
4279}
4280
4281fn put_bound(out: &mut Vec<u8>, bound: Option<&Bound>) -> Result<()> {
4282 match bound {
4283 None => out.push(0),
4284 Some(Bound::Int(value)) => {
4285 out.push(1);
4286 out.extend_from_slice(&value.to_le_bytes());
4287 }
4288 Some(Bound::Real(value)) => {
4289 out.push(2);
4290 out.extend_from_slice(&value.to_le_bytes());
4291 }
4292 Some(Bound::Bytes(value)) => {
4293 out.push(3);
4294 put_u32(out, u32::try_from(value.len()).map_err(|_| invalid("bound length overflow"))?);
4295 out.extend_from_slice(value);
4296 }
4297 Some(Bound::Scaled { unscaled, scale }) => {
4298 out.push(4);
4299 out.extend_from_slice(&unscaled.to_le_bytes());
4300 out.push(*scale);
4301 }
4302 }
4303 Ok(())
4304}
4305
4306#[derive(Debug)]
4323struct Codes;
4324
4325impl chooser::Chooser for Codes {
4326 fn name(&self) -> &'static str {
4327 "codes"
4328 }
4329
4330 fn narrow_strings(
4331 &self,
4332 _values: &[&[u8]],
4333 offered: &[string::Kind],
4334 _depth: u8,
4335 ) -> Vec<string::Kind> {
4336 offered.to_vec()
4339 }
4340
4341 fn narrow_integers(
4342 &self,
4343 _values: &[i64],
4344 offered: &[integer::Kind],
4345 depth: u8,
4346 ) -> Vec<integer::Kind> {
4347 let keep: &[integer::Kind] = if depth == 0 {
4348 &[integer::Kind::Constant, integer::Kind::Packed, integer::Kind::Rle]
4349 } else {
4350 &[integer::Kind::Constant, integer::Kind::Packed]
4351 };
4352 let narrowed: Vec<integer::Kind> =
4353 offered.iter().copied().filter(|kind| keep.contains(kind)).collect();
4354 if narrowed.is_empty() { offered.to_vec() } else { narrowed }
4357 }
4358}
4359
4360#[derive(Debug)]
4372struct Fixed;
4373
4374impl chooser::Chooser for Fixed {
4375 fn name(&self) -> &'static str {
4376 "fixed"
4377 }
4378
4379 fn narrow_strings(
4380 &self,
4381 _values: &[&[u8]],
4382 offered: &[string::Kind],
4383 _depth: u8,
4384 ) -> Vec<string::Kind> {
4385 offered.to_vec()
4386 }
4387
4388 fn narrow_integers(
4389 &self,
4390 _values: &[i64],
4391 offered: &[integer::Kind],
4392 depth: u8,
4393 ) -> Vec<integer::Kind> {
4394 let keep: &[integer::Kind] = if depth == 0 {
4395 &[
4396 integer::Kind::Constant,
4397 integer::Kind::Packed,
4398 integer::Kind::Delta,
4399 integer::Kind::Rle,
4400 integer::Kind::Sparse,
4401 integer::Kind::Strided,
4402 ]
4403 } else {
4404 &[integer::Kind::Constant, integer::Kind::Packed, integer::Kind::Delta]
4405 };
4406 let narrowed: Vec<integer::Kind> =
4407 offered.iter().copied().filter(|kind| keep.contains(kind)).collect();
4408 if narrowed.is_empty() { offered.to_vec() } else { narrowed }
4409 }
4410}
4411
4412fn widened(data: &Data) -> Option<Vec<i64>> {
4419 match data {
4420 Data::Int8(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4421 Data::UInt8(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4422 Data::Int16(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4423 Data::UInt16(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4424 Data::Int32(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4425 Data::UInt32(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4426 Data::Int64(values) => Some(values.to_vec()),
4427 _ => None,
4428 }
4429}
4430
4431trait Narrow: Copy {
4438 const BIASED: (u32, u64);
4443
4444 fn narrow(value: i64) -> Self;
4446}
4447
4448#[allow(clippy::cast_sign_loss, reason = "a residue is a bit pattern and not a number")]
4465fn residue<T: Narrow>(value: i64) -> u64 {
4466 let (bits, bias) = T::BIASED;
4467 (value as u64).wrapping_add(bias) >> bits
4468}
4469
4470macro_rules! narrows {
4475 ($($ty:ty => $bias:expr),* $(,)?) => {$(
4476 impl Narrow for $ty {
4477 const BIASED: (u32, u64) = (<$ty>::BITS, $bias);
4478
4479 #[allow(
4480 clippy::cast_possible_truncation,
4481 clippy::cast_sign_loss,
4482 reason = "the caller has checked the bits this truncates away"
4483 )]
4484 fn narrow(value: i64) -> Self {
4485 value as Self
4486 }
4487 }
4488 )*};
4489}
4490
4491narrows! {
4492 i8 => 1 << 7,
4493 u8 => 0,
4494 i16 => 1 << 15,
4495 u16 => 0,
4496 i32 => 1 << 31,
4497 u32 => 0,
4498}
4499
4500fn fit<T: Narrow>(values: &[i64]) -> Result<Vec<T>> {
4513 let mut spilled = 0u64;
4514 for value in values {
4515 spilled |= residue::<T>(*value);
4516 }
4517 if spilled != 0 {
4518 return Err(invalid("page value is not of its type"));
4519 }
4520 Ok(values.iter().map(|value| T::narrow(*value)).collect())
4521}
4522
4523fn narrowed(ty: &LogicalType, values: Vec<i64>) -> Result<Data> {
4528 Ok(match ty {
4529 LogicalType::TinyInt => Data::Int8(fit::<i8>(&values)?.into()),
4530 LogicalType::UTinyInt => Data::UInt8(fit::<u8>(&values)?.into()),
4531 LogicalType::SmallInt => Data::Int16(fit::<i16>(&values)?.into()),
4532 LogicalType::USmallInt => Data::UInt16(fit::<u16>(&values)?.into()),
4533 LogicalType::Integer | LogicalType::Date => Data::Int32(fit::<i32>(&values)?.into()),
4534 LogicalType::UInteger => Data::UInt32(fit::<u32>(&values)?.into()),
4535 LogicalType::BigInt | LogicalType::Timestamp => Data::Int64(values.into()),
4536 LogicalType::Decimal { .. } => match ty.physical() {
4539 PhysicalType::Int16 => Data::Int16(fit::<i16>(&values)?.into()),
4540 PhysicalType::Int32 => Data::Int32(fit::<i32>(&values)?.into()),
4541 PhysicalType::Int64 => Data::Int64(values.into()),
4542 _ => return Err(invalid("cascade codec belongs to a decimal that is not an integer")),
4543 },
4544 _ => return Err(invalid("cascade codec belongs to a page that is not integers")),
4545 })
4546}
4547
4548fn plain_width(ty: &LogicalType) -> Option<usize> {
4551 Some(match ty {
4552 LogicalType::TinyInt | LogicalType::UTinyInt => 1,
4553 LogicalType::SmallInt | LogicalType::USmallInt => 2,
4554 LogicalType::Integer | LogicalType::UInteger | LogicalType::Date => 4,
4555 LogicalType::BigInt | LogicalType::Timestamp => 8,
4556 LogicalType::Decimal { .. } => match ty.physical() {
4557 PhysicalType::Int16 => 2,
4558 PhysicalType::Int32 => 4,
4559 PhysicalType::Int64 => 8,
4560 _ => return None,
4563 },
4564 _ => return None,
4565 })
4566}
4567
4568fn cascaded(
4574 flat: &Vector,
4575 ty: &LogicalType,
4576 packed: Option<&Packed<'_>>,
4577) -> Result<Option<Vec<u8>>> {
4578 let (Some(width), Some(data)) = (plain_width(ty), flat.data()) else { return Ok(None) };
4579 let Some(values) = widened(data) else { return Ok(None) };
4580 let plain = values.len().saturating_mul(width);
4581 let best = match packed {
4582 Some(packed) => plain.min(21 + size_of_val(packed.words())),
4584 None => plain,
4585 };
4586 let out = integer::encode_with(&values, &Fixed)?;
4587 Ok((out.len() < best).then_some(out))
4588}
4589
4590fn text_compressed(flat: &Vector) -> Result<Option<Vec<u8>>> {
4628 let mut values: Vec<&[u8]> = Vec::with_capacity(flat.len());
4629 let mut payload = 0_usize;
4630 for row in 0..flat.len() {
4631 let text = flat.text_at(row).unwrap_or("").as_bytes();
4632 payload = payload.saturating_add(text.len());
4633 values.push(text);
4634 }
4635 let plain = (flat.len() + 1).saturating_mul(4).saturating_add(payload);
4637 let Some(out) = string::encode_only(string::Kind::Fsst, &values)? else {
4638 return Ok(None);
4639 };
4640 Ok((out.len() < plain).then_some(out))
4641}
4642
4643fn encoded_codes(codes: &[u32]) -> Result<Option<Vec<u8>>> {
4644 let wide: Vec<i64> = codes.iter().map(|code| i64::from(*code)).collect();
4645 let coded = integer::encode_with(&wide, &Codes)?;
4646 let plain = codes.len().saturating_mul(size_of::<u32>());
4647 Ok((coded.len() < plain).then_some(coded))
4648}
4649
4650fn encode(
4651 vector: &Vector,
4652 global: Option<&mut GlobalDictionary>,
4653) -> Result<(Vec<u8>, Option<Vec<u32>>)> {
4654 let ty = vector.logical_type();
4655 let flat = vector.flatten()?;
4657 let mut out = Vec::new();
4658 let mut global_codes = None;
4659 if let Some(global) = global {
4660 let mut codes = Vec::with_capacity(flat.len());
4661 for row in 0..flat.len() {
4662 let text = flat.text_at(row).unwrap_or("");
4663 let code = global.code(text)?;
4664 global.observe(code, flat.is_null_at(row))?;
4665 codes.push(code);
4666 }
4667 global_codes = Some(codes);
4668 }
4669 let membership = global_codes.as_deref().map(unique_codes);
4670 let dictionary = if global_codes.is_none() && ty == &LogicalType::Varchar {
4671 string_dictionary(&flat)?
4672 } else {
4673 None
4674 };
4675 let compressed_text =
4676 if global_codes.is_none() && dictionary.is_none() && ty == &LogicalType::Varchar {
4677 text_compressed(&flat)?
4678 } else {
4679 None
4680 };
4681 let packed_vector = if dictionary.is_none() && global_codes.is_none() {
4682 Some(flat.bit_packed()?)
4683 } else {
4684 None
4685 };
4686 let packed = packed_vector.as_ref().and_then(Vector::packed_parts);
4687 let coded = match global_codes.as_deref() {
4688 Some(codes) => encoded_codes(codes)?,
4689 None => None,
4690 };
4691 let cascade = if dictionary.is_none() && global_codes.is_none() {
4695 cascaded(&flat, ty, packed.as_ref())?
4696 } else {
4697 None
4698 };
4699 out.push(if coded.is_some() {
4700 4
4701 } else if cascade.is_some() {
4702 5
4703 } else if global_codes.is_some() {
4704 3
4705 } else if dictionary.is_some() {
4706 1
4707 } else if compressed_text.is_some() {
4708 6
4709 } else if packed.is_some() {
4710 2
4711 } else {
4712 0
4713 });
4714 let nulls = flat.validity();
4715 let flag = match nulls {
4716 Validity::AllValid => 0,
4717 Validity::AllInvalid => 1,
4718 Validity::Mask(_) => 2,
4719 };
4720 out.push(flag);
4721 if flag == 2 {
4722 for group in (0..vector.len()).step_by(8) {
4723 let mut bits = 0_u8;
4724 for bit in 0..8 {
4725 if group + bit < vector.len() && !flat.is_null_at(group + bit) {
4726 bits |= 1 << bit;
4727 }
4728 }
4729 out.push(bits);
4730 }
4731 }
4732 if let Some(coded) = coded {
4733 out.extend_from_slice(&coded);
4734 return Ok((out, membership));
4735 }
4736 if let Some(cascade) = cascade {
4737 out.extend_from_slice(&cascade);
4738 return Ok((out, membership));
4739 }
4740 if let Some(codes) = global_codes {
4741 for code in codes {
4742 put_u32(&mut out, code);
4743 }
4744 return Ok((out, membership));
4745 }
4746 if let Some(dictionary) = dictionary {
4747 out.extend_from_slice(&dictionary);
4748 return Ok((out, membership));
4749 }
4750 if let Some(compressed_text) = compressed_text {
4751 out.extend_from_slice(&compressed_text);
4752 return Ok((out, membership));
4753 }
4754 if let Some(packed) = packed {
4755 if packed.offset() != 0 {
4756 return Err(invalid("writer received a sliced packed vector"));
4757 }
4758 out.push(u8::try_from(packed.width()).map_err(|_| invalid("packed width overflow"))?);
4759 out.extend_from_slice(&packed.base().to_le_bytes());
4760 put_u32(
4761 &mut out,
4762 u32::try_from(packed.words().len()).map_err(|_| invalid("too many packed words"))?,
4763 );
4764 for word in packed.words() {
4765 put_u64(&mut out, *word);
4766 }
4767 return Ok((out, membership));
4768 }
4769 let data = flat.data().ok_or_else(|| invalid("scalar column did not flatten"))?;
4770 match (ty, data) {
4771 (LogicalType::TinyInt, Data::Int8(values)) => {
4772 for value in &**values {
4773 out.extend_from_slice(&value.to_le_bytes());
4774 }
4775 }
4776 (LogicalType::UTinyInt, Data::UInt8(values)) => {
4777 for value in &**values {
4778 out.extend_from_slice(&value.to_le_bytes());
4779 }
4780 }
4781 (LogicalType::SmallInt, Data::Int16(values)) => {
4782 for value in &**values {
4783 out.extend_from_slice(&value.to_le_bytes());
4784 }
4785 }
4786 (LogicalType::USmallInt, Data::UInt16(values)) => {
4787 for value in &**values {
4788 out.extend_from_slice(&value.to_le_bytes());
4789 }
4790 }
4791 (LogicalType::UInteger, Data::UInt32(values)) => {
4792 for value in &**values {
4793 out.extend_from_slice(&value.to_le_bytes());
4794 }
4795 }
4796 (LogicalType::UBigInt, Data::UInt64(values)) => {
4797 for value in &**values {
4798 out.extend_from_slice(&value.to_le_bytes());
4799 }
4800 }
4801 (LogicalType::Integer | LogicalType::Date, Data::Int32(values)) => {
4802 for value in &**values {
4803 out.extend_from_slice(&value.to_le_bytes());
4804 }
4805 }
4806 (LogicalType::BigInt | LogicalType::Timestamp, Data::Int64(values)) => {
4807 for value in &**values {
4808 out.extend_from_slice(&value.to_le_bytes());
4809 }
4810 }
4811 (LogicalType::Boolean, Data::Bool(values)) => {
4812 for value in &**values {
4813 out.push(u8::from(*value));
4814 }
4815 }
4816 (LogicalType::Decimal { .. }, Data::Int16(values)) => {
4819 for value in &**values {
4820 out.extend_from_slice(&value.to_le_bytes());
4821 }
4822 }
4823 (LogicalType::Decimal { .. }, Data::Int32(values)) => {
4824 for value in &**values {
4825 out.extend_from_slice(&value.to_le_bytes());
4826 }
4827 }
4828 (LogicalType::Decimal { .. }, Data::Int64(values)) => {
4829 for value in &**values {
4830 out.extend_from_slice(&value.to_le_bytes());
4831 }
4832 }
4833 (LogicalType::Decimal { .. }, Data::Int128(values)) => {
4834 for value in &**values {
4835 out.extend_from_slice(&value.to_le_bytes());
4836 }
4837 }
4838 (LogicalType::Varchar, Data::Varlen(values)) => {
4839 let mut bytes = Vec::new();
4840 put_u32(&mut out, 0);
4841 for row in 0..vector.len() {
4842 let value = values.bytes(row).ok_or_else(|| invalid("string view is invalid"))?;
4843 bytes.extend_from_slice(value);
4844 put_u32(
4845 &mut out,
4846 u32::try_from(bytes.len())
4847 .map_err(|_| invalid("string payload exceeds 4GiB"))?,
4848 );
4849 }
4850 out.extend_from_slice(&bytes);
4851 }
4852 _ => return Err(Error::not_implemented(format!("native page for {ty}"))),
4853 }
4854 Ok((out, membership))
4855}
4856
4857fn put_varint(out: &mut Vec<u8>, mut value: u32) {
4858 while value >= 0x80 {
4859 out.push((value as u8 & 0x7f) | 0x80);
4860 value >>= 7;
4861 }
4862 out.push(value as u8);
4863}
4864
4865fn unique_codes(codes: &[u32]) -> Vec<u32> {
4867 let mut unique = codes.to_vec();
4868 unique.sort_unstable();
4869 unique.dedup();
4870 unique
4871}
4872
4873fn merged_codes(lists: Vec<Vec<u32>>) -> Vec<u32> {
4879 let mut lists = lists;
4880 while lists.len() > 1 {
4881 let mut next = Vec::with_capacity(lists.len().div_ceil(2));
4882 for pair in lists.chunks(2) {
4883 match pair {
4884 [left, right] => next.push(merged_pair(left, right)),
4885 [only] => next.push(only.clone()),
4886 _ => {}
4887 }
4888 }
4889 lists = next;
4890 }
4891 lists.pop().unwrap_or_default()
4892}
4893
4894fn merged_pair(left: &[u32], right: &[u32]) -> Vec<u32> {
4895 let mut out = Vec::with_capacity(left.len().saturating_add(right.len()));
4896 let mut at = 0;
4897 let mut to = 0;
4898 while at < left.len() && to < right.len() {
4899 match left[at].cmp(&right[to]) {
4900 Ordering::Less => {
4901 out.push(left[at]);
4902 at += 1;
4903 }
4904 Ordering::Greater => {
4905 out.push(right[to]);
4906 to += 1;
4907 }
4908 Ordering::Equal => {
4909 out.push(left[at]);
4910 at += 1;
4911 to += 1;
4912 }
4913 }
4914 }
4915 out.extend_from_slice(&left[at..]);
4916 out.extend_from_slice(&right[to..]);
4917 out
4918}
4919
4920fn merged_range(ranges: impl Iterator<Item = Range>) -> Range {
4925 let mut merged = Range::default();
4926 let mut first = true;
4927 for range in ranges {
4928 merged.nulls = merged.nulls.saturating_add(range.nulls);
4929 merged.sum = match (merged.sum.take(), range.sum) {
4933 (Some(held), Some(next)) if !first => held.checked_add(next),
4934 (_, next) if first => next,
4935 _ => None,
4936 };
4937 merged.exact = if first { range.exact } else { merged.exact && range.exact };
4938 if first {
4939 merged.low = range.low;
4940 merged.high = range.high;
4941 first = false;
4942 continue;
4943 }
4944 merged.low = match (merged.low.take(), range.low) {
4945 (Some(held), Some(next)) => Some(held.smaller(next)),
4946 _ => None,
4947 };
4948 merged.high = match (merged.high.take(), range.high) {
4949 (Some(held), Some(next)) => Some(held.larger(next)),
4950 _ => None,
4951 };
4952 }
4953 merged
4954}
4955
4956fn shortened(bound: Option<Bound>, high: bool) -> Option<Bound> {
4969 match bound {
4970 Some(Bound::Bytes(mut value)) if value.len() > PART_BOUND_BYTES => {
4971 value.truncate(PART_BOUND_BYTES);
4972 if !high {
4973 return Some(Bound::Bytes(value));
4974 }
4975 while let Some(last) = value.pop() {
4976 if last < u8::MAX {
4977 value.push(last + 1);
4978 return Some(Bound::Bytes(value));
4979 }
4980 }
4981 None
4982 }
4983 other => other,
4984 }
4985}
4986
4987fn encode_part_ranges(ranges: &[Range]) -> Result<Vec<u8>> {
4995 let mut out = Vec::new();
4996 put_u32(
4997 &mut out,
4998 u32::try_from(ranges.len()).map_err(|_| invalid("too many parts in a stripe"))?,
4999 );
5000 for range in ranges {
5001 put_bound(&mut out, shortened(range.low.clone(), false).as_ref())?;
5002 put_bound(&mut out, shortened(range.high.clone(), true).as_ref())?;
5003 put_u32(&mut out, u32::try_from(range.nulls).map_err(|_| invalid("null count overflow"))?);
5004 }
5005 Ok(out)
5006}
5007
5008fn decode_part_ranges(bytes: &[u8]) -> Result<Vec<Range>> {
5010 let mut cur = Cursor { bytes, at: 0 };
5011 let parts = cur.u32()? as usize;
5012 let mut out = Vec::new();
5013 for _ in 0..parts {
5014 let low = cur.bound()?;
5015 let high = cur.bound()?;
5016 let nulls = cur.u32()? as usize;
5017 out.push(Range { low, high, nulls, exact: false, sum: None });
5018 }
5019 Ok(out)
5020}
5021
5022fn encode_sieves<'a>(sieves: impl Iterator<Item = &'a Option<Sieve>>) -> Result<Vec<u8>> {
5023 let held: Vec<&Option<Sieve>> = sieves.collect();
5024 let mut out = Vec::new();
5025 put_u32(
5026 &mut out,
5027 u32::try_from(held.len()).map_err(|_| invalid("too many parts in a stripe"))?,
5028 );
5029 for sieve in &held {
5030 let length = sieve.as_ref().map_or(0, Sieve::len);
5031 put_u32(&mut out, u32::try_from(length).map_err(|_| invalid("sieve length overflow"))?);
5032 }
5033 for sieve in held.into_iter().flatten() {
5035 out.extend_from_slice(&sieve.to_bytes());
5036 }
5037 Ok(out)
5038}
5039
5040fn decode_sieves(bytes: &[u8]) -> Result<Vec<Option<Sieve>>> {
5046 let parts = u32::from_le_bytes(
5047 bytes
5048 .get(..4)
5049 .ok_or_else(|| invalid("sieve page is truncated"))?
5050 .try_into()
5051 .map_err(|_| invalid("sieve page is truncated"))?,
5052 ) as usize;
5053 let mut lengths = Vec::with_capacity(parts);
5054 for part in 0..parts {
5055 let at = 4 + part * 4;
5056 let field = bytes.get(at..at + 4).ok_or_else(|| invalid("sieve page is truncated"))?;
5057 lengths.push(u32::from_le_bytes(
5058 field.try_into().map_err(|_| invalid("sieve page is truncated"))?,
5059 ) as usize);
5060 }
5061 let mut at = 4 + parts * 4;
5062 let mut out = Vec::with_capacity(parts);
5063 for length in lengths {
5064 if length == 0 {
5065 out.push(None);
5066 continue;
5067 }
5068 let end = at.checked_add(length).ok_or_else(|| invalid("sieve page is truncated"))?;
5069 let field = bytes.get(at..end).ok_or_else(|| invalid("sieve page is truncated"))?;
5070 out.push(Sieve::from_bytes(field));
5071 at = end;
5072 }
5073 if at != bytes.len() {
5074 return Err(invalid("sieve page has trailing bytes"));
5075 }
5076 Ok(out)
5077}
5078
5079fn encode_membership(unique: &[u32]) -> Vec<u8> {
5085 let mut out = Vec::with_capacity(unique.len().saturating_mul(2).saturating_add(5));
5086 put_varint(&mut out, u32::try_from(unique.len()).unwrap_or(u32::MAX));
5087 let mut previous = 0;
5088 for (at, &code) in unique.iter().enumerate() {
5089 put_varint(&mut out, if at == 0 { code } else { code - previous });
5090 previous = code;
5091 }
5092 out
5093}
5094
5095fn take_varint(bytes: &[u8], at: &mut usize) -> Result<u32> {
5096 let mut value = 0_u32;
5097 for shift in (0..35).step_by(7) {
5098 let byte = *bytes.get(*at).ok_or_else(|| invalid("membership varint is truncated"))?;
5099 *at += 1;
5100 let part = u32::from(byte & 0x7f);
5101 if shift == 28 && part > 0x0f {
5102 return Err(invalid("membership varint overflow"));
5103 }
5104 value = value
5105 .checked_add(
5106 part.checked_shl(shift).ok_or_else(|| invalid("membership varint overflow"))?,
5107 )
5108 .ok_or_else(|| invalid("membership varint overflow"))?;
5109 if byte & 0x80 == 0 {
5110 return Ok(value);
5111 }
5112 }
5113 Err(invalid("membership varint is too long"))
5114}
5115
5116fn decode_membership(bytes: &[u8]) -> Result<Vec<u32>> {
5117 let mut at = 0;
5118 let count = take_varint(bytes, &mut at)? as usize;
5119 let mut codes = Vec::with_capacity(count);
5120 let mut previous = 0_u32;
5121 for index in 0..count {
5122 let delta = take_varint(bytes, &mut at)?;
5123 let code = if index == 0 {
5124 delta
5125 } else {
5126 previous.checked_add(delta).ok_or_else(|| invalid("membership code overflow"))?
5127 };
5128 if index > 0 && code <= previous {
5129 return Err(invalid("membership codes are not increasing"));
5130 }
5131 codes.push(code);
5132 previous = code;
5133 }
5134 if at != bytes.len() {
5135 return Err(invalid("membership page has trailing bytes"));
5136 }
5137 Ok(codes)
5138}
5139
5140fn string_dictionary(vector: &Vector) -> Result<Option<Vec<u8>>> {
5141 let mut by_text = HashMap::new();
5142 let mut values = Vec::new();
5143 let mut codes = Vec::with_capacity(vector.len());
5144 let mut plain_bytes = 0_usize;
5145 for row in 0..vector.len() {
5146 let text = vector.text_at(row).unwrap_or("");
5147 plain_bytes = plain_bytes.saturating_add(text.len());
5148 let code = match by_text.get(text) {
5149 Some(&code) => code,
5150 None => {
5151 let code = u32::try_from(values.len())
5152 .map_err(|_| invalid("too many dictionary values"))?;
5153 by_text.insert(text, code);
5154 values.push(text);
5155 code
5156 }
5157 };
5158 codes.push(code);
5159 }
5160 let dictionary_bytes = values.iter().map(|value| value.len()).sum::<usize>();
5161 let encoded = 8_usize
5162 .saturating_add((values.len() + 1).saturating_mul(4))
5163 .saturating_add(dictionary_bytes)
5164 .saturating_add(codes.len().saturating_mul(4));
5165 let plain = (vector.len() + 1).saturating_mul(4).saturating_add(plain_bytes);
5166 if encoded >= plain {
5167 return Ok(None);
5168 }
5169 let mut out = Vec::with_capacity(encoded);
5170 put_u32(
5171 &mut out,
5172 u32::try_from(values.len()).map_err(|_| invalid("too many dictionary values"))?,
5173 );
5174 put_u32(
5175 &mut out,
5176 u32::try_from(dictionary_bytes).map_err(|_| invalid("dictionary payload exceeds 4GiB"))?,
5177 );
5178 let mut offset = 0_u32;
5179 put_u32(&mut out, offset);
5180 for value in &values {
5181 offset = offset
5182 .checked_add(
5183 u32::try_from(value.len()).map_err(|_| invalid("dictionary value is too long"))?,
5184 )
5185 .ok_or_else(|| invalid("dictionary payload exceeds 4GiB"))?;
5186 put_u32(&mut out, offset);
5187 }
5188 for value in values {
5189 out.extend_from_slice(value.as_bytes());
5190 }
5191 for code in codes {
5192 put_u32(&mut out, code);
5193 }
5194 Ok(Some(out))
5195}
5196
5197struct EncodedDictionary {
5198 index: Vec<u8>,
5199 ranks: Vec<u8>,
5200 payload: Vec<Vec<u8>>,
5203}
5204
5205fn head(bytes: &[u8]) -> u64 {
5207 let mut word = [0; 8];
5208 let take = bytes.len().min(8);
5209 word[..take].copy_from_slice(&bytes[..take]);
5210 u64::from_be_bytes(word)
5211}
5212
5213fn rankings(dictionaries: &[Option<GlobalDictionary>]) -> Result<Vec<Vec<(u64, u32)>>> {
5221 let present =
5222 dictionaries.iter().enumerate().filter(|(_, held)| held.is_some()).map(|(at, _)| at);
5223 let present = present.collect::<Vec<_>>();
5224 let mut orders = vec![Vec::new(); dictionaries.len()];
5225 let workers = std::thread::available_parallelism()
5226 .map_or(1, usize::from)
5227 .min(MAX_FREQUENCY_WORKERS)
5228 .min(present.len());
5229 if workers <= 1 {
5230 for at in present {
5231 if let Some(dictionary) = &dictionaries[at] {
5232 orders[at] = dictionary.ranked();
5233 }
5234 }
5235 return Ok(orders);
5236 }
5237 let width = present.len().div_ceil(workers);
5238 let pieces = std::thread::scope(|scope| {
5239 present
5240 .chunks(width)
5241 .map(|columns| {
5242 scope.spawn(|| {
5243 columns
5244 .iter()
5245 .filter_map(|&at| dictionaries[at].as_ref().map(|held| (at, held.ranked())))
5246 .collect::<Vec<_>>()
5247 })
5248 })
5249 .collect::<Vec<_>>()
5250 .into_iter()
5251 .map(|handle| {
5252 handle.join().map_err(|_| Error::internal("a dictionary sort worker panicked"))
5253 })
5254 .collect::<Result<Vec<_>>>()
5255 })?;
5256 for piece in pieces {
5257 for (at, order) in piece {
5258 orders[at] = order;
5259 }
5260 }
5261 Ok(orders)
5262}
5263
5264fn encode_global_dictionary(
5265 dictionary: GlobalDictionary,
5266 order: &[(u64, u32)],
5267) -> Result<EncodedDictionary> {
5268 let values = dictionary.offsets.len() - 1;
5269 if order.len() != values {
5270 return Err(invalid("global dictionary order does not cover its values"));
5271 }
5272 let blocks = values.div_ceil(TEXT_PAYLOAD_VALUES);
5273 let payload = encode_payload(&dictionary)?;
5274 if payload.len() != blocks {
5275 return Err(invalid("global dictionary payload is not the blocks it says it is"));
5276 }
5277 let (ranks, rank_ends) = encode_ranks(order, code_width(values))?;
5278 let rank_blocks = values.div_ceil(TEXT_RANK_BLOCK);
5279 let offset_bits = offset_width(&dictionary.offsets);
5280 let mut index = Vec::with_capacity(
5281 DICTIONARY_HEADER + offset_bytes(values, offset_bits) + (blocks + rank_blocks) * 16,
5282 );
5283 put_u32(
5284 &mut index,
5285 u32::try_from(values).map_err(|_| invalid("global dictionary has too many values"))?,
5286 );
5287 put_u32(&mut index, TEXT_PAYLOAD_VALUES as u32);
5288 put_u32(
5289 &mut index,
5290 u32::try_from(blocks).map_err(|_| invalid("global dictionary has too many blocks"))?,
5291 );
5292 put_u32(&mut index, offset_bits as u32);
5293 encode_offsets(&dictionary.offsets, offset_bits, &mut index)?;
5294 let mut at = 0_u64;
5298 for block in &payload {
5299 at = at
5300 .checked_add(block.len() as u64)
5301 .ok_or_else(|| invalid("global dictionary payload overflow"))?;
5302 put_u64(&mut index, at);
5303 }
5304 for block in &payload {
5305 put_u64(&mut index, checksum(block));
5306 }
5307 if rank_ends.len() != rank_blocks {
5310 return Err(invalid("global dictionary order is not the blocks it says it is"));
5311 }
5312 for end in &rank_ends {
5313 put_u64(&mut index, *end);
5314 }
5315 let mut at = 0_usize;
5316 for end in &rank_ends {
5317 let end = usize::try_from(*end).map_err(|_| invalid("global dictionary order overflow"))?;
5318 put_u64(&mut index, checksum(&ranks[at..end]));
5319 at = end;
5320 }
5321 Ok(EncodedDictionary { index, ranks, payload })
5322}
5323
5324const PAYLOAD_SAMPLE_BLOCKS: usize = 8;
5331
5332fn payload_shapes() -> Vec<chooser::Settled> {
5358 let integers = vec![integer::Kind::Packed];
5359 [
5360 vec![string::Kind::Front, string::Kind::Lz],
5361 vec![string::Kind::Lz, string::Kind::Fsst],
5362 vec![string::Kind::Lz, string::Kind::Plain],
5363 vec![string::Kind::Fsst],
5364 vec![string::Kind::Plain],
5365 ]
5366 .into_iter()
5367 .map(|strings| chooser::Settled::new(strings, integers.clone()))
5368 .collect()
5369}
5370
5371fn encode_payload(dictionary: &GlobalDictionary) -> Result<Vec<Vec<u8>>> {
5377 let values = dictionary.offsets.len() - 1;
5378 let blocks = values.div_ceil(TEXT_PAYLOAD_VALUES);
5379 let run = |block: usize| {
5380 let first = block * TEXT_PAYLOAD_VALUES;
5381 let last = (first + TEXT_PAYLOAD_VALUES).min(values);
5382 (first..last)
5383 .map(|value| {
5384 let from = dictionary.offsets[value] as usize;
5385 let to = dictionary.offsets[value + 1] as usize;
5386 &dictionary.payload[from..to]
5387 })
5388 .collect::<Vec<_>>()
5389 };
5390 let shape = (blocks > PAYLOAD_SAMPLE_BLOCKS).then(|| settle_shape(&run, blocks)).transpose()?;
5393 let one = |block: usize| match &shape {
5394 Some(shape) => string::encode_with(&run(block), shape),
5395 None => string::encode(&run(block)),
5396 };
5397 let workers = std::thread::available_parallelism()
5398 .map_or(1, usize::from)
5399 .min(MAX_FREQUENCY_WORKERS)
5400 .min(blocks);
5401 if workers <= 1 {
5402 return (0..blocks).map(one).collect();
5403 }
5404 let next = AtomicUsize::new(0);
5405 let pieces = std::thread::scope(|scope| {
5406 (0..workers)
5407 .map(|_| {
5408 scope.spawn(|| {
5409 let mut mine = Vec::new();
5410 loop {
5411 let block = next.fetch_add(1, Atomic::Relaxed);
5412 if block >= blocks {
5413 break;
5414 }
5415 mine.push((block, one(block)?));
5416 }
5417 Ok(mine)
5418 })
5419 })
5420 .collect::<Vec<_>>()
5421 .into_iter()
5422 .map(|handle| {
5423 handle.join().map_err(|_| Error::internal("a dictionary encode worker panicked"))?
5424 })
5425 .collect::<Result<Vec<_>>>()
5426 })?;
5427 let mut payload = vec![Vec::new(); blocks];
5428 for piece in pieces {
5429 for (block, bytes) in piece {
5430 payload[block] = bytes;
5431 }
5432 }
5433 Ok(payload)
5434}
5435
5436fn settle_shape<'a>(
5444 run: &dyn Fn(usize) -> Vec<&'a [u8]>,
5445 blocks: usize,
5446) -> Result<chooser::Settled> {
5447 let last = blocks - 1;
5448 let sample = (0..PAYLOAD_SAMPLE_BLOCKS)
5449 .map(|region| run(region * last / (PAYLOAD_SAMPLE_BLOCKS - 1)))
5450 .collect::<Vec<_>>();
5451 let mut best: Option<(chooser::Settled, usize)> = None;
5452 for shape in payload_shapes() {
5453 let mut size = 0;
5454 for block in &sample {
5455 size += string::encode_with(block, &shape)?.len();
5456 }
5457 if best.as_ref().is_none_or(|(_, smallest)| size < *smallest) {
5458 best = Some((shape, size));
5459 }
5460 }
5461 best.map(|(shape, _)| shape)
5462 .ok_or_else(|| invalid("no shape applies to a global dictionary payload"))
5463}
5464
5465fn encode_ranks(order: &[(u64, u32)], code_bits: usize) -> Result<(Vec<u8>, Vec<u64>)> {
5472 let mut out = Vec::with_capacity(order.len() * 4);
5473 let mut ends = Vec::with_capacity(order.len().div_ceil(TEXT_RANK_BLOCK));
5474 let mut heads = Vec::with_capacity(TEXT_RANK_BLOCK);
5475 let mut codes = Vec::with_capacity(TEXT_RANK_BLOCK);
5476 for block in order.chunks(TEXT_RANK_BLOCK) {
5477 let base = block.first().map_or(0, |&(head, _)| head);
5480 let span = block.last().map_or(0, |&(head, _)| head.wrapping_sub(base));
5481 let width = (u64::BITS - span.leading_zeros()) as usize;
5482 heads.clear();
5483 codes.clear();
5484 for &(head, code) in block {
5485 heads.push(head.wrapping_sub(base));
5486 codes.push(u64::from(code));
5487 }
5488 put_u64(&mut out, base);
5489 out.push(width as u8);
5490 bitpack::pack_tail(&heads, width, &mut out)
5491 .map_err(|_| invalid("global dictionary heads do not pack"))?;
5492 bitpack::pack_tail(&codes, code_bits, &mut out)
5493 .map_err(|_| invalid("global dictionary codes do not pack"))?;
5494 ends.push(out.len() as u64);
5495 }
5496 Ok((out, ends))
5497}
5498
5499fn open_global_dictionary(
5506 file: Arc<File>,
5507 page: Page,
5508 ty: &LogicalType,
5509 keep_budget: usize,
5510) -> Result<Vector> {
5511 if ty != &LogicalType::Varchar {
5512 return Err(invalid("global dictionary belongs to a non-string column"));
5513 }
5514 let mut header = [0; DICTIONARY_HEADER];
5515 read_at(&file, page.offset, &mut header)?;
5516 let count = u32::from_le_bytes(header[0..4].try_into().expect("four bytes")) as usize;
5517 let per_block = u32::from_le_bytes(header[4..8].try_into().expect("four bytes")) as usize;
5518 let blocks = u32::from_le_bytes(header[8..12].try_into().expect("four bytes")) as usize;
5519 let offset_bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
5520 if per_block != TEXT_PAYLOAD_VALUES {
5521 return Err(invalid("global dictionary block width differs"));
5522 }
5523 if blocks != count.div_ceil(TEXT_PAYLOAD_VALUES) {
5524 return Err(invalid("global dictionary block count differs from its value count"));
5525 }
5526 if offset_bits > u32::BITS as usize {
5527 return Err(invalid("global dictionary packs offsets past a payload"));
5528 }
5529 let offset_len = offset_bytes(count, offset_bits);
5530 let ranks = count;
5535 let rank_blocks = ranks.div_ceil(TEXT_RANK_BLOCK);
5536 let hash_len = blocks
5539 .checked_add(rank_blocks)
5540 .and_then(|words| words.checked_mul(16))
5541 .ok_or_else(|| invalid("global dictionary block count overflow"))?;
5542 let index_len = DICTIONARY_HEADER
5543 .checked_add(offset_len)
5544 .and_then(|len| len.checked_add(hash_len))
5545 .ok_or_else(|| invalid("global dictionary header overflow"))?;
5546 if index_len > page.length as usize {
5547 return Err(invalid("global dictionary offset index exceeds its page"));
5548 }
5549 let mut index = vec![0; index_len];
5550 index[..DICTIONARY_HEADER].copy_from_slice(&header);
5551 read_at(&file, page.offset + DICTIONARY_HEADER as u64, &mut index[DICTIONARY_HEADER..])?;
5552 if checksum(&index) != page.hash {
5553 return Err(invalid("global dictionary index checksum differs"));
5554 }
5555 let offsets = index[DICTIONARY_HEADER..DICTIONARY_HEADER + offset_len].to_vec();
5556 let mut words = index[DICTIONARY_HEADER + offset_len..]
5557 .chunks_exact(8)
5558 .map(|part| u64::from_le_bytes(part.try_into().expect("eight bytes")))
5559 .collect::<Vec<_>>();
5560 let mut hashes = words.split_off(blocks);
5561 let mut rank_ends = hashes.split_off(blocks);
5562 let rank_hashes = rank_ends.split_off(rank_blocks);
5563 let ends = words;
5564 if rank_ends.windows(2).any(|pair| pair[0] >= pair[1]) {
5567 return Err(invalid("global dictionary order blocks do not rise"));
5568 }
5569 let rank_len = usize::try_from(rank_ends.last().copied().unwrap_or_default())
5570 .map_err(|_| invalid("global dictionary rank overflow"))?;
5571 let body_len = index_len
5572 .checked_add(rank_len)
5573 .ok_or_else(|| invalid("global dictionary header overflow"))?;
5574 if body_len > page.length as usize {
5575 return Err(invalid("global dictionary order exceeds its page"));
5576 }
5577 let stored_len = page.length as usize - body_len;
5580 if ends.last().copied().unwrap_or_default() as usize != stored_len
5581 || ends.windows(2).any(|pair| pair[0] > pair[1])
5582 {
5583 return Err(invalid("global dictionary blocks do not bound the payload"));
5584 }
5585 Vector::external_text(
5586 LogicalType::Varchar,
5587 Arc::new(NativeText {
5588 file,
5589 values: count,
5590 offsets,
5591 offset_bits,
5592 ranks,
5593 rank_at: page.offset + index_len as u64,
5594 rank_ends,
5595 rank_hashes,
5596 rank_blocks: (0..rank_blocks).map(|_| OnceLock::new()).collect(),
5597 code_bits: code_width(count),
5598 code_ranks: OnceLock::new(),
5599 payload: page.offset + body_len as u64,
5600 ends,
5601 hashes,
5602 blocks: (0..blocks).map(|_| OnceLock::new()).collect(),
5603 keep_budget,
5604 payload_kept: AtomicUsize::new(0),
5605 searched: Mutex::new(HashMap::new()),
5606 }),
5607 )
5608}
5609
5610fn decode(
5611 ty: &LogicalType,
5612 rows: usize,
5613 bytes: &[u8],
5614 global: Option<Arc<Vector>>,
5615) -> Result<Vector> {
5616 let mut cur = Cursor { bytes, at: 0 };
5617 let codec = cur.u8()?;
5618 let flag = cur.u8()?;
5619 let validity = match flag {
5620 0 => Validity::AllValid,
5621 1 => Validity::AllInvalid,
5622 2 => {
5623 let mask = cur.take(rows.div_ceil(8))?;
5624 Validity::from_iter(rows, |row| mask[row / 8] >> (row % 8) & 1 == 1)
5625 }
5626 _ => return Err(invalid("page validity tag differs")),
5627 };
5628 if codec == 1 {
5629 if ty != &LogicalType::Varchar {
5630 return Err(invalid("dictionary codec belongs to a non-string page"));
5631 }
5632 let count = cur.u32()? as usize;
5633 let payload_len = cur.u32()? as usize;
5634 let offset_bytes = cur.take(
5635 (count + 1)
5636 .checked_mul(4)
5637 .ok_or_else(|| invalid("dictionary offset count overflow"))?,
5638 )?;
5639 let offsets = offset_bytes
5640 .chunks_exact(4)
5641 .map(|part| u32::from_le_bytes(part.try_into().expect("four bytes")))
5642 .collect::<Vec<_>>();
5643 let payload = cur.take(payload_len)?.to_vec();
5644 if offsets.first() != Some(&0)
5645 || offsets.last().copied().map(|last| last as usize) != Some(payload.len())
5646 || offsets.windows(2).any(|pair| pair[0] > pair[1])
5647 {
5648 return Err(invalid("dictionary offsets do not bound the payload"));
5649 }
5650 let mut strings = StringColumn::over(Buffer::from_vec(payload).into_page());
5653 for pair in offsets.windows(2) {
5654 strings.push_in_place(pair[0] as usize, (pair[1] - pair[0]) as usize)?;
5655 }
5656 let mut codes = Vec::with_capacity(rows);
5657 for _ in 0..rows {
5658 codes.push(cur.u32()?);
5659 }
5660 if codes.iter().any(|code| *code as usize >= count) {
5661 return Err(invalid("dictionary code is out of range"));
5662 }
5663 if cur.at != bytes.len() {
5664 return Err(invalid("dictionary page has trailing bytes"));
5665 }
5666 let dictionary = Vector::flat(LogicalType::Varchar, Data::Varlen(strings))?;
5667 return Ok(Vector::dictionary(codes, dictionary)?.with_validity(validity));
5668 }
5669 if codec == 3 || codec == 4 {
5670 let dictionary = global.ok_or_else(|| invalid("global code page has no dictionary"))?;
5671 let codes = if codec == 4 {
5672 let wide = integer::decode(&bytes[cur.at..])?;
5675 if wide.len() != rows {
5676 return Err(invalid("encoded code page holds the wrong number of rows"));
5677 }
5678 let mut codes = Vec::with_capacity(wide.len());
5685 let mut seen = 0_i64;
5686 for &code in &wide {
5687 seen |= code;
5688 codes.push(code as u32);
5689 }
5690 if seen < 0 || seen > i64::from(u32::MAX) {
5691 return Err(invalid("code is not a code"));
5692 }
5693 codes
5694 } else {
5695 let mut codes = Vec::with_capacity(rows);
5696 for _ in 0..rows {
5697 codes.push(cur.u32()?);
5698 }
5699 if cur.at != bytes.len() {
5700 return Err(invalid("global code page has trailing bytes"));
5701 }
5702 codes
5703 };
5704 let highest = codes.iter().copied().max();
5705 return Ok(Vector::stable_dictionary_validated(codes, dictionary, highest)?
5706 .with_validity(validity));
5707 }
5708 if codec == 6 {
5709 if ty != &LogicalType::Varchar {
5710 return Err(invalid("compressed text codec belongs to a non-string page"));
5711 }
5712 let (payload, ends) = string::decode_flat(&bytes[cur.at..])?.into_parts();
5716 if ends.len() != rows {
5717 return Err(invalid("compressed text page holds the wrong number of rows"));
5718 }
5719 let mut values = StringColumn::over(Buffer::from_vec(payload).into_page());
5722 let mut start = 0;
5723 for end in ends {
5724 let len = end
5725 .checked_sub(start)
5726 .ok_or_else(|| invalid("compressed text value ends before it starts"))?;
5727 values.push_in_place(start, len)?;
5728 start = end;
5729 }
5730 return Ok(Vector::flat(ty.clone(), Data::Varlen(values))?.with_validity(validity));
5731 }
5732 if codec == 5 {
5733 let values = integer::decode(&bytes[cur.at..])?;
5735 if values.len() != rows {
5736 return Err(invalid("cascade page holds the wrong number of rows"));
5737 }
5738 let data = narrowed(ty, values)?;
5739 return Ok(Vector::flat(ty.clone(), data)?.with_validity(validity));
5740 }
5741 if codec == 2 {
5742 let width = u32::from(cur.u8()?);
5743 let base = i128::from_le_bytes(cur.take(16)?.try_into().expect("sixteen bytes"));
5744 let count = cur.u32()? as usize;
5745 let mut words = Vec::with_capacity(count);
5746 for _ in 0..count {
5747 words.push(cur.u64()?);
5748 }
5749 if cur.at != bytes.len() {
5750 return Err(invalid("packed page has trailing bytes"));
5751 }
5752 return Ok(Vector::packed(ty.clone(), words, width, base, rows)?.with_validity(validity));
5753 }
5754 if codec != 0 {
5755 return Err(invalid("page codec is unknown"));
5756 }
5757 let data = match ty {
5758 LogicalType::TinyInt => {
5759 let values = cur.take(rows)?;
5760 Data::Int8(values.iter().map(|item| *item as i8).collect::<Vec<_>>().into())
5761 }
5762 LogicalType::UTinyInt => Data::UInt8(cur.take(rows)?.to_vec().into()),
5763 LogicalType::SmallInt => {
5764 let values =
5765 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5766 Data::Int16(
5767 values
5768 .chunks_exact(2)
5769 .map(|item| i16::from_le_bytes(item.try_into().expect("two bytes")))
5770 .collect::<Vec<_>>()
5771 .into(),
5772 )
5773 }
5774 LogicalType::USmallInt => {
5775 let values =
5776 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5777 Data::UInt16(
5778 values
5779 .chunks_exact(2)
5780 .map(|item| u16::from_le_bytes(item.try_into().expect("two bytes")))
5781 .collect::<Vec<_>>()
5782 .into(),
5783 )
5784 }
5785 LogicalType::UInteger => {
5786 let values =
5787 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5788 Data::UInt32(
5789 values
5790 .chunks_exact(4)
5791 .map(|item| u32::from_le_bytes(item.try_into().expect("four bytes")))
5792 .collect::<Vec<_>>()
5793 .into(),
5794 )
5795 }
5796 LogicalType::UBigInt => {
5797 let values =
5798 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5799 Data::UInt64(
5800 values
5801 .chunks_exact(8)
5802 .map(|item| u64::from_le_bytes(item.try_into().expect("eight bytes")))
5803 .collect::<Vec<_>>()
5804 .into(),
5805 )
5806 }
5807 LogicalType::Integer | LogicalType::Date => {
5808 let values =
5809 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5810 Data::Int32(
5811 values
5812 .chunks_exact(4)
5813 .map(|item| i32::from_le_bytes(item.try_into().expect("four bytes")))
5814 .collect::<Vec<_>>()
5815 .into(),
5816 )
5817 }
5818 LogicalType::BigInt | LogicalType::Timestamp => {
5819 let values =
5820 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5821 Data::Int64(
5822 values
5823 .chunks_exact(8)
5824 .map(|item| i64::from_le_bytes(item.try_into().expect("eight bytes")))
5825 .collect::<Vec<_>>()
5826 .into(),
5827 )
5828 }
5829 LogicalType::Boolean => {
5830 let values = cur.take(rows)?;
5831 if values.iter().any(|value| *value > 1) {
5832 return Err(invalid("boolean page has another value"));
5833 }
5834 Data::Bool(values.iter().map(|value| *value == 1).collect::<Vec<_>>().into())
5835 }
5836 LogicalType::Decimal { .. } => match ty.physical() {
5839 PhysicalType::Int16 => {
5840 let values =
5841 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5842 Data::Int16(
5843 values
5844 .chunks_exact(2)
5845 .map(|item| i16::from_le_bytes(item.try_into().expect("two bytes")))
5846 .collect::<Vec<_>>()
5847 .into(),
5848 )
5849 }
5850 PhysicalType::Int32 => {
5851 let values =
5852 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5853 Data::Int32(
5854 values
5855 .chunks_exact(4)
5856 .map(|item| i32::from_le_bytes(item.try_into().expect("four bytes")))
5857 .collect::<Vec<_>>()
5858 .into(),
5859 )
5860 }
5861 PhysicalType::Int64 => {
5862 let values =
5863 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5864 Data::Int64(
5865 values
5866 .chunks_exact(8)
5867 .map(|item| i64::from_le_bytes(item.try_into().expect("eight bytes")))
5868 .collect::<Vec<_>>()
5869 .into(),
5870 )
5871 }
5872 _ => {
5873 let values =
5874 cur.take(rows.checked_mul(16).ok_or_else(|| invalid("page size overflow"))?)?;
5875 Data::Int128(
5876 values
5877 .chunks_exact(16)
5878 .map(|item| i128::from_le_bytes(item.try_into().expect("sixteen bytes")))
5879 .collect::<Vec<_>>()
5880 .into(),
5881 )
5882 }
5883 },
5884 LogicalType::Varchar => {
5885 let offset_bytes = cur
5886 .take((rows + 1).checked_mul(4).ok_or_else(|| invalid("offset count overflow"))?)?;
5887 let offsets = offset_bytes
5888 .chunks_exact(4)
5889 .map(|part| u32::from_le_bytes(part.try_into().expect("four bytes")))
5890 .collect::<Vec<_>>();
5891 let payload = cur.take(bytes.len() - cur.at)?.to_vec();
5892 if offsets.first() != Some(&0)
5893 || offsets.last().copied().map(|last| last as usize) != Some(payload.len())
5894 || offsets.windows(2).any(|pair| pair[0] > pair[1])
5895 {
5896 return Err(invalid("string offsets do not bound the payload"));
5897 }
5898 let mut values = StringColumn::over(Buffer::from_vec(payload).into_page());
5902 for pair in offsets.windows(2) {
5903 values.push_in_place(pair[0] as usize, (pair[1] - pair[0]) as usize)?;
5904 }
5905 Data::Varlen(values)
5906 }
5907 _ => return Err(Error::not_implemented(format!("native page for {ty}"))),
5908 };
5909 if cur.at != bytes.len() {
5910 return Err(invalid("page has trailing bytes"));
5911 }
5912 Ok(Vector::flat(ty.clone(), data)?.with_validity(validity))
5913}
5914
5915#[cfg(test)]
5916mod tests {
5917 use std::fs;
5918 use std::io::{Seek, SeekFrom, Write};
5919 use std::path::PathBuf;
5920 use std::time::{SystemTime, UNIX_EPOCH};
5921
5922 use rudb_common::Stat;
5923 use rudb_common::Value;
5924 use rudb_common::bounds::{Frequencies, Op, Zones};
5925 use rudb_common::stat::Provenance;
5926
5927 use super::*;
5928
5929 #[test]
5930 fn checksum_matches_fixed_vectors() {
5931 assert_eq!(checksum(b""), 0xef46_db37_51d8_e999);
5932 assert_eq!(checksum(b"a"), 0xd24e_c4f1_a98c_6e5b);
5933 assert_eq!(checksum(b"abc"), 0x44bc_2cf5_ad77_0999);
5934 }
5935
5936 fn path(label: &str) -> PathBuf {
5937 let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
5938 std::env::temp_dir().join(format!("rudb-native-{label}-{}-{stamp}.rdb", std::process::id()))
5939 }
5940
5941 #[test]
5943 fn a_read_at_an_offset_ignores_where_another_thread_left_the_cursor() {
5944 const SPANS: usize = 64;
5945 const SPAN: usize = 512;
5946 let path = path("positional");
5947 let content: Vec<u8> =
5948 (0..SPANS).flat_map(|span| std::iter::repeat_n(span as u8, SPAN)).collect();
5949 fs::write(&path, &content).expect("the file is written");
5950 let file = Arc::new(File::open(&path).expect("the file opens"));
5951 std::thread::scope(|scope| {
5952 for _ in 0..8 {
5953 let file = Arc::clone(&file);
5954 scope.spawn(move || {
5955 for _ in 0..64 {
5956 for span in 0..SPANS {
5957 let mut bytes = [0_u8; SPAN];
5958 read_at(&file, (span * SPAN) as u64, &mut bytes)
5959 .expect("the span reads");
5960 assert!(
5961 bytes.iter().all(|byte| *byte == span as u8),
5962 "span {span} came back as {}",
5963 bytes[0],
5964 );
5965 }
5966 }
5967 });
5968 }
5969 });
5970 let mut past = [0_u8; SPAN];
5971 let end = (SPANS * SPAN) as u64;
5972 let error = read_at(&file, end, &mut past).expect_err("a read past the end is refused");
5973 assert!(error.message().contains("ends before its declared length"), "{error}");
5974 drop(file);
5975 let _ = fs::remove_file(&path);
5976 }
5977
5978 #[test]
5984 fn a_writer_puts_a_page_where_it_said_it_did_wherever_the_cursor_has_got_to() {
5985 let path = path("cursor");
5986 let mut writer = Writer::create(
5987 &path,
5988 "items",
5989 vec![
5990 Field::required("id", LogicalType::Integer),
5991 Field::new("text", LogicalType::Varchar),
5992 ],
5993 )
5994 .expect("new file");
5995 writer.append(&sample()).expect("first part");
5996 writer.file.seek(SeekFrom::Start(0)).expect("the cursor goes back to the header");
5997 writer.append(&sample()).expect("second part");
5998 writer.file.seek(SeekFrom::Start(1)).expect("and somewhere useless again");
5999 writer.finish().expect("commit");
6000 let reader = Reader::open(&path).expect("reopen from disk");
6001 assert_eq!(reader.table().rows(), 6);
6002 let ids = reader.read(0, &[0]).expect("the integer page reads back");
6003 assert_eq!(ids.value_at(0, 0), Value::Integer(4));
6004 assert_eq!(ids.value_at(2, 0), Value::Integer(-2));
6005 let text = reader.read(1, &[1]).expect("the text page reads back");
6006 assert_eq!(text.value_at(1, 0), Value::Null);
6007 assert_eq!(text.value_at(2, 0), Value::Varchar("long text after a slash".into()));
6008 let end = reader.table().stripes().iter().flat_map(|stripe| {
6011 stripe
6012 .pages
6013 .iter()
6014 .map(|page| page.offset + u64::from(page.length))
6015 .chain(std::iter::once(stripe.index.offset + u64::from(stripe.index.length)))
6016 });
6017 let last = end.fold(HEADER, u64::max);
6018 let directory = fs::metadata(&path).expect("the file is there").len();
6019 assert!(last <= directory, "a page runs to {last} in a file of {directory} bytes");
6020 fs::remove_file(path).expect("remove scratch file");
6021 }
6022
6023 fn dictionary_index_len(header: &[u8; DICTIONARY_HEADER]) -> u64 {
6029 let count = u64::from(u32::from_le_bytes(header[0..4].try_into().expect("four bytes")));
6030 let blocks = u64::from(u32::from_le_bytes(header[8..12].try_into().expect("four bytes")));
6031 let bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
6032 let rank_blocks = count.div_ceil(TEXT_RANK_BLOCK as u64);
6033 DICTIONARY_HEADER as u64
6034 + offset_bytes(count as usize, bits) as u64
6035 + (blocks + rank_blocks) * 16
6036 }
6037
6038 fn last_rank_end(file: &File, offset: u64, header: &[u8; DICTIONARY_HEADER]) -> u64 {
6040 let count = u64::from(u32::from_le_bytes(header[0..4].try_into().expect("four bytes")));
6041 let blocks = u64::from(u32::from_le_bytes(header[8..12].try_into().expect("four bytes")));
6042 let bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
6043 let rank_blocks = count.div_ceil(TEXT_RANK_BLOCK as u64);
6044 let at = offset
6045 + DICTIONARY_HEADER as u64
6046 + offset_bytes(count as usize, bits) as u64
6047 + blocks * 16
6048 + (rank_blocks - 1) * 8;
6049 let mut end = [0; 8];
6050 read_at(file, at, &mut end).expect("the last rank block end");
6051 u64::from_le_bytes(end)
6052 }
6053
6054 fn sample() -> Chunk {
6055 Chunk::new(vec![
6056 Vector::from_values(
6057 LogicalType::Integer,
6058 &[Value::Integer(4), Value::Integer(9), Value::Integer(-2)],
6059 )
6060 .expect("integers"),
6061 Vector::from_values(
6062 LogicalType::Varchar,
6063 &[
6064 Value::Varchar("alpha".into()),
6065 Value::Null,
6066 Value::Varchar("long text after a slash".into()),
6067 ],
6068 )
6069 .expect("strings"),
6070 ])
6071 .expect("matching rows")
6072 }
6073
6074 fn sample_ids() -> Chunk {
6075 Chunk::new(vec![
6076 Vector::flat(LogicalType::Integer, Data::Int32(vec![7, 8, 9].into()))
6077 .expect("integers"),
6078 ])
6079 .expect("one column")
6080 }
6081
6082 #[test]
6083 fn the_planner_gets_the_null_count_off_the_same_directory_the_bounds_are_in() {
6084 let path = path("nulls_for_the_planner");
6087 let mut writer =
6088 Writer::create(&path, "items", vec![Field::new("a", LogicalType::Integer)])
6089 .expect("new file");
6090 let rows = Chunk::new(vec![
6091 Vector::from_values(
6092 LogicalType::Integer,
6093 &[
6094 Value::Integer(4),
6095 Value::Null,
6096 Value::Integer(9),
6097 Value::Null,
6098 Value::Integer(1),
6099 Value::Integer(2),
6100 ],
6101 )
6102 .expect("integers"),
6103 ])
6104 .expect("one column");
6105 writer.append(&rows).expect("the only part");
6106 writer.finish().expect("commit");
6107 let reader = Reader::open(&path).expect("reopen from disk");
6108 let stripes = Stripes::new(reader);
6109 let column = stripes.column("a").expect("the file has that column");
6110 assert_eq!(stripes.nulls(column), Stat::exact(2, Provenance::NullCount));
6111 assert_eq!(stripes.nulls(column + 1), Stat::Unknown);
6114 fs::remove_file(&path).expect("clean up");
6115 }
6116
6117 #[test]
6118 fn the_planner_gets_a_row_count_per_value_off_a_complete_synopsis() {
6119 let path = path("frequencies_for_the_planner");
6124 let mut writer =
6125 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6126 .expect("new file");
6127 let rows = Chunk::new(vec![
6128 Vector::from_values(
6129 LogicalType::Integer,
6130 &[
6131 Value::Integer(4),
6132 Value::Integer(4),
6133 Value::Integer(4),
6134 Value::Integer(9),
6135 Value::Integer(9),
6136 Value::Integer(1),
6137 ],
6138 )
6139 .expect("integers"),
6140 ])
6141 .expect("one column");
6142 writer.append(&rows).expect("the only part");
6143 writer.finish().expect("commit");
6144 let reader = Reader::open(&path).expect("reopen from disk");
6145 let common = Common::new(reader);
6146 assert_eq!(common.rows(), 6);
6147 let column = common.column("id").expect("the file has that column");
6148 assert_eq!(common.column("nothing"), None);
6149 assert_eq!(
6150 common.rows_with(column, &Bound::Int(4)),
6151 Stat::exact(3, Provenance::FrequencySynopsis)
6152 );
6153 assert_eq!(
6155 common.rows_with(column, &Bound::Int(7)),
6156 Stat::exact(0, Provenance::FrequencySynopsis)
6157 );
6158 assert_eq!(common.rows_with(column, &Bound::Bytes(b"four".to_vec())), Stat::Unknown);
6161 fs::remove_file(&path).expect("clean up");
6162 }
6163
6164 #[test]
6165 fn committed_file_reopens_and_reads_only_requested_columns() {
6166 let path = path("reopen");
6167 let mut writer = Writer::create(
6168 &path,
6169 "items",
6170 vec![
6171 Field::required("id", LogicalType::Integer),
6172 Field::new("text", LogicalType::Varchar),
6173 ],
6174 )
6175 .expect("new file");
6176 writer.append(&sample()).expect("first part");
6177 writer.append(&sample()).expect("second part");
6178 writer.finish().expect("commit");
6179 let reader = Reader::open(&path).expect("reopen from disk");
6180 assert_eq!(reader.table().rows(), 6);
6181 assert_eq!(reader.table().stripes().len(), 1);
6184 assert_eq!(reader.parts(), 2);
6185 assert_eq!(reader.part_rows(0), 3);
6186 assert_eq!(reader.part_rows(1), 3);
6187 let text = reader.read(1, &[1]).expect("only text page");
6188 assert_eq!(text.width(), 1);
6189 assert_eq!(text.value_at(1, 0), Value::Null);
6190 assert_eq!(text.value_at(2, 0), Value::Varchar("long text after a slash".into()));
6191 let sparse = reader.read_sparse(1, &[1]).expect("one part without its whole page");
6192 assert_eq!(sparse.width(), 1);
6193 assert_eq!(sparse.value_at(1, 0), Value::Null);
6194 assert_eq!(sparse.value_at(2, 0), Value::Varchar("long text after a slash".into()));
6195 assert!(!reader.skips_codes(0, 1, &[0]).expect("alpha is in the stripe"));
6196 assert!(!reader.skips_codes(0, 1, &[2]).expect("long text is in the stripe"));
6197 assert!(reader.skips_codes(0, 1, &[3]).expect("unknown code is absent"));
6198 let count = reader.read(0, &[]).expect("no page is needed for count");
6199 assert_eq!(count.len(), 3);
6200 assert!(reader.skips(0, &[Probe { column: 0, op: Op::Greater, value: Bound::Int(100) }]));
6201 assert!(!reader.skips(0, &[Probe { column: 0, op: Op::Greater, value: Bound::Int(0) }]));
6202 let integers = reader.top_frequencies(0, 1).expect("valid integer synopsis").expect("kept");
6203 assert_eq!(
6204 integers,
6205 vec![(Value::Integer(-2), 2), (Value::Integer(4), 2), (Value::Integer(9), 2),]
6206 );
6207 let strings = reader.top_frequencies(1, 1).expect("valid string synopsis").expect("kept");
6208 assert_eq!(strings.len(), 3);
6209 assert!(strings.contains(&(Value::Null, 2)));
6210 assert!(strings.contains(&(Value::Varchar("alpha".into()), 2)));
6211 assert!(strings.contains(&(Value::Varchar("long text after a slash".into()), 2)));
6212 fs::remove_file(path).expect("remove scratch file");
6213 }
6214
6215 #[test]
6223 fn runs_handed_over_out_of_order_still_read_back_in_source_order() {
6224 let path = path("interleaved-runs");
6225 let mut writer =
6226 Writer::create(&path, "interleaved", vec![Field::new("v", LogicalType::BigInt)])
6227 .expect("new file");
6228 for morsel in [2_u64, 0, 3, 1] {
6229 let parts = (0..4_u64)
6230 .map(|chunk| {
6231 let first = i64::try_from(morsel * 32 + chunk * 8).expect("small");
6232 let values =
6233 (0..8_i64).map(|row| Value::BigInt(first + row)).collect::<Vec<_>>();
6234 let column =
6235 Vector::from_values(LogicalType::BigInt, &values).expect("a column");
6236 ((morsel, chunk), Chunk::new(vec![column]).expect("one column"))
6237 })
6238 .collect::<Vec<_>>();
6239 writer.append_stripe(parts).expect("a stripe");
6240 }
6241 writer.finish().expect("commit");
6242
6243 let reader = Reader::open(&path).expect("valid directory");
6244 assert_eq!(reader.table().stripes().len(), 4, "a run is a stripe of its own");
6245 assert_eq!(reader.table().rows(), 128);
6246 for part in 0..16_usize {
6247 let read = reader.read(part, &[0]).expect("a part back");
6248 for row in 0..8_usize {
6249 let want = i64::try_from(part * 8 + row).expect("small");
6250 assert_eq!(read.value_at(row, 0), Value::BigInt(want), "part {part} row {row}");
6251 }
6252 }
6253 fs::remove_file(path).expect("remove scratch file");
6254 }
6255
6256 #[test]
6259 fn runs_that_overlap_each_other_are_refused_at_commit() {
6260 let path = path("overlapping-runs");
6261 let mut writer =
6262 Writer::create(&path, "overlapping", vec![Field::new("v", LogicalType::BigInt)])
6263 .expect("new file");
6264 let one = |order: (u64, u64)| {
6265 let column =
6266 Vector::from_values(LogicalType::BigInt, &[Value::BigInt(1)]).expect("a column");
6267 (order, Chunk::new(vec![column]).expect("one column"))
6268 };
6269 writer.append_stripe(vec![one((0, 0)), one((0, 2))]).expect("a stripe");
6272 writer.append_stripe(vec![one((0, 1))]).expect("a stripe");
6273 let error = writer.finish().expect_err("the runs overlap");
6274 assert!(error.message().contains("source order"), "{error}");
6275 fs::remove_file(path).expect("remove scratch file");
6276 }
6277
6278 #[test]
6281 fn a_run_longer_than_a_stripe_is_refused() {
6282 let path = path("overlong-run");
6283 let mut writer =
6284 Writer::create(&path, "overlong", vec![Field::new("v", LogicalType::BigInt)])
6285 .expect("new file");
6286 let parts = (0..=STRIPE_PARTS)
6287 .map(|at| {
6288 let column = Vector::from_values(LogicalType::BigInt, &[Value::BigInt(1)])
6289 .expect("a column");
6290 let chunk = Chunk::new(vec![column]).expect("one column");
6291 ((0, u64::try_from(at).expect("small")), chunk)
6292 })
6293 .collect::<Vec<_>>();
6294 let error = writer.append_stripe(parts).expect_err("one part too many");
6295 assert!(error.message().contains("more parts than it holds"), "{error}");
6296 fs::remove_file(path).expect("remove scratch file");
6297 }
6298
6299 #[test]
6305 fn parts_past_the_stripe_bound_start_a_new_stripe() {
6306 let path = path("stripe-bound");
6307 let mut writer = Writer::create(
6308 &path,
6309 "items",
6310 vec![
6311 Field::required("id", LogicalType::Integer),
6312 Field::new("text", LogicalType::Varchar),
6313 ],
6314 )
6315 .expect("new file");
6316 let parts = STRIPE_PARTS * 2 + 3;
6317 for part in 0..parts {
6318 let id = part as i32;
6319 let chunk = Chunk::new(vec![
6320 Vector::from_values(
6321 LogicalType::Integer,
6322 &[Value::Integer(id), Value::Integer(-id)],
6323 )
6324 .expect("integers"),
6325 Vector::from_values(
6326 LogicalType::Varchar,
6327 &[Value::Varchar(format!("value {part}")), Value::Null],
6328 )
6329 .expect("strings"),
6330 ])
6331 .expect("matching rows");
6332 writer.append(&chunk).expect("one part");
6333 }
6334 writer.finish().expect("commit");
6335
6336 let reader = Reader::open(&path).expect("reopen from disk");
6337 assert_eq!(reader.parts(), parts);
6338 assert_eq!(reader.table().rows(), parts * 2);
6339 assert_eq!(reader.table().stripes().len(), parts.div_ceil(STRIPE_PARTS));
6340 assert_eq!(reader.table().stripes()[0].parts(), STRIPE_PARTS);
6341 assert_eq!(reader.table().stripes()[0].rows(), STRIPE_PARTS * 2);
6342 assert_eq!(reader.table().stripes()[2].parts(), 3);
6343 for part in (0..parts).rev() {
6346 let dense = reader.read(part, &[0, 1]).expect("a whole page read");
6347 let sparse = reader.read_sparse(part, &[0, 1]).expect("one part read");
6348 for chunk in [&dense, &sparse] {
6349 assert_eq!(chunk.len(), 2, "part {part} has its own row count");
6350 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6351 assert_eq!(chunk.value_at(1, 0), Value::Integer(-(part as i32)));
6352 assert_eq!(chunk.value_at(0, 1), Value::Varchar(format!("value {part}")));
6353 assert_eq!(chunk.value_at(1, 1), Value::Null);
6354 }
6355 }
6356 let above = [Probe { column: 0, op: Op::Greater, value: Bound::Int(100) }];
6359 assert!(reader.skips(0, &above), "the first stripe stops at 63");
6360 assert!(!reader.skips(STRIPE_PARTS * 2, &above), "the third stripe reaches 130");
6361 fs::remove_file(path).expect("remove scratch file");
6362 }
6363
6364 fn scattered(n: i64) -> i64 {
6366 n.wrapping_mul(-7_046_029_254_386_353_131)
6367 }
6368
6369 #[test]
6375 fn a_part_is_skipped_when_its_sieve_does_not_hold_the_constant() {
6376 let path = path("sieve-skip");
6377 let mut writer =
6378 Writer::create(&path, "hits", vec![Field::required("id", LogicalType::BigInt)])
6379 .expect("new file");
6380 let parts = STRIPE_PARTS + 3;
6381 let per_part = 128;
6385 for part in 0..parts {
6386 let held: Vec<Value> = (0..per_part)
6387 .map(|row| Value::BigInt(scattered((part * per_part + row) as i64)))
6388 .collect();
6389 let chunk =
6390 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6391 .expect("one column");
6392 writer.append(&chunk).expect("one part");
6393 }
6394 writer.finish().expect("commit");
6395
6396 let reader = Reader::open(&path).expect("reopen from disk");
6397 let probe = |value: i64| Probe {
6398 column: 0,
6399 op: Op::Equal,
6400 value: Bound::Int(i128::from(scattered(value))),
6401 };
6402 for wanted in [0_i64, (per_part + 1) as i64, (parts * per_part - 1) as i64] {
6403 let tests = [probe(wanted)];
6404 let kept: Vec<usize> = (0..parts).filter(|&part| !reader.skips(part, &tests)).collect();
6405 let home = wanted as usize / per_part;
6406 assert!(kept.contains(&home), "the part holding {wanted} is read");
6407 assert!(kept.len() <= 2, "{wanted} keeps {kept:?}, which is more than one stray part");
6411 }
6412 let absent = [probe((parts * per_part) as i64 + 1)];
6413 let kept = (0..parts).filter(|&part| !reader.skips(part, &absent)).count();
6414 assert!(kept <= 1, "{kept} parts of {parts} kept a value no part holds");
6415 let tests = [probe(0)];
6418 assert!(
6419 reader.table().stripes().iter().all(|stripe| !stripe.zone.skips(&tests)),
6420 "the bounds rule out no stripe at all"
6421 );
6422 fs::remove_file(path).expect("remove scratch file");
6423 }
6424
6425 #[test]
6431 fn a_part_is_skipped_when_its_own_bounds_rule_out_a_comparison_the_stripe_keeps() {
6432 let path = path("part-range-skip");
6433 let mut writer =
6434 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
6435 .expect("new file");
6436 let parts = STRIPE_PARTS + 3;
6437 let per_part = 128;
6438 for part in 0..parts {
6439 let held: Vec<Value> = (0..per_part)
6443 .map(|row| {
6444 Value::BigInt((part * 1_000) as i64 + (scattered(row as i64).rem_euclid(900)))
6445 })
6446 .collect();
6447 let chunk =
6448 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6449 .expect("one column");
6450 writer.append(&chunk).expect("one part");
6451 }
6452 writer.finish().expect("commit");
6453
6454 let reader = Reader::open(&path).expect("reopen from disk");
6455 let under = [Probe { column: 0, op: Op::Less, value: Bound::Int(3_000) }];
6456 let kept: Vec<usize> = (0..parts).filter(|&part| !reader.skips(part, &under)).collect();
6457 assert_eq!(kept, vec![0, 1, 2], "only the three parts that start under three thousand");
6458 assert!(!reader.stripe_skips(0, &under), "the stripe reaches from zero and keeps itself");
6460 fs::remove_file(path).expect("remove scratch file");
6461 }
6462
6463 #[test]
6467 fn a_part_is_waved_through_when_its_own_bounds_pass_a_comparison_the_stripe_cannot() {
6468 let path = path("part-range-certain");
6469 let mut writer =
6470 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
6471 .expect("new file");
6472 let parts = STRIPE_PARTS + 3;
6473 let per_part = 128;
6474 for part in 0..parts {
6475 let held: Vec<Value> = (0..per_part)
6476 .map(|row| {
6477 Value::BigInt((part * 1_000) as i64 + (scattered(row as i64).rem_euclid(900)))
6478 })
6479 .collect();
6480 let chunk =
6481 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6482 .expect("one column");
6483 writer.append(&chunk).expect("one part");
6484 }
6485 writer.finish().expect("commit");
6486
6487 let reader = Reader::open(&path).expect("reopen from disk");
6488 let under = [Probe { column: 0, op: Op::Less, value: Bound::Int(3_000) }];
6489 let waved: Vec<usize> = (0..parts).filter(|&part| reader.certain(part, &under)).collect();
6490 assert_eq!(waved, vec![0, 1, 2], "the three parts that end under three thousand");
6491 assert!(!reader.stripe_skips(0, &under), "the stripe straddles the comparison");
6494 fs::remove_file(path).expect("remove scratch file");
6495 }
6496
6497 #[test]
6500 fn a_stripe_of_one_part_writes_no_range_page_and_a_stripe_of_many_does() {
6501 for (parts, wanted) in [(1_usize, false), (STRIPE_PARTS, true)] {
6502 let path = path("part-range-page");
6503 let mut writer =
6504 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
6505 .expect("new file");
6506 for part in 0..parts {
6507 let held: Vec<Value> = (0..128)
6508 .map(|row| {
6509 Value::BigInt((part * 1_000) as i64 + scattered(row as i64).rem_euclid(900))
6510 })
6511 .collect();
6512 let chunk = Chunk::new(vec![
6513 Vector::from_values(LogicalType::BigInt, &held).expect("numbers"),
6514 ])
6515 .expect("one column");
6516 writer.append(&chunk).expect("one part");
6517 }
6518 writer.finish().expect("commit");
6519 let reader = Reader::open(&path).expect("reopen from disk");
6520 let bytes = reader.layout().columns[0].part_ranges;
6521 assert_eq!(bytes > 0, wanted, "{parts} parts wrote {bytes} bytes of ranges");
6522 fs::remove_file(path).expect("remove scratch file");
6523 }
6524 }
6525
6526 #[test]
6529 fn a_string_end_that_is_cut_down_still_covers_the_value_it_came_from() {
6530 let long = vec![b'a'; PART_BOUND_BYTES * 2];
6531 let low = shortened(Some(Bound::Bytes(long.clone())), false).expect("a low end");
6532 let high = shortened(Some(Bound::Bytes(long.clone())), true).expect("a high end");
6533 let Bound::Bytes(low) = low else { panic!("a string stays a string") };
6534 let Bound::Bytes(high) = high else { panic!("a string stays a string") };
6535 assert!(low.len() <= PART_BOUND_BYTES && high.len() <= PART_BOUND_BYTES);
6536 assert!(low.as_slice() <= long.as_slice(), "the low end is at or under the value");
6537 assert!(high.as_slice() >= long.as_slice(), "the high end is at or over the value");
6538 }
6539
6540 #[test]
6543 fn a_string_end_with_no_room_to_step_up_gives_up_the_bound() {
6544 let long = vec![u8::MAX; PART_BOUND_BYTES * 2];
6545 assert_eq!(shortened(Some(Bound::Bytes(long.clone())), true), None);
6546 let low = shortened(Some(Bound::Bytes(long)), false).expect("a low end is still a prefix");
6547 assert_eq!(low, Bound::Bytes(vec![u8::MAX; PART_BOUND_BYTES]));
6548 }
6549
6550 #[test]
6560 fn a_sieve_larger_than_the_part_it_indexes_is_not_written() {
6561 let path = path("sieve-pays");
6562 let fields = vec![
6563 Field::required("spread", LogicalType::BigInt),
6564 Field::required("repeated", LogicalType::BigInt),
6565 ];
6566 let mut writer = Writer::create(&path, "hits", fields).expect("new file");
6567 let parts = 3;
6568 let per_part = 1024;
6569 for part in 0..parts {
6570 let base = (part * per_part) as i64;
6571 let spread: Vec<Value> =
6572 (0..per_part).map(|row| Value::BigInt(scattered(base + row as i64))).collect();
6573 let repeated: Vec<Value> =
6574 (0..per_part).map(|row| Value::BigInt(scattered((row / 256) as i64))).collect();
6575 let chunk = Chunk::new(vec![
6576 Vector::from_values(LogicalType::BigInt, &spread).expect("numbers"),
6577 Vector::from_values(LogicalType::BigInt, &repeated).expect("numbers"),
6578 ])
6579 .expect("two columns");
6580 writer.append(&chunk).expect("one part");
6581 }
6582 writer.finish().expect("commit");
6583
6584 let reader = Reader::open(&path).expect("reopen from disk");
6585 let layout = reader.layout();
6586 let spread = &layout.columns[0];
6587 let repeated = &layout.columns[1];
6588 assert!(spread.sieves > 0, "a column whose parts are worth a filter keeps one");
6589 assert_eq!(
6590 repeated.sieves, 0,
6591 "a column whose filter costs more than its parts keeps none"
6592 );
6593 for column in &layout.columns {
6596 assert!(
6597 column.sieves < column.pages,
6598 "{} spends {} on sieves over {} of data",
6599 column.name,
6600 column.sieves,
6601 column.pages
6602 );
6603 }
6604 let absent = [Probe {
6606 column: 0,
6607 op: Op::Equal,
6608 value: Bound::Int(i128::from(scattered((parts * per_part) as i64 + 1))),
6609 }];
6610 assert!((0..parts).all(|part| reader.skips(part, &absent)), "no part holds it");
6611 fs::remove_file(path).expect("remove scratch file");
6612 }
6613
6614 #[test]
6620 fn a_damaged_sieve_page_is_read_through_rather_than_refused() {
6621 let path = path("sieve-damaged");
6622 let mut writer =
6623 Writer::create(&path, "hits", vec![Field::required("id", LogicalType::BigInt)])
6624 .expect("new file");
6625 let rows = 128;
6626 let held: Vec<Value> = (0..rows).map(|row| Value::BigInt(scattered(row))).collect();
6627 let chunk =
6628 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6629 .expect("one column");
6630 writer.append(&chunk).expect("one part");
6631 writer.finish().expect("commit");
6632
6633 let page =
6634 Reader::open(&path).expect("reopen").table.stripes[0].sieves[0].expect("a sieve page");
6635 let mut file = OpenOptions::new().write(true).open(&path).expect("open the sieve page");
6636 file.seek(SeekFrom::Start(page.offset + u64::from(page.length) - 1)).expect("seek");
6637 file.write_all(&[0xff]).expect("damage one byte");
6638 drop(file);
6639
6640 let reader = Reader::open(&path).expect("reopen the damaged file");
6641 let absent =
6642 [Probe { column: 0, op: Op::Equal, value: Bound::Int(i128::from(scattered(99))) }];
6643 assert!(!reader.skips(0, &absent), "a sieve that cannot be read skips nothing");
6644 assert_eq!(
6645 reader.read(0, &[0]).expect("the rows are untouched").len(),
6646 usize::try_from(rows).expect("a small count")
6647 );
6648 fs::remove_file(path).expect("remove scratch file");
6649 }
6650
6651 #[test]
6662 fn workers_that_want_the_same_stripe_read_it_once() {
6663 let path = path("single-flight");
6664 let mut writer =
6665 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6666 .expect("new file");
6667 for part in 0..STRIPE_PARTS {
6668 let id = part as i32;
6669 let chunk = Chunk::new(vec![
6670 Vector::from_values(
6671 LogicalType::Integer,
6672 &[Value::Integer(id), Value::Integer(-id)],
6673 )
6674 .expect("integers"),
6675 ])
6676 .expect("matching rows");
6677 writer.append(&chunk).expect("one part");
6678 }
6679 writer.finish().expect("commit");
6680
6681 let reader = Reader::open(&path).expect("reopen from disk");
6682 assert_eq!(reader.table().stripes().len(), 1, "one stripe is the point of the test");
6683 let barrier = std::sync::Barrier::new(8);
6684 std::thread::scope(|scope| {
6685 for worker in 0..8 {
6686 let reader = &reader;
6687 let barrier = &barrier;
6688 scope.spawn(move || {
6689 barrier.wait();
6690 for part in (worker..STRIPE_PARTS).step_by(8) {
6691 let chunk = reader.read(part, &[0]).expect("a whole page read");
6692 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6693 assert_eq!(chunk.value_at(1, 0), Value::Integer(-(part as i32)));
6694 }
6695 });
6696 }
6697 });
6698 assert_eq!(reader.pages.load(Atomic::Relaxed), 1, "one stripe, one page read, whoever won");
6699 fs::remove_file(path).expect("remove scratch file");
6700 }
6701
6702 #[test]
6715 fn opening_costs_the_same_over_a_thousand_times_the_rows() {
6716 let opened = |label: &str, rows_per_part: i32| {
6717 let path = path(label);
6718 let mut writer =
6719 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6720 .expect("new file");
6721 for part in 0..STRIPE_PARTS * 3 {
6722 let values = (0..rows_per_part)
6726 .map(|row| {
6727 Value::Integer((part as i32 * rows_per_part + row).wrapping_mul(2_654_435))
6728 })
6729 .collect::<Vec<_>>();
6730 let chunk = Chunk::new(vec![
6731 Vector::from_values(LogicalType::Integer, &values).expect("integers"),
6732 ])
6733 .expect("matching rows");
6734 writer.append(&chunk).expect("one part");
6735 }
6736 writer.finish().expect("commit");
6737 let reader = Reader::open(&path).expect("reopen from disk");
6738 let size = fs::metadata(&path).expect("the file is there").len();
6739 let out = (reader.reads(), reader.table().stripes().len(), size);
6740 fs::remove_file(path).expect("remove scratch file");
6741 out
6742 };
6743
6744 let (thin, thin_stripes, thin_size) = opened("open-thin", 1);
6745 let (fat, fat_stripes, fat_size) = opened("open-fat", 1000);
6746 assert_eq!(
6747 thin_stripes, fat_stripes,
6748 "the same stripe count is what makes this a fair ask"
6749 );
6750 assert!(
6751 fat_size > thin_size * 50,
6752 "the fat file has to actually be larger, and it is {fat_size} against {thin_size}"
6753 );
6754
6755 assert_eq!(thin.opening.reads, fat.opening.reads, "the same reads either way");
6756 assert_eq!(thin.pages, 0, "opening read a page");
6757 assert_eq!(fat.pages, 0, "opening read a page");
6758 assert_eq!(thin.indexes, 0, "opening read an index");
6759 assert_eq!(fat.indexes, 0, "opening read an index");
6760 assert!(
6763 fat.opening.bytes < thin.opening.bytes * 2,
6764 "opening the thin file read {} bytes and the fat one read {}",
6765 thin.opening.bytes,
6766 fat.opening.bytes
6767 );
6768 }
6769
6770 #[test]
6778 fn two_opens_of_one_file_cost_the_same_and_the_second_is_not_cheaper() {
6779 let path = path("open-twice");
6780 let mut writer =
6781 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6782 .expect("new file");
6783 for part in 0..STRIPE_PARTS * 3 {
6784 let chunk = Chunk::new(vec![
6785 Vector::from_values(LogicalType::Integer, &[Value::Integer(part as i32)])
6786 .expect("integers"),
6787 ])
6788 .expect("matching rows");
6789 writer.append(&chunk).expect("one part");
6790 }
6791 writer.finish().expect("commit");
6792
6793 let first = Reader::open(&path).expect("open");
6794 for part in 0..first.parts() {
6797 first.read(part, &[0]).expect("a part");
6798 }
6799 assert!(first.reads().pages > 0, "the scan has to have read something");
6800 let second = Reader::open(&path).expect("open again");
6801
6802 assert_eq!(first.reads().opening, second.reads().opening);
6803 assert_eq!(
6804 second.reads().pages,
6805 0,
6806 "the second open read a page off the back of the first"
6807 );
6808 assert_eq!(second.reads().indexes, 0, "the second open read an index it inherited");
6809 fs::remove_file(path).expect("remove scratch file");
6810 }
6811
6812 #[test]
6820 fn an_index_is_read_once_per_stripe_however_often_the_page_is_evicted() {
6821 let path = path("index-cache");
6822 let mut writer =
6823 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6824 .expect("new file");
6825 let parts = STRIPE_PARTS * (CACHED_STRIPES_PER_COLUMN + 2);
6826 for part in 0..parts {
6827 let id = part as i32;
6828 let chunk = Chunk::new(vec![
6829 Vector::from_values(LogicalType::Integer, &[Value::Integer(id)]).expect("integers"),
6830 ])
6831 .expect("matching rows");
6832 writer.append(&chunk).expect("one part");
6833 }
6834 writer.finish().expect("commit");
6835
6836 let reader = Reader::open(&path).expect("reopen from disk");
6837 let stripes = reader.table().stripes().len();
6838 assert!(stripes > CACHED_STRIPES_PER_COLUMN, "the page cache has to be too small for this");
6839 for _ in 0..2 {
6841 for part in 0..parts {
6842 let chunk = reader.read(part, &[0]).expect("a part");
6843 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6844 }
6845 }
6846 assert_eq!(reader.indexes.load(Atomic::Relaxed), stripes, "one index read per stripe");
6847 assert!(
6848 reader.pages.load(Atomic::Relaxed) > stripes,
6849 "the pages are the ones that get read again, which is what makes the index count mean \
6850 something"
6851 );
6852 fs::remove_file(path).expect("remove scratch file");
6853 }
6854
6855 #[test]
6864 fn a_worker_per_stripe_reads_its_page_once_when_the_cache_was_told_to_expect_it() {
6865 let workers = CACHED_STRIPES_PER_COLUMN + 4;
6866 let path = path("stripe-per-worker");
6867 let mut writer =
6868 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6869 .expect("new file");
6870 for part in 0..STRIPE_PARTS * workers {
6871 let chunk = Chunk::new(vec![
6872 Vector::from_values(LogicalType::Integer, &[Value::Integer(part as i32)])
6873 .expect("integers"),
6874 ])
6875 .expect("matching rows");
6876 writer.append(&chunk).expect("one part");
6877 }
6878 writer.finish().expect("commit");
6879
6880 let read = |told: bool| {
6881 let reader = Reader::open(&path).expect("reopen from disk");
6882 assert_eq!(reader.table().stripes().len(), workers, "a stripe per worker");
6883 if told {
6884 reader.keep_stripes(workers);
6885 }
6886 let barrier = std::sync::Barrier::new(workers);
6887 std::thread::scope(|scope| {
6888 for (worker, run) in reader.stripe_parts().into_iter().enumerate() {
6889 let reader = &reader;
6890 let barrier = &barrier;
6891 scope.spawn(move || {
6892 for part in run {
6893 barrier.wait();
6894 let chunk = reader.read(part, &[0]).expect("a part of my own stripe");
6895 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6896 }
6897 assert!(worker < workers);
6898 });
6899 }
6900 });
6901 reader.pages.load(Atomic::Relaxed)
6902 };
6903
6904 assert_eq!(read(true), workers, "one page read per stripe and no more");
6905 assert!(read(false) > workers, "a cache that small is read again on every part");
6906 fs::remove_file(path).expect("remove scratch file");
6907 }
6908
6909 #[test]
6914 fn a_damaged_index_page_is_an_error() {
6915 let path = path("damaged-index");
6916 let mut writer =
6917 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6918 .expect("new file");
6919 writer.append(&sample_ids()).expect("first part");
6920 writer.append(&sample_ids()).expect("second part");
6921 writer.finish().expect("commit");
6922
6923 let reader = Reader::open(&path).expect("valid directory");
6924 let index = reader.table.stripes[0].index;
6925 let mut byte = [0; 1];
6926 read_at(&reader.file, index.offset, &mut byte).expect("the first part length");
6927 let mut file = OpenOptions::new().write(true).open(&path).expect("open index page");
6928 file.seek(SeekFrom::Start(index.offset)).expect("index start");
6929 file.write_all(&[!byte[0]]).expect("damage the first part length");
6930 let error = reader.read(1, &[0]).expect_err("a damaged index must not be used");
6931 assert!(error.message().contains("index page section checksum differs"), "{error}");
6932 fs::remove_file(path).expect("remove scratch file");
6933 }
6934
6935 #[test]
6942 fn every_integer_width_round_trips_through_a_page() {
6943 let path = path("integer-widths");
6944 let columns = [
6945 (LogicalType::TinyInt, vec![Value::TinyInt(i8::MIN), Value::TinyInt(i8::MAX)]),
6946 (LogicalType::UTinyInt, vec![Value::UTinyInt(0), Value::UTinyInt(u8::MAX)]),
6947 (LogicalType::SmallInt, vec![Value::SmallInt(i16::MIN), Value::SmallInt(i16::MAX)]),
6948 (LogicalType::USmallInt, vec![Value::USmallInt(0), Value::USmallInt(u16::MAX)]),
6949 (LogicalType::Integer, vec![Value::Integer(i32::MIN), Value::Integer(i32::MAX)]),
6950 (LogicalType::UInteger, vec![Value::UInteger(0), Value::UInteger(u32::MAX)]),
6951 (LogicalType::BigInt, vec![Value::BigInt(i64::MIN), Value::BigInt(i64::MAX)]),
6952 (LogicalType::UBigInt, vec![Value::UBigInt(0), Value::UBigInt(u64::MAX)]),
6953 ];
6954 let fields = columns
6955 .iter()
6956 .enumerate()
6957 .map(|(at, (ty, _))| Field::required(format!("c{at}"), ty.clone()))
6958 .collect::<Vec<_>>();
6959 let vectors = columns
6960 .iter()
6961 .map(|(ty, values)| Vector::from_values(ty.clone(), values).expect("a vector"))
6962 .collect::<Vec<_>>();
6963 let mut writer = Writer::create(&path, "widths", fields).expect("new file");
6964 writer.append(&Chunk::new(vectors).expect("matching rows")).expect("one stripe");
6965 writer.finish().expect("commit");
6966
6967 let reader = Reader::open(&path).expect("reopen from disk");
6968 let wanted = (0..columns.len()).collect::<Vec<_>>();
6969 let read = reader.read(0, &wanted).expect("every column");
6970 assert_eq!(read.len(), 2);
6971 for (at, (ty, values)) in columns.iter().enumerate() {
6973 assert_eq!(read.value_at(0, at), values[0], "the low end of {ty}");
6974 assert_eq!(read.value_at(1, at), values[1], "the high end of {ty}");
6975 }
6976 fs::remove_file(path).expect("remove scratch file");
6977 }
6978
6979 #[test]
6980 fn numeric_frequency_candidates_keep_bounded_row_ordinals() {
6981 let path = path("frequency-ordinals");
6982 let mut writer =
6983 Writer::create(&path, "items", vec![Field::required("id", LogicalType::BigInt)])
6984 .expect("new file");
6985 let mut values = Vec::new();
6986 for leader in 0..10_i64 {
6987 values.extend(std::iter::repeat_n(leader, 100));
6988 }
6989 values.extend(1_000_i64..41_000);
6990 for part in values.chunks(1_024) {
6991 let vector = Vector::flat(LogicalType::BigInt, Data::Int64(part.to_vec().into()))
6992 .expect("big integers");
6993 writer.append(&Chunk::new(vec![vector]).expect("one column")).expect("one stripe");
6994 }
6995 writer.finish().expect("commit");
6996
6997 let reader = Reader::open(&path).expect("reopen from disk");
6998 let occurrences =
6999 reader.frequency_occurrences(0).expect("valid metadata").expect("bounded ordinals");
7000 assert!(occurrences.omitted_max < 100);
7001 assert!(occurrences.ordinals.len() <= FREQUENCY_ORDINALS);
7002 assert!(occurrences.ordinals.windows(2).all(|pair| pair[0] < pair[1]));
7003 assert_eq!(&occurrences.ordinals[..1_000], &(0_u64..1_000).collect::<Vec<_>>());
7004 fs::remove_file(path).expect("remove scratch file");
7005 }
7006
7007 #[test]
7013 fn a_file_from_another_format_says_which_format_it_is() {
7014 let older = path("older-format");
7015 let mut writer =
7016 Writer::create(&older, "items", vec![Field::new("id", LogicalType::Integer)])
7017 .expect("new file");
7018 let chunk = Chunk::new(vec![
7019 Vector::flat(LogicalType::Integer, Data::Int32(vec![1, 2, 3].into()))
7020 .expect("integers"),
7021 ])
7022 .expect("chunk");
7023 writer.append(&chunk).expect("page written");
7024 writer.finish().expect("commit");
7025
7026 let mut file = OpenOptions::new().write(true).open(&older).expect("open for the header");
7027 file.seek(SeekFrom::Start(8)).expect("the version follows the magic");
7028 file.write_all(&(FORMAT - 1).to_le_bytes()).expect("write an older version");
7029 drop(file);
7030 let complaint = Reader::open(&older).expect_err("an older format is refused").to_string();
7031 assert!(complaint.contains(&format!("format {}", FORMAT - 1)), "{complaint}");
7032 assert!(complaint.contains(&format!("format {FORMAT}")), "{complaint}");
7033
7034 let mut file = OpenOptions::new().write(true).open(&older).expect("open for the header");
7035 file.seek(SeekFrom::Start(0)).expect("the magic is first");
7036 file.write_all(b"NOTRUDB!").expect("write another engine's magic");
7037 drop(file);
7038 let complaint = Reader::open(&older).expect_err("a foreign file is refused").to_string();
7039 assert!(complaint.contains("magic"), "{complaint}");
7040 assert!(!complaint.contains("format"), "a version has nothing to do with it: {complaint}");
7041 fs::remove_file(older).expect("remove scratch file");
7042 }
7043
7044 #[test]
7045 fn an_unfinished_or_damaged_file_does_not_answer_with_partial_rows() {
7046 let unfinished = path("unfinished");
7047 let mut writer =
7048 Writer::create(&unfinished, "items", vec![Field::new("id", LogicalType::Integer)])
7049 .expect("new file");
7050 let chunk = Chunk::new(vec![
7051 Vector::flat(LogicalType::Integer, Data::Int32(vec![1, 2, 3].into()))
7052 .expect("integers"),
7053 ])
7054 .expect("chunk");
7055 writer.append(&chunk).expect("page written");
7056 drop(writer);
7057 assert!(Reader::open(&unfinished).is_err(), "no directory was committed");
7058 fs::remove_file(unfinished).expect("remove scratch file");
7059
7060 let damaged = path("damaged");
7061 let mut writer =
7062 Writer::create(&damaged, "items", vec![Field::new("id", LogicalType::Integer)])
7063 .expect("new file");
7064 writer.append(&chunk).expect("page written");
7065 writer.finish().expect("commit");
7066 let reader = Reader::open(&damaged).expect("valid directory");
7067 let mut file =
7068 OpenOptions::new().write(true).open(&damaged).expect("open for a damaged page");
7069 file.seek(SeekFrom::Start(HEADER + 1)).expect("inside first page");
7070 file.write_all(&[255]).expect("damage one byte");
7071 assert!(reader.read(0, &[0]).is_err(), "page checksum rejects corruption");
7072 fs::remove_file(damaged).expect("remove scratch file");
7073 }
7074
7075 #[test]
7076 fn damaged_lazy_dictionary_payload_is_an_error() {
7077 let path = path("damaged-dictionary");
7078 let mut writer = Writer::create(
7079 &path,
7080 "items",
7081 vec![
7082 Field::required("id", LogicalType::Integer),
7083 Field::new("text", LogicalType::Varchar),
7084 ],
7085 )
7086 .expect("new file");
7087 writer.append(&sample()).expect("stripe written");
7088 writer.finish().expect("commit");
7089
7090 let reader = Reader::open(&path).expect("valid directory");
7091 let dictionary = reader.table.dictionaries[1].expect("string dictionary page");
7092 let mut header = [0; DICTIONARY_HEADER];
7095 read_at(&reader.file, dictionary.offset, &mut header).expect("dictionary header");
7096 let index_len = dictionary_index_len(&header);
7097 let rank_len = last_rank_end(&reader.file, dictionary.offset, &header);
7098 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
7099 file.seek(SeekFrom::Start(dictionary.offset + index_len + rank_len))
7100 .expect("inside dictionary payload");
7101 file.write_all(&[255]).expect("damage dictionary payload");
7102
7103 let chunk = reader.read(0, &[1]).expect("code page and dictionary index remain valid");
7104 let error =
7105 chunk.validate_external().expect_err("payload corruption must reach the caller");
7106 assert!(error.message().contains("payload checksum differs"), "{error}");
7107 fs::remove_file(path).expect("remove scratch file");
7108 }
7109
7110 #[test]
7120 fn a_column_of_all_different_values_is_written_without_a_dictionary() {
7121 let path = path("dictionary-decide");
7122 let rows = 20_000;
7123 let unique =
7125 |row: usize| format!("{row:09} a value that appears exactly once in the table");
7126 let repeated = |row: usize| unique(row / 40);
7128 let mut writer = Writer::create(
7129 &path,
7130 "items",
7131 vec![
7132 Field::required("unique", LogicalType::Varchar),
7133 Field::required("repeated", LogicalType::Varchar),
7134 ],
7135 )
7136 .expect("new file");
7137 for part in (0..rows).step_by(1_000) {
7138 let span = part..(part + 1_000).min(rows);
7139 let left = span.clone().map(|row| Value::Varchar(unique(row))).collect::<Vec<_>>();
7140 let right = span.map(|row| Value::Varchar(repeated(row))).collect::<Vec<_>>();
7141 writer
7142 .append(
7143 &Chunk::new(vec![
7144 Vector::from_values(LogicalType::Varchar, &left).expect("strings"),
7145 Vector::from_values(LogicalType::Varchar, &right).expect("strings"),
7146 ])
7147 .expect("two columns"),
7148 )
7149 .expect("a part");
7150 }
7151 writer.finish().expect("commit");
7152
7153 let reader = Reader::open(&path).expect("reopen from disk");
7154 assert!(
7155 reader.table.dictionaries[0].is_none(),
7156 "a column with no repeats has nothing to say twice"
7157 );
7158 assert!(
7159 reader.table.dictionaries[1].is_some(),
7160 "a column whose values come round again keeps its dictionary"
7161 );
7162 let mut first = 0;
7163 for part in 0..reader.parts() {
7164 let chunk = reader.read(part, &[0, 1]).expect("a part");
7165 for row in 0..chunk.len() {
7166 assert_eq!(chunk.value_at(row, 0), Value::Varchar(unique(first + row)));
7167 assert_eq!(chunk.value_at(row, 1), Value::Varchar(repeated(first + row)));
7168 }
7169 first += chunk.len();
7170 }
7171 assert_eq!(first, rows, "every row was read back");
7172 let raw = (0..rows).map(|row| unique(row).len()).sum::<usize>();
7173 let size = fs::metadata(&path).expect("the file is there").len() as usize;
7174 assert!(size < raw, "a column without a dictionary is still encoded: {size} against {raw}");
7175 fs::remove_file(path).expect("remove scratch file");
7176 }
7177
7178 #[test]
7191 fn a_dictionary_over_many_blocks_checks_every_block_of_it() {
7192 let path = path("dictionary-blocks");
7193 let value = |row: usize| {
7194 let row = row.saturating_sub(8_000);
7195 format!("{row:07} a value long enough to be worth a payload block")
7196 };
7197 let parts = 40;
7198 let per_part = 1000;
7199 let mut writer =
7200 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
7201 .expect("new file");
7202 for part in 0..parts {
7203 let values = (0..per_part)
7204 .map(|row| Value::Varchar(value(part * per_part + row)))
7205 .collect::<Vec<_>>();
7206 let chunk = Chunk::new(vec![
7207 Vector::from_values(LogicalType::Varchar, &values).expect("strings"),
7208 ])
7209 .expect("matching rows");
7210 writer.append(&chunk).expect("a part");
7211 }
7212 writer.finish().expect("commit");
7213
7214 let reader = Reader::open(&path).expect("reopen from disk");
7215 let dictionary = reader.table.dictionaries[0].expect("string dictionary page");
7216 assert!(
7217 parts * per_part > TEXT_PAYLOAD_VALUES * 4,
7218 "the dictionary has to be several blocks for this to be testing anything"
7219 );
7220 for part in [0, parts - 1] {
7221 let chunk = reader.read(part, &[0]).expect("a part");
7222 chunk.validate_external().expect("every payload block checks out");
7223 assert_eq!(chunk.value_at(0, 0), Value::Varchar(value(part * per_part)));
7224 }
7225
7226 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
7227 file.seek(SeekFrom::Start(dictionary.offset + u64::from(dictionary.length) - 4))
7228 .expect("the last bytes of the page are payload");
7229 file.write_all(&[255]).expect("damage the last payload block");
7230 let reader = Reader::open(&path).expect("the directory and the index are untouched");
7231 let chunk = reader.read(parts - 1, &[0]).expect("the code page remains valid");
7232 let error = chunk.validate_external().expect_err("the damage must reach the caller");
7233 assert!(error.message().contains("payload checksum differs"), "{error}");
7234 fs::remove_file(path).expect("remove scratch file");
7235 }
7236
7237 #[test]
7251 fn values_of_different_lengths_read_back_out_of_packed_offsets() {
7252 let path = path("dictionary-offsets");
7253 let value = |row: usize| {
7254 let row = row % 5_000;
7255 if row % 511 == 3 { String::new() } else { "x".repeat(row % 97) + &format!("{row:05}") }
7256 };
7257 let rows = 6_000;
7258 let mut writer =
7259 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
7260 .expect("new file");
7261 let values = (0..rows).map(|row| Value::Varchar(value(row))).collect::<Vec<_>>();
7262 for part in values.chunks(1_000) {
7263 let chunk =
7264 Chunk::new(vec![Vector::from_values(LogicalType::Varchar, part).expect("strings")])
7265 .expect("matching rows");
7266 writer.append(&chunk).expect("a part");
7267 }
7268 writer.finish().expect("commit");
7269
7270 let reader = Reader::open(&path).expect("reopen from disk");
7271 assert!(
7272 rows > TEXT_PAYLOAD_VALUES * 4,
7273 "the dictionary has to be several blocks for this to be testing anything"
7274 );
7275 for part in 0..rows / 1_000 {
7276 let chunk = reader.read(part, &[0]).expect("a part");
7277 for row in 0..1_000 {
7278 let row = part * 1_000 + row;
7279 assert_eq!(
7280 chunk.value_at(row % 1_000, 0),
7281 Value::Varchar(value(row)),
7282 "value {row}"
7283 );
7284 }
7285 }
7286 fs::remove_file(path).expect("remove scratch file");
7287 }
7288
7289 #[test]
7301 fn a_global_dictionary_is_opened_once_however_many_workers_ask_at_once() {
7302 let path = path("dictionary-once");
7303 let parts = 8;
7304 let per_part = 500;
7305 let value =
7306 |row: usize| format!("{row:07} a value long enough to be worth a payload block");
7307 let mut writer =
7308 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
7309 .expect("new file");
7310 for part in 0..parts {
7311 let values = (0..per_part)
7312 .map(|row| Value::Varchar(value(part * per_part + row)))
7313 .collect::<Vec<_>>();
7314 let chunk = Chunk::new(vec![
7315 Vector::from_values(LogicalType::Varchar, &values).expect("strings"),
7316 ])
7317 .expect("matching rows");
7318 writer.append(&chunk).expect("a part");
7319 }
7320 writer.finish().expect("commit");
7321
7322 let reader = Reader::open(&path).expect("reopen from disk");
7323 assert!(reader.table.dictionaries[0].is_some(), "the column has to have one to share");
7324 assert_eq!(reader.reads().dictionaries, 0, "opening the file does not open a dictionary");
7325
7326 let workers = 16;
7327 let gate = std::sync::Barrier::new(workers);
7328 std::thread::scope(|scope| {
7329 for worker in 0..workers {
7330 let reader = reader.clone();
7331 let gate = &gate;
7332 scope.spawn(move || {
7333 gate.wait();
7334 let chunk = reader.read(worker % parts, &[0]).expect("a part");
7335 assert_eq!(
7336 chunk.value_at(0, 0),
7337 Value::Varchar(value((worker % parts) * per_part))
7338 );
7339 });
7340 }
7341 });
7342
7343 assert_eq!(reader.reads().dictionaries, 1, "sixteen workers, one dictionary, one open");
7344 fs::remove_file(path).expect("remove scratch file");
7345 }
7346
7347 #[test]
7352 fn a_damaged_sorted_order_is_an_error() {
7353 let path = path("damaged-order");
7354 let mut writer = Writer::create(
7355 &path,
7356 "items",
7357 vec![
7358 Field::required("id", LogicalType::Integer),
7359 Field::new("text", LogicalType::Varchar),
7360 ],
7361 )
7362 .expect("new file");
7363 writer.append(&sample()).expect("stripe written");
7364 writer.finish().expect("commit");
7365
7366 let reader = Reader::open(&path).expect("valid directory");
7367 let page = reader.table.dictionaries[1].expect("string dictionary page");
7368 let mut header = [0; DICTIONARY_HEADER];
7369 read_at(&reader.file, page.offset, &mut header).expect("dictionary header");
7370 let index_len = dictionary_index_len(&header);
7371 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
7372 file.seek(SeekFrom::Start(page.offset + index_len)).expect("the first head");
7373 file.write_all(&[255]).expect("damage the order");
7374
7375 let dictionary = reader.dictionary(1).expect("read").expect("a string column has one");
7376 let error = dictionary.compare_rank(0, b"anything").expect_err("a damaged order is caught");
7377 assert!(error.message().contains("rank checksum differs"), "{error}");
7378 fs::remove_file(path).expect("remove scratch file");
7379 }
7380
7381 #[test]
7385 fn a_global_dictionary_carries_the_sorted_order_of_its_values() {
7386 let spellings = ["overlong1z", "b", "", "overlong1a", "overlong", "ab", "a", "overlong1"];
7389 let path = path("dictionary-order");
7390 let mut writer =
7391 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7392 .expect("new file");
7393 writer
7394 .append(
7395 &Chunk::new(vec![
7396 Vector::from_values(
7397 LogicalType::Varchar,
7398 &spellings.map(|text| Value::Varchar(text.into())),
7399 )
7400 .expect("strings"),
7401 ])
7402 .expect("one column"),
7403 )
7404 .expect("stripe written");
7405 writer.finish().expect("commit");
7406
7407 let reader = Reader::open(&path).expect("valid directory");
7408 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
7409 let count = dictionary.ranks().expect("a v10 file stores one");
7410 assert_eq!(count, spellings.len(), "every distinct value has a rank");
7411 let order = (0..count)
7412 .map(|rank| dictionary.code_at_rank(rank).expect("a code"))
7413 .collect::<Vec<_>>();
7414 let mut seen = order.clone();
7415 seen.sort_unstable();
7416 assert_eq!(seen, (0..spellings.len() as u32).collect::<Vec<_>>(), "a permutation of codes");
7417
7418 let ranked = order
7419 .iter()
7420 .map(|&code| {
7421 dictionary.try_bytes_at(code as usize).expect("read").expect("a value").to_vec()
7422 })
7423 .collect::<Vec<_>>();
7424 let mut expected = spellings.map(|text| text.as_bytes().to_vec()).to_vec();
7425 expected.sort();
7426 assert_eq!(ranked, expected, "rank order is value order");
7427
7428 for (rank, value) in expected.iter().enumerate() {
7431 assert_eq!(
7432 dictionary.compare_rank(rank, value).expect("compare"),
7433 Ordering::Equal,
7434 "rank {rank} is its own value"
7435 );
7436 if rank > 0 {
7437 assert_eq!(
7438 dictionary.compare_rank(rank - 1, value).expect("compare"),
7439 Ordering::Less,
7440 "rank {rank} follows the one before it"
7441 );
7442 }
7443 }
7444 fs::remove_file(path).expect("remove scratch file");
7445 }
7446
7447 #[test]
7455 fn a_dictionary_sweep_reads_every_value_and_keeps_it_under_the_budget() {
7456 let path = path("dictionary-sweep");
7457 let spellings = (0..2_500)
7460 .map(|index| Value::Varchar(format!("value {index:08} {}", "x".repeat(index % 40))))
7461 .collect::<Vec<_>>();
7462 let mut writer =
7463 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7464 .expect("new file");
7465 for part in spellings.chunks(1_024) {
7468 writer
7469 .append(
7470 &Chunk::new(vec![
7471 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
7472 ])
7473 .expect("one column"),
7474 )
7475 .expect("stripe written");
7476 }
7477 writer.finish().expect("commit");
7478
7479 let reader = Reader::open(&path).expect("valid directory");
7480 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
7481 assert_eq!(dictionary.len(), spellings.len(), "every value is distinct");
7482
7483 let resting = dictionary.footprint();
7484 let mut swept: Vec<Vec<u8>> = Vec::new();
7485 let mut at = 0;
7486 let mut calls = 0;
7487 while at < dictionary.len() {
7488 let stopped = dictionary
7489 .sweep_text(at, dictionary.len(), &mut |index: usize, text: &[u8]| {
7490 assert_eq!(index, swept.len(), "a sweep hands its values over in order");
7491 swept.push(text.to_vec());
7492 Ok(())
7493 })
7494 .expect("a sweep reads");
7495 assert!(stopped > at, "a sweep moves");
7496 at = stopped;
7497 calls += 1;
7498 }
7499 assert_eq!(calls, 3, "a sweep hands over one block at a time");
7500 let after = dictionary.footprint();
7501 assert!(after > resting, "a sweep under the budget keeps what it decoded");
7502
7503 let read = (0..dictionary.len())
7504 .map(|code| dictionary.try_bytes_at(code).expect("read").expect("a value").to_vec())
7505 .collect::<Vec<_>>();
7506 assert_eq!(swept, read, "a sweep answers what a point read answers");
7507 assert_eq!(dictionary.footprint(), after, "a point read of a kept block decodes nothing");
7508 fs::remove_file(path).expect("remove scratch file");
7509 }
7510
7511 #[test]
7522 fn a_sweep_over_a_block_with_a_short_second_run_reads_what_a_point_read_reads() {
7523 let path = path("dictionary-sweep-short-run");
7524 let spellings = (0..2_800)
7525 .map(|index| Value::Varchar(format!("value {index:08} {}", "x".repeat(index % 40))))
7526 .collect::<Vec<_>>();
7527 let mut writer =
7528 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7529 .expect("new file");
7530 for part in spellings.chunks(1_024) {
7531 writer
7532 .append(
7533 &Chunk::new(vec![
7534 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
7535 ])
7536 .expect("one column"),
7537 )
7538 .expect("stripe written");
7539 }
7540 writer.finish().expect("commit");
7541
7542 let reader = Reader::open(&path).expect("valid directory");
7543 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
7544 assert_eq!(dictionary.len(), spellings.len(), "every value is distinct");
7545 let last = dictionary.len() % TEXT_PAYLOAD_VALUES;
7546 assert!(last > TEXT_OFFSET_RUN, "the last block has to reach into a second run of offsets");
7547 assert!(last < TEXT_PAYLOAD_VALUES, "and that second run has to be short of a whole one");
7548
7549 let mut swept: Vec<Vec<u8>> = Vec::new();
7550 let mut at = 0;
7551 while at < dictionary.len() {
7552 let stopped = dictionary
7553 .sweep_text(at, dictionary.len(), &mut |index: usize, text: &[u8]| {
7554 assert_eq!(index, swept.len(), "a sweep hands its values over in order");
7555 swept.push(text.to_vec());
7556 Ok(())
7557 })
7558 .expect("a sweep reads");
7559 assert!(stopped > at, "a sweep moves");
7560 at = stopped;
7561 }
7562 let read = (0..dictionary.len())
7563 .map(|code| dictionary.try_bytes_at(code).expect("read").expect("a value").to_vec())
7564 .collect::<Vec<_>>();
7565 assert_eq!(swept, read, "a sweep answers what a point read answers");
7566 fs::remove_file(path).expect("remove scratch file");
7567 }
7568
7569 #[test]
7579 fn narrowing_a_page_takes_what_fits_and_refuses_what_does_not() {
7580 assert_eq!(fit::<i8>(&[]).expect("an empty page fits anything"), Vec::<i8>::new());
7581 assert_eq!(fit::<i8>(&[-128, 0, 127]).expect("the edges fit"), vec![-128_i8, 0, 127]);
7582 fit::<i8>(&[128]).expect_err("one past the top does not fit");
7583 fit::<i8>(&[-129]).expect_err("one past the bottom does not fit");
7584 assert_eq!(fit::<u8>(&[0, 255]).expect("the edges fit"), vec![0_u8, 255]);
7585 fit::<u8>(&[256]).expect_err("one past the top does not fit");
7586 fit::<u8>(&[-1]).expect_err("a negative does not fit an unsigned page");
7587 assert_eq!(
7588 fit::<i16>(&[-32_768, 0, 32_767]).expect("the edges fit"),
7589 vec![-32_768_i16, 0, 32_767]
7590 );
7591 fit::<i16>(&[32_768]).expect_err("one past the top does not fit");
7592 fit::<i16>(&[-32_769]).expect_err("one past the bottom does not fit");
7593 assert_eq!(fit::<u16>(&[0, 65_535]).expect("the edges fit"), vec![0_u16, 65_535]);
7594 fit::<u16>(&[65_536]).expect_err("one past the top does not fit");
7595 fit::<u16>(&[-1]).expect_err("a negative does not fit an unsigned page");
7596 assert_eq!(
7597 fit::<i32>(&[i64::from(i32::MIN), 0, i64::from(i32::MAX)]).expect("the edges fit"),
7598 vec![i32::MIN, 0, i32::MAX]
7599 );
7600 fit::<i32>(&[i64::from(i32::MAX) + 1]).expect_err("one past the top does not fit");
7601 fit::<i32>(&[i64::from(i32::MIN) - 1]).expect_err("one past the bottom does not fit");
7602 assert_eq!(
7603 fit::<u32>(&[0, 4_294_967_295]).expect("the edges fit"),
7604 vec![0_u32, 4_294_967_295]
7605 );
7606 fit::<u32>(&[4_294_967_296]).expect_err("one past the top does not fit");
7607 fit::<u32>(&[-1]).expect_err("a negative does not fit an unsigned page");
7608
7609 fit::<i8>(&[0, 1, 2, 128, 3]).expect_err("one bad value spoils the page");
7612 }
7613
7614 #[test]
7621 fn the_residue_agrees_with_a_checked_conversion_everywhere() {
7622 for value in -70_000_i64..70_000 {
7623 assert_eq!(fit::<i8>(&[value]).is_ok(), i8::try_from(value).is_ok(), "{value} as i8");
7624 assert_eq!(fit::<u8>(&[value]).is_ok(), u8::try_from(value).is_ok(), "{value} as u8");
7625 assert_eq!(fit::<i16>(&[value]).is_ok(), i16::try_from(value).is_ok(), "{value} i16");
7626 assert_eq!(fit::<u16>(&[value]).is_ok(), u16::try_from(value).is_ok(), "{value} u16");
7627 }
7628 let wide = [i64::MIN, i64::MIN + 1, i64::from(i32::MIN), 0, i64::from(u32::MAX), i64::MAX];
7629 for edge in wide {
7630 for step in -2_i64..=2 {
7631 let value = edge.saturating_add(step);
7632 assert_eq!(
7633 fit::<i32>(&[value]).is_ok(),
7634 i32::try_from(value).is_ok(),
7635 "{value} as i32"
7636 );
7637 assert_eq!(
7638 fit::<u32>(&[value]).is_ok(),
7639 u32::try_from(value).is_ok(),
7640 "{value} as u32"
7641 );
7642 }
7643 }
7644 }
7645
7646 #[test]
7654 fn a_dictionary_at_its_budget_sweeps_without_keeping() {
7655 let path = path("dictionary-budget");
7656 let spellings = (0..2_500)
7657 .map(|index| Value::Varchar(format!("value {index:08} {}", "y".repeat(index % 40))))
7658 .collect::<Vec<_>>();
7659 let mut writer =
7660 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7661 .expect("new file");
7662 for part in spellings.chunks(1_024) {
7663 writer
7664 .append(
7665 &Chunk::new(vec![
7666 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
7667 ])
7668 .expect("one column"),
7669 )
7670 .expect("stripe written");
7671 }
7672 writer.finish().expect("commit");
7673
7674 let reader = Reader::open(&path).expect("valid directory");
7675 let page = reader.table.dictionaries[0].expect("a string column has one");
7676 let file = Arc::clone(&reader.file);
7677 let starved = open_global_dictionary(file, page, &LogicalType::Varchar, 0)
7678 .expect("a dictionary opens whatever it may keep");
7679
7680 let resting = starved.footprint();
7681 let mut swept: Vec<Vec<u8>> = Vec::new();
7682 let mut at = 0;
7683 while at < starved.len() {
7684 at = starved
7685 .sweep_text(at, starved.len(), &mut |_index: usize, text: &[u8]| {
7686 swept.push(text.to_vec());
7687 Ok(())
7688 })
7689 .expect("a sweep reads");
7690 }
7691 assert_eq!(swept.len(), spellings.len(), "a starved sweep still reads every value");
7692 assert_eq!(starved.footprint(), resting, "and keeps no block it decoded");
7693
7694 let generous = reader.dictionary(0).expect("read").expect("a string column has one");
7695 let read = (0..generous.len())
7696 .map(|code| generous.try_bytes_at(code).expect("read").expect("a value").to_vec())
7697 .collect::<Vec<_>>();
7698 assert_eq!(swept, read, "a starved sweep answers what a point read answers");
7699 fs::remove_file(path).expect("remove scratch file");
7700 }
7701
7702 #[test]
7703 fn damaged_membership_cannot_skip_a_string_page() {
7704 let path = path("damaged-membership");
7705 let mut writer = Writer::create(
7706 &path,
7707 "items",
7708 vec![
7709 Field::required("id", LogicalType::Integer),
7710 Field::new("text", LogicalType::Varchar),
7711 ],
7712 )
7713 .expect("new file");
7714 writer.append(&sample()).expect("stripe written");
7715 writer.finish().expect("commit");
7716
7717 let reader = Reader::open(&path).expect("valid directory");
7718 let membership = reader.table.stripes[0].memberships[1].expect("string membership");
7719 let mut file = OpenOptions::new().write(true).open(&path).expect("open membership page");
7720 file.seek(SeekFrom::Start(membership.offset)).expect("membership start");
7721 file.write_all(&[255]).expect("damage membership");
7722 let error = reader.skips_codes(0, 1, &[3]).expect_err("corruption must not skip rows");
7723 assert!(error.message().contains("membership page checksum differs"), "{error}");
7724 fs::remove_file(path).expect("remove scratch file");
7725 }
7726
7727 #[test]
7728 fn membership_delta_stream_is_sorted_exact_and_bounded() {
7729 let unique = unique_codes(&[900, 4, 4, 72, 9, u32::MAX]);
7730 assert_eq!(unique, [4, 9, 72, 900, u32::MAX]);
7731 let encoded = encode_membership(&unique);
7732 assert_eq!(
7733 decode_membership(&encoded).expect("valid membership"),
7734 [4, 9, 72, 900, u32::MAX]
7735 );
7736 let merged = merged_codes(vec![vec![4, 900], vec![9, 900, u32::MAX], vec![72]]);
7739 assert_eq!(merged, [4, 9, 72, 900, u32::MAX]);
7740 assert_eq!(
7741 decode_membership(&encode_membership(&merged)).expect("valid membership"),
7742 unique
7743 );
7744 assert!(decode_membership(&[1, 0x80]).is_err(), "a truncated varint is invalid");
7745 assert!(
7746 decode_membership(&[1, 0xff, 0xff, 0xff, 0xff, 0x10]).is_err(),
7747 "a value past u32 is invalid"
7748 );
7749 }
7750
7751 #[test]
7752 fn a_global_dictionary_may_be_larger_than_one_column_page() {
7753 let dictionary = Page {
7754 offset: HEADER,
7755 length: u32::try_from(MAX_PAGE + 1).expect("the page bound fits on disk"),
7756 hash: 0,
7757 };
7758 let table = Table {
7759 name: "items".to_owned(),
7760 fields: vec![Field::new("text", LogicalType::Varchar)],
7761 stripes: Vec::new(),
7762 rows: 0,
7763 dictionaries: vec![Some(dictionary)],
7764 distincts: vec![None],
7765 frequencies: vec![None],
7766 };
7767 let directory = encode_directory(&table).expect("directory");
7768 let file_size = dictionary.offset + u64::from(dictionary.length) + 1;
7769
7770 let decoded = decode_directory(&directory, file_size).expect("large lazy dictionary");
7771 assert_eq!(decoded.dictionaries[0].expect("dictionary").length, dictionary.length);
7772 }
7773
7774 #[test]
7775 fn a_column_with_one_value_everywhere_costs_almost_nothing_a_row() {
7776 let path = path("constant-codes");
7777 let mut writer =
7778 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7779 .expect("new file");
7780 let empty = vec![Value::Varchar(String::new()); 1024];
7781 for _ in 0..4 {
7782 let column = Vector::from_values(LogicalType::Varchar, &empty).expect("strings");
7783 writer.append(&Chunk::new(vec![column]).expect("one column")).expect("a part");
7784 }
7785 writer.finish().expect("commit");
7786
7787 let reader = Reader::open(&path).expect("valid directory");
7788 let pages = reader.layout().columns.first().expect("one column").pages;
7789 assert!(pages < 256, "{pages} bytes of pages for 4,096 rows of one value");
7793 let read = reader.read(3, &[0]).expect("the last part back");
7794 assert_eq!(read.value_at(0, 0), Value::Varchar(String::new()));
7795 assert_eq!(read.value_at(1023, 0), Value::Varchar(String::new()));
7796 fs::remove_file(path).expect("remove scratch file");
7797 }
7798
7799 #[test]
7800 fn a_cascade_value_too_wide_for_its_column_is_refused_rather_than_cut() {
7801 let over = vec![i64::from(i32::MAX) + 1];
7804 let error = narrowed(&LogicalType::Integer, over).expect_err("a page that disagrees");
7805 assert!(format!("{error}").contains("not of its type"), "{error}");
7806 assert!(narrowed(&LogicalType::BigInt, vec![i64::MIN]).is_ok(), "bigint holds all of i64");
7807 assert!(narrowed(&LogicalType::Varchar, vec![0]).is_err(), "strings are not integers");
7808 }
7809
7810 #[test]
7811 fn a_code_stream_the_cascade_cannot_shrink_is_left_alone() {
7812 let mut state: u32 = 0x9e37_79b9;
7816 let spread: Vec<u32> = (0..1024)
7817 .map(|_| {
7818 state ^= state << 13;
7819 state ^= state >> 17;
7820 state ^= state << 5;
7821 state
7822 })
7823 .collect();
7824 assert_eq!(encoded_codes(&spread).expect("no failure"), None);
7825 let near: Vec<u32> = (0..1024).collect();
7826 let coded = encoded_codes(&near).expect("no failure").expect("counting up is packable");
7827 assert!(coded.len() < near.len() * 4, "{} bytes for a run of 1,024", coded.len());
7828 }
7829
7830 #[test]
7836 fn two_writes_of_the_same_rows_give_the_same_bytes() {
7837 fn written(path: &PathBuf) {
7838 let fields = (0..40)
7839 .map(|column| {
7840 let ty =
7841 if column % 4 == 0 { LogicalType::Varchar } else { LogicalType::BigInt };
7842 Field::new(format!("c{column}"), ty)
7843 })
7844 .collect::<Vec<_>>();
7845 let mut writer = Writer::create(path, "wide", fields).expect("new file");
7846 for part in 0..70_u64 {
7847 let columns = (0..40)
7848 .map(|column| {
7849 let values = (0..64_u64)
7850 .map(|row| {
7851 let seed = part.wrapping_mul(31).wrapping_add(row);
7852 if column % 4 == 0 {
7853 Value::Varchar(format!("v{}", seed % 17))
7854 } else {
7855 Value::BigInt(i64::try_from(seed % 97).expect("small"))
7856 }
7857 })
7858 .collect::<Vec<_>>();
7859 let ty = if column % 4 == 0 {
7860 LogicalType::Varchar
7861 } else {
7862 LogicalType::BigInt
7863 };
7864 Vector::from_values(ty, &values).expect("a column")
7865 })
7866 .collect::<Vec<_>>();
7867 writer.append(&Chunk::new(columns).expect("forty columns")).expect("a part");
7868 }
7869 writer.finish().expect("commit");
7870 }
7871
7872 let first = path("repeatable-one");
7873 let second = path("repeatable-two");
7874 written(&first);
7875 written(&second);
7876 let left = fs::read(&first).expect("the first file");
7877 let right = fs::read(&second).expect("the second file");
7878 assert_eq!(left.len(), right.len(), "two writes of the same rows differ in length");
7879 assert!(left == right, "two writes of the same rows differ in their bytes");
7880
7881 let reader = Reader::open(&first).expect("valid directory");
7884 assert_eq!(reader.table().rows(), 70 * 64);
7885 let read = reader.read(0, &[0, 1]).expect("the first part back");
7886 assert_eq!(read.value_at(0, 0), Value::Varchar("v0".to_owned()));
7887 assert_eq!(read.value_at(0, 1), Value::BigInt(0));
7888 fs::remove_file(first).expect("remove scratch file");
7889 fs::remove_file(second).expect("remove scratch file");
7890 }
7891
7892 fn three_tables(path: &PathBuf) {
7894 let writer = Writer::create(
7895 path,
7896 "region",
7897 vec![
7898 Field::new("r_key", LogicalType::Integer),
7899 Field::new("r_name", LogicalType::Varchar),
7900 ],
7901 )
7902 .expect("new file");
7903 let mut writer = writer;
7904 writer
7905 .append(
7906 &Chunk::new(vec![
7907 Vector::from_values(
7908 LogicalType::Integer,
7909 &[Value::Integer(0), Value::Integer(1)],
7910 )
7911 .expect("keys"),
7912 Vector::from_values(
7913 LogicalType::Varchar,
7914 &[Value::Varchar("AFRICA".to_owned()), Value::Varchar("ASIA".to_owned())],
7915 )
7916 .expect("names"),
7917 ])
7918 .expect("two columns"),
7919 )
7920 .expect("a part");
7921 let mut writer = writer
7922 .next("empty", vec![Field::new("nothing", LogicalType::BigInt)])
7923 .expect("a second table");
7924 writer
7925 .append(
7926 &Chunk::new(vec![
7927 Vector::from_values(LogicalType::BigInt, &[Value::BigInt(7)]).expect("a row"),
7928 ])
7929 .expect("one column"),
7930 )
7931 .expect("a part");
7932 let mut writer =
7933 writer.next("wide", vec![Field::new("n", LogicalType::BigInt)]).expect("a third table");
7934 for part in 0..70_i64 {
7935 let values = (0..64).map(|row| Value::BigInt(part * 64 + row)).collect::<Vec<_>>();
7936 writer
7937 .append(
7938 &Chunk::new(vec![
7939 Vector::from_values(LogicalType::BigInt, &values).expect("a column"),
7940 ])
7941 .expect("one column"),
7942 )
7943 .expect("a part");
7944 }
7945 writer.finish().expect("commit");
7946 }
7947
7948 #[test]
7949 fn three_tables_in_one_file_read_back_by_name() {
7950 let file = path("three-tables");
7951 three_tables(&file);
7952 let catalog = Catalog::open(&file).expect("a committed catalog");
7953 assert_eq!(catalog.names().collect::<Vec<_>>(), ["region", "empty", "wide"]);
7954
7955 let region = catalog.table("region").expect("the first table");
7956 assert_eq!(region.table().rows(), 2);
7957 assert_eq!(
7958 region.read(0, &[1]).expect("names").value_at(1, 0),
7959 Value::Varchar("ASIA".to_owned())
7960 );
7961
7962 let wide = catalog.table("wide").expect("the third table");
7963 assert_eq!(wide.table().rows(), 70 * 64);
7964 assert_eq!(wide.read(0, &[0]).expect("the first part").value_at(0, 0), Value::BigInt(0));
7965
7966 let empty = catalog.table("empty").expect("the second table");
7969 assert_eq!(empty.table().rows(), 1);
7970 assert_eq!(empty.read(0, &[0]).expect("the row").value_at(0, 0), Value::BigInt(7));
7971
7972 fs::remove_file(file).expect("remove scratch file");
7973 }
7974
7975 #[test]
7976 fn a_name_the_file_does_not_hold_is_an_error_rather_than_the_first_table() {
7977 let file = path("three-tables-missing");
7978 three_tables(&file);
7979 let catalog = Catalog::open(&file).expect("a committed catalog");
7980 let error = catalog.table("nation").expect_err("no such table");
7981 assert!(error.message().contains("nation"), "{}", error.message());
7982 fs::remove_file(file).expect("remove scratch file");
7983 }
7984
7985 #[test]
7986 fn a_file_of_three_tables_will_not_open_as_one() {
7987 let file = path("three-tables-unnamed");
7988 three_tables(&file);
7989 let error = Reader::open(&file).expect_err("more than one table");
7990 assert!(error.message().contains("more than one table"), "{}", error.message());
7991 fs::remove_file(file).expect("remove scratch file");
7992 }
7993
7994 #[test]
7996 fn decimals_of_every_storage_width_round_trip() {
7997 let file = path("decimals");
7998 let widths = [(4_u8, 2_u8), (9, 2), (18, 4), (38, 6)];
7999 let fields = widths
8000 .iter()
8001 .enumerate()
8002 .map(|(index, (width, scale))| {
8003 Field::new(
8004 format!("d{index}"),
8005 LogicalType::decimal(*width, *scale).expect("a decimal type"),
8006 )
8007 })
8008 .collect::<Vec<_>>();
8009 let mut writer = Writer::create(&file, "money", fields).expect("new file");
8010 let rows: [i128; 3] = [-1234, 0, 999];
8011 let columns = widths
8012 .iter()
8013 .map(|(width, scale)| {
8014 let values = rows
8015 .iter()
8016 .map(|unscaled| Value::Decimal {
8017 unscaled: *unscaled,
8018 width: *width,
8019 scale: *scale,
8020 })
8021 .collect::<Vec<_>>();
8022 Vector::from_values(
8023 LogicalType::decimal(*width, *scale).expect("a decimal type"),
8024 &values,
8025 )
8026 .expect("a decimal column")
8027 })
8028 .collect::<Vec<_>>();
8029 writer.append(&Chunk::new(columns).expect("four columns")).expect("a part");
8030 writer.finish().expect("commit");
8031
8032 let reader = Reader::open(&file).expect("a committed file");
8033 for (index, (width, scale)) in widths.iter().enumerate() {
8034 assert_eq!(
8035 reader.table().fields()[index].ty,
8036 LogicalType::decimal(*width, *scale).expect("a decimal type"),
8037 "column {index} came back as another type"
8038 );
8039 let column = reader.read(0, &[index]).expect("the column");
8040 for (row, unscaled) in rows.iter().enumerate() {
8041 assert_eq!(
8042 column.value_at(row, 0),
8043 Value::Decimal { unscaled: *unscaled, width: *width, scale: *scale },
8044 "column {index} row {row}"
8045 );
8046 }
8047 }
8048 fs::remove_file(file).expect("remove scratch file");
8049 }
8050
8051 #[test]
8052 fn two_tables_of_one_name_are_refused_before_anything_is_committed() {
8053 let file = path("two-of-a-name");
8054 let writer = Writer::create(&file, "t", vec![Field::new("a", LogicalType::BigInt)])
8055 .expect("new file");
8056 let error = writer
8057 .next("t", vec![Field::new("a", LogicalType::BigInt)])
8058 .expect_err("the same name twice");
8059 assert!(error.message().contains("same name"), "{}", error.message());
8060 fs::remove_file(file).expect("remove scratch file");
8061 }
8062
8063 #[test]
8064 fn opening_the_catalog_reads_no_table_directory() {
8065 let file = path("catalog-only");
8066 three_tables(&file);
8067 let catalog = Catalog::open(&file).expect("a committed catalog");
8068 assert_eq!(catalog.opening.reads, 2, "opening the catalog read more than the slot");
8071 assert_eq!(catalog.names().len(), 3);
8072 fs::remove_file(file).expect("remove scratch file");
8073 }
8074}