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::{Clustering, Error, Field, LogicalType, PhysicalType, Result, Value, Width};
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 CLUSTERING: &[u8; 8] = b"RUDBCL1\0";
74const FREQUENCY_CANDIDATES: usize = 32_768;
75const FREQUENCY_ENTRIES: usize = 512;
76const FREQUENCY_BUILD_RANK: usize = 10;
77const FREQUENCY_ORDINALS: usize = 65_536;
78const MAX_FREQUENCY_WORKERS: usize = 32;
85
86const MAX_ENCODE_WORKERS: usize = 32;
93
94const SIEVE_BUDGET: usize = 8 * 1024;
102
103const PART_BOUND_BYTES: usize = 24;
112
113fn io(error: std::io::Error) -> Error {
114 Error::io(error.to_string())
115}
116
117fn invalid(message: &str) -> Error {
118 Error::invalid_input(format!("invalid rudb native file: {message}"))
119}
120
121fn sum(counts: impl Iterator<Item = u64>) -> u64 {
123 counts.fold(0, u64::saturating_add)
124}
125
126fn span_bytes(spans: &[Span], at: usize) -> u64 {
128 spans.get(at).map_or(0, |span| u64::from(span.length))
129}
130
131fn page_bytes(pages: &[Option<Page>], at: usize) -> u64 {
133 pages.get(at).and_then(Option::as_ref).map_or(0, Page::bytes)
134}
135
136fn checksum(bytes: &[u8]) -> u64 {
146 const P1: u64 = 11_400_714_785_074_694_791;
147 const P2: u64 = 14_029_467_366_897_019_727;
148 const P3: u64 = 1_609_587_929_392_839_161;
149 const P4: u64 = 9_650_029_242_287_828_579;
150 const P5: u64 = 2_870_177_450_012_600_261;
151 let round = |state: u64, word: u64| {
152 state.wrapping_add(word.wrapping_mul(P2)).rotate_left(31).wrapping_mul(P1)
153 };
154 let merge = |state: u64, lane: u64| (state ^ round(0, lane)).wrapping_mul(P1).wrapping_add(P4);
155 let word = |chunk: &[u8]| u64::from_le_bytes(chunk.try_into().expect("eight checksum bytes"));
156
157 let mut blocks = bytes.chunks_exact(32);
160 let mut rest = blocks.remainder();
161 let mut hash = if bytes.len() >= 32 {
162 let mut one = P1.wrapping_add(P2);
163 let mut two = P2;
164 let mut three = 0;
165 let mut four = 0_u64.wrapping_sub(P1);
166 for block in blocks.by_ref() {
167 one = round(one, word(&block[..8]));
168 two = round(two, word(&block[8..16]));
169 three = round(three, word(&block[16..24]));
170 four = round(four, word(&block[24..]));
171 }
172 let combined = one
173 .rotate_left(1)
174 .wrapping_add(two.rotate_left(7))
175 .wrapping_add(three.rotate_left(12))
176 .wrapping_add(four.rotate_left(18));
177 merge(merge(merge(merge(combined, one), two), three), four)
178 } else {
179 P5
180 };
181 hash = hash.wrapping_add(bytes.len() as u64);
182 let mut words = rest.chunks_exact(8);
183 for chunk in words.by_ref() {
184 hash ^= round(0, word(chunk));
185 hash = hash.rotate_left(27).wrapping_mul(P1).wrapping_add(P4);
186 }
187 rest = words.remainder();
188 if rest.len() >= 4 {
189 let (head, tail) = rest.split_at(4);
190 let quarter = u32::from_le_bytes(head.try_into().expect("four checksum bytes"));
191 hash ^= u64::from(quarter).wrapping_mul(P1);
192 hash = hash.rotate_left(23).wrapping_mul(P2).wrapping_add(P3);
193 rest = tail;
194 }
195 for &byte in rest {
196 hash ^= u64::from(byte).wrapping_mul(P5);
197 hash = hash.rotate_left(11).wrapping_mul(P1);
198 }
199 hash ^= hash >> 33;
200 hash = hash.wrapping_mul(P2);
201 hash ^= hash >> 29;
202 hash = hash.wrapping_mul(P3);
203 hash ^ (hash >> 32)
204}
205
206#[derive(Debug, Clone, Copy)]
207struct Slot {
208 offset: u64,
209 length: u32,
210 generation: u64,
211 hash: u64,
212}
213
214impl Slot {
215 fn bytes(self) -> [u8; SLOT_BYTES] {
216 let mut result = [0; SLOT_BYTES];
217 result[..8].copy_from_slice(&self.offset.to_le_bytes());
218 result[8..12].copy_from_slice(&self.length.to_le_bytes());
219 result[12..20].copy_from_slice(&self.generation.to_le_bytes());
220 result[20..28].copy_from_slice(&self.hash.to_le_bytes());
221 result
222 }
223
224 fn read(bytes: &[u8]) -> Self {
225 Self {
226 offset: u64::from_le_bytes(bytes[..8].try_into().expect("eight bytes")),
227 length: u32::from_le_bytes(bytes[8..12].try_into().expect("four bytes")),
228 generation: u64::from_le_bytes(bytes[12..20].try_into().expect("eight bytes")),
229 hash: u64::from_le_bytes(bytes[20..28].try_into().expect("eight bytes")),
230 }
231 }
232}
233
234#[derive(Debug, Clone, Copy)]
235struct Page {
236 offset: u64,
237 length: u32,
238 hash: u64,
239}
240
241impl Page {
242 fn bytes(&self) -> u64 {
244 u64::from(self.length)
245 }
246}
247
248#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
249enum FrequencyValue {
250 Null,
251 Integer(i128),
252 Code(u32),
253}
254
255#[derive(Debug, Clone)]
256struct FrequencyEntry {
257 value: FrequencyValue,
258 count: u64,
259}
260
261#[derive(Debug, Clone)]
266struct FrequencySummary {
267 entries: Vec<FrequencyEntry>,
268 omitted_max: u64,
269 ordinals: Vec<u64>,
270}
271
272#[derive(Debug, Clone)]
277pub struct FrequencyPrefix {
278 pub entries: Vec<(Value, u64)>,
280 pub omitted_max: u64,
282}
283
284#[derive(Debug, Clone, PartialEq, Eq)]
286pub struct FrequencyOccurrences {
287 pub omitted_max: u64,
289 pub ordinals: Vec<u64>,
291}
292
293#[derive(Debug, Clone, Copy, Default)]
300struct Span {
301 offset: u64,
302 length: u32,
303}
304
305#[derive(Debug, Clone)]
307pub struct Stripe {
308 rows: usize,
309 parts: Vec<u32>,
312 index: Span,
316 pages: Vec<Span>,
317 memberships: Vec<Option<Page>>,
318 sieves: Vec<Option<Page>>,
321 part_ranges: Vec<Option<Page>>,
332 zone: Zone,
333}
334
335impl Stripe {
336 #[must_use]
338 pub fn rows(&self) -> usize {
339 self.rows
340 }
341
342 #[must_use]
344 pub fn parts(&self) -> usize {
345 self.parts.len()
346 }
347
348 #[must_use]
354 pub fn zone(&self) -> &Zone {
355 &self.zone
356 }
357}
358
359#[derive(Debug, Clone)]
361pub struct Table {
362 name: String,
363 fields: Vec<Field>,
364 stripes: Vec<Stripe>,
365 rows: usize,
366 dictionaries: Vec<Option<Page>>,
367 frequencies: Vec<Option<FrequencySummary>>,
368 distincts: Vec<Option<u64>>,
378 clustering: Option<Clustering>,
386}
387
388impl Table {
389 #[must_use]
391 pub fn name(&self) -> &str {
392 &self.name
393 }
394
395 #[must_use]
397 pub fn fields(&self) -> &[Field] {
398 &self.fields
399 }
400
401 #[must_use]
403 pub fn rows(&self) -> usize {
404 self.rows
405 }
406
407 #[must_use]
409 pub fn stripes(&self) -> &[Stripe] {
410 &self.stripes
411 }
412
413 #[must_use]
415 pub fn clustering(&self) -> Option<&Clustering> {
416 self.clustering.as_ref()
417 }
418}
419
420#[derive(Debug, Clone)]
432struct Entry {
433 name: String,
434 fields: Vec<Field>,
435 rows: usize,
436 directory: Page,
438}
439
440#[derive(Debug, Clone)]
442pub struct ColumnLayout {
443 pub name: String,
445 pub kind: String,
447 pub pages: u64,
449 pub memberships: u64,
451 pub sieves: u64,
453 pub part_ranges: u64,
455 pub dictionary: u64,
457}
458
459impl ColumnLayout {
460 #[must_use]
462 pub fn total(&self) -> u64 {
463 self.pages
464 .saturating_add(self.memberships)
465 .saturating_add(self.sieves)
466 .saturating_add(self.part_ranges)
467 .saturating_add(self.dictionary)
468 }
469}
470
471#[derive(Debug, Clone)]
482pub struct Layout {
483 pub file: u64,
485 pub rows: usize,
487 pub stripes: usize,
489 pub parts: usize,
491 pub columns: Vec<ColumnLayout>,
493 pub indexes: u64,
496 pub directory: u64,
498 pub header: u64,
500}
501
502impl Layout {
503 #[must_use]
505 pub fn columns_total(&self) -> u64 {
506 self.columns.iter().map(ColumnLayout::total).fold(0, u64::saturating_add)
507 }
508
509 #[must_use]
515 pub fn unaccounted(&self) -> u64 {
516 self.file
517 .saturating_sub(self.columns_total())
518 .saturating_sub(self.indexes)
519 .saturating_sub(self.directory)
520 .saturating_sub(self.header)
521 }
522}
523
524#[derive(Debug)]
526struct GlobalDictionary {
527 primary: HashMap<u64, u32>,
528 collisions: HashMap<u64, Vec<u32>>,
529 offsets: Vec<u32>,
530 payload: Vec<u8>,
531 counts: Vec<u64>,
532 nulls: u64,
533}
534
535impl GlobalDictionary {
536 fn new() -> Self {
537 Self {
538 primary: HashMap::new(),
539 collisions: HashMap::new(),
540 offsets: vec![0],
541 payload: Vec::new(),
542 counts: Vec::new(),
543 nulls: 0,
544 }
545 }
546
547 fn bytes(&self, code: u32) -> Option<&[u8]> {
548 let start = *self.offsets.get(code as usize)? as usize;
549 let end = *self.offsets.get(code as usize + 1)? as usize;
550 self.payload.get(start..end)
551 }
552
553 fn code(&mut self, text: &str) -> Result<u32> {
554 let hash = checksum(text.as_bytes());
555 if let Some(&code) = self.primary.get(&hash) {
556 if self.bytes(code) == Some(text.as_bytes()) {
557 return Ok(code);
558 }
559 if let Some(codes) = self.collisions.get(&hash) {
560 if let Some(code) =
561 codes.iter().copied().find(|&code| self.bytes(code) == Some(text.as_bytes()))
562 {
563 return Ok(code);
564 }
565 }
566 let code = self.insert(text)?;
567 self.collisions.entry(hash).or_default().push(code);
568 return Ok(code);
569 }
570 let code = self.insert(text)?;
571 self.primary.insert(hash, code);
572 Ok(code)
573 }
574
575 fn insert(&mut self, text: &str) -> Result<u32> {
576 let code = u32::try_from(self.offsets.len() - 1)
577 .map_err(|_| invalid("global dictionary has too many values"))?;
578 self.payload.extend_from_slice(text.as_bytes());
579 self.offsets.push(
580 u32::try_from(self.payload.len())
581 .map_err(|_| invalid("global dictionary payload exceeds 4 GiB"))?,
582 );
583 self.counts.push(0);
584 Ok(code)
585 }
586
587 fn ranked(&self) -> Vec<(u64, u32)> {
607 let count = self.offsets.len() - 1;
608 let mut ranked = (0..count)
609 .map(|code| {
610 let code = code as u32;
611 (head(self.bytes(code).unwrap_or_default()), code)
612 })
613 .collect::<Vec<_>>();
614 ranked.sort_unstable_by(|left, right| {
615 left.0.cmp(&right.0).then_with(|| self.bytes(left.1).cmp(&self.bytes(right.1)))
616 });
617 ranked
618 }
619
620 fn observe(&mut self, code: u32, null: bool) -> Result<()> {
621 if null {
622 self.nulls = self.nulls.saturating_add(1);
623 return Ok(());
624 }
625 let count = self
626 .counts
627 .get_mut(code as usize)
628 .ok_or_else(|| invalid("global dictionary count code is out of range"))?;
629 *count = count.saturating_add(1);
630 Ok(())
631 }
632}
633
634#[derive(Debug)]
642pub struct Writer {
643 file: File,
644 at: u64,
652 table: Table,
653 generation: u64,
654 order: Vec<((u64, u64), (u64, u64))>,
657 next_order: u64,
658 dictionaries: Vec<Option<GlobalDictionary>>,
659 pending: Vec<PendingChunk>,
660 closed: Vec<Entry>,
662}
663
664#[derive(Debug)]
672struct PendingChunk {
673 order: (u64, u64),
674 chunk: Chunk,
675}
676
677#[derive(Debug)]
683struct ColumnStripe {
684 pages: Vec<Vec<u8>>,
685 codes: Vec<Option<Vec<u32>>>,
686 sieves: Vec<Option<Sieve>>,
687 ranges: Vec<Range>,
688}
689
690fn weight(ty: &LogicalType) -> usize {
698 match ty {
699 LogicalType::Varchar | LogicalType::Blob => 64,
700 LogicalType::BigInt
701 | LogicalType::UBigInt
702 | LogicalType::Timestamp
703 | LogicalType::Double
704 | LogicalType::Decimal { .. } => 8,
705 LogicalType::Integer | LogicalType::UInteger | LogicalType::Date | LogicalType::Float => 4,
706 LogicalType::SmallInt | LogicalType::USmallInt => 2,
707 _ => 1,
708 }
709}
710
711pub const STRIPE_PARTS: usize = 64;
718
719const DICTIONARY_DECIDE_ROWS: usize = 4_096;
727
728const DICTIONARY_DISTINCT_IN_TEN: usize = 9;
744
745const INDEX_ENTRY: usize = size_of::<u32>() + size_of::<u64>();
747
748fn index_section(parts: usize) -> Result<usize> {
750 parts
751 .checked_mul(INDEX_ENTRY)
752 .and_then(|bytes| bytes.checked_add(size_of::<u64>()))
753 .ok_or_else(|| invalid("index page length overflow"))
754}
755
756impl Writer {
757 pub fn open(
775 path: impl AsRef<Path>,
776 name: impl Into<String>,
777 fields: Vec<Field>,
778 ) -> Result<Self> {
779 for field in &fields {
780 type_tag(&field.ty)?;
781 }
782 let name = name.into();
783 let path = path.as_ref();
784 let (_, size, slot, bytes, _) = slot_bytes(path)?;
785 let closed = decode_catalog(&bytes, size)?;
786 if closed.iter().any(|held| held.name == name) {
787 return Err(invalid("two tables in one native file have the same name"));
788 }
789 let generation = slot
794 .generation
795 .checked_add(1)
796 .ok_or_else(|| invalid("native file generation overflow"))?;
797 let file = OpenOptions::new().write(true).read(true).open(path).map_err(io)?;
798 Ok(Self {
799 file,
800 at: size,
803 dictionaries: fields
804 .iter()
805 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
806 .collect(),
807 table: Table {
808 name,
809 dictionaries: vec![None; fields.len()],
810 distincts: vec![None; fields.len()],
811 fields,
812 stripes: Vec::new(),
813 rows: 0,
814 frequencies: Vec::new(),
815 clustering: None,
816 },
817 generation,
818 order: Vec::new(),
819 next_order: 0,
820 pending: Vec::with_capacity(STRIPE_PARTS),
821 closed,
822 })
823 }
824
825 pub fn create(
831 path: impl AsRef<Path>,
832 name: impl Into<String>,
833 fields: Vec<Field>,
834 ) -> Result<Self> {
835 for field in &fields {
836 type_tag(&field.ty)?;
837 }
838 let file =
839 OpenOptions::new().write(true).read(true).create_new(true).open(path).map_err(io)?;
840 let mut header = [0; HEADER as usize];
841 header[..8].copy_from_slice(MAGIC);
842 header[8..12].copy_from_slice(&FORMAT.to_le_bytes());
843 write_at(&file, 0, &header)?;
844 Ok(Self {
845 file,
846 at: HEADER,
847 dictionaries: fields
848 .iter()
849 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
850 .collect(),
851 table: Table {
852 name: name.into(),
853 dictionaries: vec![None; fields.len()],
854 distincts: vec![None; fields.len()],
855 fields,
856 stripes: Vec::new(),
857 rows: 0,
858 frequencies: Vec::new(),
859 clustering: None,
860 },
861 generation: 1,
862 order: Vec::new(),
863 next_order: 0,
864 pending: Vec::with_capacity(STRIPE_PARTS),
865 closed: Vec::new(),
866 })
867 }
868
869 pub fn empty(path: impl AsRef<Path>) -> Result<()> {
887 let file =
888 OpenOptions::new().write(true).read(true).create_new(true).open(path).map_err(io)?;
889 let mut header = [0; HEADER as usize];
890 header[..8].copy_from_slice(MAGIC);
891 header[8..12].copy_from_slice(&FORMAT.to_le_bytes());
892 write_at(&file, 0, &header)?;
893 let catalog = encode_catalog(&[])?;
894 write_at(&file, HEADER, &catalog)?;
895 file.sync_all().map_err(io)?;
899 let slot = Slot {
900 offset: HEADER,
901 length: u32::try_from(catalog.len()).map_err(|_| invalid("catalog length overflow"))?,
902 generation: 1,
903 hash: checksum(&catalog),
904 };
905 write_at(&file, slot_offset(1), &slot.bytes())?;
906 file.sync_all().map_err(io)?;
907 Ok(())
908 }
909
910 pub fn next(mut self, name: impl Into<String>, fields: Vec<Field>) -> Result<Self> {
921 for field in &fields {
922 type_tag(&field.ty)?;
923 }
924 let name = name.into();
925 let entry = self.close()?;
926 if self.closed.iter().chain(std::iter::once(&entry)).any(|held| held.name == name) {
927 return Err(invalid("two tables in one native file have the same name"));
928 }
929 let Self { file, at, generation, mut closed, .. } = self;
930 closed.push(entry);
931 Ok(Self {
932 file,
933 at,
934 generation,
935 closed,
936 dictionaries: fields
937 .iter()
938 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
939 .collect(),
940 table: Table {
941 name,
942 dictionaries: vec![None; fields.len()],
943 distincts: vec![None; fields.len()],
944 fields,
945 stripes: Vec::new(),
946 rows: 0,
947 frequencies: Vec::new(),
948 clustering: None,
949 },
950 order: Vec::new(),
951 next_order: 0,
952 pending: Vec::with_capacity(STRIPE_PARTS),
953 })
954 }
955
956 pub fn declare(mut self, clustering: Clustering) -> Result<Self> {
971 self.table.clustering = Some(Clustering::new(
974 clustering.columns().to_vec(),
975 clustering.width(),
976 &self.table.fields,
977 )?);
978 Ok(self)
979 }
980
981 fn put(&mut self, bytes: &[u8]) -> Result<()> {
986 write_at(&self.file, self.at, bytes)?;
987 self.at = self
988 .at
989 .checked_add(bytes.len() as u64)
990 .ok_or_else(|| invalid("native file length overflow"))?;
991 Ok(())
992 }
993
994 pub fn append(&mut self, chunk: &Chunk) -> Result<()> {
1000 let order = (self.next_order, 0);
1001 self.next_order = self.next_order.saturating_add(1);
1002 self.append_at(order, chunk)
1003 }
1004
1005 pub fn append_at(&mut self, order: (u64, u64), chunk: &Chunk) -> Result<()> {
1016 if chunk.is_empty() {
1017 return Ok(());
1018 }
1019 self.admit(chunk)?;
1020 if self.pending.last().is_some_and(|last| last.order > order) {
1021 self.flush_pending()?;
1022 }
1023 self.pending.push(PendingChunk { order, chunk: chunk.clone() });
1028 if self.pending.len() == STRIPE_PARTS {
1029 self.flush_pending()?;
1030 }
1031 Ok(())
1032 }
1033
1034 pub fn append_stripe(&mut self, parts: Vec<((u64, u64), Chunk)>) -> Result<()> {
1050 if parts.len() > STRIPE_PARTS {
1051 return Err(invalid("a stripe was handed more parts than it holds"));
1052 }
1053 self.flush_pending()?;
1056 for (order, chunk) in parts {
1057 if chunk.is_empty() {
1058 continue;
1059 }
1060 self.admit(&chunk)?;
1061 self.pending.push(PendingChunk { order, chunk });
1062 }
1063 self.flush_pending()
1064 }
1065
1066 fn admit(&mut self, chunk: &Chunk) -> Result<()> {
1068 if chunk.width() != self.table.fields.len() {
1069 return Err(invalid("chunk width differs from table schema"));
1070 }
1071 for (index, field) in self.table.fields.iter().enumerate() {
1072 if chunk.column(index)?.logical_type() != &field.ty {
1073 return Err(invalid("chunk type differs from table schema"));
1074 }
1075 }
1076 self.table.rows = self
1077 .table
1078 .rows
1079 .checked_add(chunk.len())
1080 .ok_or_else(|| invalid("row count overflow"))?;
1081 Ok(())
1082 }
1083
1084 fn encode_column(
1115 index: usize,
1116 held: &[PendingChunk],
1117 dictionary: &mut Option<GlobalDictionary>,
1118 ) -> Result<ColumnStripe> {
1119 let deciding = dictionary.as_ref().is_some_and(|held| held.offsets.len() == 1);
1122 let stripe = Self::encode_pages(index, held, dictionary.as_mut())?;
1123 if !deciding {
1124 return Ok(stripe);
1125 }
1126 let rows: usize = held.iter().map(|pending| pending.chunk.len()).sum();
1127 let distinct = dictionary.as_ref().map_or(0, |held| held.offsets.len() - 1);
1128 if rows < DICTIONARY_DECIDE_ROWS
1129 || distinct.saturating_mul(10) <= rows.saturating_mul(DICTIONARY_DISTINCT_IN_TEN)
1130 {
1131 return Ok(stripe);
1132 }
1133 *dictionary = None;
1134 Self::encode_pages(index, held, None)
1135 }
1136
1137 fn encode_pages(
1139 index: usize,
1140 held: &[PendingChunk],
1141 mut dictionary: Option<&mut GlobalDictionary>,
1142 ) -> Result<ColumnStripe> {
1143 let mut stripe = ColumnStripe {
1144 pages: Vec::with_capacity(held.len()),
1145 codes: Vec::with_capacity(held.len()),
1146 sieves: Vec::with_capacity(held.len()),
1147 ranges: Vec::with_capacity(held.len()),
1148 };
1149 for pending in held {
1150 let column = pending.chunk.column(index)?;
1151 let (bytes, unique) = encode(column, dictionary.as_deref_mut())?;
1152 if bytes.len() > MAX_PAGE {
1153 return Err(invalid("column page exceeds the configured bound"));
1154 }
1155 let range = Range::of(column);
1158 let sieve = match dictionary {
1171 Some(_) => None,
1172 None => Sieve::of(column, &range, SIEVE_BUDGET)
1173 .filter(|sieve| sieve.len() < bytes.len()),
1174 };
1175 stripe.pages.push(bytes);
1176 stripe.codes.push(unique);
1177 stripe.sieves.push(sieve);
1178 stripe.ranges.push(range);
1179 }
1180 Ok(stripe)
1181 }
1182
1183 fn encode_columns(&mut self, held: &[PendingChunk]) -> Result<Vec<ColumnStripe>> {
1192 let width = self.table.fields.len();
1193 let workers = std::thread::available_parallelism()
1194 .map_or(1, usize::from)
1195 .min(MAX_ENCODE_WORKERS)
1196 .min(width);
1197 if workers <= 1 || held.len() <= 1 {
1198 return self
1199 .dictionaries
1200 .iter_mut()
1201 .enumerate()
1202 .map(|(index, dictionary)| Self::encode_column(index, held, dictionary))
1203 .collect();
1204 }
1205 let mut jobs: Vec<(usize, Option<GlobalDictionary>)> =
1208 std::mem::take(&mut self.dictionaries).into_iter().enumerate().collect();
1209 jobs.sort_by_key(|(index, _)| weight(&self.table.fields[*index].ty));
1211 let queue = Mutex::new(jobs);
1212 let pieces = std::thread::scope(|scope| {
1213 (0..workers)
1214 .map(|_| {
1215 scope.spawn(|| {
1216 let mut mine = Vec::new();
1217 loop {
1218 let taken = queue
1219 .lock()
1220 .map_err(|_| Error::internal("a native encode worker panicked"))?
1221 .pop();
1222 let Some((index, mut dictionary)) = taken else { break };
1223 let encoded = Self::encode_column(index, held, &mut dictionary)?;
1224 mine.push((index, dictionary, encoded));
1225 }
1226 Ok(mine)
1227 })
1228 })
1229 .collect::<Vec<_>>()
1230 .into_iter()
1231 .map(|handle| {
1232 handle.join().map_err(|_| Error::internal("a native encode worker panicked"))?
1233 })
1234 .collect::<Result<Vec<_>>>()
1235 })?;
1236 let mut dictionaries: Vec<Option<GlobalDictionary>> = (0..width).map(|_| None).collect();
1237 let mut encoded: Vec<Option<ColumnStripe>> = (0..width).map(|_| None).collect();
1238 for piece in pieces {
1239 for (index, dictionary, stripe) in piece {
1240 dictionaries[index] = dictionary;
1241 encoded[index] = Some(stripe);
1242 }
1243 }
1244 self.dictionaries = dictionaries;
1245 encoded
1246 .into_iter()
1247 .map(|stripe| stripe.ok_or_else(|| Error::internal("a column was never encoded")))
1248 .collect()
1249 }
1250
1251 fn flush_pending(&mut self) -> Result<()> {
1253 if self.pending.is_empty() {
1254 return Ok(());
1255 }
1256 let width = self.table.fields.len();
1257 let mut held = std::mem::take(&mut self.pending);
1260 let parts = held.len();
1261 let encoded = self.encode_columns(&held)?;
1262 let mut pages = Vec::with_capacity(width);
1263 let mut memberships = vec![None; width];
1264 let mut ranges = Vec::with_capacity(width);
1265 let mut index = Vec::with_capacity(width.saturating_mul(index_section(parts)?));
1266 for stripe in &encoded {
1267 let offset = self.at;
1268 let section = index.len();
1269 let mut length = 0_usize;
1270 for bytes in &stripe.pages {
1271 write_at(&self.file, self.at + length as u64, bytes)?;
1272 put_u32(
1273 &mut index,
1274 u32::try_from(bytes.len()).map_err(|_| invalid("part length overflow"))?,
1275 );
1276 put_u64(&mut index, checksum(bytes));
1277 length = length
1278 .checked_add(bytes.len())
1279 .ok_or_else(|| invalid("column page length overflow"))?;
1280 }
1281 let hash = checksum(&index[section..]);
1282 put_u64(&mut index, hash);
1283 if length > MAX_PAGE {
1284 return Err(invalid("column page exceeds the configured bound"));
1285 }
1286 self.at = self
1287 .at
1288 .checked_add(length as u64)
1289 .ok_or_else(|| invalid("native file length overflow"))?;
1290 pages.push(Span {
1291 offset,
1292 length: u32::try_from(length).map_err(|_| invalid("page length overflow"))?,
1293 });
1294 ranges.push(merged_range(stripe.ranges.iter().cloned()));
1295 }
1296 for (membership, stripe) in memberships.iter_mut().zip(&encoded) {
1297 if stripe.codes.iter().all(Option::is_none) {
1298 continue;
1299 }
1300 let lists = stripe
1301 .codes
1302 .iter()
1303 .map(|codes| codes.clone().unwrap_or_default())
1304 .collect::<Vec<_>>();
1305 let bytes = encode_membership(&merged_codes(lists));
1306 let offset = self.at;
1307 self.put(&bytes)?;
1308 *membership = Some(Page {
1309 offset,
1310 length: u32::try_from(bytes.len())
1311 .map_err(|_| invalid("membership page length overflow"))?,
1312 hash: checksum(&bytes),
1313 });
1314 }
1315 let mut sieves = vec![None; width];
1316 for (page, stripe) in sieves.iter_mut().zip(&encoded) {
1317 if stripe.sieves.iter().all(Option::is_none) {
1318 continue;
1319 }
1320 let bytes = encode_sieves(stripe.sieves.iter())?;
1321 let offset = self.at;
1322 self.put(&bytes)?;
1323 *page = Some(Page {
1324 offset,
1325 length: u32::try_from(bytes.len())
1326 .map_err(|_| invalid("sieve page length overflow"))?,
1327 hash: checksum(&bytes),
1328 });
1329 }
1330 let mut part_ranges = vec![None; width];
1336 if parts > 1 {
1337 for ((page, stripe), span) in part_ranges.iter_mut().zip(&encoded).zip(&pages) {
1338 let bytes = encode_part_ranges(&stripe.ranges)?;
1339 if bytes.len() >= span.length as usize {
1340 continue;
1341 }
1342 let offset = self.at;
1343 self.put(&bytes)?;
1344 *page = Some(Page {
1345 offset,
1346 length: u32::try_from(bytes.len())
1347 .map_err(|_| invalid("part range page length overflow"))?,
1348 hash: checksum(&bytes),
1349 });
1350 }
1351 }
1352 let offset = self.at;
1353 self.put(&index)?;
1354 let index = Span {
1355 offset,
1356 length: u32::try_from(index.len())
1357 .map_err(|_| invalid("index page length overflow"))?,
1358 };
1359 let mut rows = 0_usize;
1360 let mut lengths = Vec::with_capacity(parts);
1361 let mut span = None;
1362 for pending in held.drain(..) {
1363 let part = pending.chunk.len();
1364 rows = rows.checked_add(part).ok_or_else(|| invalid("row count overflow"))?;
1365 lengths.push(u32::try_from(part).map_err(|_| invalid("part row count overflow"))?);
1366 span = Some(
1367 span.map_or((pending.order, pending.order), |(first, _)| (first, pending.order)),
1368 );
1369 }
1370 self.order.push(span.ok_or_else(|| invalid("a stripe was flushed with no parts"))?);
1371 self.table.stripes.push(Stripe {
1372 rows,
1373 parts: lengths,
1374 index,
1375 pages,
1376 memberships,
1377 sieves,
1378 part_ranges,
1379 zone: Zone::from_ranges(ranges),
1380 });
1381 self.pending = held;
1383 Ok(())
1384 }
1385
1386 fn numeric_frequency(&self, column: usize) -> Result<Option<FrequencySummary>> {
1390 let ty = &self.table.fields[column].ty;
1391 if !matches!(
1392 ty,
1393 LogicalType::TinyInt
1394 | LogicalType::SmallInt
1395 | LogicalType::Integer
1396 | LogicalType::BigInt
1397 | LogicalType::UTinyInt
1398 | LogicalType::USmallInt
1399 | LogicalType::UInteger
1400 | LogicalType::UBigInt
1401 | LogicalType::Date
1402 | LogicalType::Timestamp
1403 ) {
1404 return Ok(None);
1405 }
1406 let mut candidates: HashMap<FrequencyValue, u32> = HashMap::new();
1407 let mut decrements = 0_u64;
1408 self.visit_numeric(column, |_, value| {
1409 if let Some(count) = candidates.get_mut(&value) {
1410 *count = count.saturating_add(1);
1411 } else if candidates.len() < FREQUENCY_CANDIDATES {
1412 candidates.insert(value, 1);
1413 } else {
1414 candidates.retain(|_, count| {
1415 *count -= 1;
1416 *count != 0
1417 });
1418 decrements = decrements.saturating_add(1);
1419 }
1420 })?;
1421 let (exact, ordinals) = if decrements == 0 {
1422 (
1423 candidates
1424 .into_iter()
1425 .map(|(value, count)| (value, u64::from(count)))
1426 .collect::<HashMap<_, _>>(),
1427 Vec::new(),
1428 )
1429 } else {
1430 let mut lower = candidates.values().copied().collect::<Vec<_>>();
1431 lower.sort_unstable_by(|left, right| right.cmp(left));
1432 if lower.len() < FREQUENCY_BUILD_RANK
1433 || u64::from(lower[FREQUENCY_BUILD_RANK - 1]) <= decrements
1434 {
1435 return Ok(None);
1436 }
1437 let mut exact =
1438 candidates.into_keys().map(|value| (value, 0_u64)).collect::<HashMap<_, _>>();
1439 let mut ordinals = Vec::new();
1440 let mut exceeded = false;
1441 self.visit_numeric(column, |ordinal, value| {
1442 if let Some(count) = exact.get_mut(&value) {
1443 *count = count.saturating_add(1);
1444 if !exceeded {
1445 if ordinals.len() < FREQUENCY_ORDINALS {
1446 ordinals.push(ordinal);
1447 } else {
1448 ordinals.clear();
1449 exceeded = true;
1450 }
1451 }
1452 }
1453 })?;
1454 (exact, ordinals)
1455 };
1456 let mut entries = exact
1457 .into_iter()
1458 .map(|(value, count)| FrequencyEntry { value, count })
1459 .collect::<Vec<_>>();
1460 entries.sort_unstable_by(|left, right| {
1461 right.count.cmp(&left.count).then_with(|| frequency_order(left.value, right.value))
1462 });
1463 let omitted_max =
1464 entries.get(FREQUENCY_ENTRIES).map_or(decrements, |entry| decrements.max(entry.count));
1465 entries.truncate(FREQUENCY_ENTRIES);
1466 Ok(Some(FrequencySummary { entries, omitted_max, ordinals }))
1467 }
1468
1469 fn visit_numeric(
1470 &self,
1471 column: usize,
1472 mut visit: impl FnMut(u64, FrequencyValue),
1473 ) -> Result<()> {
1474 let ty = &self.table.fields[column].ty;
1475 let mut start = 0_u64;
1476 for stripe in &self.table.stripes {
1477 let spans = read_index(&self.file, stripe, column)?;
1478 let page = stripe.pages[column];
1479 let mut bytes = vec![0; page.length as usize];
1480 read_at(&self.file, page.offset, &mut bytes)?;
1481 for (span, &rows) in spans.iter().zip(&stripe.parts) {
1482 let part = part_bytes(&bytes, *span)?;
1483 if checksum(part) != span.hash {
1484 return Err(invalid("column page checksum differs while building frequencies"));
1485 }
1486 let rows = rows as usize;
1487 let vector = decode(ty, rows, part, None)?;
1488 for row in 0..rows {
1490 let value = if vector.is_null_at(row) {
1491 FrequencyValue::Null
1492 } else {
1493 let widened = match vector.signed_at(row) {
1497 Some(value) => Some(value),
1498 None => match vector.value_at(row) {
1499 Value::UTinyInt(value) => Some(i128::from(value)),
1500 Value::USmallInt(value) => Some(i128::from(value)),
1501 Value::UInteger(value) => Some(i128::from(value)),
1502 Value::UBigInt(value) => Some(i128::from(value)),
1503 _ => None,
1504 },
1505 };
1506 FrequencyValue::Integer(widened.ok_or_else(|| {
1507 invalid("numeric frequency page did not contain an integer value")
1508 })?)
1509 };
1510 visit(start.saturating_add(row as u64), value);
1511 }
1512 start = start.saturating_add(rows as u64);
1513 }
1514 }
1515 Ok(())
1516 }
1517
1518 fn numeric_frequencies(&self) -> Result<Vec<Option<FrequencySummary>>> {
1526 let mut columns = self
1527 .table
1528 .fields
1529 .iter()
1530 .enumerate()
1531 .filter_map(|(column, field)| {
1532 matches!(
1533 field.ty,
1534 LogicalType::TinyInt
1535 | LogicalType::SmallInt
1536 | LogicalType::Integer
1537 | LogicalType::BigInt
1538 | LogicalType::UTinyInt
1539 | LogicalType::USmallInt
1540 | LogicalType::UInteger
1541 | LogicalType::UBigInt
1542 | LogicalType::Date
1543 | LogicalType::Timestamp
1544 )
1545 .then_some(column)
1546 })
1547 .collect::<Vec<_>>();
1548 let workers = std::thread::available_parallelism()
1549 .map_or(1, usize::from)
1550 .min(MAX_FREQUENCY_WORKERS)
1551 .min(columns.len());
1552 if workers <= 1 {
1553 let mut frequencies = vec![None; self.table.fields.len()];
1554 for column in columns {
1555 frequencies[column] = self.numeric_frequency(column)?;
1556 }
1557 return Ok(frequencies);
1558 }
1559 columns.sort_by_key(|&column| weight(&self.table.fields[column].ty));
1562 let queue = Mutex::new(columns);
1563 let pieces = std::thread::scope(|scope| {
1564 (0..workers)
1565 .map(|_| {
1566 scope.spawn(|| {
1567 let mut mine = Vec::new();
1568 loop {
1569 let taken = queue
1570 .lock()
1571 .map_err(|_| Error::internal("a native frequency worker panicked"))?
1572 .pop();
1573 let Some(column) = taken else { break };
1574 mine.push((column, self.numeric_frequency(column)?));
1575 }
1576 Ok(mine)
1577 })
1578 })
1579 .collect::<Vec<_>>()
1580 .into_iter()
1581 .map(|handle| {
1582 handle
1583 .join()
1584 .map_err(|_| Error::internal("a native frequency worker panicked"))?
1585 })
1586 .collect::<Result<Vec<_>>>()
1587 })?;
1588 let mut frequencies = vec![None; self.table.fields.len()];
1589 for piece in pieces {
1590 for (column, summary) in piece {
1591 frequencies[column] = summary;
1592 }
1593 }
1594 Ok(frequencies)
1595 }
1596
1597 fn close(&mut self) -> Result<Entry> {
1608 self.flush_pending()?;
1609 let mut stripes = std::mem::take(&mut self.order)
1610 .into_iter()
1611 .zip(std::mem::take(&mut self.table.stripes))
1612 .collect::<Vec<_>>();
1613 stripes.sort_by_key(|(order, _)| order.0);
1614 let mut previous: Option<(u64, u64)> = None;
1615 for ((first, last), _) in &stripes {
1616 if previous.is_some_and(|previous| previous >= *first) {
1617 return Err(invalid("chunks did not arrive in source order"));
1618 }
1619 previous = Some(*last);
1620 }
1621 self.table.stripes = stripes.into_iter().map(|(_, stripe)| stripe).collect();
1622 self.table.frequencies = self.numeric_frequencies()?;
1623 let dictionaries = std::mem::take(&mut self.dictionaries);
1624 let orders = rankings(&dictionaries)?;
1625 for (index, (dictionary, order)) in dictionaries.into_iter().zip(orders).enumerate() {
1626 let Some(dictionary) = dictionary else { continue };
1627 self.table.distincts[index] =
1631 Some(dictionary.counts.iter().filter(|count| **count != 0).count() as u64);
1632 self.table.frequencies[index] = Some(code_frequency(&dictionary));
1633 let encoded = encode_global_dictionary(dictionary, &order)?;
1634 let offset = self.at;
1635 self.put(&encoded.index)?;
1636 self.put(&encoded.ranks)?;
1637 for block in &encoded.payload {
1638 self.put(block)?;
1639 }
1640 let payload_len =
1641 encoded.payload.iter().try_fold(0_usize, |len, block| len.checked_add(block.len()));
1642 let length = payload_len
1643 .and_then(|len| len.checked_add(encoded.index.len()))
1644 .and_then(|len| len.checked_add(encoded.ranks.len()))
1645 .ok_or_else(|| invalid("dictionary page length overflow"))?;
1646 self.table.dictionaries[index] = Some(Page {
1647 offset,
1648 length: u32::try_from(length)
1649 .map_err(|_| invalid("dictionary page length overflow"))?,
1650 hash: checksum(&encoded.index),
1651 });
1652 }
1653 let directory = encode_directory(&self.table)?;
1654 if directory.len() > MAX_DIRECTORY {
1655 return Err(invalid("directory exceeds the configured bound"));
1656 }
1657 let offset = self.at;
1658 self.put(&directory)?;
1659 Ok(Entry {
1660 name: self.table.name.clone(),
1661 fields: self.table.fields.clone(),
1662 rows: self.table.rows,
1663 directory: Page {
1664 offset,
1665 length: u32::try_from(directory.len())
1666 .map_err(|_| invalid("directory length overflow"))?,
1667 hash: checksum(&directory),
1668 },
1669 })
1670 }
1671
1672 pub fn finish(mut self) -> Result<Table> {
1682 let entry = self.close()?;
1683 let mut tables = std::mem::take(&mut self.closed);
1684 tables.push(entry);
1685 let catalog = encode_catalog(&tables)?;
1686 if catalog.len() > MAX_DIRECTORY {
1687 return Err(invalid("catalog exceeds the configured bound"));
1688 }
1689 let offset = self.at;
1690 self.put(&catalog)?;
1691 self.file.sync_all().map_err(io)?;
1695 let slot = Slot {
1696 offset,
1697 length: u32::try_from(catalog.len()).map_err(|_| invalid("catalog length overflow"))?,
1698 generation: self.generation,
1699 hash: checksum(&catalog),
1700 };
1701 write_at(&self.file, slot_offset(self.generation), &slot.bytes())?;
1706 self.file.sync_all().map_err(io)?;
1707 Ok(self.table)
1708 }
1709}
1710
1711#[derive(Debug, Clone)]
1713pub struct Reader {
1714 file: Arc<File>,
1715 table: Arc<Table>,
1716 dictionaries: Arc<Vec<OnceLock<Arc<Vector>>>>,
1717 loading: Arc<Vec<Mutex<()>>>,
1726 opened: Arc<AtomicUsize>,
1730 sieves: Arc<Vec<Vec<SieveSlot>>>,
1734 part_ranges: Arc<Vec<Vec<RangeSlot>>>,
1737 places: Arc<Vec<Place>>,
1739 cache: Arc<Vec<Mutex<Cached>>>,
1740 pages: Arc<AtomicUsize>,
1743 indexes: Arc<AtomicUsize>,
1746 kept: Arc<AtomicUsize>,
1749 size: u64,
1751 directory: u64,
1753 opening: Opening,
1755}
1756
1757#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1769pub struct Opening {
1770 pub reads: u32,
1773 pub bytes: u64,
1775}
1776
1777#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1779pub struct Reads {
1780 pub opening: Opening,
1782 pub pages: usize,
1784 pub indexes: usize,
1786 pub dictionaries: usize,
1789}
1790
1791#[derive(Debug, Clone, Copy)]
1793struct Place {
1794 stripe: u32,
1795 part: u32,
1796 rows: u32,
1797}
1798
1799#[derive(Debug, Clone, Copy)]
1801struct PartSpan {
1802 start: usize,
1803 length: usize,
1804 hash: u64,
1805}
1806
1807#[derive(Debug, Clone)]
1813struct CachedColumn {
1814 stripe: usize,
1815 index: Arc<Vec<PartSpan>>,
1816 page: Option<Arc<Vec<u8>>>,
1817}
1818
1819#[derive(Debug, Default)]
1839struct Cached {
1840 pages: Vec<Option<Arc<Vec<u8>>>>,
1841 order: VecDeque<usize>,
1842 loading: Vec<usize>,
1843 index: Vec<Option<Arc<Vec<PartSpan>>>>,
1844}
1845
1846const CACHED_STRIPES_PER_COLUMN: usize = 4;
1858
1859type SieveSlot = OnceLock<Arc<Vec<Option<Sieve>>>>;
1861
1862type RangeSlot = OnceLock<Arc<Vec<Range>>>;
1863
1864#[derive(Debug)]
1865struct NativeText {
1866 file: Arc<File>,
1867 values: usize,
1869 offsets: Vec<u8>,
1878 offset_bits: usize,
1881 ranks: usize,
1883 rank_at: u64,
1887 rank_ends: Vec<u64>,
1891 rank_hashes: Vec<u64>,
1892 rank_blocks: Vec<OnceLock<Result<Vec<u8>>>>,
1893 code_bits: usize,
1896 code_ranks: OnceLock<Option<Vec<u32>>>,
1903 payload: u64,
1904 ends: Vec<u64>,
1907 hashes: Vec<u64>,
1908 blocks: Vec<OnceLock<Result<Vec<u8>>>>,
1910 keep_budget: usize,
1913 payload_kept: AtomicUsize,
1921 searched: Mutex<HashMap<Vec<u8>, (usize, bool)>>,
1938}
1939
1940const TEXT_SEARCH_MEMO: usize = 64;
1945
1946const TEXT_PAYLOAD_VALUES: usize = 1024;
1962
1963const TEXT_KEEP_BUDGET: usize = 256 * 1024 * 1024;
1984
1985const TEXT_OFFSET_RUN: usize = 512;
1992
1993const DICTIONARY_HEADER: usize = 16;
1996
1997const TEXT_RANK_BLOCK: usize = 512;
2008
2009const RANK_BLOCK_HEADER: usize = size_of::<u64>() + 1;
2023
2024impl NativeText {
2025 fn payload_block(&self, block: usize) -> Result<Option<&[u8]>> {
2032 let Some(slot) = self.blocks.get(block) else { return Ok(None) };
2033 let bytes = slot.get_or_init(|| self.decode_block(block)).as_ref().map_err(Clone::clone)?;
2034 Ok(Some(bytes.as_slice()))
2035 }
2036
2037 fn decode_block(&self, block: usize) -> Result<Vec<u8>> {
2042 let start = if block == 0 { 0 } else { self.ends[block - 1] };
2043 let end = self.ends[block];
2044 let len = end
2045 .checked_sub(start)
2046 .ok_or_else(|| invalid("global dictionary block ends before it starts"))?;
2047 let mut stored = vec![
2048 0;
2049 usize::try_from(len).map_err(|_| invalid(
2050 "global dictionary block does not fit in memory"
2051 ))?
2052 ];
2053 read_at(&self.file, self.payload + start, &mut stored)?;
2054 if checksum(&stored) != self.hashes[block] {
2055 return Err(invalid("global dictionary payload checksum differs"));
2056 }
2057 let first = block * TEXT_PAYLOAD_VALUES;
2058 let last = (first + TEXT_PAYLOAD_VALUES).min(self.values);
2059 let want = self.end_within(last - 1)? as usize;
2060 let values = string::decode_flat(&stored)?;
2061 if values.len() != last - first {
2062 return Err(invalid("global dictionary block holds the wrong value count"));
2063 }
2064 let bytes = values.into_bytes();
2065 if bytes.len() != want {
2066 return Err(invalid("global dictionary block decodes to the wrong length"));
2067 }
2068 Ok(bytes)
2069 }
2070
2071 fn end_within(&self, index: usize) -> Result<u32> {
2073 let run = index / TEXT_OFFSET_RUN;
2074 let bytes = self
2075 .offsets
2076 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
2077 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
2078 let end = bitpack::tail_at(bytes, self.offset_bits, index % TEXT_OFFSET_RUN)
2079 .map_err(|_| invalid("global dictionary offsets are short"))?;
2080 u32::try_from(end).map_err(|_| invalid("global dictionary offset is past the payload"))
2081 }
2082
2083 fn ends_within(&self, first: usize, last: usize) -> Result<Vec<u64>> {
2096 let mut ends = Vec::with_capacity(last.saturating_sub(first));
2097 let mut at = first;
2098 while at < last {
2099 let run = at / TEXT_OFFSET_RUN;
2100 let stop = ((run + 1) * TEXT_OFFSET_RUN).min(last);
2101 let held = self.values.saturating_sub(run * TEXT_OFFSET_RUN).min(TEXT_OFFSET_RUN);
2102 let bytes = self
2103 .offsets
2104 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
2105 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
2106 let run_ends = bitpack::unpack_tail(bytes, self.offset_bits, held)
2107 .map_err(|_| invalid("global dictionary offsets are short"))?;
2108 let within = run_ends
2109 .get(at % TEXT_OFFSET_RUN..stop - run * TEXT_OFFSET_RUN)
2110 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
2111 ends.extend_from_slice(within);
2112 at = stop;
2113 }
2114 Ok(ends)
2115 }
2116
2117 fn start_within(&self, index: usize) -> Result<u32> {
2120 if index % TEXT_PAYLOAD_VALUES == 0 { Ok(0) } else { self.end_within(index - 1) }
2121 }
2122
2123 fn span_within(&self, index: usize) -> Result<(u32, u32)> {
2131 let within = index % TEXT_OFFSET_RUN;
2132 let (start, end) = if within == 0 {
2133 (self.start_within(index)?, self.end_within(index)?)
2134 } else {
2135 let run = index / TEXT_OFFSET_RUN;
2136 let bytes = self
2137 .offsets
2138 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
2139 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
2140 let (start, end) = bitpack::tail_pair(bytes, self.offset_bits, within)
2141 .map_err(|_| invalid("global dictionary offsets are short"))?;
2142 let ends = u32::try_from(end)
2143 .map_err(|_| invalid("global dictionary offset is past the payload"))?;
2144 let starts = u32::try_from(start)
2145 .map_err(|_| invalid("global dictionary offset is past the payload"))?;
2146 (starts, ends)
2147 };
2148 if start > end {
2149 return Err(invalid("global dictionary value ends before it starts"));
2150 }
2151 Ok((start, end))
2152 }
2153
2154 fn rank_parts(&self, rank: usize) -> Result<(&[u8], usize)> {
2161 let slot = self
2162 .rank_blocks
2163 .get(rank / TEXT_RANK_BLOCK)
2164 .ok_or_else(|| invalid("global dictionary rank is past the order"))?;
2165 let block = slot
2166 .get_or_init(|| {
2167 let which = rank / TEXT_RANK_BLOCK;
2168 let start = if which == 0 { 0 } else { self.rank_ends[which - 1] };
2169 let end = self.rank_ends[which];
2170 let mut bytes = vec![0; (end - start) as usize];
2171 read_at(&self.file, self.rank_at + start, &mut bytes)?;
2172 if checksum(&bytes)
2173 != *self
2174 .rank_hashes
2175 .get(rank / TEXT_RANK_BLOCK)
2176 .ok_or_else(|| invalid("global dictionary rank block has no checksum"))?
2177 {
2178 return Err(invalid("global dictionary rank checksum differs"));
2179 }
2180 Ok(bytes)
2181 })
2182 .as_ref()
2183 .map_err(Clone::clone)?;
2184 Ok((block.as_slice(), rank % TEXT_RANK_BLOCK))
2185 }
2186
2187 fn head_at(&self, rank: usize) -> Result<u64> {
2189 let (block, within) = self.rank_parts(rank)?;
2190 let (base, width, packed) = rank_heads(block)?;
2191 let above = bitpack::tail_at(packed, width, within)
2192 .map_err(|_| invalid("global dictionary rank block is short of heads"))?;
2193 Ok(base.wrapping_add(above))
2194 }
2195
2196 fn rank_codes<'block>(&self, block: &'block [u8], count: usize) -> Result<&'block [u8]> {
2198 let (_, width, packed) = rank_heads(block)?;
2199 packed
2200 .get(bitpack::tail_len(count, width)..)
2201 .ok_or_else(|| invalid("global dictionary rank block is short of codes"))
2202 }
2203
2204 fn rank_block_len(&self, rank: usize) -> usize {
2206 let first = rank / TEXT_RANK_BLOCK * TEXT_RANK_BLOCK;
2207 TEXT_RANK_BLOCK.min(self.ranks - first)
2208 }
2209}
2210
2211fn rank_heads(block: &[u8]) -> Result<(u64, usize, &[u8])> {
2213 let header = block
2214 .get(..RANK_BLOCK_HEADER)
2215 .ok_or_else(|| invalid("global dictionary rank block is short"))?;
2216 let base = u64::from_le_bytes(header[..8].try_into().expect("eight bytes"));
2217 let width = header[8] as usize;
2218 if width > 64 {
2219 return Err(invalid("global dictionary rank block packs heads past a word"));
2220 }
2221 Ok((base, width, &block[RANK_BLOCK_HEADER..]))
2222}
2223
2224fn offset_width(offsets: &[u32]) -> usize {
2231 let values = offsets.len() - 1;
2232 let mut span = 0;
2233 for first in (0..values).step_by(TEXT_PAYLOAD_VALUES) {
2234 let last = (first + TEXT_PAYLOAD_VALUES).min(values);
2235 span = span.max(offsets[last] - offsets[first]);
2236 }
2237 (u32::BITS - span.leading_zeros()) as usize
2238}
2239
2240fn offset_bytes(values: usize, bits: usize) -> usize {
2243 let full = values / TEXT_OFFSET_RUN;
2244 let rest = values % TEXT_OFFSET_RUN;
2245 full * TEXT_OFFSET_RUN / 8 * bits + bitpack::tail_len(rest, bits)
2246}
2247
2248fn encode_offsets(offsets: &[u32], bits: usize, out: &mut Vec<u8>) -> Result<()> {
2250 let values = offsets.len() - 1;
2251 let mut run = Vec::with_capacity(TEXT_OFFSET_RUN);
2252 for first in (0..values).step_by(TEXT_OFFSET_RUN) {
2253 let last = (first + TEXT_OFFSET_RUN).min(values);
2254 let base = offsets[first / TEXT_PAYLOAD_VALUES * TEXT_PAYLOAD_VALUES];
2255 run.clear();
2256 run.extend((first..last).map(|value| u64::from(offsets[value + 1] - base)));
2257 bitpack::pack_tail(&run, bits, out)
2258 .map_err(|_| invalid("global dictionary offsets do not pack"))?;
2259 }
2260 Ok(())
2261}
2262
2263fn code_width(values: usize) -> usize {
2265 match u64::try_from(values).unwrap_or(u64::MAX) {
2266 0 | 1 => 0,
2267 last => (u64::BITS - (last - 1).leading_zeros()) as usize,
2268 }
2269}
2270
2271impl TextSource for NativeText {
2272 fn len(&self) -> usize {
2273 self.values
2274 }
2275
2276 fn bytes_at(&self, index: usize) -> Result<Option<&[u8]>> {
2277 if index >= self.values {
2278 return Ok(None);
2279 }
2280 let (start, end) = self.span_within(index)?;
2281 if start == end {
2282 return Ok(Some(&[]));
2283 }
2284 let block = index / TEXT_PAYLOAD_VALUES;
2287 let Some(bytes) = self.payload_block(block)? else { return Ok(None) };
2288 Ok(bytes.get(start as usize..end as usize))
2289 }
2290
2291 fn bytes_len_at(&self, index: usize) -> Result<Option<usize>> {
2292 if index >= self.values {
2293 return Ok(None);
2294 }
2295 let (start, end) = self.span_within(index)?;
2296 Ok(Some((end - start) as usize))
2297 }
2298
2299 fn sweep(
2312 &self,
2313 first: usize,
2314 limit: usize,
2315 body: &mut dyn FnMut(usize, &[u8]) -> Result<()>,
2316 ) -> Result<usize> {
2317 let limit = limit.min(self.values);
2318 if first >= limit {
2319 return Ok(first);
2320 }
2321 let block = first / TEXT_PAYLOAD_VALUES;
2322 let last = ((block + 1) * TEXT_PAYLOAD_VALUES).min(limit);
2323 let decoded;
2324 let bytes: &[u8] = match self.blocks.get(block).and_then(OnceLock::get) {
2325 Some(Ok(kept)) => kept,
2326 _ if self.payload_kept.load(Atomic::Relaxed) < self.keep_budget => {
2327 let kept = self
2328 .payload_block(block)?
2329 .ok_or_else(|| invalid("global dictionary block is past the payload"))?;
2330 self.payload_kept.fetch_add(kept.len(), Atomic::Relaxed);
2331 kept
2332 }
2333 _ => {
2334 decoded = self.decode_block(block)?;
2335 &decoded
2336 }
2337 };
2338 let ends = self.ends_within(first, last)?;
2339 if ends.len() != last - first {
2340 return Err(invalid("global dictionary offsets are short"));
2341 }
2342 let mut start = u64::from(self.start_within(first)?);
2343 for (index, &end) in (first..last).zip(&ends) {
2346 let value = usize::try_from(start)
2347 .ok()
2348 .zip(usize::try_from(end).ok())
2349 .and_then(|(from, to)| bytes.get(from..to))
2350 .ok_or_else(|| invalid("global dictionary value is past its block"))?;
2351 body(index, value)?;
2352 start = end;
2353 }
2354 Ok(last)
2355 }
2356
2357 fn ranks(&self) -> Option<usize> {
2358 (self.ranks > 0).then_some(self.ranks)
2359 }
2360
2361 fn below(&self, ranks: usize, wanted: &[u8]) -> Result<(usize, bool)> {
2369 let mut memo = self.searched.lock().map_err(|_| invalid("a poisoned dictionary search"))?;
2370 if let Some(&answer) = memo.get(wanted) {
2371 return Ok(answer);
2372 }
2373 let answer = search_below(self, ranks, wanted)?;
2374 if memo.len() >= TEXT_SEARCH_MEMO {
2375 memo.clear();
2376 }
2377 memo.insert(wanted.to_vec(), answer);
2378 Ok(answer)
2379 }
2380
2381 fn compare_rank(&self, rank: usize, wanted: &[u8]) -> Result<Ordering> {
2382 let settled = self.head_at(rank)?.cmp(&head(wanted));
2386 if settled != Ordering::Equal {
2387 return Ok(settled);
2388 }
2389 let code = self.code_at_rank(rank)?;
2390 let bytes = self
2391 .bytes_at(code as usize)?
2392 .ok_or_else(|| invalid("global dictionary order names a code it does not have"))?;
2393 Ok(bytes.cmp(wanted))
2394 }
2395
2396 fn code_at_rank(&self, rank: usize) -> Result<u32> {
2397 let (block, within) = self.rank_parts(rank)?;
2398 let codes = self.rank_codes(block, self.rank_block_len(rank))?;
2399 let code = bitpack::tail_at(codes, self.code_bits, within)
2400 .map_err(|_| invalid("global dictionary rank block is short of codes"))?;
2401 let code = u32::try_from(code)
2402 .map_err(|_| invalid("global dictionary order names a code it does not have"))?;
2403 if code as usize >= self.len() {
2404 return Err(invalid("global dictionary order names a code it does not have"));
2405 }
2406 Ok(code)
2407 }
2408
2409 fn code_ranks(&self) -> Option<&[u32]> {
2410 if self.ranks == 0 || self.ranks != self.len() {
2414 return None;
2415 }
2416 self.code_ranks
2417 .get_or_init(|| {
2418 let mut ranks = vec![u32::MAX; self.ranks];
2419 for first in (0..self.ranks).step_by(TEXT_RANK_BLOCK) {
2422 let (block, _) = self.rank_parts(first).ok()?;
2423 let count = self.rank_block_len(first);
2424 let codes = self.rank_codes(block, count).ok()?;
2425 for (within, code) in bitpack::unpack_tail(codes, self.code_bits, count)
2426 .ok()?
2427 .into_iter()
2428 .enumerate()
2429 {
2430 let code = usize::try_from(code).ok()?;
2431 *ranks.get_mut(code)? = u32::try_from(first + within).ok()?;
2432 }
2433 }
2434 if ranks.contains(&u32::MAX) {
2435 return None;
2436 }
2437 Some(ranks)
2438 })
2439 .as_deref()
2440 }
2441
2442 fn footprint(&self) -> usize {
2443 self.offsets.capacity()
2444 + self
2445 .code_ranks
2446 .get()
2447 .and_then(Option::as_ref)
2448 .map_or(0, |ranks| ranks.capacity() * size_of::<u32>())
2449 + self.rank_hashes.capacity() * size_of::<u64>()
2450 + self.rank_ends.capacity() * size_of::<u64>()
2451 + self.rank_blocks.capacity() * size_of::<OnceLock<Result<Vec<u8>>>>()
2452 + self
2453 .rank_blocks
2454 .iter()
2455 .filter_map(OnceLock::get)
2456 .filter_map(|result| result.as_ref().ok())
2457 .map(Vec::capacity)
2458 .sum::<usize>()
2459 + self.blocks.capacity() * size_of::<OnceLock<Result<Vec<u8>>>>()
2460 + self.hashes.capacity() * size_of::<u64>()
2461 + self.ends.capacity() * size_of::<u64>()
2462 + self
2463 .blocks
2464 .iter()
2465 .filter_map(OnceLock::get)
2466 .filter_map(|result| result.as_ref().ok())
2467 .map(Vec::capacity)
2468 .sum::<usize>()
2469 }
2470}
2471
2472fn places(table: &Table) -> Result<Vec<Place>> {
2474 let mut places = Vec::with_capacity(table.stripes.len().saturating_mul(STRIPE_PARTS));
2475 for (at, stripe) in table.stripes.iter().enumerate() {
2476 let index = u32::try_from(at).map_err(|_| invalid("too many stripes"))?;
2477 for (part, &rows) in stripe.parts.iter().enumerate() {
2478 places.push(Place {
2479 stripe: index,
2480 part: u32::try_from(part).map_err(|_| invalid("too many parts in a stripe"))?,
2481 rows,
2482 });
2483 }
2484 }
2485 Ok(places)
2486}
2487
2488fn read_index(file: &File, stripe: &Stripe, column: usize) -> Result<Vec<PartSpan>> {
2493 let parts = stripe.parts.len();
2494 let section = index_section(parts)?;
2495 let at = column.checked_mul(section).ok_or_else(|| invalid("index page offset overflow"))?;
2496 let end = at.checked_add(section).ok_or_else(|| invalid("index page offset overflow"))?;
2497 if end > stripe.index.length as usize {
2498 return Err(invalid("index page is shorter than its columns"));
2499 }
2500 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
2501 let mut bytes = vec![0; section];
2502 let offset = stripe
2503 .index
2504 .offset
2505 .checked_add(at as u64)
2506 .ok_or_else(|| invalid("index page offset overflow"))?;
2507 read_at(file, offset, &mut bytes)?;
2508 let entries = section - size_of::<u64>();
2509 let stored = u64::from_le_bytes(bytes[entries..].try_into().expect("eight bytes"));
2510 if checksum(&bytes[..entries]) != stored {
2511 return Err(invalid(&format!(
2514 "index page section checksum differs, column {column} of {parts} parts at {offset}, \
2515 wanted {stored:016x} and got {:016x}",
2516 checksum(&bytes[..entries]),
2517 )));
2518 }
2519 let mut spans = Vec::with_capacity(parts);
2520 let mut start = 0_usize;
2521 for part in 0..parts {
2522 let at = part * INDEX_ENTRY;
2523 let length = u32::from_le_bytes(bytes[at..at + 4].try_into().expect("four bytes")) as usize;
2524 let hash = u64::from_le_bytes(bytes[at + 4..at + 12].try_into().expect("eight bytes"));
2525 spans.push(PartSpan { start, length, hash });
2526 start = start.checked_add(length).ok_or_else(|| invalid("column page length overflow"))?;
2527 }
2528 if start != page.length as usize {
2529 return Err(invalid("column page length differs from its index"));
2530 }
2531 Ok(spans)
2532}
2533
2534fn part_bytes(page: &[u8], span: PartSpan) -> Result<&[u8]> {
2536 let end = span.start.checked_add(span.length).ok_or_else(|| invalid("part range overflow"))?;
2537 page.get(span.start..end).ok_or_else(|| invalid("part exceeds its column page"))
2538}
2539
2540fn remember(cached: &mut Cached, held: &CachedColumn, kept: usize) {
2545 if let Some(slot) = cached.index.get_mut(held.stripe) {
2546 if slot.is_none() {
2547 *slot = Some(Arc::clone(&held.index));
2548 }
2549 }
2550 let Some(page) = held.page.clone() else { return };
2551 let Some(slot) = cached.pages.get_mut(held.stripe) else { return };
2552 if slot.is_none() {
2553 cached.order.push_back(held.stripe);
2554 }
2555 *slot = Some(page);
2556 while cached.order.len() > kept.max(1) {
2557 let Some(oldest) = cached.order.pop_front() else { break };
2558 if let Some(slot) = cached.pages.get_mut(oldest) {
2559 *slot = None;
2560 }
2561 }
2562}
2563
2564#[derive(Debug, Clone)]
2573pub struct Catalog {
2574 file: Arc<File>,
2575 size: u64,
2576 entries: Arc<Vec<Entry>>,
2577 opening: Opening,
2578}
2579
2580impl Catalog {
2581 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
2587 let (file, size, _, bytes, opening) = slot_bytes(path)?;
2588 let entries = decode_catalog(&bytes, size)?;
2589 Ok(Self { file: Arc::new(file), size, entries: Arc::new(entries), opening })
2590 }
2591
2592 pub fn names(&self) -> impl ExactSizeIterator<Item = &str> {
2594 self.entries.iter().map(|entry| entry.name.as_str())
2595 }
2596
2597 #[must_use]
2599 pub fn len(&self) -> usize {
2600 self.entries.len()
2601 }
2602
2603 #[must_use]
2606 pub fn is_empty(&self) -> bool {
2607 self.entries.is_empty()
2608 }
2609
2610 pub fn table(&self, name: &str) -> Result<Reader> {
2616 let entry = self
2617 .entries
2618 .iter()
2619 .find(|entry| entry.name == name)
2620 .ok_or_else(|| invalid(&format!("the file holds no table called {name}")))?;
2621 let mut bytes = vec![0; entry.directory.length as usize];
2622 read_at(&self.file, entry.directory.offset, &mut bytes)?;
2623 if checksum(&bytes) != entry.directory.hash {
2624 return Err(invalid(&format!("the directory of table {name} does not checksum")));
2625 }
2626 let mut opening = self.opening;
2627 opening.reads += 1;
2628 opening.bytes += u64::from(entry.directory.length);
2629 Reader::build(
2630 Arc::clone(&self.file),
2631 self.size,
2632 decode_directory(&bytes, self.size)?,
2633 u64::from(entry.directory.length),
2634 opening,
2635 )
2636 }
2637}
2638
2639fn slot_offset(generation: u64) -> u64 {
2644 16 + (generation - 1) % 2 * SLOT_BYTES as u64
2645}
2646
2647fn slot_bytes(path: impl AsRef<Path>) -> Result<(File, u64, Slot, Vec<u8>, Opening)> {
2652 let mut file = File::open(path).map_err(io)?;
2653 let size = file.metadata().map_err(io)?.len();
2654 if size < HEADER {
2655 return Err(invalid("file is shorter than its header"));
2656 }
2657 let mut header = [0; HEADER as usize];
2658 file.read_exact(&mut header).map_err(io)?;
2659 let mut opening = Opening { reads: 1, bytes: HEADER };
2660 let version = u32::from_le_bytes([header[8], header[9], header[10], header[11]]);
2661 if &header[..8] != MAGIC {
2666 return Err(invalid("the header does not begin with a rudb native magic"));
2667 }
2668 if version != FORMAT {
2669 return Err(invalid(&format!(
2670 "the file is format {version} and this build reads format {FORMAT}, so it has to \
2671 be written again"
2672 )));
2673 }
2674 let mut selected = None;
2675 for start in [16, 16 + SLOT_BYTES] {
2676 let slot = Slot::read(&header[start..start + SLOT_BYTES]);
2677 if slot.generation == 0 || slot.length == 0 || slot.length as usize > MAX_DIRECTORY {
2678 continue;
2679 }
2680 let Some(end) = slot.offset.checked_add(u64::from(slot.length)) else { continue };
2681 if slot.offset < HEADER || end > size {
2682 continue;
2683 }
2684 let mut bytes = vec![0; slot.length as usize];
2685 file.seek(SeekFrom::Start(slot.offset)).map_err(io)?;
2686 file.read_exact(&mut bytes).map_err(io)?;
2687 opening.reads += 1;
2688 opening.bytes += u64::from(slot.length);
2689 if checksum(&bytes) == slot.hash
2690 && selected
2691 .as_ref()
2692 .is_none_or(|(old, _): &(Slot, Vec<u8>)| old.generation < slot.generation)
2693 {
2694 selected = Some((slot, bytes));
2695 }
2696 }
2697 let (slot, bytes) = selected.ok_or_else(|| invalid("no committed directory slot is valid"))?;
2698 Ok((file, size, slot, bytes, opening))
2699}
2700
2701impl Reader {
2702 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
2709 let catalog = Catalog::open(path)?;
2710 let mut names = catalog.names();
2711 let name = names.next().ok_or_else(|| invalid("the file holds no table"))?.to_string();
2712 if names.next().is_some() {
2713 return Err(invalid(
2714 "the file holds more than one table, so it has to be opened by name",
2715 ));
2716 }
2717 catalog.table(&name)
2718 }
2719
2720 fn build(
2722 file: Arc<File>,
2723 size: u64,
2724 table: Table,
2725 directory: u64,
2726 opening: Opening,
2727 ) -> Result<Self> {
2728 let places = places(&table)?;
2729 let dictionaries = (0..table.fields.len()).map(|_| OnceLock::new()).collect();
2730 let table_fields = table.fields.len();
2731 let stripes = table.stripes.len();
2732 let cache = (0..table.fields.len())
2733 .map(|_| {
2734 Mutex::new(Cached {
2735 pages: (0..stripes).map(|_| None).collect(),
2736 index: (0..stripes).map(|_| None).collect(),
2737 ..Cached::default()
2738 })
2739 })
2740 .collect::<Vec<_>>();
2741 let sieves: Vec<Vec<SieveSlot>> = (0..table.fields.len())
2742 .map(|_| table.stripes.iter().map(|_| OnceLock::new()).collect())
2743 .collect();
2744 let part_ranges: Vec<Vec<RangeSlot>> = (0..table.fields.len())
2745 .map(|_| table.stripes.iter().map(|_| OnceLock::new()).collect())
2746 .collect();
2747 Ok(Self {
2748 file,
2749 table: Arc::new(table),
2750 dictionaries: Arc::new(dictionaries),
2751 loading: Arc::new((0..table_fields).map(|_| Mutex::new(())).collect()),
2752 opened: Arc::new(AtomicUsize::new(0)),
2753 sieves: Arc::new(sieves),
2754 part_ranges: Arc::new(part_ranges),
2755 places: Arc::new(places),
2756 cache: Arc::new(cache),
2757 pages: Arc::new(AtomicUsize::new(0)),
2758 indexes: Arc::new(AtomicUsize::new(0)),
2759 kept: Arc::new(AtomicUsize::new(CACHED_STRIPES_PER_COLUMN)),
2760 size,
2761 directory,
2762 opening,
2763 })
2764 }
2765
2766 #[must_use]
2773 pub fn reads(&self) -> Reads {
2774 Reads {
2775 opening: self.opening,
2776 pages: self.pages.load(Atomic::Relaxed),
2777 indexes: self.indexes.load(Atomic::Relaxed),
2778 dictionaries: self.opened.load(Atomic::Relaxed),
2779 }
2780 }
2781
2782 #[must_use]
2787 pub fn layout(&self) -> Layout {
2788 let table = &self.table;
2789 let stripes = table.stripes.as_slice();
2790 let columns = table
2791 .fields
2792 .iter()
2793 .enumerate()
2794 .map(|(at, field)| ColumnLayout {
2795 name: field.name.clone(),
2796 kind: field.ty.to_string(),
2797 pages: sum(stripes.iter().map(|stripe| span_bytes(&stripe.pages, at))),
2798 memberships: sum(stripes.iter().map(|stripe| page_bytes(&stripe.memberships, at))),
2799 sieves: sum(stripes.iter().map(|stripe| page_bytes(&stripe.sieves, at))),
2800 part_ranges: sum(stripes.iter().map(|stripe| page_bytes(&stripe.part_ranges, at))),
2801 dictionary: page_bytes(&table.dictionaries, at),
2802 })
2803 .collect();
2804 Layout {
2805 file: self.size,
2806 rows: table.rows,
2807 stripes: stripes.len(),
2808 parts: self.places.len(),
2809 columns,
2810 indexes: sum(stripes.iter().map(|stripe| u64::from(stripe.index.length))),
2811 directory: self.directory,
2812 header: HEADER,
2813 }
2814 }
2815
2816 #[must_use]
2818 pub fn parts(&self) -> usize {
2819 self.places.len()
2820 }
2821
2822 #[must_use]
2829 pub fn stripe_parts(&self) -> Vec<std::ops::Range<usize>> {
2830 let mut runs = Vec::with_capacity(self.table.stripes.len());
2831 let mut start = 0;
2832 for stripe in &self.table.stripes {
2833 let end = start + stripe.parts.len();
2834 runs.push(start..end);
2835 start = end;
2836 }
2837 runs
2838 }
2839
2840 #[must_use]
2845 pub fn stripe_rows(&self, stripe: usize) -> usize {
2846 self.table.stripes.get(stripe).map_or(0, |held| held.rows)
2847 }
2848
2849 pub fn keep_stripes(&self, stripes: usize) {
2856 self.kept.fetch_max(stripes, Atomic::Relaxed);
2857 }
2858
2859 #[must_use]
2861 pub fn part_rows(&self, at: usize) -> usize {
2862 self.places.get(at).map_or(0, |place| place.rows as usize)
2863 }
2864
2865 #[must_use]
2867 pub fn table(&self) -> &Table {
2868 &self.table
2869 }
2870
2871 pub fn top_frequencies(&self, column: usize, top: usize) -> Result<Option<Vec<(Value, u64)>>> {
2880 let field = self
2881 .table
2882 .fields
2883 .get(column)
2884 .ok_or_else(|| invalid("frequency column index out of range"))?;
2885 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2886 return Ok(None);
2887 };
2888 if top == 0 || summary.entries.len() < top {
2889 return Ok(None);
2890 }
2891 let boundary = summary.entries[top - 1].count;
2892 if boundary <= summary.omitted_max {
2893 return Ok(None);
2894 }
2895 self.decode_frequencies(column, &field.ty, &summary.entries).map(Some)
2896 }
2897
2898 pub fn exact_frequencies(&self, column: usize) -> Result<Option<Vec<(Value, u64)>>> {
2918 let Some(prefix) = self.frequency_prefix(column)? else {
2919 return Ok(None);
2920 };
2921 Ok((prefix.omitted_max == 0).then_some(prefix.entries))
2922 }
2923
2924 pub fn frequency_prefix(&self, column: usize) -> Result<Option<FrequencyPrefix>> {
2947 let field = self
2948 .table
2949 .fields
2950 .get(column)
2951 .ok_or_else(|| invalid("frequency column index out of range"))?;
2952 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2953 return Ok(None);
2954 };
2955 let entries = self.decode_frequencies(column, &field.ty, &summary.entries)?;
2956 Ok(Some(FrequencyPrefix { entries, omitted_max: summary.omitted_max }))
2957 }
2958
2959 fn decode_frequencies(
2961 &self,
2962 column: usize,
2963 ty: &LogicalType,
2964 entries: &[FrequencyEntry],
2965 ) -> Result<Vec<(Value, u64)>> {
2966 let dictionary = if *ty == LogicalType::Varchar { self.dictionary(column)? } else { None };
2967 let mut out = Vec::with_capacity(entries.len());
2968 for entry in entries {
2969 let value = match entry.value {
2970 FrequencyValue::Null => Value::Null,
2971 FrequencyValue::Integer(value) => match *ty {
2972 LogicalType::TinyInt => Value::TinyInt(
2973 i8::try_from(value)
2974 .map_err(|_| invalid("frequency TINYINT is out of range"))?,
2975 ),
2976 LogicalType::UTinyInt => Value::UTinyInt(
2977 u8::try_from(value)
2978 .map_err(|_| invalid("frequency UTINYINT is out of range"))?,
2979 ),
2980 LogicalType::USmallInt => Value::USmallInt(
2981 u16::try_from(value)
2982 .map_err(|_| invalid("frequency USMALLINT is out of range"))?,
2983 ),
2984 LogicalType::UInteger => Value::UInteger(
2985 u32::try_from(value)
2986 .map_err(|_| invalid("frequency UINTEGER is out of range"))?,
2987 ),
2988 LogicalType::UBigInt => Value::UBigInt(
2989 u64::try_from(value)
2990 .map_err(|_| invalid("frequency UBIGINT is out of range"))?,
2991 ),
2992 LogicalType::SmallInt => Value::SmallInt(
2993 i16::try_from(value)
2994 .map_err(|_| invalid("frequency SMALLINT is out of range"))?,
2995 ),
2996 LogicalType::Integer => Value::Integer(
2997 i32::try_from(value)
2998 .map_err(|_| invalid("frequency INTEGER is out of range"))?,
2999 ),
3000 LogicalType::BigInt => Value::BigInt(
3001 i64::try_from(value)
3002 .map_err(|_| invalid("frequency BIGINT is out of range"))?,
3003 ),
3004 LogicalType::Date => Value::Date(
3005 i32::try_from(value)
3006 .map_err(|_| invalid("frequency DATE is out of range"))?,
3007 ),
3008 LogicalType::Timestamp => Value::Timestamp(
3009 i64::try_from(value)
3010 .map_err(|_| invalid("frequency TIMESTAMP is out of range"))?,
3011 ),
3012 _ => return Err(invalid("integer frequency belongs to another type")),
3013 },
3014 FrequencyValue::Code(code) => dictionary
3015 .as_ref()
3016 .ok_or_else(|| invalid("frequency code has no dictionary"))?
3017 .try_value_at(code as usize)?,
3018 };
3019 out.push((value, entry.count));
3020 }
3021 Ok(out)
3022 }
3023
3024 pub fn frequency_occurrences(&self, column: usize) -> Result<Option<FrequencyOccurrences>> {
3034 self.table
3035 .fields
3036 .get(column)
3037 .ok_or_else(|| invalid("frequency column index out of range"))?;
3038 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
3039 return Ok(None);
3040 };
3041 if summary.ordinals.is_empty() {
3042 return Ok(None);
3043 }
3044 Ok(Some(FrequencyOccurrences {
3045 omitted_max: summary.omitted_max,
3046 ordinals: summary.ordinals.clone(),
3047 }))
3048 }
3049
3050 pub fn distinct_values(&self, column: usize) -> Result<Option<u64>> {
3074 self.table
3075 .distincts
3076 .get(column)
3077 .copied()
3078 .ok_or_else(|| invalid("distinct column index out of range"))
3079 }
3080
3081 pub fn null_count(&self, column: usize) -> Result<u64> {
3092 if column >= self.table.fields.len() {
3093 return Err(invalid("null count column index out of range"));
3094 }
3095 let mut nulls = 0_u64;
3096 for stripe in &self.table.stripes {
3097 let range = stripe
3098 .zone
3099 .column(column)
3100 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
3101 nulls = nulls
3102 .checked_add(range.nulls as u64)
3103 .ok_or_else(|| invalid("null count overflow"))?;
3104 }
3105 Ok(nulls)
3106 }
3107
3108 pub fn text_extremes(&self, column: usize) -> Result<Option<(Value, Value)>> {
3123 if self.null_count(column)? > 0 {
3124 return Ok(None);
3125 }
3126 let Some(dictionary) = self.dictionary(column)? else { return Ok(None) };
3127 let Some(ranks) = dictionary.ranks() else { return Ok(None) };
3128 if ranks == 0 {
3129 return Ok(None);
3130 }
3131 let low = text_at_rank(&dictionary, 0)?;
3132 let high = text_at_rank(&dictionary, ranks - 1)?;
3133 Ok(Some((low, high)))
3134 }
3135
3136 pub fn exact_extremes(&self, column: usize) -> Result<Option<(Bound, Bound)>> {
3159 if column >= self.table.fields.len() {
3160 return Err(invalid("extremes column index out of range"));
3161 }
3162 let mut low: Option<Bound> = None;
3163 let mut high: Option<Bound> = None;
3164 for stripe in &self.table.stripes {
3165 let range = stripe
3166 .zone
3167 .column(column)
3168 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
3169 if !range.exact {
3170 return Ok(None);
3171 }
3172 let (Some(small), Some(large)) = (range.low.as_ref(), range.high.as_ref()) else {
3177 if stripe.rows > range.nulls {
3178 return Ok(None);
3179 }
3180 continue;
3181 };
3182 low = Some(low.map_or_else(|| small.clone(), |held| held.smaller(small.clone())));
3183 high = Some(high.map_or_else(|| large.clone(), |held| held.larger(large.clone())));
3184 }
3185 Ok(low.zip(high))
3186 }
3187
3188 pub fn exact_sum(&self, column: usize) -> Result<Option<(i128, u64)>> {
3201 if column >= self.table.fields.len() {
3202 return Err(invalid("sum column index out of range"));
3203 }
3204 let mut total = 0_i128;
3205 let mut rows = 0_u64;
3206 for stripe in &self.table.stripes {
3207 let range = stripe
3208 .zone
3209 .column(column)
3210 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
3211 let Some(part) = range.sum else { return Ok(None) };
3212 let Some(sum) = total.checked_add(part) else { return Ok(None) };
3213 total = sum;
3214 rows = rows.saturating_add(stripe.rows as u64 - range.nulls as u64);
3215 }
3216 Ok(Some((total, rows)))
3217 }
3218
3219 fn dictionary(&self, column: usize) -> Result<Option<Arc<Vector>>> {
3228 let Some(page) = self.table.dictionaries[column] else { return Ok(None) };
3229 if let Some(dictionary) = self.dictionaries[column].get() {
3230 return Ok(Some(Arc::clone(dictionary)));
3231 }
3232 let _queued = self.loading[column].lock().map_err(|_| invalid("a poisoned dictionary"))?;
3233 if let Some(dictionary) = self.dictionaries[column].get() {
3234 return Ok(Some(Arc::clone(dictionary)));
3235 }
3236 self.opened.fetch_add(1, Atomic::Relaxed);
3237 let dictionary = Arc::new(open_global_dictionary(
3238 Arc::clone(&self.file),
3239 page,
3240 &self.table.fields[column].ty,
3241 TEXT_KEEP_BUDGET,
3242 )?);
3243 let _ = self.dictionaries[column].set(Arc::clone(&dictionary));
3244 Ok(Some(dictionary))
3245 }
3246
3247 pub fn read(&self, part: usize, columns: &[usize]) -> Result<Chunk> {
3256 self.read_impl(part, columns, true)
3257 }
3258
3259 pub fn read_sparse(&self, part: usize, columns: &[usize]) -> Result<Chunk> {
3269 self.read_impl(part, columns, false)
3270 }
3271
3272 pub fn skips_codes(&self, part: usize, column: usize, candidates: &[u32]) -> Result<bool> {
3279 if candidates.is_empty() {
3280 return Ok(true);
3281 }
3282 if candidates.windows(2).any(|pair| pair[0] >= pair[1]) {
3283 return Err(Error::internal("native code candidates are not sorted and unique"));
3284 }
3285 let stripe = self.stripe_of(part)?;
3286 let Some(page) = stripe.memberships.get(column).copied().flatten() else {
3287 return Ok(false);
3288 };
3289 let mut bytes = vec![0; page.length as usize];
3290 read_at(&self.file, page.offset, &mut bytes)?;
3291 if checksum(&bytes) != page.hash {
3292 return Err(invalid("membership page checksum differs"));
3293 }
3294 let codes = decode_membership(&bytes)?;
3295 let mut left = 0;
3296 let mut right = 0;
3297 while left < codes.len() && right < candidates.len() {
3298 match codes[left].cmp(&candidates[right]) {
3299 Ordering::Less => left += 1,
3300 Ordering::Greater => right += 1,
3301 Ordering::Equal => return Ok(false),
3302 }
3303 }
3304 Ok(true)
3305 }
3306
3307 fn stripe_of(&self, part: usize) -> Result<&Stripe> {
3308 let place = self.places.get(part).ok_or_else(|| invalid("part index out of range"))?;
3309 self.table
3310 .stripes
3311 .get(place.stripe as usize)
3312 .ok_or_else(|| invalid("stripe index out of range"))
3313 }
3314
3315 fn held(&self, at: usize, stripe: &Stripe, column: usize, whole: bool) -> Result<CachedColumn> {
3332 let cache = self.cache.get(column).ok_or_else(|| invalid("column index out of range"))?;
3333 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3334 let known = cached.index.get(at).and_then(Clone::clone);
3335 let page = cached.pages.get(at).and_then(Clone::clone);
3336 if let Some(index) = known.clone() {
3337 if !whole || page.is_some() {
3338 return Ok(CachedColumn { stripe: at, index, page });
3339 }
3340 }
3341 if cached.loading.contains(&at) {
3342 drop(cached);
3343 if let Some(index) = known {
3347 return Ok(CachedColumn { stripe: at, index, page: None });
3348 }
3349 let held = self.page_of(stripe, column, at, false, None)?;
3350 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3351 remember(&mut cached, &held, self.kept.load(Atomic::Relaxed));
3352 return Ok(held);
3353 }
3354 cached.loading.push(at);
3355 drop(cached);
3356
3357 let read = self.page_of(stripe, column, at, whole, known);
3358
3359 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3363 if let Some(position) = cached.loading.iter().position(|loading| *loading == at) {
3364 cached.loading.remove(position);
3365 }
3366 let held = read?;
3367 remember(&mut cached, &held, self.kept.load(Atomic::Relaxed));
3368 Ok(held)
3369 }
3370
3371 fn page_of(
3377 &self,
3378 stripe: &Stripe,
3379 column: usize,
3380 at: usize,
3381 whole: bool,
3382 known: Option<Arc<Vec<PartSpan>>>,
3383 ) -> Result<CachedColumn> {
3384 let index = match known {
3385 Some(index) => index,
3386 None => {
3387 self.indexes.fetch_add(1, Atomic::Relaxed);
3388 Arc::new(read_index(&self.file, stripe, column)?)
3389 }
3390 };
3391 let page = if whole {
3392 self.pages.fetch_add(1, Atomic::Relaxed);
3393 let span = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
3394 let mut bytes = vec![0; span.length as usize];
3395 read_at(&self.file, span.offset, &mut bytes)?;
3396 Some(Arc::new(bytes))
3397 } else {
3398 None
3399 };
3400 Ok(CachedColumn { stripe: at, index, page })
3401 }
3402
3403 fn read_impl(&self, at: usize, columns: &[usize], whole: bool) -> Result<Chunk> {
3404 let place = *self.places.get(at).ok_or_else(|| invalid("part index out of range"))?;
3405 let index = place.stripe as usize;
3406 let stripe =
3407 self.table.stripes.get(index).ok_or_else(|| invalid("stripe index out of range"))?;
3408 let rows = place.rows as usize;
3409 let mut picked = Vec::with_capacity(columns.len());
3410 for &column in columns {
3411 let field = self
3412 .table
3413 .fields
3414 .get(column)
3415 .ok_or_else(|| invalid("column index out of range"))?;
3416 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
3417 let held = self.held(index, stripe, column, whole)?;
3418 let span = *held
3419 .index
3420 .get(place.part as usize)
3421 .ok_or_else(|| invalid("part index out of range"))?;
3422 let owned;
3423 let bytes = match &held.page {
3424 Some(held) => part_bytes(held, span)?,
3425 None => {
3426 let offset = page
3427 .offset
3428 .checked_add(span.start as u64)
3429 .ok_or_else(|| invalid("part range overflow"))?;
3430 let mut bytes = vec![0; span.length];
3431 read_at(&self.file, offset, &mut bytes)?;
3432 owned = bytes;
3433 &owned
3434 }
3435 };
3436 if checksum(bytes) != span.hash {
3437 return Err(invalid(&format!(
3438 "column page checksum differs, column {column} part {} at {}+{} of {} bytes, \
3439 wanted {:016x} and got {:016x}",
3440 place.part,
3441 page.offset,
3442 span.start,
3443 span.length,
3444 span.hash,
3445 checksum(bytes),
3446 )));
3447 }
3448 let dictionary = self.dictionary(column)?;
3449 picked.push(decode(&field.ty, rows, bytes, dictionary)?.into_pages());
3455 }
3456 Chunk::with_rows(picked, rows)
3457 }
3458
3459 #[must_use]
3475 pub fn skips(&self, part: usize, probes: &[Probe]) -> bool {
3476 let Some(place) = self.places.get(part).copied() else { return false };
3477 let Some(stripe) = self.table.stripes.get(place.stripe as usize) else { return false };
3478 if stripe.zone.skips(probes) {
3479 return true;
3480 }
3481 probes.iter().any(|probe| self.outside(place, probe) || self.sifted(place, probe))
3482 }
3483
3484 fn outside(&self, place: Place, probe: &Probe) -> bool {
3490 match self.stripe_part_ranges(place.stripe as usize, probe.column) {
3491 Some(ranges) => ranges
3492 .get(place.part as usize)
3493 .is_some_and(|range| range.excludes(probe.op, &probe.value)),
3494 None => false,
3495 }
3496 }
3497
3498 fn stripe_part_ranges(&self, stripe: usize, column: usize) -> Option<&[Range]> {
3504 let slot = self.part_ranges.get(column)?.get(stripe)?;
3505 if let Some(held) = slot.get() {
3506 return Some(held);
3507 }
3508 let page = self.table.stripes.get(stripe)?.part_ranges.get(column).copied().flatten()?;
3509 let mut bytes = vec![0; page.length as usize];
3510 read_at(&self.file, page.offset, &mut bytes).ok()?;
3511 if checksum(&bytes) != page.hash {
3512 return None;
3513 }
3514 let ranges = Arc::new(decode_part_ranges(&bytes).ok()?);
3515 let _ = slot.set(ranges);
3516 slot.get().map(|held| held.as_slice())
3517 }
3518
3519 #[must_use]
3536 pub fn certain(&self, part: usize, probes: &[Probe]) -> bool {
3537 let Some(place) = self.places.get(part).copied() else { return false };
3538 let Some(stripe) = self.table.stripes.get(place.stripe as usize) else { return false };
3539 if stripe.zone.certain(probes) {
3540 return true;
3541 }
3542 probes
3543 .iter()
3544 .all(|probe| stripe.zone.certain(slice::from_ref(probe)) || self.inside(place, probe))
3545 }
3546
3547 fn inside(&self, place: Place, probe: &Probe) -> bool {
3553 match self.stripe_part_ranges(place.stripe as usize, probe.column) {
3554 Some(ranges) => ranges
3555 .get(place.part as usize)
3556 .is_some_and(|range| range.certain(probe.op, &probe.value)),
3557 None => false,
3558 }
3559 }
3560
3561 #[must_use]
3572 pub fn stripe_skips(&self, stripe: usize, probes: &[Probe]) -> bool {
3573 self.table.stripes.get(stripe).is_some_and(|held| held.zone.skips(probes))
3574 }
3575
3576 fn sifted(&self, place: Place, probe: &Probe) -> bool {
3582 if probe.op != Op::Equal {
3583 return false;
3584 }
3585 match self.stripe_sieves(place.stripe as usize, probe.column) {
3586 Some(sieves) => sieves
3587 .get(place.part as usize)
3588 .and_then(Option::as_ref)
3589 .is_some_and(|sieve| sieve.excludes(&probe.value)),
3590 None => false,
3591 }
3592 }
3593
3594 fn stripe_sieves(&self, stripe: usize, column: usize) -> Option<&[Option<Sieve>]> {
3601 let slot = self.sieves.get(column)?.get(stripe)?;
3602 if let Some(held) = slot.get() {
3603 return Some(held);
3604 }
3605 let page = self.table.stripes.get(stripe)?.sieves.get(column).copied().flatten()?;
3606 let mut bytes = vec![0; page.length as usize];
3607 read_at(&self.file, page.offset, &mut bytes).ok()?;
3608 if checksum(&bytes) != page.hash {
3609 return None;
3610 }
3611 let sieves = Arc::new(decode_sieves(&bytes).ok()?);
3612 let _ = slot.set(sieves);
3613 slot.get().map(|held| held.as_slice())
3614 }
3615}
3616
3617fn text_at_rank(dictionary: &Vector, rank: usize) -> Result<Value> {
3619 let code = dictionary.code_at_rank(rank)? as usize;
3620 let text = dictionary
3621 .try_text_at(code)?
3622 .ok_or_else(|| invalid("global dictionary order names a code it does not have"))?;
3623 Ok(Value::Varchar(text.into()))
3624}
3625
3626#[cfg(unix)]
3631fn write_at(file: &File, mut offset: u64, mut bytes: &[u8]) -> Result<()> {
3632 use std::os::unix::fs::FileExt;
3633 while !bytes.is_empty() {
3634 let written = file.write_at(bytes, offset).map_err(io)?;
3635 if written == 0 {
3636 return Err(invalid("a write to the native file wrote nothing"));
3637 }
3638 offset += written as u64;
3639 bytes = &bytes[written..];
3640 }
3641 Ok(())
3642}
3643
3644#[cfg(windows)]
3646fn write_at(file: &File, mut offset: u64, mut bytes: &[u8]) -> Result<()> {
3647 use std::os::windows::fs::FileExt;
3648 while !bytes.is_empty() {
3649 let written = file.seek_write(bytes, offset).map_err(io)?;
3650 if written == 0 {
3651 return Err(invalid("a write to the native file wrote nothing"));
3652 }
3653 offset += written as u64;
3654 bytes = &bytes[written..];
3655 }
3656 Ok(())
3657}
3658
3659#[cfg(not(any(unix, windows)))]
3661fn write_at(file: &File, offset: u64, bytes: &[u8]) -> Result<()> {
3662 use std::io::Write;
3663 let mut file = file.try_clone().map_err(io)?;
3664 file.seek(SeekFrom::Start(offset)).map_err(io)?;
3665 file.write_all(bytes).map_err(io)
3666}
3667
3668#[cfg(unix)]
3678fn read_at(file: &File, mut offset: u64, mut bytes: &mut [u8]) -> Result<()> {
3679 use std::os::unix::fs::FileExt;
3680 while !bytes.is_empty() {
3681 let read = file.read_at(bytes, offset).map_err(io)?;
3682 if read == 0 {
3683 return Err(invalid("column page ends before its declared length"));
3684 }
3685 offset += read as u64;
3686 bytes = &mut bytes[read..];
3687 }
3688 Ok(())
3689}
3690
3691#[cfg(windows)]
3697fn read_at(file: &File, mut offset: u64, mut bytes: &mut [u8]) -> Result<()> {
3698 use std::os::windows::fs::FileExt;
3699 while !bytes.is_empty() {
3700 let read = file.seek_read(bytes, offset).map_err(io)?;
3701 if read == 0 {
3702 return Err(invalid("column page ends before its declared length"));
3703 }
3704 offset += read as u64;
3705 bytes = &mut bytes[read..];
3706 }
3707 Ok(())
3708}
3709
3710#[cfg(not(any(unix, windows)))]
3715fn read_at(file: &File, offset: u64, bytes: &mut [u8]) -> Result<()> {
3716 let mut file = file.try_clone().map_err(io)?;
3717 file.seek(SeekFrom::Start(offset)).map_err(io)?;
3718 file.read_exact(bytes).map_err(io)
3719}
3720
3721fn type_tag(ty: &LogicalType) -> Result<u8> {
3722 match ty {
3723 LogicalType::SmallInt => Ok(1),
3724 LogicalType::Integer => Ok(2),
3725 LogicalType::BigInt => Ok(3),
3726 LogicalType::Varchar => Ok(4),
3727 LogicalType::Date => Ok(5),
3728 LogicalType::Timestamp => Ok(6),
3729 LogicalType::Boolean => Ok(7),
3730 LogicalType::TinyInt => Ok(8),
3731 LogicalType::UTinyInt => Ok(9),
3732 LogicalType::USmallInt => Ok(10),
3733 LogicalType::UInteger => Ok(11),
3734 LogicalType::UBigInt => Ok(12),
3735 LogicalType::Decimal { .. } => Ok(13),
3736 _ => Err(Error::not_implemented(format!("native storage for {ty}"))),
3737 }
3738}
3739
3740fn put_type(out: &mut Vec<u8>, ty: &LogicalType) -> Result<()> {
3746 out.push(type_tag(ty)?);
3747 if let LogicalType::Decimal { width, scale } = ty {
3748 out.push(*width);
3749 out.push(*scale);
3750 }
3751 Ok(())
3752}
3753
3754fn read_type(cur: &mut Cursor<'_>) -> Result<LogicalType> {
3756 let tag = cur.u8()?;
3757 if tag == 13 {
3758 let width = cur.u8()?;
3759 let scale = cur.u8()?;
3760 return LogicalType::decimal(width, scale)
3761 .map_err(|_| invalid("decimal column width and scale are not a decimal"));
3762 }
3763 tag_type(tag)
3764}
3765
3766fn tag_type(tag: u8) -> Result<LogicalType> {
3767 match tag {
3768 1 => Ok(LogicalType::SmallInt),
3769 2 => Ok(LogicalType::Integer),
3770 3 => Ok(LogicalType::BigInt),
3771 4 => Ok(LogicalType::Varchar),
3772 5 => Ok(LogicalType::Date),
3773 6 => Ok(LogicalType::Timestamp),
3774 7 => Ok(LogicalType::Boolean),
3775 8 => Ok(LogicalType::TinyInt),
3776 9 => Ok(LogicalType::UTinyInt),
3777 10 => Ok(LogicalType::USmallInt),
3778 11 => Ok(LogicalType::UInteger),
3779 12 => Ok(LogicalType::UBigInt),
3780 _ => Err(invalid("column type tag is unknown")),
3781 }
3782}
3783
3784fn put_u16(out: &mut Vec<u8>, value: u16) {
3785 out.extend_from_slice(&value.to_le_bytes());
3786}
3787fn put_u32(out: &mut Vec<u8>, value: u32) {
3788 out.extend_from_slice(&value.to_le_bytes());
3789}
3790fn put_u64(out: &mut Vec<u8>, value: u64) {
3791 out.extend_from_slice(&value.to_le_bytes());
3792}
3793fn put_var_u64(out: &mut Vec<u8>, mut value: u64) {
3794 while value >= 0x80 {
3795 out.push((value as u8 & 0x7f) | 0x80);
3796 value >>= 7;
3797 }
3798 out.push(value as u8);
3799}
3800
3801fn frequency_order(left: FrequencyValue, right: FrequencyValue) -> Ordering {
3802 match (left, right) {
3803 (FrequencyValue::Null, FrequencyValue::Null) => Ordering::Equal,
3804 (FrequencyValue::Null, _) => Ordering::Less,
3805 (_, FrequencyValue::Null) => Ordering::Greater,
3806 (FrequencyValue::Integer(left), FrequencyValue::Integer(right)) => left.cmp(&right),
3807 (FrequencyValue::Code(left), FrequencyValue::Code(right)) => left.cmp(&right),
3808 (FrequencyValue::Integer(_), FrequencyValue::Code(_)) => Ordering::Less,
3809 (FrequencyValue::Code(_), FrequencyValue::Integer(_)) => Ordering::Greater,
3810 }
3811}
3812
3813fn code_frequency(dictionary: &GlobalDictionary) -> FrequencySummary {
3814 let mut entries = dictionary
3815 .counts
3816 .iter()
3817 .enumerate()
3818 .filter(|(_, count)| **count != 0)
3819 .map(|(code, &count)| FrequencyEntry { value: FrequencyValue::Code(code as u32), count })
3820 .collect::<Vec<_>>();
3821 if dictionary.nulls != 0 {
3822 entries.push(FrequencyEntry { value: FrequencyValue::Null, count: dictionary.nulls });
3823 }
3824 entries.sort_unstable_by(|left, right| {
3825 right.count.cmp(&left.count).then_with(|| frequency_order(left.value, right.value))
3826 });
3827 let omitted_max = entries.get(FREQUENCY_ENTRIES).map_or(0, |entry| entry.count);
3828 entries.truncate(FREQUENCY_ENTRIES);
3829 FrequencySummary { entries, omitted_max, ordinals: Vec::new() }
3830}
3831
3832fn encode_directory(table: &Table) -> Result<Vec<u8>> {
3833 let mut out = DIRECTORY.to_vec();
3834 let name = table.name.as_bytes();
3835 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("table name too long"))?);
3836 out.extend_from_slice(name);
3837 put_u16(&mut out, u16::try_from(table.fields.len()).map_err(|_| invalid("too many columns"))?);
3838 for field in &table.fields {
3839 let name = field.name.as_bytes();
3840 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("column name too long"))?);
3841 out.extend_from_slice(name);
3842 put_type(&mut out, &field.ty)?;
3843 out.push(u8::from(field.not_null));
3844 }
3845 for dictionary in &table.dictionaries {
3846 match dictionary {
3847 None => out.push(0),
3848 Some(page) => {
3849 out.push(1);
3850 put_u64(&mut out, page.offset);
3851 put_u32(&mut out, page.length);
3852 put_u64(&mut out, page.hash);
3853 }
3854 }
3855 }
3856 for distinct in &table.distincts {
3857 match distinct {
3858 None => out.push(0),
3859 Some(count) => {
3860 out.push(1);
3861 put_u64(&mut out, *count);
3862 }
3863 }
3864 }
3865 put_u64(&mut out, u64::try_from(table.rows).map_err(|_| invalid("row count overflow"))?);
3866 put_u32(&mut out, u32::try_from(table.stripes.len()).map_err(|_| invalid("too many stripes"))?);
3867 for stripe in &table.stripes {
3868 put_u32(
3869 &mut out,
3870 u32::try_from(stripe.parts.len()).map_err(|_| invalid("too many parts in a stripe"))?,
3871 );
3872 for &rows in &stripe.parts {
3873 put_u32(&mut out, rows);
3874 }
3875 put_u64(&mut out, stripe.index.offset);
3876 put_u32(&mut out, stripe.index.length);
3877 for page in &stripe.pages {
3878 put_u64(&mut out, page.offset);
3879 put_u32(&mut out, page.length);
3880 }
3881 for ((field, dictionary), membership) in
3886 table.fields.iter().zip(&table.dictionaries).zip(&stripe.memberships)
3887 {
3888 if field.ty != LogicalType::Varchar || dictionary.is_none() {
3889 continue;
3890 }
3891 let page =
3892 membership.ok_or_else(|| invalid("string page has no code membership index"))?;
3893 put_u64(&mut out, page.offset);
3894 put_u32(&mut out, page.length);
3895 put_u64(&mut out, page.hash);
3896 }
3897 for sieve in &stripe.sieves {
3898 match sieve {
3899 None => out.push(0),
3900 Some(page) => {
3901 out.push(1);
3902 put_u64(&mut out, page.offset);
3903 put_u32(&mut out, page.length);
3904 put_u64(&mut out, page.hash);
3905 }
3906 }
3907 }
3908 for held in &stripe.part_ranges {
3909 match held {
3910 None => out.push(0),
3911 Some(page) => {
3912 out.push(1);
3913 put_u64(&mut out, page.offset);
3914 put_u32(&mut out, page.length);
3915 put_u64(&mut out, page.hash);
3916 }
3917 }
3918 }
3919 for range in stripe.zone.columns() {
3920 put_bound(&mut out, range.low.as_ref())?;
3921 put_bound(&mut out, range.high.as_ref())?;
3922 put_u32(
3923 &mut out,
3924 u32::try_from(range.nulls).map_err(|_| invalid("null count overflow"))?,
3925 );
3926 out.push(u8::from(range.exact));
3927 match range.sum {
3928 None => out.push(0),
3929 Some(total) => {
3930 out.push(1);
3931 out.extend_from_slice(&total.to_le_bytes());
3932 }
3933 }
3934 }
3935 }
3936 out.extend_from_slice(FREQUENCIES);
3937 put_u16(
3938 &mut out,
3939 u16::try_from(table.frequencies.len())
3940 .map_err(|_| invalid("too many frequency columns"))?,
3941 );
3942 for summary in &table.frequencies {
3943 let Some(summary) = summary else {
3944 out.push(0);
3945 continue;
3946 };
3947 out.push(1);
3948 put_u64(&mut out, summary.omitted_max);
3949 put_u32(
3950 &mut out,
3951 u32::try_from(summary.entries.len())
3952 .map_err(|_| invalid("too many frequency entries"))?,
3953 );
3954 for entry in &summary.entries {
3955 match entry.value {
3956 FrequencyValue::Null => out.push(0),
3957 FrequencyValue::Integer(value) => {
3958 out.push(1);
3959 out.extend_from_slice(&value.to_le_bytes());
3960 }
3961 FrequencyValue::Code(value) => {
3962 out.push(2);
3963 put_u32(&mut out, value);
3964 }
3965 }
3966 put_u64(&mut out, entry.count);
3967 }
3968 put_u32(
3969 &mut out,
3970 u32::try_from(summary.ordinals.len())
3971 .map_err(|_| invalid("too many frequency ordinals"))?,
3972 );
3973 let mut previous = 0_u64;
3974 for (at, &ordinal) in summary.ordinals.iter().enumerate() {
3975 let delta = if at == 0 {
3976 ordinal
3977 } else {
3978 ordinal
3979 .checked_sub(previous)
3980 .ok_or_else(|| invalid("frequency ordinals are not ordered"))?
3981 };
3982 if at != 0 && delta == 0 {
3983 return Err(invalid("frequency ordinals are not unique"));
3984 }
3985 put_var_u64(&mut out, delta);
3986 previous = ordinal;
3987 }
3988 }
3989 if let Some(clustering) = &table.clustering {
3992 out.extend_from_slice(CLUSTERING);
3993 out.push(clustering.width().tag());
3994 put_u16(
3995 &mut out,
3996 u16::try_from(clustering.columns().len())
3997 .map_err(|_| invalid("too many clustering columns"))?,
3998 );
3999 for &column in clustering.columns() {
4000 put_u16(
4001 &mut out,
4002 u16::try_from(column).map_err(|_| invalid("clustering column index overflow"))?,
4003 );
4004 }
4005 }
4006 Ok(out)
4007}
4008
4009fn encode_catalog(entries: &[Entry]) -> Result<Vec<u8>> {
4015 let mut out = CATALOG.to_vec();
4016 put_u32(&mut out, u32::try_from(entries.len()).map_err(|_| invalid("too many tables"))?);
4017 for entry in entries {
4018 let name = entry.name.as_bytes();
4019 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("table name too long"))?);
4020 out.extend_from_slice(name);
4021 put_u64(&mut out, u64::try_from(entry.rows).map_err(|_| invalid("row count overflow"))?);
4022 put_u16(
4023 &mut out,
4024 u16::try_from(entry.fields.len()).map_err(|_| invalid("too many columns"))?,
4025 );
4026 for field in &entry.fields {
4027 let name = field.name.as_bytes();
4028 put_u16(
4029 &mut out,
4030 u16::try_from(name.len()).map_err(|_| invalid("column name too long"))?,
4031 );
4032 out.extend_from_slice(name);
4033 put_type(&mut out, &field.ty)?;
4034 out.push(u8::from(field.not_null));
4035 }
4036 put_u64(&mut out, entry.directory.offset);
4037 put_u32(&mut out, entry.directory.length);
4038 put_u64(&mut out, entry.directory.hash);
4039 }
4040 Ok(out)
4041}
4042
4043fn decode_catalog(bytes: &[u8], size: u64) -> Result<Vec<Entry>> {
4046 let mut cur = Cursor { bytes, at: 0 };
4047 if cur.take(8)? != CATALOG {
4048 return Err(invalid("catalog magic differs"));
4049 }
4050 let count = cur.u32()? as usize;
4051 let mut entries: Vec<Entry> = Vec::with_capacity(count.min(1024));
4052 for _ in 0..count {
4053 let name = cur.text()?;
4054 let rows = usize::try_from(cur.u64()?).map_err(|_| invalid("row count does not fit"))?;
4055 let width = cur.u16()? as usize;
4056 let mut fields = Vec::with_capacity(width);
4057 for _ in 0..width {
4058 let name = cur.text()?;
4059 let ty = read_type(&mut cur)?;
4060 let not_null = match cur.u8()? {
4061 0 => false,
4062 1 => true,
4063 _ => return Err(invalid("nullability flag differs")),
4064 };
4065 fields.push(Field { name, ty, not_null });
4066 }
4067 let directory = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4068 let end = directory
4069 .offset
4070 .checked_add(u64::from(directory.length))
4071 .ok_or_else(|| invalid("table directory offset overflow"))?;
4072 if directory.offset < HEADER
4073 || end > size
4074 || directory.length as usize > MAX_DIRECTORY
4075 || directory.length == 0
4076 {
4077 return Err(invalid("table directory range is outside the file"));
4078 }
4079 if entries.iter().any(|held| held.name == name) {
4080 return Err(invalid("two tables in the catalog have the same name"));
4081 }
4082 entries.push(Entry { name, fields, rows, directory });
4083 }
4084 Ok(entries)
4085}
4086
4087struct Cursor<'a> {
4088 bytes: &'a [u8],
4089 at: usize,
4090}
4091impl<'a> Cursor<'a> {
4092 fn take(&mut self, len: usize) -> Result<&'a [u8]> {
4093 let end = self.at.checked_add(len).ok_or_else(|| invalid("directory offset overflow"))?;
4094 let bytes =
4095 self.bytes.get(self.at..end).ok_or_else(|| invalid("directory is truncated"))?;
4096 self.at = end;
4097 Ok(bytes)
4098 }
4099 fn u8(&mut self) -> Result<u8> {
4100 Ok(self.take(1)?[0])
4101 }
4102 fn u16(&mut self) -> Result<u16> {
4103 Ok(u16::from_le_bytes(self.take(2)?.try_into().expect("two bytes")))
4104 }
4105 fn u32(&mut self) -> Result<u32> {
4106 Ok(u32::from_le_bytes(self.take(4)?.try_into().expect("four bytes")))
4107 }
4108 fn u64(&mut self) -> Result<u64> {
4109 Ok(u64::from_le_bytes(self.take(8)?.try_into().expect("eight bytes")))
4110 }
4111 fn var_u64(&mut self) -> Result<u64> {
4112 let mut value = 0_u64;
4113 for shift in (0..=63).step_by(7) {
4114 let byte = self.u8()?;
4115 let part = u64::from(byte & 0x7f);
4116 if shift == 63 && part > 1 {
4117 return Err(invalid("frequency ordinal varint overflows"));
4118 }
4119 value |= part << shift;
4120 if byte & 0x80 == 0 {
4121 return Ok(value);
4122 }
4123 }
4124 Err(invalid("frequency ordinal varint is too long"))
4125 }
4126 fn bound(&mut self) -> Result<Option<Bound>> {
4127 Ok(match self.u8()? {
4128 0 => None,
4129 1 => Some(Bound::Int(i128::from_le_bytes(
4130 self.take(16)?.try_into().expect("sixteen bytes"),
4131 ))),
4132 2 => Some(Bound::Real(f64::from_le_bytes(
4133 self.take(8)?.try_into().expect("eight bytes"),
4134 ))),
4135 3 => {
4136 let length = self.u32()? as usize;
4137 Some(Bound::Bytes(self.take(length)?.to_vec()))
4138 }
4139 4 => {
4140 let unscaled =
4141 i128::from_le_bytes(self.take(16)?.try_into().expect("sixteen bytes"));
4142 Some(Bound::Scaled { unscaled, scale: self.u8()? })
4143 }
4144 _ => return Err(invalid("bound tag differs")),
4145 })
4146 }
4147 fn text(&mut self) -> Result<String> {
4148 let len = self.u16()? as usize;
4149 String::from_utf8(self.take(len)?.to_vec()).map_err(|_| invalid("name is not UTF-8"))
4150 }
4151}
4152
4153fn decode_directory(bytes: &[u8], size: u64) -> Result<Table> {
4154 let mut cur = Cursor { bytes, at: 0 };
4155 if cur.take(8)? != DIRECTORY {
4156 return Err(invalid("directory magic differs"));
4157 }
4158 let name = cur.text()?;
4159 let width = cur.u16()? as usize;
4160 let mut fields = Vec::with_capacity(width);
4161 for _ in 0..width {
4162 let name = cur.text()?;
4163 let ty = read_type(&mut cur)?;
4164 let not_null = match cur.u8()? {
4165 0 => false,
4166 1 => true,
4167 _ => return Err(invalid("nullability flag differs")),
4168 };
4169 fields.push(Field { name, ty, not_null });
4170 }
4171 let mut dictionaries = Vec::with_capacity(width);
4172 for _ in 0..width {
4173 dictionaries.push(match cur.u8()? {
4174 0 => None,
4175 1 => {
4176 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4177 let end = page
4178 .offset
4179 .checked_add(u64::from(page.length))
4180 .ok_or_else(|| invalid("dictionary page offset overflow"))?;
4181 if page.offset < HEADER || end > size {
4186 return Err(invalid("dictionary page range is outside the file"));
4187 }
4188 Some(page)
4189 }
4190 _ => return Err(invalid("dictionary page tag differs")),
4191 });
4192 }
4193 let mut distincts = Vec::with_capacity(width);
4194 for _ in 0..width {
4195 distincts.push(match cur.u8()? {
4196 0 => None,
4197 1 => Some(cur.u64()?),
4198 _ => return Err(invalid("distinct count tag differs")),
4199 });
4200 }
4201 let rows = usize::try_from(cur.u64()?).map_err(|_| invalid("row count does not fit"))?;
4202 let count = cur.u32()? as usize;
4203 let mut stripes = Vec::with_capacity(count);
4204 let mut total = 0_usize;
4205 for _ in 0..count {
4206 let count = cur.u32()? as usize;
4207 if count == 0 || count > STRIPE_PARTS {
4208 return Err(invalid("stripe part count is outside its bound"));
4209 }
4210 let mut parts = Vec::with_capacity(count);
4211 let mut stripe_rows = 0_usize;
4212 for _ in 0..count {
4213 let rows = cur.u32()?;
4214 if rows == 0 {
4215 return Err(invalid("empty part"));
4216 }
4217 parts.push(rows);
4218 stripe_rows = stripe_rows
4219 .checked_add(rows as usize)
4220 .ok_or_else(|| invalid("stripe row count overflow"))?;
4221 }
4222 total =
4223 total.checked_add(stripe_rows).ok_or_else(|| invalid("stripe row count overflow"))?;
4224 let index = Span { offset: cur.u64()?, length: cur.u32()? };
4225 let section = index_section(count)?;
4226 let wanted = section
4227 .checked_mul(width)
4228 .and_then(|bytes| u32::try_from(bytes).ok())
4229 .ok_or_else(|| invalid("index page length overflow"))?;
4230 let end = index
4231 .offset
4232 .checked_add(u64::from(index.length))
4233 .ok_or_else(|| invalid("index page offset overflow"))?;
4234 if index.offset < HEADER || end > size || index.length != wanted {
4235 return Err(invalid("index page range is outside the file"));
4236 }
4237 let mut pages = Vec::with_capacity(width);
4238 for _ in 0..width {
4239 let offset = cur.u64()?;
4240 let length = cur.u32()?;
4241 let end = offset
4242 .checked_add(u64::from(length))
4243 .ok_or_else(|| invalid("page offset overflow"))?;
4244 if offset < HEADER || end > size || length as usize > MAX_PAGE {
4245 return Err(invalid("page range is outside the file"));
4246 }
4247 pages.push(Span { offset, length });
4248 }
4249 let mut memberships = vec![None; width];
4250 for (column, field) in fields.iter().enumerate() {
4251 if field.ty != LogicalType::Varchar || dictionaries[column].is_none() {
4252 continue;
4253 }
4254 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4255 let end = page
4256 .offset
4257 .checked_add(u64::from(page.length))
4258 .ok_or_else(|| invalid("membership page offset overflow"))?;
4259 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
4260 return Err(invalid("membership page range is outside the file"));
4261 }
4262 memberships[column] = Some(page);
4263 }
4264 let mut sieves = vec![None; width];
4265 for sieve in sieves.iter_mut().take(width) {
4266 match cur.u8()? {
4267 0 => continue,
4268 1 => {}
4269 _ => return Err(invalid("a sieve page has an unknown tag")),
4270 }
4271 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4272 let end = page
4273 .offset
4274 .checked_add(u64::from(page.length))
4275 .ok_or_else(|| invalid("sieve page offset overflow"))?;
4276 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
4277 return Err(invalid("sieve page range is outside the file"));
4278 }
4279 *sieve = Some(page);
4280 }
4281 let mut part_ranges = vec![None; width];
4282 for held in part_ranges.iter_mut().take(width) {
4283 match cur.u8()? {
4284 0 => continue,
4285 1 => {}
4286 _ => return Err(invalid("a part range page has an unknown tag")),
4287 }
4288 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4289 let end = page
4290 .offset
4291 .checked_add(u64::from(page.length))
4292 .ok_or_else(|| invalid("part range page offset overflow"))?;
4293 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
4294 return Err(invalid("part range page range is outside the file"));
4295 }
4296 *held = Some(page);
4297 }
4298 let mut ranges = Vec::with_capacity(width);
4299 for column in 0..width {
4300 let low = cur.bound()?;
4301 let high = cur.bound()?;
4302 let nulls = cur.u32()? as usize;
4303 if nulls > stripe_rows {
4304 return Err(invalid("null count exceeds stripe rows"));
4305 }
4306 let exact = cur.u8()? != 0;
4307 let sum = match cur.u8()? {
4308 0 => None,
4309 1 => Some(i128::from_le_bytes(
4310 cur.take(16)?.try_into().map_err(|_| invalid("a stripe sum is truncated"))?,
4311 )),
4312 _ => return Err(invalid("a stripe sum has an unknown tag")),
4313 };
4314 let ty = &fields.get(column).ok_or_else(|| invalid("a stripe range has no column"))?.ty;
4320 let low = low.map(|bound| scaled_as(bound, ty));
4321 let high = high.map(|bound| scaled_as(bound, ty));
4322 ranges.push(Range { low, high, nulls, exact, sum });
4323 }
4324 stripes.push(Stripe {
4325 rows: stripe_rows,
4326 parts,
4327 index,
4328 pages,
4329 memberships,
4330 sieves,
4331 part_ranges,
4332 zone: Zone::from_ranges(ranges),
4333 });
4334 }
4335 if total != rows {
4336 return Err(invalid("table row count differs from stripes"));
4337 }
4338 let frequencies = if cur.at == bytes.len() {
4339 vec![None; width]
4340 } else {
4341 if cur.take(8)? != FREQUENCIES {
4342 return Err(invalid("directory extension magic differs"));
4343 }
4344 if cur.u16()? as usize != width {
4345 return Err(invalid("frequency column count differs"));
4346 }
4347 let mut frequencies = Vec::with_capacity(width);
4348 for field in &fields {
4349 let summary = match cur.u8()? {
4350 0 => None,
4351 1 => {
4352 let omitted_max = cur.u64()?;
4353 let count = cur.u32()? as usize;
4354 if count > FREQUENCY_ENTRIES {
4355 return Err(invalid("frequency entry count exceeds its bound"));
4356 }
4357 let mut entries = Vec::with_capacity(count);
4358 for _ in 0..count {
4360 let value = match cur.u8()? {
4361 0 => FrequencyValue::Null,
4362 1 => FrequencyValue::Integer(i128::from_le_bytes(
4363 cur.take(16)?.try_into().expect("sixteen bytes"),
4364 )),
4365 2 => FrequencyValue::Code(cur.u32()?),
4366 _ => return Err(invalid("frequency value tag differs")),
4367 };
4368 let valid = matches!(
4369 (&field.ty, value),
4370 (_, FrequencyValue::Null)
4371 | (LogicalType::Varchar, FrequencyValue::Code(_))
4372 | (
4373 LogicalType::TinyInt
4374 | LogicalType::SmallInt
4375 | LogicalType::Integer
4376 | LogicalType::BigInt
4377 | LogicalType::UTinyInt
4378 | LogicalType::USmallInt
4379 | LogicalType::UInteger
4380 | LogicalType::UBigInt
4381 | LogicalType::Date
4382 | LogicalType::Timestamp,
4383 FrequencyValue::Integer(_),
4384 )
4385 );
4386 if !valid {
4387 return Err(invalid("frequency value does not match its column"));
4388 }
4389 let count = cur.u64()?;
4390 if count == 0 || count > rows as u64 {
4391 return Err(invalid("frequency count is outside the table"));
4392 }
4393 entries.push(FrequencyEntry { value, count });
4394 }
4395 if entries.windows(2).any(|pair| pair[0].count < pair[1].count) {
4396 return Err(invalid("frequency entries are not descending"));
4397 }
4398 let ordinals = {
4399 let ordinal_count = cur.u32()? as usize;
4400 if ordinal_count > FREQUENCY_ORDINALS || ordinal_count > rows {
4401 return Err(invalid("frequency ordinal count exceeds its bound"));
4402 }
4403 let mut ordinals = Vec::with_capacity(ordinal_count);
4404 let mut previous = 0_u64;
4405 for at in 0..ordinal_count {
4406 let delta = cur.var_u64()?;
4407 if at != 0 && delta == 0 {
4408 return Err(invalid("frequency ordinals are not increasing"));
4409 }
4410 let ordinal = if at == 0 {
4411 delta
4412 } else {
4413 previous
4414 .checked_add(delta)
4415 .ok_or_else(|| invalid("frequency ordinal overflows"))?
4416 };
4417 if ordinal >= rows as u64 {
4418 return Err(invalid("frequency ordinal is outside the table"));
4419 }
4420 ordinals.push(ordinal);
4421 previous = ordinal;
4422 }
4423 ordinals
4424 };
4425 Some(FrequencySummary { entries, omitted_max, ordinals })
4426 }
4427 _ => return Err(invalid("frequency summary tag differs")),
4428 };
4429 frequencies.push(summary);
4430 }
4431 frequencies
4432 };
4433 let clustering = if cur.at == bytes.len() {
4434 None
4435 } else {
4436 if cur.take(8)? != CLUSTERING {
4437 return Err(invalid("directory extension magic differs"));
4438 }
4439 let bucket =
4440 Width::from_tag(cur.u8()?).ok_or_else(|| invalid("clustering width tag differs"))?;
4441 let count = cur.u16()? as usize;
4442 let mut columns = Vec::with_capacity(count.min(fields.len()));
4443 for _ in 0..count {
4444 columns.push(u32::from(cur.u16()?));
4445 }
4446 Some(Clustering::new(columns, bucket, &fields).map_err(|_| {
4449 invalid("stored clustering declaration does not match the table it is on")
4450 })?)
4451 };
4452 if cur.at != bytes.len() {
4453 return Err(invalid("directory has trailing bytes"));
4454 }
4455 Ok(Table { name, fields, stripes, rows, dictionaries, distincts, frequencies, clustering })
4456}
4457
4458fn put_bound(out: &mut Vec<u8>, bound: Option<&Bound>) -> Result<()> {
4459 match bound {
4460 None => out.push(0),
4461 Some(Bound::Int(value)) => {
4462 out.push(1);
4463 out.extend_from_slice(&value.to_le_bytes());
4464 }
4465 Some(Bound::Real(value)) => {
4466 out.push(2);
4467 out.extend_from_slice(&value.to_le_bytes());
4468 }
4469 Some(Bound::Bytes(value)) => {
4470 out.push(3);
4471 put_u32(out, u32::try_from(value.len()).map_err(|_| invalid("bound length overflow"))?);
4472 out.extend_from_slice(value);
4473 }
4474 Some(Bound::Scaled { unscaled, scale }) => {
4475 out.push(4);
4476 out.extend_from_slice(&unscaled.to_le_bytes());
4477 out.push(*scale);
4478 }
4479 }
4480 Ok(())
4481}
4482
4483#[derive(Debug)]
4500struct Codes;
4501
4502impl chooser::Chooser for Codes {
4503 fn name(&self) -> &'static str {
4504 "codes"
4505 }
4506
4507 fn narrow_strings(
4508 &self,
4509 _values: &[&[u8]],
4510 offered: &[string::Kind],
4511 _depth: u8,
4512 ) -> Vec<string::Kind> {
4513 offered.to_vec()
4516 }
4517
4518 fn narrow_integers(
4519 &self,
4520 _values: &[i64],
4521 offered: &[integer::Kind],
4522 depth: u8,
4523 ) -> Vec<integer::Kind> {
4524 let keep: &[integer::Kind] = if depth == 0 {
4525 &[integer::Kind::Constant, integer::Kind::Packed, integer::Kind::Rle]
4526 } else {
4527 &[integer::Kind::Constant, integer::Kind::Packed]
4528 };
4529 let narrowed: Vec<integer::Kind> =
4530 offered.iter().copied().filter(|kind| keep.contains(kind)).collect();
4531 if narrowed.is_empty() { offered.to_vec() } else { narrowed }
4534 }
4535}
4536
4537#[derive(Debug)]
4549struct Fixed;
4550
4551impl chooser::Chooser for Fixed {
4552 fn name(&self) -> &'static str {
4553 "fixed"
4554 }
4555
4556 fn narrow_strings(
4557 &self,
4558 _values: &[&[u8]],
4559 offered: &[string::Kind],
4560 _depth: u8,
4561 ) -> Vec<string::Kind> {
4562 offered.to_vec()
4563 }
4564
4565 fn narrow_integers(
4566 &self,
4567 _values: &[i64],
4568 offered: &[integer::Kind],
4569 depth: u8,
4570 ) -> Vec<integer::Kind> {
4571 let keep: &[integer::Kind] = if depth == 0 {
4572 &[
4573 integer::Kind::Constant,
4574 integer::Kind::Packed,
4575 integer::Kind::Delta,
4576 integer::Kind::Rle,
4577 integer::Kind::Sparse,
4578 integer::Kind::Strided,
4579 ]
4580 } else {
4581 &[integer::Kind::Constant, integer::Kind::Packed, integer::Kind::Delta]
4582 };
4583 let narrowed: Vec<integer::Kind> =
4584 offered.iter().copied().filter(|kind| keep.contains(kind)).collect();
4585 if narrowed.is_empty() { offered.to_vec() } else { narrowed }
4586 }
4587}
4588
4589fn widened(data: &Data) -> Option<Vec<i64>> {
4596 match data {
4597 Data::Int8(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4598 Data::UInt8(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4599 Data::Int16(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4600 Data::UInt16(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4601 Data::Int32(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4602 Data::UInt32(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4603 Data::Int64(values) => Some(values.to_vec()),
4604 _ => None,
4605 }
4606}
4607
4608trait Narrow: Copy {
4615 const BIASED: (u32, u64);
4620
4621 fn narrow(value: i64) -> Self;
4623}
4624
4625#[allow(clippy::cast_sign_loss, reason = "a residue is a bit pattern and not a number")]
4642fn residue<T: Narrow>(value: i64) -> u64 {
4643 let (bits, bias) = T::BIASED;
4644 (value as u64).wrapping_add(bias) >> bits
4645}
4646
4647macro_rules! narrows {
4652 ($($ty:ty => $bias:expr),* $(,)?) => {$(
4653 impl Narrow for $ty {
4654 const BIASED: (u32, u64) = (<$ty>::BITS, $bias);
4655
4656 #[allow(
4657 clippy::cast_possible_truncation,
4658 clippy::cast_sign_loss,
4659 reason = "the caller has checked the bits this truncates away"
4660 )]
4661 fn narrow(value: i64) -> Self {
4662 value as Self
4663 }
4664 }
4665 )*};
4666}
4667
4668narrows! {
4669 i8 => 1 << 7,
4670 u8 => 0,
4671 i16 => 1 << 15,
4672 u16 => 0,
4673 i32 => 1 << 31,
4674 u32 => 0,
4675}
4676
4677fn fit<T: Narrow>(values: &[i64]) -> Result<Vec<T>> {
4690 let mut spilled = 0u64;
4691 for value in values {
4692 spilled |= residue::<T>(*value);
4693 }
4694 if spilled != 0 {
4695 return Err(invalid("page value is not of its type"));
4696 }
4697 Ok(values.iter().map(|value| T::narrow(*value)).collect())
4698}
4699
4700fn narrowed(ty: &LogicalType, values: Vec<i64>) -> Result<Data> {
4705 Ok(match ty {
4706 LogicalType::TinyInt => Data::Int8(fit::<i8>(&values)?.into()),
4707 LogicalType::UTinyInt => Data::UInt8(fit::<u8>(&values)?.into()),
4708 LogicalType::SmallInt => Data::Int16(fit::<i16>(&values)?.into()),
4709 LogicalType::USmallInt => Data::UInt16(fit::<u16>(&values)?.into()),
4710 LogicalType::Integer | LogicalType::Date => Data::Int32(fit::<i32>(&values)?.into()),
4711 LogicalType::UInteger => Data::UInt32(fit::<u32>(&values)?.into()),
4712 LogicalType::BigInt | LogicalType::Timestamp => Data::Int64(values.into()),
4713 LogicalType::Decimal { .. } => match ty.physical() {
4716 PhysicalType::Int16 => Data::Int16(fit::<i16>(&values)?.into()),
4717 PhysicalType::Int32 => Data::Int32(fit::<i32>(&values)?.into()),
4718 PhysicalType::Int64 => Data::Int64(values.into()),
4719 _ => return Err(invalid("cascade codec belongs to a decimal that is not an integer")),
4720 },
4721 _ => return Err(invalid("cascade codec belongs to a page that is not integers")),
4722 })
4723}
4724
4725fn plain_width(ty: &LogicalType) -> Option<usize> {
4728 Some(match ty {
4729 LogicalType::TinyInt | LogicalType::UTinyInt => 1,
4730 LogicalType::SmallInt | LogicalType::USmallInt => 2,
4731 LogicalType::Integer | LogicalType::UInteger | LogicalType::Date => 4,
4732 LogicalType::BigInt | LogicalType::Timestamp => 8,
4733 LogicalType::Decimal { .. } => match ty.physical() {
4734 PhysicalType::Int16 => 2,
4735 PhysicalType::Int32 => 4,
4736 PhysicalType::Int64 => 8,
4737 _ => return None,
4740 },
4741 _ => return None,
4742 })
4743}
4744
4745fn cascaded(
4751 flat: &Vector,
4752 ty: &LogicalType,
4753 packed: Option<&Packed<'_>>,
4754) -> Result<Option<Vec<u8>>> {
4755 let (Some(width), Some(data)) = (plain_width(ty), flat.data()) else { return Ok(None) };
4756 let Some(values) = widened(data) else { return Ok(None) };
4757 let plain = values.len().saturating_mul(width);
4758 let best = match packed {
4759 Some(packed) => plain.min(21 + size_of_val(packed.words())),
4761 None => plain,
4762 };
4763 let out = integer::encode_with(&values, &Fixed)?;
4764 Ok((out.len() < best).then_some(out))
4765}
4766
4767fn text_compressed(flat: &Vector) -> Result<Option<Vec<u8>>> {
4805 let mut values: Vec<&[u8]> = Vec::with_capacity(flat.len());
4806 let mut payload = 0_usize;
4807 for row in 0..flat.len() {
4808 let text = flat.text_at(row).unwrap_or("").as_bytes();
4809 payload = payload.saturating_add(text.len());
4810 values.push(text);
4811 }
4812 let plain = (flat.len() + 1).saturating_mul(4).saturating_add(payload);
4814 let Some(out) = string::encode_only(string::Kind::Fsst, &values)? else {
4815 return Ok(None);
4816 };
4817 Ok((out.len() < plain).then_some(out))
4818}
4819
4820fn encoded_codes(codes: &[u32]) -> Result<Option<Vec<u8>>> {
4821 let wide: Vec<i64> = codes.iter().map(|code| i64::from(*code)).collect();
4822 let coded = integer::encode_with(&wide, &Codes)?;
4823 let plain = codes.len().saturating_mul(size_of::<u32>());
4824 Ok((coded.len() < plain).then_some(coded))
4825}
4826
4827fn encode(
4828 vector: &Vector,
4829 global: Option<&mut GlobalDictionary>,
4830) -> Result<(Vec<u8>, Option<Vec<u32>>)> {
4831 let ty = vector.logical_type();
4832 let flat = vector.flatten()?;
4834 let mut out = Vec::new();
4835 let mut global_codes = None;
4836 if let Some(global) = global {
4837 let mut codes = Vec::with_capacity(flat.len());
4838 for row in 0..flat.len() {
4839 let text = flat.text_at(row).unwrap_or("");
4840 let code = global.code(text)?;
4841 global.observe(code, flat.is_null_at(row))?;
4842 codes.push(code);
4843 }
4844 global_codes = Some(codes);
4845 }
4846 let membership = global_codes.as_deref().map(unique_codes);
4847 let dictionary = if global_codes.is_none() && ty == &LogicalType::Varchar {
4848 string_dictionary(&flat)?
4849 } else {
4850 None
4851 };
4852 let compressed_text =
4853 if global_codes.is_none() && dictionary.is_none() && ty == &LogicalType::Varchar {
4854 text_compressed(&flat)?
4855 } else {
4856 None
4857 };
4858 let packed_vector = if dictionary.is_none() && global_codes.is_none() {
4859 Some(flat.bit_packed()?)
4860 } else {
4861 None
4862 };
4863 let packed = packed_vector.as_ref().and_then(Vector::packed_parts);
4864 let coded = match global_codes.as_deref() {
4865 Some(codes) => encoded_codes(codes)?,
4866 None => None,
4867 };
4868 let cascade = if dictionary.is_none() && global_codes.is_none() {
4872 cascaded(&flat, ty, packed.as_ref())?
4873 } else {
4874 None
4875 };
4876 out.push(if coded.is_some() {
4877 4
4878 } else if cascade.is_some() {
4879 5
4880 } else if global_codes.is_some() {
4881 3
4882 } else if dictionary.is_some() {
4883 1
4884 } else if compressed_text.is_some() {
4885 6
4886 } else if packed.is_some() {
4887 2
4888 } else {
4889 0
4890 });
4891 let nulls = flat.validity();
4892 let flag = match nulls {
4893 Validity::AllValid => 0,
4894 Validity::AllInvalid => 1,
4895 Validity::Mask(_) => 2,
4896 };
4897 out.push(flag);
4898 if flag == 2 {
4899 for group in (0..vector.len()).step_by(8) {
4900 let mut bits = 0_u8;
4901 for bit in 0..8 {
4902 if group + bit < vector.len() && !flat.is_null_at(group + bit) {
4903 bits |= 1 << bit;
4904 }
4905 }
4906 out.push(bits);
4907 }
4908 }
4909 if let Some(coded) = coded {
4910 out.extend_from_slice(&coded);
4911 return Ok((out, membership));
4912 }
4913 if let Some(cascade) = cascade {
4914 out.extend_from_slice(&cascade);
4915 return Ok((out, membership));
4916 }
4917 if let Some(codes) = global_codes {
4918 for code in codes {
4919 put_u32(&mut out, code);
4920 }
4921 return Ok((out, membership));
4922 }
4923 if let Some(dictionary) = dictionary {
4924 out.extend_from_slice(&dictionary);
4925 return Ok((out, membership));
4926 }
4927 if let Some(compressed_text) = compressed_text {
4928 out.extend_from_slice(&compressed_text);
4929 return Ok((out, membership));
4930 }
4931 if let Some(packed) = packed {
4932 if packed.offset() != 0 {
4933 return Err(invalid("writer received a sliced packed vector"));
4934 }
4935 out.push(u8::try_from(packed.width()).map_err(|_| invalid("packed width overflow"))?);
4936 out.extend_from_slice(&packed.base().to_le_bytes());
4937 put_u32(
4938 &mut out,
4939 u32::try_from(packed.words().len()).map_err(|_| invalid("too many packed words"))?,
4940 );
4941 for word in packed.words() {
4942 put_u64(&mut out, *word);
4943 }
4944 return Ok((out, membership));
4945 }
4946 let data = flat.data().ok_or_else(|| invalid("scalar column did not flatten"))?;
4947 match (ty, data) {
4948 (LogicalType::TinyInt, Data::Int8(values)) => {
4949 for value in &**values {
4950 out.extend_from_slice(&value.to_le_bytes());
4951 }
4952 }
4953 (LogicalType::UTinyInt, Data::UInt8(values)) => {
4954 for value in &**values {
4955 out.extend_from_slice(&value.to_le_bytes());
4956 }
4957 }
4958 (LogicalType::SmallInt, Data::Int16(values)) => {
4959 for value in &**values {
4960 out.extend_from_slice(&value.to_le_bytes());
4961 }
4962 }
4963 (LogicalType::USmallInt, Data::UInt16(values)) => {
4964 for value in &**values {
4965 out.extend_from_slice(&value.to_le_bytes());
4966 }
4967 }
4968 (LogicalType::UInteger, Data::UInt32(values)) => {
4969 for value in &**values {
4970 out.extend_from_slice(&value.to_le_bytes());
4971 }
4972 }
4973 (LogicalType::UBigInt, Data::UInt64(values)) => {
4974 for value in &**values {
4975 out.extend_from_slice(&value.to_le_bytes());
4976 }
4977 }
4978 (LogicalType::Integer | LogicalType::Date, Data::Int32(values)) => {
4979 for value in &**values {
4980 out.extend_from_slice(&value.to_le_bytes());
4981 }
4982 }
4983 (LogicalType::BigInt | LogicalType::Timestamp, Data::Int64(values)) => {
4984 for value in &**values {
4985 out.extend_from_slice(&value.to_le_bytes());
4986 }
4987 }
4988 (LogicalType::Boolean, Data::Bool(values)) => {
4989 for value in &**values {
4990 out.push(u8::from(*value));
4991 }
4992 }
4993 (LogicalType::Decimal { .. }, Data::Int16(values)) => {
4996 for value in &**values {
4997 out.extend_from_slice(&value.to_le_bytes());
4998 }
4999 }
5000 (LogicalType::Decimal { .. }, Data::Int32(values)) => {
5001 for value in &**values {
5002 out.extend_from_slice(&value.to_le_bytes());
5003 }
5004 }
5005 (LogicalType::Decimal { .. }, Data::Int64(values)) => {
5006 for value in &**values {
5007 out.extend_from_slice(&value.to_le_bytes());
5008 }
5009 }
5010 (LogicalType::Decimal { .. }, Data::Int128(values)) => {
5011 for value in &**values {
5012 out.extend_from_slice(&value.to_le_bytes());
5013 }
5014 }
5015 (LogicalType::Varchar, Data::Varlen(values)) => {
5016 let mut bytes = Vec::new();
5017 put_u32(&mut out, 0);
5018 for row in 0..vector.len() {
5019 let value = values.bytes(row).ok_or_else(|| invalid("string view is invalid"))?;
5020 bytes.extend_from_slice(value);
5021 put_u32(
5022 &mut out,
5023 u32::try_from(bytes.len())
5024 .map_err(|_| invalid("string payload exceeds 4GiB"))?,
5025 );
5026 }
5027 out.extend_from_slice(&bytes);
5028 }
5029 _ => return Err(Error::not_implemented(format!("native page for {ty}"))),
5030 }
5031 Ok((out, membership))
5032}
5033
5034fn put_varint(out: &mut Vec<u8>, mut value: u32) {
5035 while value >= 0x80 {
5036 out.push((value as u8 & 0x7f) | 0x80);
5037 value >>= 7;
5038 }
5039 out.push(value as u8);
5040}
5041
5042fn unique_codes(codes: &[u32]) -> Vec<u32> {
5044 let mut unique = codes.to_vec();
5045 unique.sort_unstable();
5046 unique.dedup();
5047 unique
5048}
5049
5050fn merged_codes(lists: Vec<Vec<u32>>) -> Vec<u32> {
5056 let mut lists = lists;
5057 while lists.len() > 1 {
5058 let mut next = Vec::with_capacity(lists.len().div_ceil(2));
5059 for pair in lists.chunks(2) {
5060 match pair {
5061 [left, right] => next.push(merged_pair(left, right)),
5062 [only] => next.push(only.clone()),
5063 _ => {}
5064 }
5065 }
5066 lists = next;
5067 }
5068 lists.pop().unwrap_or_default()
5069}
5070
5071fn merged_pair(left: &[u32], right: &[u32]) -> Vec<u32> {
5072 let mut out = Vec::with_capacity(left.len().saturating_add(right.len()));
5073 let mut at = 0;
5074 let mut to = 0;
5075 while at < left.len() && to < right.len() {
5076 match left[at].cmp(&right[to]) {
5077 Ordering::Less => {
5078 out.push(left[at]);
5079 at += 1;
5080 }
5081 Ordering::Greater => {
5082 out.push(right[to]);
5083 to += 1;
5084 }
5085 Ordering::Equal => {
5086 out.push(left[at]);
5087 at += 1;
5088 to += 1;
5089 }
5090 }
5091 }
5092 out.extend_from_slice(&left[at..]);
5093 out.extend_from_slice(&right[to..]);
5094 out
5095}
5096
5097fn merged_range(ranges: impl Iterator<Item = Range>) -> Range {
5102 let mut merged = Range::default();
5103 let mut first = true;
5104 for range in ranges {
5105 merged.nulls = merged.nulls.saturating_add(range.nulls);
5106 merged.sum = match (merged.sum.take(), range.sum) {
5110 (Some(held), Some(next)) if !first => held.checked_add(next),
5111 (_, next) if first => next,
5112 _ => None,
5113 };
5114 merged.exact = if first { range.exact } else { merged.exact && range.exact };
5115 if first {
5116 merged.low = range.low;
5117 merged.high = range.high;
5118 first = false;
5119 continue;
5120 }
5121 merged.low = match (merged.low.take(), range.low) {
5122 (Some(held), Some(next)) => Some(held.smaller(next)),
5123 _ => None,
5124 };
5125 merged.high = match (merged.high.take(), range.high) {
5126 (Some(held), Some(next)) => Some(held.larger(next)),
5127 _ => None,
5128 };
5129 }
5130 merged
5131}
5132
5133fn shortened(bound: Option<Bound>, high: bool) -> Option<Bound> {
5146 match bound {
5147 Some(Bound::Bytes(mut value)) if value.len() > PART_BOUND_BYTES => {
5148 value.truncate(PART_BOUND_BYTES);
5149 if !high {
5150 return Some(Bound::Bytes(value));
5151 }
5152 while let Some(last) = value.pop() {
5153 if last < u8::MAX {
5154 value.push(last + 1);
5155 return Some(Bound::Bytes(value));
5156 }
5157 }
5158 None
5159 }
5160 other => other,
5161 }
5162}
5163
5164fn encode_part_ranges(ranges: &[Range]) -> Result<Vec<u8>> {
5172 let mut out = Vec::new();
5173 put_u32(
5174 &mut out,
5175 u32::try_from(ranges.len()).map_err(|_| invalid("too many parts in a stripe"))?,
5176 );
5177 for range in ranges {
5178 put_bound(&mut out, shortened(range.low.clone(), false).as_ref())?;
5179 put_bound(&mut out, shortened(range.high.clone(), true).as_ref())?;
5180 put_u32(&mut out, u32::try_from(range.nulls).map_err(|_| invalid("null count overflow"))?);
5181 }
5182 Ok(out)
5183}
5184
5185fn decode_part_ranges(bytes: &[u8]) -> Result<Vec<Range>> {
5187 let mut cur = Cursor { bytes, at: 0 };
5188 let parts = cur.u32()? as usize;
5189 let mut out = Vec::new();
5190 for _ in 0..parts {
5191 let low = cur.bound()?;
5192 let high = cur.bound()?;
5193 let nulls = cur.u32()? as usize;
5194 out.push(Range { low, high, nulls, exact: false, sum: None });
5195 }
5196 Ok(out)
5197}
5198
5199fn encode_sieves<'a>(sieves: impl Iterator<Item = &'a Option<Sieve>>) -> Result<Vec<u8>> {
5200 let held: Vec<&Option<Sieve>> = sieves.collect();
5201 let mut out = Vec::new();
5202 put_u32(
5203 &mut out,
5204 u32::try_from(held.len()).map_err(|_| invalid("too many parts in a stripe"))?,
5205 );
5206 for sieve in &held {
5207 let length = sieve.as_ref().map_or(0, Sieve::len);
5208 put_u32(&mut out, u32::try_from(length).map_err(|_| invalid("sieve length overflow"))?);
5209 }
5210 for sieve in held.into_iter().flatten() {
5212 out.extend_from_slice(&sieve.to_bytes());
5213 }
5214 Ok(out)
5215}
5216
5217fn decode_sieves(bytes: &[u8]) -> Result<Vec<Option<Sieve>>> {
5223 let parts = u32::from_le_bytes(
5224 bytes
5225 .get(..4)
5226 .ok_or_else(|| invalid("sieve page is truncated"))?
5227 .try_into()
5228 .map_err(|_| invalid("sieve page is truncated"))?,
5229 ) as usize;
5230 let mut lengths = Vec::with_capacity(parts);
5231 for part in 0..parts {
5232 let at = 4 + part * 4;
5233 let field = bytes.get(at..at + 4).ok_or_else(|| invalid("sieve page is truncated"))?;
5234 lengths.push(u32::from_le_bytes(
5235 field.try_into().map_err(|_| invalid("sieve page is truncated"))?,
5236 ) as usize);
5237 }
5238 let mut at = 4 + parts * 4;
5239 let mut out = Vec::with_capacity(parts);
5240 for length in lengths {
5241 if length == 0 {
5242 out.push(None);
5243 continue;
5244 }
5245 let end = at.checked_add(length).ok_or_else(|| invalid("sieve page is truncated"))?;
5246 let field = bytes.get(at..end).ok_or_else(|| invalid("sieve page is truncated"))?;
5247 out.push(Sieve::from_bytes(field));
5248 at = end;
5249 }
5250 if at != bytes.len() {
5251 return Err(invalid("sieve page has trailing bytes"));
5252 }
5253 Ok(out)
5254}
5255
5256fn encode_membership(unique: &[u32]) -> Vec<u8> {
5262 let mut out = Vec::with_capacity(unique.len().saturating_mul(2).saturating_add(5));
5263 put_varint(&mut out, u32::try_from(unique.len()).unwrap_or(u32::MAX));
5264 let mut previous = 0;
5265 for (at, &code) in unique.iter().enumerate() {
5266 put_varint(&mut out, if at == 0 { code } else { code - previous });
5267 previous = code;
5268 }
5269 out
5270}
5271
5272fn take_varint(bytes: &[u8], at: &mut usize) -> Result<u32> {
5273 let mut value = 0_u32;
5274 for shift in (0..35).step_by(7) {
5275 let byte = *bytes.get(*at).ok_or_else(|| invalid("membership varint is truncated"))?;
5276 *at += 1;
5277 let part = u32::from(byte & 0x7f);
5278 if shift == 28 && part > 0x0f {
5279 return Err(invalid("membership varint overflow"));
5280 }
5281 value = value
5282 .checked_add(
5283 part.checked_shl(shift).ok_or_else(|| invalid("membership varint overflow"))?,
5284 )
5285 .ok_or_else(|| invalid("membership varint overflow"))?;
5286 if byte & 0x80 == 0 {
5287 return Ok(value);
5288 }
5289 }
5290 Err(invalid("membership varint is too long"))
5291}
5292
5293fn decode_membership(bytes: &[u8]) -> Result<Vec<u32>> {
5294 let mut at = 0;
5295 let count = take_varint(bytes, &mut at)? as usize;
5296 let mut codes = Vec::with_capacity(count);
5297 let mut previous = 0_u32;
5298 for index in 0..count {
5299 let delta = take_varint(bytes, &mut at)?;
5300 let code = if index == 0 {
5301 delta
5302 } else {
5303 previous.checked_add(delta).ok_or_else(|| invalid("membership code overflow"))?
5304 };
5305 if index > 0 && code <= previous {
5306 return Err(invalid("membership codes are not increasing"));
5307 }
5308 codes.push(code);
5309 previous = code;
5310 }
5311 if at != bytes.len() {
5312 return Err(invalid("membership page has trailing bytes"));
5313 }
5314 Ok(codes)
5315}
5316
5317fn string_dictionary(vector: &Vector) -> Result<Option<Vec<u8>>> {
5318 let mut by_text = HashMap::new();
5319 let mut values = Vec::new();
5320 let mut codes = Vec::with_capacity(vector.len());
5321 let mut plain_bytes = 0_usize;
5322 for row in 0..vector.len() {
5323 let text = vector.text_at(row).unwrap_or("");
5324 plain_bytes = plain_bytes.saturating_add(text.len());
5325 let code = match by_text.get(text) {
5326 Some(&code) => code,
5327 None => {
5328 let code = u32::try_from(values.len())
5329 .map_err(|_| invalid("too many dictionary values"))?;
5330 by_text.insert(text, code);
5331 values.push(text);
5332 code
5333 }
5334 };
5335 codes.push(code);
5336 }
5337 let dictionary_bytes = values.iter().map(|value| value.len()).sum::<usize>();
5338 let encoded = 8_usize
5339 .saturating_add((values.len() + 1).saturating_mul(4))
5340 .saturating_add(dictionary_bytes)
5341 .saturating_add(codes.len().saturating_mul(4));
5342 let plain = (vector.len() + 1).saturating_mul(4).saturating_add(plain_bytes);
5343 if encoded >= plain {
5344 return Ok(None);
5345 }
5346 let mut out = Vec::with_capacity(encoded);
5347 put_u32(
5348 &mut out,
5349 u32::try_from(values.len()).map_err(|_| invalid("too many dictionary values"))?,
5350 );
5351 put_u32(
5352 &mut out,
5353 u32::try_from(dictionary_bytes).map_err(|_| invalid("dictionary payload exceeds 4GiB"))?,
5354 );
5355 let mut offset = 0_u32;
5356 put_u32(&mut out, offset);
5357 for value in &values {
5358 offset = offset
5359 .checked_add(
5360 u32::try_from(value.len()).map_err(|_| invalid("dictionary value is too long"))?,
5361 )
5362 .ok_or_else(|| invalid("dictionary payload exceeds 4GiB"))?;
5363 put_u32(&mut out, offset);
5364 }
5365 for value in values {
5366 out.extend_from_slice(value.as_bytes());
5367 }
5368 for code in codes {
5369 put_u32(&mut out, code);
5370 }
5371 Ok(Some(out))
5372}
5373
5374struct EncodedDictionary {
5375 index: Vec<u8>,
5376 ranks: Vec<u8>,
5377 payload: Vec<Vec<u8>>,
5380}
5381
5382fn head(bytes: &[u8]) -> u64 {
5384 let mut word = [0; 8];
5385 let take = bytes.len().min(8);
5386 word[..take].copy_from_slice(&bytes[..take]);
5387 u64::from_be_bytes(word)
5388}
5389
5390fn rankings(dictionaries: &[Option<GlobalDictionary>]) -> Result<Vec<Vec<(u64, u32)>>> {
5398 let present =
5399 dictionaries.iter().enumerate().filter(|(_, held)| held.is_some()).map(|(at, _)| at);
5400 let present = present.collect::<Vec<_>>();
5401 let mut orders = vec![Vec::new(); dictionaries.len()];
5402 let workers = std::thread::available_parallelism()
5403 .map_or(1, usize::from)
5404 .min(MAX_FREQUENCY_WORKERS)
5405 .min(present.len());
5406 if workers <= 1 {
5407 for at in present {
5408 if let Some(dictionary) = &dictionaries[at] {
5409 orders[at] = dictionary.ranked();
5410 }
5411 }
5412 return Ok(orders);
5413 }
5414 let width = present.len().div_ceil(workers);
5415 let pieces = std::thread::scope(|scope| {
5416 present
5417 .chunks(width)
5418 .map(|columns| {
5419 scope.spawn(|| {
5420 columns
5421 .iter()
5422 .filter_map(|&at| dictionaries[at].as_ref().map(|held| (at, held.ranked())))
5423 .collect::<Vec<_>>()
5424 })
5425 })
5426 .collect::<Vec<_>>()
5427 .into_iter()
5428 .map(|handle| {
5429 handle.join().map_err(|_| Error::internal("a dictionary sort worker panicked"))
5430 })
5431 .collect::<Result<Vec<_>>>()
5432 })?;
5433 for piece in pieces {
5434 for (at, order) in piece {
5435 orders[at] = order;
5436 }
5437 }
5438 Ok(orders)
5439}
5440
5441fn encode_global_dictionary(
5442 dictionary: GlobalDictionary,
5443 order: &[(u64, u32)],
5444) -> Result<EncodedDictionary> {
5445 let values = dictionary.offsets.len() - 1;
5446 if order.len() != values {
5447 return Err(invalid("global dictionary order does not cover its values"));
5448 }
5449 let blocks = values.div_ceil(TEXT_PAYLOAD_VALUES);
5450 let payload = encode_payload(&dictionary)?;
5451 if payload.len() != blocks {
5452 return Err(invalid("global dictionary payload is not the blocks it says it is"));
5453 }
5454 let (ranks, rank_ends) = encode_ranks(order, code_width(values))?;
5455 let rank_blocks = values.div_ceil(TEXT_RANK_BLOCK);
5456 let offset_bits = offset_width(&dictionary.offsets);
5457 let mut index = Vec::with_capacity(
5458 DICTIONARY_HEADER + offset_bytes(values, offset_bits) + (blocks + rank_blocks) * 16,
5459 );
5460 put_u32(
5461 &mut index,
5462 u32::try_from(values).map_err(|_| invalid("global dictionary has too many values"))?,
5463 );
5464 put_u32(&mut index, TEXT_PAYLOAD_VALUES as u32);
5465 put_u32(
5466 &mut index,
5467 u32::try_from(blocks).map_err(|_| invalid("global dictionary has too many blocks"))?,
5468 );
5469 put_u32(&mut index, offset_bits as u32);
5470 encode_offsets(&dictionary.offsets, offset_bits, &mut index)?;
5471 let mut at = 0_u64;
5475 for block in &payload {
5476 at = at
5477 .checked_add(block.len() as u64)
5478 .ok_or_else(|| invalid("global dictionary payload overflow"))?;
5479 put_u64(&mut index, at);
5480 }
5481 for block in &payload {
5482 put_u64(&mut index, checksum(block));
5483 }
5484 if rank_ends.len() != rank_blocks {
5487 return Err(invalid("global dictionary order is not the blocks it says it is"));
5488 }
5489 for end in &rank_ends {
5490 put_u64(&mut index, *end);
5491 }
5492 let mut at = 0_usize;
5493 for end in &rank_ends {
5494 let end = usize::try_from(*end).map_err(|_| invalid("global dictionary order overflow"))?;
5495 put_u64(&mut index, checksum(&ranks[at..end]));
5496 at = end;
5497 }
5498 Ok(EncodedDictionary { index, ranks, payload })
5499}
5500
5501const PAYLOAD_SAMPLE_BLOCKS: usize = 8;
5508
5509fn payload_shapes() -> Vec<chooser::Settled> {
5535 let integers = vec![integer::Kind::Packed];
5536 [
5537 vec![string::Kind::Front, string::Kind::Lz],
5538 vec![string::Kind::Lz, string::Kind::Fsst],
5539 vec![string::Kind::Lz, string::Kind::Plain],
5540 vec![string::Kind::Fsst],
5541 vec![string::Kind::Plain],
5542 ]
5543 .into_iter()
5544 .map(|strings| chooser::Settled::new(strings, integers.clone()))
5545 .collect()
5546}
5547
5548fn encode_payload(dictionary: &GlobalDictionary) -> Result<Vec<Vec<u8>>> {
5554 let values = dictionary.offsets.len() - 1;
5555 let blocks = values.div_ceil(TEXT_PAYLOAD_VALUES);
5556 let run = |block: usize| {
5557 let first = block * TEXT_PAYLOAD_VALUES;
5558 let last = (first + TEXT_PAYLOAD_VALUES).min(values);
5559 (first..last)
5560 .map(|value| {
5561 let from = dictionary.offsets[value] as usize;
5562 let to = dictionary.offsets[value + 1] as usize;
5563 &dictionary.payload[from..to]
5564 })
5565 .collect::<Vec<_>>()
5566 };
5567 let shape = (blocks > PAYLOAD_SAMPLE_BLOCKS).then(|| settle_shape(&run, blocks)).transpose()?;
5570 let one = |block: usize| match &shape {
5571 Some(shape) => string::encode_with(&run(block), shape),
5572 None => string::encode(&run(block)),
5573 };
5574 let workers = std::thread::available_parallelism()
5575 .map_or(1, usize::from)
5576 .min(MAX_FREQUENCY_WORKERS)
5577 .min(blocks);
5578 if workers <= 1 {
5579 return (0..blocks).map(one).collect();
5580 }
5581 let next = AtomicUsize::new(0);
5582 let pieces = std::thread::scope(|scope| {
5583 (0..workers)
5584 .map(|_| {
5585 scope.spawn(|| {
5586 let mut mine = Vec::new();
5587 loop {
5588 let block = next.fetch_add(1, Atomic::Relaxed);
5589 if block >= blocks {
5590 break;
5591 }
5592 mine.push((block, one(block)?));
5593 }
5594 Ok(mine)
5595 })
5596 })
5597 .collect::<Vec<_>>()
5598 .into_iter()
5599 .map(|handle| {
5600 handle.join().map_err(|_| Error::internal("a dictionary encode worker panicked"))?
5601 })
5602 .collect::<Result<Vec<_>>>()
5603 })?;
5604 let mut payload = vec![Vec::new(); blocks];
5605 for piece in pieces {
5606 for (block, bytes) in piece {
5607 payload[block] = bytes;
5608 }
5609 }
5610 Ok(payload)
5611}
5612
5613fn settle_shape<'a>(
5621 run: &dyn Fn(usize) -> Vec<&'a [u8]>,
5622 blocks: usize,
5623) -> Result<chooser::Settled> {
5624 let last = blocks - 1;
5625 let sample = (0..PAYLOAD_SAMPLE_BLOCKS)
5626 .map(|region| run(region * last / (PAYLOAD_SAMPLE_BLOCKS - 1)))
5627 .collect::<Vec<_>>();
5628 let mut best: Option<(chooser::Settled, usize)> = None;
5629 for shape in payload_shapes() {
5630 let mut size = 0;
5631 for block in &sample {
5632 size += string::encode_with(block, &shape)?.len();
5633 }
5634 if best.as_ref().is_none_or(|(_, smallest)| size < *smallest) {
5635 best = Some((shape, size));
5636 }
5637 }
5638 best.map(|(shape, _)| shape)
5639 .ok_or_else(|| invalid("no shape applies to a global dictionary payload"))
5640}
5641
5642fn encode_ranks(order: &[(u64, u32)], code_bits: usize) -> Result<(Vec<u8>, Vec<u64>)> {
5649 let mut out = Vec::with_capacity(order.len() * 4);
5650 let mut ends = Vec::with_capacity(order.len().div_ceil(TEXT_RANK_BLOCK));
5651 let mut heads = Vec::with_capacity(TEXT_RANK_BLOCK);
5652 let mut codes = Vec::with_capacity(TEXT_RANK_BLOCK);
5653 for block in order.chunks(TEXT_RANK_BLOCK) {
5654 let base = block.first().map_or(0, |&(head, _)| head);
5657 let span = block.last().map_or(0, |&(head, _)| head.wrapping_sub(base));
5658 let width = (u64::BITS - span.leading_zeros()) as usize;
5659 heads.clear();
5660 codes.clear();
5661 for &(head, code) in block {
5662 heads.push(head.wrapping_sub(base));
5663 codes.push(u64::from(code));
5664 }
5665 put_u64(&mut out, base);
5666 out.push(width as u8);
5667 bitpack::pack_tail(&heads, width, &mut out)
5668 .map_err(|_| invalid("global dictionary heads do not pack"))?;
5669 bitpack::pack_tail(&codes, code_bits, &mut out)
5670 .map_err(|_| invalid("global dictionary codes do not pack"))?;
5671 ends.push(out.len() as u64);
5672 }
5673 Ok((out, ends))
5674}
5675
5676fn open_global_dictionary(
5683 file: Arc<File>,
5684 page: Page,
5685 ty: &LogicalType,
5686 keep_budget: usize,
5687) -> Result<Vector> {
5688 if ty != &LogicalType::Varchar {
5689 return Err(invalid("global dictionary belongs to a non-string column"));
5690 }
5691 let mut header = [0; DICTIONARY_HEADER];
5692 read_at(&file, page.offset, &mut header)?;
5693 let count = u32::from_le_bytes(header[0..4].try_into().expect("four bytes")) as usize;
5694 let per_block = u32::from_le_bytes(header[4..8].try_into().expect("four bytes")) as usize;
5695 let blocks = u32::from_le_bytes(header[8..12].try_into().expect("four bytes")) as usize;
5696 let offset_bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
5697 if per_block != TEXT_PAYLOAD_VALUES {
5698 return Err(invalid("global dictionary block width differs"));
5699 }
5700 if blocks != count.div_ceil(TEXT_PAYLOAD_VALUES) {
5701 return Err(invalid("global dictionary block count differs from its value count"));
5702 }
5703 if offset_bits > u32::BITS as usize {
5704 return Err(invalid("global dictionary packs offsets past a payload"));
5705 }
5706 let offset_len = offset_bytes(count, offset_bits);
5707 let ranks = count;
5712 let rank_blocks = ranks.div_ceil(TEXT_RANK_BLOCK);
5713 let hash_len = blocks
5716 .checked_add(rank_blocks)
5717 .and_then(|words| words.checked_mul(16))
5718 .ok_or_else(|| invalid("global dictionary block count overflow"))?;
5719 let index_len = DICTIONARY_HEADER
5720 .checked_add(offset_len)
5721 .and_then(|len| len.checked_add(hash_len))
5722 .ok_or_else(|| invalid("global dictionary header overflow"))?;
5723 if index_len > page.length as usize {
5724 return Err(invalid("global dictionary offset index exceeds its page"));
5725 }
5726 let mut index = vec![0; index_len];
5727 index[..DICTIONARY_HEADER].copy_from_slice(&header);
5728 read_at(&file, page.offset + DICTIONARY_HEADER as u64, &mut index[DICTIONARY_HEADER..])?;
5729 if checksum(&index) != page.hash {
5730 return Err(invalid("global dictionary index checksum differs"));
5731 }
5732 let offsets = index[DICTIONARY_HEADER..DICTIONARY_HEADER + offset_len].to_vec();
5733 let mut words = index[DICTIONARY_HEADER + offset_len..]
5734 .chunks_exact(8)
5735 .map(|part| u64::from_le_bytes(part.try_into().expect("eight bytes")))
5736 .collect::<Vec<_>>();
5737 let mut hashes = words.split_off(blocks);
5738 let mut rank_ends = hashes.split_off(blocks);
5739 let rank_hashes = rank_ends.split_off(rank_blocks);
5740 let ends = words;
5741 if rank_ends.windows(2).any(|pair| pair[0] >= pair[1]) {
5744 return Err(invalid("global dictionary order blocks do not rise"));
5745 }
5746 let rank_len = usize::try_from(rank_ends.last().copied().unwrap_or_default())
5747 .map_err(|_| invalid("global dictionary rank overflow"))?;
5748 let body_len = index_len
5749 .checked_add(rank_len)
5750 .ok_or_else(|| invalid("global dictionary header overflow"))?;
5751 if body_len > page.length as usize {
5752 return Err(invalid("global dictionary order exceeds its page"));
5753 }
5754 let stored_len = page.length as usize - body_len;
5757 if ends.last().copied().unwrap_or_default() as usize != stored_len
5758 || ends.windows(2).any(|pair| pair[0] > pair[1])
5759 {
5760 return Err(invalid("global dictionary blocks do not bound the payload"));
5761 }
5762 Vector::external_text(
5763 LogicalType::Varchar,
5764 Arc::new(NativeText {
5765 file,
5766 values: count,
5767 offsets,
5768 offset_bits,
5769 ranks,
5770 rank_at: page.offset + index_len as u64,
5771 rank_ends,
5772 rank_hashes,
5773 rank_blocks: (0..rank_blocks).map(|_| OnceLock::new()).collect(),
5774 code_bits: code_width(count),
5775 code_ranks: OnceLock::new(),
5776 payload: page.offset + body_len as u64,
5777 ends,
5778 hashes,
5779 blocks: (0..blocks).map(|_| OnceLock::new()).collect(),
5780 keep_budget,
5781 payload_kept: AtomicUsize::new(0),
5782 searched: Mutex::new(HashMap::new()),
5783 }),
5784 )
5785}
5786
5787fn decode(
5788 ty: &LogicalType,
5789 rows: usize,
5790 bytes: &[u8],
5791 global: Option<Arc<Vector>>,
5792) -> Result<Vector> {
5793 let mut cur = Cursor { bytes, at: 0 };
5794 let codec = cur.u8()?;
5795 let flag = cur.u8()?;
5796 let validity = match flag {
5797 0 => Validity::AllValid,
5798 1 => Validity::AllInvalid,
5799 2 => {
5800 let mask = cur.take(rows.div_ceil(8))?;
5801 Validity::from_iter(rows, |row| mask[row / 8] >> (row % 8) & 1 == 1)
5802 }
5803 _ => return Err(invalid("page validity tag differs")),
5804 };
5805 if codec == 1 {
5806 if ty != &LogicalType::Varchar {
5807 return Err(invalid("dictionary codec belongs to a non-string page"));
5808 }
5809 let count = cur.u32()? as usize;
5810 let payload_len = cur.u32()? as usize;
5811 let offset_bytes = cur.take(
5812 (count + 1)
5813 .checked_mul(4)
5814 .ok_or_else(|| invalid("dictionary offset count overflow"))?,
5815 )?;
5816 let offsets = offset_bytes
5817 .chunks_exact(4)
5818 .map(|part| u32::from_le_bytes(part.try_into().expect("four bytes")))
5819 .collect::<Vec<_>>();
5820 let payload = cur.take(payload_len)?.to_vec();
5821 if offsets.first() != Some(&0)
5822 || offsets.last().copied().map(|last| last as usize) != Some(payload.len())
5823 || offsets.windows(2).any(|pair| pair[0] > pair[1])
5824 {
5825 return Err(invalid("dictionary offsets do not bound the payload"));
5826 }
5827 let mut strings = StringColumn::over(Buffer::from_vec(payload).into_page());
5830 for pair in offsets.windows(2) {
5831 strings.push_in_place(pair[0] as usize, (pair[1] - pair[0]) as usize)?;
5832 }
5833 let mut codes = Vec::with_capacity(rows);
5834 for _ in 0..rows {
5835 codes.push(cur.u32()?);
5836 }
5837 if codes.iter().any(|code| *code as usize >= count) {
5838 return Err(invalid("dictionary code is out of range"));
5839 }
5840 if cur.at != bytes.len() {
5841 return Err(invalid("dictionary page has trailing bytes"));
5842 }
5843 let dictionary = Vector::flat(LogicalType::Varchar, Data::Varlen(strings))?;
5844 return Ok(Vector::dictionary(codes, dictionary)?.with_validity(validity));
5845 }
5846 if codec == 3 || codec == 4 {
5847 let dictionary = global.ok_or_else(|| invalid("global code page has no dictionary"))?;
5848 let codes = if codec == 4 {
5849 let wide = integer::decode(&bytes[cur.at..])?;
5852 if wide.len() != rows {
5853 return Err(invalid("encoded code page holds the wrong number of rows"));
5854 }
5855 let mut codes = Vec::with_capacity(wide.len());
5862 let mut seen = 0_i64;
5863 for &code in &wide {
5864 seen |= code;
5865 codes.push(code as u32);
5866 }
5867 if seen < 0 || seen > i64::from(u32::MAX) {
5868 return Err(invalid("code is not a code"));
5869 }
5870 codes
5871 } else {
5872 let mut codes = Vec::with_capacity(rows);
5873 for _ in 0..rows {
5874 codes.push(cur.u32()?);
5875 }
5876 if cur.at != bytes.len() {
5877 return Err(invalid("global code page has trailing bytes"));
5878 }
5879 codes
5880 };
5881 let highest = codes.iter().copied().max();
5882 return Ok(Vector::stable_dictionary_validated(codes, dictionary, highest)?
5883 .with_validity(validity));
5884 }
5885 if codec == 6 {
5886 if ty != &LogicalType::Varchar {
5887 return Err(invalid("compressed text codec belongs to a non-string page"));
5888 }
5889 let (payload, ends) = string::decode_flat(&bytes[cur.at..])?.into_parts();
5893 if ends.len() != rows {
5894 return Err(invalid("compressed text page holds the wrong number of rows"));
5895 }
5896 let mut values = StringColumn::over(Buffer::from_vec(payload).into_page());
5899 let mut start = 0;
5900 for end in ends {
5901 let len = end
5902 .checked_sub(start)
5903 .ok_or_else(|| invalid("compressed text value ends before it starts"))?;
5904 values.push_in_place(start, len)?;
5905 start = end;
5906 }
5907 return Ok(Vector::flat(ty.clone(), Data::Varlen(values))?.with_validity(validity));
5908 }
5909 if codec == 5 {
5910 let values = integer::decode(&bytes[cur.at..])?;
5912 if values.len() != rows {
5913 return Err(invalid("cascade page holds the wrong number of rows"));
5914 }
5915 let data = narrowed(ty, values)?;
5916 return Ok(Vector::flat(ty.clone(), data)?.with_validity(validity));
5917 }
5918 if codec == 2 {
5919 let width = u32::from(cur.u8()?);
5920 let base = i128::from_le_bytes(cur.take(16)?.try_into().expect("sixteen bytes"));
5921 let count = cur.u32()? as usize;
5922 let mut words = Vec::with_capacity(count);
5923 for _ in 0..count {
5924 words.push(cur.u64()?);
5925 }
5926 if cur.at != bytes.len() {
5927 return Err(invalid("packed page has trailing bytes"));
5928 }
5929 return Ok(Vector::packed(ty.clone(), words, width, base, rows)?.with_validity(validity));
5930 }
5931 if codec != 0 {
5932 return Err(invalid("page codec is unknown"));
5933 }
5934 let data = match ty {
5935 LogicalType::TinyInt => {
5936 let values = cur.take(rows)?;
5937 Data::Int8(values.iter().map(|item| *item as i8).collect::<Vec<_>>().into())
5938 }
5939 LogicalType::UTinyInt => Data::UInt8(cur.take(rows)?.to_vec().into()),
5940 LogicalType::SmallInt => {
5941 let values =
5942 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5943 Data::Int16(
5944 values
5945 .chunks_exact(2)
5946 .map(|item| i16::from_le_bytes(item.try_into().expect("two bytes")))
5947 .collect::<Vec<_>>()
5948 .into(),
5949 )
5950 }
5951 LogicalType::USmallInt => {
5952 let values =
5953 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5954 Data::UInt16(
5955 values
5956 .chunks_exact(2)
5957 .map(|item| u16::from_le_bytes(item.try_into().expect("two bytes")))
5958 .collect::<Vec<_>>()
5959 .into(),
5960 )
5961 }
5962 LogicalType::UInteger => {
5963 let values =
5964 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5965 Data::UInt32(
5966 values
5967 .chunks_exact(4)
5968 .map(|item| u32::from_le_bytes(item.try_into().expect("four bytes")))
5969 .collect::<Vec<_>>()
5970 .into(),
5971 )
5972 }
5973 LogicalType::UBigInt => {
5974 let values =
5975 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5976 Data::UInt64(
5977 values
5978 .chunks_exact(8)
5979 .map(|item| u64::from_le_bytes(item.try_into().expect("eight bytes")))
5980 .collect::<Vec<_>>()
5981 .into(),
5982 )
5983 }
5984 LogicalType::Integer | LogicalType::Date => {
5985 let values =
5986 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5987 Data::Int32(
5988 values
5989 .chunks_exact(4)
5990 .map(|item| i32::from_le_bytes(item.try_into().expect("four bytes")))
5991 .collect::<Vec<_>>()
5992 .into(),
5993 )
5994 }
5995 LogicalType::BigInt | LogicalType::Timestamp => {
5996 let values =
5997 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5998 Data::Int64(
5999 values
6000 .chunks_exact(8)
6001 .map(|item| i64::from_le_bytes(item.try_into().expect("eight bytes")))
6002 .collect::<Vec<_>>()
6003 .into(),
6004 )
6005 }
6006 LogicalType::Boolean => {
6007 let values = cur.take(rows)?;
6008 if values.iter().any(|value| *value > 1) {
6009 return Err(invalid("boolean page has another value"));
6010 }
6011 Data::Bool(values.iter().map(|value| *value == 1).collect::<Vec<_>>().into())
6012 }
6013 LogicalType::Decimal { .. } => match ty.physical() {
6016 PhysicalType::Int16 => {
6017 let values =
6018 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
6019 Data::Int16(
6020 values
6021 .chunks_exact(2)
6022 .map(|item| i16::from_le_bytes(item.try_into().expect("two bytes")))
6023 .collect::<Vec<_>>()
6024 .into(),
6025 )
6026 }
6027 PhysicalType::Int32 => {
6028 let values =
6029 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
6030 Data::Int32(
6031 values
6032 .chunks_exact(4)
6033 .map(|item| i32::from_le_bytes(item.try_into().expect("four bytes")))
6034 .collect::<Vec<_>>()
6035 .into(),
6036 )
6037 }
6038 PhysicalType::Int64 => {
6039 let values =
6040 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
6041 Data::Int64(
6042 values
6043 .chunks_exact(8)
6044 .map(|item| i64::from_le_bytes(item.try_into().expect("eight bytes")))
6045 .collect::<Vec<_>>()
6046 .into(),
6047 )
6048 }
6049 _ => {
6050 let values =
6051 cur.take(rows.checked_mul(16).ok_or_else(|| invalid("page size overflow"))?)?;
6052 Data::Int128(
6053 values
6054 .chunks_exact(16)
6055 .map(|item| i128::from_le_bytes(item.try_into().expect("sixteen bytes")))
6056 .collect::<Vec<_>>()
6057 .into(),
6058 )
6059 }
6060 },
6061 LogicalType::Varchar => {
6062 let offset_bytes = cur
6063 .take((rows + 1).checked_mul(4).ok_or_else(|| invalid("offset count overflow"))?)?;
6064 let offsets = offset_bytes
6065 .chunks_exact(4)
6066 .map(|part| u32::from_le_bytes(part.try_into().expect("four bytes")))
6067 .collect::<Vec<_>>();
6068 let payload = cur.take(bytes.len() - cur.at)?.to_vec();
6069 if offsets.first() != Some(&0)
6070 || offsets.last().copied().map(|last| last as usize) != Some(payload.len())
6071 || offsets.windows(2).any(|pair| pair[0] > pair[1])
6072 {
6073 return Err(invalid("string offsets do not bound the payload"));
6074 }
6075 let mut values = StringColumn::over(Buffer::from_vec(payload).into_page());
6079 for pair in offsets.windows(2) {
6080 values.push_in_place(pair[0] as usize, (pair[1] - pair[0]) as usize)?;
6081 }
6082 Data::Varlen(values)
6083 }
6084 _ => return Err(Error::not_implemented(format!("native page for {ty}"))),
6085 };
6086 if cur.at != bytes.len() {
6087 return Err(invalid("page has trailing bytes"));
6088 }
6089 Ok(Vector::flat(ty.clone(), data)?.with_validity(validity))
6090}
6091
6092#[cfg(test)]
6093mod tests {
6094 use std::fs;
6095 use std::io::{Seek, SeekFrom, Write};
6096 use std::path::PathBuf;
6097 use std::time::{SystemTime, UNIX_EPOCH};
6098
6099 use rudb_common::Stat;
6100 use rudb_common::Value;
6101 use rudb_common::bounds::{Frequencies, Op, Remainder, Zones};
6102 use rudb_common::stat::Provenance;
6103
6104 use super::*;
6105
6106 #[test]
6107 fn checksum_matches_fixed_vectors() {
6108 assert_eq!(checksum(b""), 0xef46_db37_51d8_e999);
6109 assert_eq!(checksum(b"a"), 0xd24e_c4f1_a98c_6e5b);
6110 assert_eq!(checksum(b"abc"), 0x44bc_2cf5_ad77_0999);
6111 }
6112
6113 fn path(label: &str) -> PathBuf {
6114 let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
6115 std::env::temp_dir().join(format!("rudb-native-{label}-{}-{stamp}.rdb", std::process::id()))
6116 }
6117
6118 #[test]
6120 fn a_read_at_an_offset_ignores_where_another_thread_left_the_cursor() {
6121 const SPANS: usize = 64;
6122 const SPAN: usize = 512;
6123 let path = path("positional");
6124 let content: Vec<u8> =
6125 (0..SPANS).flat_map(|span| std::iter::repeat_n(span as u8, SPAN)).collect();
6126 fs::write(&path, &content).expect("the file is written");
6127 let file = Arc::new(File::open(&path).expect("the file opens"));
6128 std::thread::scope(|scope| {
6129 for _ in 0..8 {
6130 let file = Arc::clone(&file);
6131 scope.spawn(move || {
6132 for _ in 0..64 {
6133 for span in 0..SPANS {
6134 let mut bytes = [0_u8; SPAN];
6135 read_at(&file, (span * SPAN) as u64, &mut bytes)
6136 .expect("the span reads");
6137 assert!(
6138 bytes.iter().all(|byte| *byte == span as u8),
6139 "span {span} came back as {}",
6140 bytes[0],
6141 );
6142 }
6143 }
6144 });
6145 }
6146 });
6147 let mut past = [0_u8; SPAN];
6148 let end = (SPANS * SPAN) as u64;
6149 let error = read_at(&file, end, &mut past).expect_err("a read past the end is refused");
6150 assert!(error.message().contains("ends before its declared length"), "{error}");
6151 drop(file);
6152 let _ = fs::remove_file(&path);
6153 }
6154
6155 #[test]
6161 fn a_writer_puts_a_page_where_it_said_it_did_wherever_the_cursor_has_got_to() {
6162 let path = path("cursor");
6163 let mut writer = Writer::create(
6164 &path,
6165 "items",
6166 vec![
6167 Field::required("id", LogicalType::Integer),
6168 Field::new("text", LogicalType::Varchar),
6169 ],
6170 )
6171 .expect("new file");
6172 writer.append(&sample()).expect("first part");
6173 writer.file.seek(SeekFrom::Start(0)).expect("the cursor goes back to the header");
6174 writer.append(&sample()).expect("second part");
6175 writer.file.seek(SeekFrom::Start(1)).expect("and somewhere useless again");
6176 writer.finish().expect("commit");
6177 let reader = Reader::open(&path).expect("reopen from disk");
6178 assert_eq!(reader.table().rows(), 6);
6179 let ids = reader.read(0, &[0]).expect("the integer page reads back");
6180 assert_eq!(ids.value_at(0, 0), Value::Integer(4));
6181 assert_eq!(ids.value_at(2, 0), Value::Integer(-2));
6182 let text = reader.read(1, &[1]).expect("the text page reads back");
6183 assert_eq!(text.value_at(1, 0), Value::Null);
6184 assert_eq!(text.value_at(2, 0), Value::Varchar("long text after a slash".into()));
6185 let end = reader.table().stripes().iter().flat_map(|stripe| {
6188 stripe
6189 .pages
6190 .iter()
6191 .map(|page| page.offset + u64::from(page.length))
6192 .chain(std::iter::once(stripe.index.offset + u64::from(stripe.index.length)))
6193 });
6194 let last = end.fold(HEADER, u64::max);
6195 let directory = fs::metadata(&path).expect("the file is there").len();
6196 assert!(last <= directory, "a page runs to {last} in a file of {directory} bytes");
6197 fs::remove_file(path).expect("remove scratch file");
6198 }
6199
6200 fn dictionary_index_len(header: &[u8; DICTIONARY_HEADER]) -> u64 {
6206 let count = u64::from(u32::from_le_bytes(header[0..4].try_into().expect("four bytes")));
6207 let blocks = u64::from(u32::from_le_bytes(header[8..12].try_into().expect("four bytes")));
6208 let bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
6209 let rank_blocks = count.div_ceil(TEXT_RANK_BLOCK as u64);
6210 DICTIONARY_HEADER as u64
6211 + offset_bytes(count as usize, bits) as u64
6212 + (blocks + rank_blocks) * 16
6213 }
6214
6215 fn last_rank_end(file: &File, offset: u64, header: &[u8; DICTIONARY_HEADER]) -> u64 {
6217 let count = u64::from(u32::from_le_bytes(header[0..4].try_into().expect("four bytes")));
6218 let blocks = u64::from(u32::from_le_bytes(header[8..12].try_into().expect("four bytes")));
6219 let bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
6220 let rank_blocks = count.div_ceil(TEXT_RANK_BLOCK as u64);
6221 let at = offset
6222 + DICTIONARY_HEADER as u64
6223 + offset_bytes(count as usize, bits) as u64
6224 + blocks * 16
6225 + (rank_blocks - 1) * 8;
6226 let mut end = [0; 8];
6227 read_at(file, at, &mut end).expect("the last rank block end");
6228 u64::from_le_bytes(end)
6229 }
6230
6231 fn sample() -> Chunk {
6232 Chunk::new(vec![
6233 Vector::from_values(
6234 LogicalType::Integer,
6235 &[Value::Integer(4), Value::Integer(9), Value::Integer(-2)],
6236 )
6237 .expect("integers"),
6238 Vector::from_values(
6239 LogicalType::Varchar,
6240 &[
6241 Value::Varchar("alpha".into()),
6242 Value::Null,
6243 Value::Varchar("long text after a slash".into()),
6244 ],
6245 )
6246 .expect("strings"),
6247 ])
6248 .expect("matching rows")
6249 }
6250
6251 fn sample_ids() -> Chunk {
6252 Chunk::new(vec![
6253 Vector::flat(LogicalType::Integer, Data::Int32(vec![7, 8, 9].into()))
6254 .expect("integers"),
6255 ])
6256 .expect("one column")
6257 }
6258
6259 #[test]
6260 fn the_planner_gets_the_null_count_off_the_same_directory_the_bounds_are_in() {
6261 let path = path("nulls_for_the_planner");
6264 let mut writer =
6265 Writer::create(&path, "items", vec![Field::new("a", LogicalType::Integer)])
6266 .expect("new file");
6267 let rows = Chunk::new(vec![
6268 Vector::from_values(
6269 LogicalType::Integer,
6270 &[
6271 Value::Integer(4),
6272 Value::Null,
6273 Value::Integer(9),
6274 Value::Null,
6275 Value::Integer(1),
6276 Value::Integer(2),
6277 ],
6278 )
6279 .expect("integers"),
6280 ])
6281 .expect("one column");
6282 writer.append(&rows).expect("the only part");
6283 writer.finish().expect("commit");
6284 let reader = Reader::open(&path).expect("reopen from disk");
6285 let stripes = Stripes::new(reader);
6286 let column = stripes.column("a").expect("the file has that column");
6287 assert_eq!(stripes.nulls(column), Stat::exact(2, Provenance::NullCount));
6288 assert_eq!(stripes.nulls(column + 1), Stat::Unknown);
6291 fs::remove_file(&path).expect("clean up");
6292 }
6293
6294 #[test]
6295 fn the_planner_gets_a_row_count_per_value_off_a_complete_synopsis() {
6296 let path = path("frequencies_for_the_planner");
6301 let mut writer =
6302 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6303 .expect("new file");
6304 let rows = Chunk::new(vec![
6305 Vector::from_values(
6306 LogicalType::Integer,
6307 &[
6308 Value::Integer(4),
6309 Value::Integer(4),
6310 Value::Integer(4),
6311 Value::Integer(9),
6312 Value::Integer(9),
6313 Value::Integer(1),
6314 ],
6315 )
6316 .expect("integers"),
6317 ])
6318 .expect("one column");
6319 writer.append(&rows).expect("the only part");
6320 writer.finish().expect("commit");
6321 let reader = Reader::open(&path).expect("reopen from disk");
6322 let common = Common::new(reader);
6323 assert_eq!(common.rows(), 6);
6324 let column = common.column("id").expect("the file has that column");
6325 assert_eq!(common.column("nothing"), None);
6326 assert_eq!(
6327 common.rows_with(column, &Bound::Int(4)),
6328 Stat::exact(3, Provenance::FrequencySynopsis)
6329 );
6330 assert_eq!(
6332 common.rows_with(column, &Bound::Int(7)),
6333 Stat::exact(0, Provenance::FrequencySynopsis)
6334 );
6335 assert_eq!(common.rows_with(column, &Bound::Bytes(b"four".to_vec())), Stat::Unknown);
6338 assert_eq!(common.remainder(column), None);
6341 fs::remove_file(&path).expect("clean up");
6342 }
6343
6344 #[test]
6345 fn the_planner_gets_an_exact_count_for_a_leading_value_of_an_incomplete_synopsis() {
6346 let path = path("frequency_prefix_for_the_planner");
6353 let mut writer =
6354 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6355 .expect("new file");
6356 let mut values = vec![Value::Integer(1); 10_000];
6357 for _ in 0..10 {
6358 values.extend((0..600).map(|tail| Value::Integer(1_000 + tail)));
6359 }
6360 for part in values.chunks(8_000) {
6363 let rows = Chunk::new(vec![
6364 Vector::from_values(LogicalType::Integer, part).expect("integers"),
6365 ])
6366 .expect("one column");
6367 writer.append(&rows).expect("a part");
6368 }
6369 writer.finish().expect("commit");
6370 let reader = Reader::open(&path).expect("reopen from disk");
6371 let prefix =
6372 reader.frequency_prefix(0).expect("a readable synopsis").expect("the column has one");
6373 assert_eq!(prefix.entries.len(), 512);
6376 assert_eq!(prefix.omitted_max, 10);
6377 let common = Common::new(reader);
6378 assert_eq!(common.rows(), 16_000);
6379 let column = common.column("id").expect("the file has that column");
6380 assert_eq!(
6381 common.rows_with(column, &Bound::Int(1)),
6382 Stat::exact(10_000, Provenance::FrequencySynopsis)
6383 );
6384 assert_eq!(
6386 common.rows_with(column, &Bound::Int(1_100)),
6387 Stat::exact(10, Provenance::FrequencySynopsis)
6388 );
6389 assert_eq!(common.rows_with(column, &Bound::Int(1_550)), Stat::Unknown);
6392 assert_eq!(common.rows_with(column, &Bound::Int(9_999)), Stat::Unknown);
6395 let remainder = common.remainder(column).expect("the list is a prefix");
6399 assert_eq!(remainder, Remainder { rows: 890, listed: 512, most: 10 });
6400 assert_eq!(remainder.rows / (601 - remainder.listed), 10);
6401 fs::remove_file(&path).expect("clean up");
6402 }
6403
6404 #[test]
6406 fn a_file_holding_no_table_commits_and_opens_and_a_table_can_be_added_to_it() {
6407 let path = path("empty");
6408 Writer::empty(&path).expect("a file with nothing in it");
6409 let catalog = Catalog::open(&path).expect("the empty file opens");
6410 assert_eq!(catalog.len(), 0);
6411 assert!(catalog.is_empty());
6412 assert_eq!(catalog.names().count(), 0);
6413 let mut writer =
6416 Writer::open(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6417 .expect("a table goes into the empty file");
6418 writer.append(&sample_ids()).expect("rows");
6419 writer.finish().expect("commit");
6420 let catalog = Catalog::open(&path).expect("the file opens again");
6421 assert_eq!(catalog.names().collect::<Vec<_>>(), vec!["items"]);
6422 fs::remove_file(&path).expect("clean up");
6423 }
6424
6425 #[test]
6426 fn committed_file_reopens_and_reads_only_requested_columns() {
6427 let path = path("reopen");
6428 let mut writer = Writer::create(
6429 &path,
6430 "items",
6431 vec![
6432 Field::required("id", LogicalType::Integer),
6433 Field::new("text", LogicalType::Varchar),
6434 ],
6435 )
6436 .expect("new file");
6437 writer.append(&sample()).expect("first part");
6438 writer.append(&sample()).expect("second part");
6439 writer.finish().expect("commit");
6440 let reader = Reader::open(&path).expect("reopen from disk");
6441 assert_eq!(reader.table().rows(), 6);
6442 assert_eq!(reader.table().stripes().len(), 1);
6445 assert_eq!(reader.parts(), 2);
6446 assert_eq!(reader.part_rows(0), 3);
6447 assert_eq!(reader.part_rows(1), 3);
6448 let text = reader.read(1, &[1]).expect("only text page");
6449 assert_eq!(text.width(), 1);
6450 assert_eq!(text.value_at(1, 0), Value::Null);
6451 assert_eq!(text.value_at(2, 0), Value::Varchar("long text after a slash".into()));
6452 let sparse = reader.read_sparse(1, &[1]).expect("one part without its whole page");
6453 assert_eq!(sparse.width(), 1);
6454 assert_eq!(sparse.value_at(1, 0), Value::Null);
6455 assert_eq!(sparse.value_at(2, 0), Value::Varchar("long text after a slash".into()));
6456 assert!(!reader.skips_codes(0, 1, &[0]).expect("alpha is in the stripe"));
6457 assert!(!reader.skips_codes(0, 1, &[2]).expect("long text is in the stripe"));
6458 assert!(reader.skips_codes(0, 1, &[3]).expect("unknown code is absent"));
6459 let count = reader.read(0, &[]).expect("no page is needed for count");
6460 assert_eq!(count.len(), 3);
6461 assert!(reader.skips(0, &[Probe { column: 0, op: Op::Greater, value: Bound::Int(100) }]));
6462 assert!(!reader.skips(0, &[Probe { column: 0, op: Op::Greater, value: Bound::Int(0) }]));
6463 let integers = reader.top_frequencies(0, 1).expect("valid integer synopsis").expect("kept");
6464 assert_eq!(
6465 integers,
6466 vec![(Value::Integer(-2), 2), (Value::Integer(4), 2), (Value::Integer(9), 2),]
6467 );
6468 let strings = reader.top_frequencies(1, 1).expect("valid string synopsis").expect("kept");
6469 assert_eq!(strings.len(), 3);
6470 assert!(strings.contains(&(Value::Null, 2)));
6471 assert!(strings.contains(&(Value::Varchar("alpha".into()), 2)));
6472 assert!(strings.contains(&(Value::Varchar("long text after a slash".into()), 2)));
6473 fs::remove_file(path).expect("remove scratch file");
6474 }
6475
6476 #[test]
6484 fn runs_handed_over_out_of_order_still_read_back_in_source_order() {
6485 let path = path("interleaved-runs");
6486 let mut writer =
6487 Writer::create(&path, "interleaved", vec![Field::new("v", LogicalType::BigInt)])
6488 .expect("new file");
6489 for morsel in [2_u64, 0, 3, 1] {
6490 let parts = (0..4_u64)
6491 .map(|chunk| {
6492 let first = i64::try_from(morsel * 32 + chunk * 8).expect("small");
6493 let values =
6494 (0..8_i64).map(|row| Value::BigInt(first + row)).collect::<Vec<_>>();
6495 let column =
6496 Vector::from_values(LogicalType::BigInt, &values).expect("a column");
6497 ((morsel, chunk), Chunk::new(vec![column]).expect("one column"))
6498 })
6499 .collect::<Vec<_>>();
6500 writer.append_stripe(parts).expect("a stripe");
6501 }
6502 writer.finish().expect("commit");
6503
6504 let reader = Reader::open(&path).expect("valid directory");
6505 assert_eq!(reader.table().stripes().len(), 4, "a run is a stripe of its own");
6506 assert_eq!(reader.table().rows(), 128);
6507 for part in 0..16_usize {
6508 let read = reader.read(part, &[0]).expect("a part back");
6509 for row in 0..8_usize {
6510 let want = i64::try_from(part * 8 + row).expect("small");
6511 assert_eq!(read.value_at(row, 0), Value::BigInt(want), "part {part} row {row}");
6512 }
6513 }
6514 fs::remove_file(path).expect("remove scratch file");
6515 }
6516
6517 #[test]
6520 fn runs_that_overlap_each_other_are_refused_at_commit() {
6521 let path = path("overlapping-runs");
6522 let mut writer =
6523 Writer::create(&path, "overlapping", vec![Field::new("v", LogicalType::BigInt)])
6524 .expect("new file");
6525 let one = |order: (u64, u64)| {
6526 let column =
6527 Vector::from_values(LogicalType::BigInt, &[Value::BigInt(1)]).expect("a column");
6528 (order, Chunk::new(vec![column]).expect("one column"))
6529 };
6530 writer.append_stripe(vec![one((0, 0)), one((0, 2))]).expect("a stripe");
6533 writer.append_stripe(vec![one((0, 1))]).expect("a stripe");
6534 let error = writer.finish().expect_err("the runs overlap");
6535 assert!(error.message().contains("source order"), "{error}");
6536 fs::remove_file(path).expect("remove scratch file");
6537 }
6538
6539 #[test]
6542 fn a_run_longer_than_a_stripe_is_refused() {
6543 let path = path("overlong-run");
6544 let mut writer =
6545 Writer::create(&path, "overlong", vec![Field::new("v", LogicalType::BigInt)])
6546 .expect("new file");
6547 let parts = (0..=STRIPE_PARTS)
6548 .map(|at| {
6549 let column = Vector::from_values(LogicalType::BigInt, &[Value::BigInt(1)])
6550 .expect("a column");
6551 let chunk = Chunk::new(vec![column]).expect("one column");
6552 ((0, u64::try_from(at).expect("small")), chunk)
6553 })
6554 .collect::<Vec<_>>();
6555 let error = writer.append_stripe(parts).expect_err("one part too many");
6556 assert!(error.message().contains("more parts than it holds"), "{error}");
6557 fs::remove_file(path).expect("remove scratch file");
6558 }
6559
6560 #[test]
6566 fn parts_past_the_stripe_bound_start_a_new_stripe() {
6567 let path = path("stripe-bound");
6568 let mut writer = Writer::create(
6569 &path,
6570 "items",
6571 vec![
6572 Field::required("id", LogicalType::Integer),
6573 Field::new("text", LogicalType::Varchar),
6574 ],
6575 )
6576 .expect("new file");
6577 let parts = STRIPE_PARTS * 2 + 3;
6578 for part in 0..parts {
6579 let id = part as i32;
6580 let chunk = Chunk::new(vec![
6581 Vector::from_values(
6582 LogicalType::Integer,
6583 &[Value::Integer(id), Value::Integer(-id)],
6584 )
6585 .expect("integers"),
6586 Vector::from_values(
6587 LogicalType::Varchar,
6588 &[Value::Varchar(format!("value {part}")), Value::Null],
6589 )
6590 .expect("strings"),
6591 ])
6592 .expect("matching rows");
6593 writer.append(&chunk).expect("one part");
6594 }
6595 writer.finish().expect("commit");
6596
6597 let reader = Reader::open(&path).expect("reopen from disk");
6598 assert_eq!(reader.parts(), parts);
6599 assert_eq!(reader.table().rows(), parts * 2);
6600 assert_eq!(reader.table().stripes().len(), parts.div_ceil(STRIPE_PARTS));
6601 assert_eq!(reader.table().stripes()[0].parts(), STRIPE_PARTS);
6602 assert_eq!(reader.table().stripes()[0].rows(), STRIPE_PARTS * 2);
6603 assert_eq!(reader.table().stripes()[2].parts(), 3);
6604 for part in (0..parts).rev() {
6607 let dense = reader.read(part, &[0, 1]).expect("a whole page read");
6608 let sparse = reader.read_sparse(part, &[0, 1]).expect("one part read");
6609 for chunk in [&dense, &sparse] {
6610 assert_eq!(chunk.len(), 2, "part {part} has its own row count");
6611 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6612 assert_eq!(chunk.value_at(1, 0), Value::Integer(-(part as i32)));
6613 assert_eq!(chunk.value_at(0, 1), Value::Varchar(format!("value {part}")));
6614 assert_eq!(chunk.value_at(1, 1), Value::Null);
6615 }
6616 }
6617 let above = [Probe { column: 0, op: Op::Greater, value: Bound::Int(100) }];
6620 assert!(reader.skips(0, &above), "the first stripe stops at 63");
6621 assert!(!reader.skips(STRIPE_PARTS * 2, &above), "the third stripe reaches 130");
6622 fs::remove_file(path).expect("remove scratch file");
6623 }
6624
6625 fn scattered(n: i64) -> i64 {
6627 n.wrapping_mul(-7_046_029_254_386_353_131)
6628 }
6629
6630 #[test]
6636 fn a_part_is_skipped_when_its_sieve_does_not_hold_the_constant() {
6637 let path = path("sieve-skip");
6638 let mut writer =
6639 Writer::create(&path, "hits", vec![Field::required("id", LogicalType::BigInt)])
6640 .expect("new file");
6641 let parts = STRIPE_PARTS + 3;
6642 let per_part = 128;
6646 for part in 0..parts {
6647 let held: Vec<Value> = (0..per_part)
6648 .map(|row| Value::BigInt(scattered((part * per_part + row) as i64)))
6649 .collect();
6650 let chunk =
6651 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6652 .expect("one column");
6653 writer.append(&chunk).expect("one part");
6654 }
6655 writer.finish().expect("commit");
6656
6657 let reader = Reader::open(&path).expect("reopen from disk");
6658 let probe = |value: i64| Probe {
6659 column: 0,
6660 op: Op::Equal,
6661 value: Bound::Int(i128::from(scattered(value))),
6662 };
6663 for wanted in [0_i64, (per_part + 1) as i64, (parts * per_part - 1) as i64] {
6664 let tests = [probe(wanted)];
6665 let kept: Vec<usize> = (0..parts).filter(|&part| !reader.skips(part, &tests)).collect();
6666 let home = wanted as usize / per_part;
6667 assert!(kept.contains(&home), "the part holding {wanted} is read");
6668 assert!(kept.len() <= 2, "{wanted} keeps {kept:?}, which is more than one stray part");
6672 }
6673 let absent = [probe((parts * per_part) as i64 + 1)];
6674 let kept = (0..parts).filter(|&part| !reader.skips(part, &absent)).count();
6675 assert!(kept <= 1, "{kept} parts of {parts} kept a value no part holds");
6676 let tests = [probe(0)];
6679 assert!(
6680 reader.table().stripes().iter().all(|stripe| !stripe.zone.skips(&tests)),
6681 "the bounds rule out no stripe at all"
6682 );
6683 fs::remove_file(path).expect("remove scratch file");
6684 }
6685
6686 #[test]
6692 fn a_part_is_skipped_when_its_own_bounds_rule_out_a_comparison_the_stripe_keeps() {
6693 let path = path("part-range-skip");
6694 let mut writer =
6695 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
6696 .expect("new file");
6697 let parts = STRIPE_PARTS + 3;
6698 let per_part = 128;
6699 for part in 0..parts {
6700 let held: Vec<Value> = (0..per_part)
6704 .map(|row| {
6705 Value::BigInt((part * 1_000) as i64 + (scattered(row as i64).rem_euclid(900)))
6706 })
6707 .collect();
6708 let chunk =
6709 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6710 .expect("one column");
6711 writer.append(&chunk).expect("one part");
6712 }
6713 writer.finish().expect("commit");
6714
6715 let reader = Reader::open(&path).expect("reopen from disk");
6716 let under = [Probe { column: 0, op: Op::Less, value: Bound::Int(3_000) }];
6717 let kept: Vec<usize> = (0..parts).filter(|&part| !reader.skips(part, &under)).collect();
6718 assert_eq!(kept, vec![0, 1, 2], "only the three parts that start under three thousand");
6719 assert!(!reader.stripe_skips(0, &under), "the stripe reaches from zero and keeps itself");
6721 fs::remove_file(path).expect("remove scratch file");
6722 }
6723
6724 #[test]
6728 fn a_part_is_waved_through_when_its_own_bounds_pass_a_comparison_the_stripe_cannot() {
6729 let path = path("part-range-certain");
6730 let mut writer =
6731 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
6732 .expect("new file");
6733 let parts = STRIPE_PARTS + 3;
6734 let per_part = 128;
6735 for part in 0..parts {
6736 let held: Vec<Value> = (0..per_part)
6737 .map(|row| {
6738 Value::BigInt((part * 1_000) as i64 + (scattered(row as i64).rem_euclid(900)))
6739 })
6740 .collect();
6741 let chunk =
6742 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6743 .expect("one column");
6744 writer.append(&chunk).expect("one part");
6745 }
6746 writer.finish().expect("commit");
6747
6748 let reader = Reader::open(&path).expect("reopen from disk");
6749 let under = [Probe { column: 0, op: Op::Less, value: Bound::Int(3_000) }];
6750 let waved: Vec<usize> = (0..parts).filter(|&part| reader.certain(part, &under)).collect();
6751 assert_eq!(waved, vec![0, 1, 2], "the three parts that end under three thousand");
6752 assert!(!reader.stripe_skips(0, &under), "the stripe straddles the comparison");
6755 fs::remove_file(path).expect("remove scratch file");
6756 }
6757
6758 #[test]
6761 fn a_stripe_of_one_part_writes_no_range_page_and_a_stripe_of_many_does() {
6762 for (parts, wanted) in [(1_usize, false), (STRIPE_PARTS, true)] {
6763 let path = path("part-range-page");
6764 let mut writer =
6765 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
6766 .expect("new file");
6767 for part in 0..parts {
6768 let held: Vec<Value> = (0..128)
6769 .map(|row| {
6770 Value::BigInt((part * 1_000) as i64 + scattered(row as i64).rem_euclid(900))
6771 })
6772 .collect();
6773 let chunk = Chunk::new(vec![
6774 Vector::from_values(LogicalType::BigInt, &held).expect("numbers"),
6775 ])
6776 .expect("one column");
6777 writer.append(&chunk).expect("one part");
6778 }
6779 writer.finish().expect("commit");
6780 let reader = Reader::open(&path).expect("reopen from disk");
6781 let bytes = reader.layout().columns[0].part_ranges;
6782 assert_eq!(bytes > 0, wanted, "{parts} parts wrote {bytes} bytes of ranges");
6783 fs::remove_file(path).expect("remove scratch file");
6784 }
6785 }
6786
6787 #[test]
6790 fn a_string_end_that_is_cut_down_still_covers_the_value_it_came_from() {
6791 let long = vec![b'a'; PART_BOUND_BYTES * 2];
6792 let low = shortened(Some(Bound::Bytes(long.clone())), false).expect("a low end");
6793 let high = shortened(Some(Bound::Bytes(long.clone())), true).expect("a high end");
6794 let Bound::Bytes(low) = low else { panic!("a string stays a string") };
6795 let Bound::Bytes(high) = high else { panic!("a string stays a string") };
6796 assert!(low.len() <= PART_BOUND_BYTES && high.len() <= PART_BOUND_BYTES);
6797 assert!(low.as_slice() <= long.as_slice(), "the low end is at or under the value");
6798 assert!(high.as_slice() >= long.as_slice(), "the high end is at or over the value");
6799 }
6800
6801 #[test]
6804 fn a_string_end_with_no_room_to_step_up_gives_up_the_bound() {
6805 let long = vec![u8::MAX; PART_BOUND_BYTES * 2];
6806 assert_eq!(shortened(Some(Bound::Bytes(long.clone())), true), None);
6807 let low = shortened(Some(Bound::Bytes(long)), false).expect("a low end is still a prefix");
6808 assert_eq!(low, Bound::Bytes(vec![u8::MAX; PART_BOUND_BYTES]));
6809 }
6810
6811 #[test]
6821 fn a_sieve_larger_than_the_part_it_indexes_is_not_written() {
6822 let path = path("sieve-pays");
6823 let fields = vec![
6824 Field::required("spread", LogicalType::BigInt),
6825 Field::required("repeated", LogicalType::BigInt),
6826 ];
6827 let mut writer = Writer::create(&path, "hits", fields).expect("new file");
6828 let parts = 3;
6829 let per_part = 1024;
6830 for part in 0..parts {
6831 let base = (part * per_part) as i64;
6832 let spread: Vec<Value> =
6833 (0..per_part).map(|row| Value::BigInt(scattered(base + row as i64))).collect();
6834 let repeated: Vec<Value> =
6835 (0..per_part).map(|row| Value::BigInt(scattered((row / 256) as i64))).collect();
6836 let chunk = Chunk::new(vec![
6837 Vector::from_values(LogicalType::BigInt, &spread).expect("numbers"),
6838 Vector::from_values(LogicalType::BigInt, &repeated).expect("numbers"),
6839 ])
6840 .expect("two columns");
6841 writer.append(&chunk).expect("one part");
6842 }
6843 writer.finish().expect("commit");
6844
6845 let reader = Reader::open(&path).expect("reopen from disk");
6846 let layout = reader.layout();
6847 let spread = &layout.columns[0];
6848 let repeated = &layout.columns[1];
6849 assert!(spread.sieves > 0, "a column whose parts are worth a filter keeps one");
6850 assert_eq!(
6851 repeated.sieves, 0,
6852 "a column whose filter costs more than its parts keeps none"
6853 );
6854 for column in &layout.columns {
6857 assert!(
6858 column.sieves < column.pages,
6859 "{} spends {} on sieves over {} of data",
6860 column.name,
6861 column.sieves,
6862 column.pages
6863 );
6864 }
6865 let absent = [Probe {
6867 column: 0,
6868 op: Op::Equal,
6869 value: Bound::Int(i128::from(scattered((parts * per_part) as i64 + 1))),
6870 }];
6871 assert!((0..parts).all(|part| reader.skips(part, &absent)), "no part holds it");
6872 fs::remove_file(path).expect("remove scratch file");
6873 }
6874
6875 #[test]
6881 fn a_damaged_sieve_page_is_read_through_rather_than_refused() {
6882 let path = path("sieve-damaged");
6883 let mut writer =
6884 Writer::create(&path, "hits", vec![Field::required("id", LogicalType::BigInt)])
6885 .expect("new file");
6886 let rows = 128;
6887 let held: Vec<Value> = (0..rows).map(|row| Value::BigInt(scattered(row))).collect();
6888 let chunk =
6889 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6890 .expect("one column");
6891 writer.append(&chunk).expect("one part");
6892 writer.finish().expect("commit");
6893
6894 let page =
6895 Reader::open(&path).expect("reopen").table.stripes[0].sieves[0].expect("a sieve page");
6896 let mut file = OpenOptions::new().write(true).open(&path).expect("open the sieve page");
6897 file.seek(SeekFrom::Start(page.offset + u64::from(page.length) - 1)).expect("seek");
6898 file.write_all(&[0xff]).expect("damage one byte");
6899 drop(file);
6900
6901 let reader = Reader::open(&path).expect("reopen the damaged file");
6902 let absent =
6903 [Probe { column: 0, op: Op::Equal, value: Bound::Int(i128::from(scattered(99))) }];
6904 assert!(!reader.skips(0, &absent), "a sieve that cannot be read skips nothing");
6905 assert_eq!(
6906 reader.read(0, &[0]).expect("the rows are untouched").len(),
6907 usize::try_from(rows).expect("a small count")
6908 );
6909 fs::remove_file(path).expect("remove scratch file");
6910 }
6911
6912 #[test]
6923 fn workers_that_want_the_same_stripe_read_it_once() {
6924 let path = path("single-flight");
6925 let mut writer =
6926 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6927 .expect("new file");
6928 for part in 0..STRIPE_PARTS {
6929 let id = part as i32;
6930 let chunk = Chunk::new(vec![
6931 Vector::from_values(
6932 LogicalType::Integer,
6933 &[Value::Integer(id), Value::Integer(-id)],
6934 )
6935 .expect("integers"),
6936 ])
6937 .expect("matching rows");
6938 writer.append(&chunk).expect("one part");
6939 }
6940 writer.finish().expect("commit");
6941
6942 let reader = Reader::open(&path).expect("reopen from disk");
6943 assert_eq!(reader.table().stripes().len(), 1, "one stripe is the point of the test");
6944 let barrier = std::sync::Barrier::new(8);
6945 std::thread::scope(|scope| {
6946 for worker in 0..8 {
6947 let reader = &reader;
6948 let barrier = &barrier;
6949 scope.spawn(move || {
6950 barrier.wait();
6951 for part in (worker..STRIPE_PARTS).step_by(8) {
6952 let chunk = reader.read(part, &[0]).expect("a whole page read");
6953 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6954 assert_eq!(chunk.value_at(1, 0), Value::Integer(-(part as i32)));
6955 }
6956 });
6957 }
6958 });
6959 assert_eq!(reader.pages.load(Atomic::Relaxed), 1, "one stripe, one page read, whoever won");
6960 fs::remove_file(path).expect("remove scratch file");
6961 }
6962
6963 #[test]
6976 fn opening_costs_the_same_over_a_thousand_times_the_rows() {
6977 let opened = |label: &str, rows_per_part: i32| {
6978 let path = path(label);
6979 let mut writer =
6980 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6981 .expect("new file");
6982 for part in 0..STRIPE_PARTS * 3 {
6983 let values = (0..rows_per_part)
6987 .map(|row| {
6988 Value::Integer((part as i32 * rows_per_part + row).wrapping_mul(2_654_435))
6989 })
6990 .collect::<Vec<_>>();
6991 let chunk = Chunk::new(vec![
6992 Vector::from_values(LogicalType::Integer, &values).expect("integers"),
6993 ])
6994 .expect("matching rows");
6995 writer.append(&chunk).expect("one part");
6996 }
6997 writer.finish().expect("commit");
6998 let reader = Reader::open(&path).expect("reopen from disk");
6999 let size = fs::metadata(&path).expect("the file is there").len();
7000 let out = (reader.reads(), reader.table().stripes().len(), size);
7001 fs::remove_file(path).expect("remove scratch file");
7002 out
7003 };
7004
7005 let (thin, thin_stripes, thin_size) = opened("open-thin", 1);
7006 let (fat, fat_stripes, fat_size) = opened("open-fat", 1000);
7007 assert_eq!(
7008 thin_stripes, fat_stripes,
7009 "the same stripe count is what makes this a fair ask"
7010 );
7011 assert!(
7012 fat_size > thin_size * 50,
7013 "the fat file has to actually be larger, and it is {fat_size} against {thin_size}"
7014 );
7015
7016 assert_eq!(thin.opening.reads, fat.opening.reads, "the same reads either way");
7017 assert_eq!(thin.pages, 0, "opening read a page");
7018 assert_eq!(fat.pages, 0, "opening read a page");
7019 assert_eq!(thin.indexes, 0, "opening read an index");
7020 assert_eq!(fat.indexes, 0, "opening read an index");
7021 assert!(
7024 fat.opening.bytes < thin.opening.bytes * 2,
7025 "opening the thin file read {} bytes and the fat one read {}",
7026 thin.opening.bytes,
7027 fat.opening.bytes
7028 );
7029 }
7030
7031 #[test]
7039 fn two_opens_of_one_file_cost_the_same_and_the_second_is_not_cheaper() {
7040 let path = path("open-twice");
7041 let mut writer =
7042 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
7043 .expect("new file");
7044 for part in 0..STRIPE_PARTS * 3 {
7045 let chunk = Chunk::new(vec![
7046 Vector::from_values(LogicalType::Integer, &[Value::Integer(part as i32)])
7047 .expect("integers"),
7048 ])
7049 .expect("matching rows");
7050 writer.append(&chunk).expect("one part");
7051 }
7052 writer.finish().expect("commit");
7053
7054 let first = Reader::open(&path).expect("open");
7055 for part in 0..first.parts() {
7058 first.read(part, &[0]).expect("a part");
7059 }
7060 assert!(first.reads().pages > 0, "the scan has to have read something");
7061 let second = Reader::open(&path).expect("open again");
7062
7063 assert_eq!(first.reads().opening, second.reads().opening);
7064 assert_eq!(
7065 second.reads().pages,
7066 0,
7067 "the second open read a page off the back of the first"
7068 );
7069 assert_eq!(second.reads().indexes, 0, "the second open read an index it inherited");
7070 fs::remove_file(path).expect("remove scratch file");
7071 }
7072
7073 #[test]
7081 fn an_index_is_read_once_per_stripe_however_often_the_page_is_evicted() {
7082 let path = path("index-cache");
7083 let mut writer =
7084 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
7085 .expect("new file");
7086 let parts = STRIPE_PARTS * (CACHED_STRIPES_PER_COLUMN + 2);
7087 for part in 0..parts {
7088 let id = part as i32;
7089 let chunk = Chunk::new(vec![
7090 Vector::from_values(LogicalType::Integer, &[Value::Integer(id)]).expect("integers"),
7091 ])
7092 .expect("matching rows");
7093 writer.append(&chunk).expect("one part");
7094 }
7095 writer.finish().expect("commit");
7096
7097 let reader = Reader::open(&path).expect("reopen from disk");
7098 let stripes = reader.table().stripes().len();
7099 assert!(stripes > CACHED_STRIPES_PER_COLUMN, "the page cache has to be too small for this");
7100 for _ in 0..2 {
7102 for part in 0..parts {
7103 let chunk = reader.read(part, &[0]).expect("a part");
7104 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
7105 }
7106 }
7107 assert_eq!(reader.indexes.load(Atomic::Relaxed), stripes, "one index read per stripe");
7108 assert!(
7109 reader.pages.load(Atomic::Relaxed) > stripes,
7110 "the pages are the ones that get read again, which is what makes the index count mean \
7111 something"
7112 );
7113 fs::remove_file(path).expect("remove scratch file");
7114 }
7115
7116 #[test]
7125 fn a_worker_per_stripe_reads_its_page_once_when_the_cache_was_told_to_expect_it() {
7126 let workers = CACHED_STRIPES_PER_COLUMN + 4;
7127 let path = path("stripe-per-worker");
7128 let mut writer =
7129 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
7130 .expect("new file");
7131 for part in 0..STRIPE_PARTS * workers {
7132 let chunk = Chunk::new(vec![
7133 Vector::from_values(LogicalType::Integer, &[Value::Integer(part as i32)])
7134 .expect("integers"),
7135 ])
7136 .expect("matching rows");
7137 writer.append(&chunk).expect("one part");
7138 }
7139 writer.finish().expect("commit");
7140
7141 let read = |told: bool| {
7142 let reader = Reader::open(&path).expect("reopen from disk");
7143 assert_eq!(reader.table().stripes().len(), workers, "a stripe per worker");
7144 if told {
7145 reader.keep_stripes(workers);
7146 }
7147 let barrier = std::sync::Barrier::new(workers);
7148 std::thread::scope(|scope| {
7149 for (worker, run) in reader.stripe_parts().into_iter().enumerate() {
7150 let reader = &reader;
7151 let barrier = &barrier;
7152 scope.spawn(move || {
7153 for part in run {
7154 barrier.wait();
7155 let chunk = reader.read(part, &[0]).expect("a part of my own stripe");
7156 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
7157 }
7158 assert!(worker < workers);
7159 });
7160 }
7161 });
7162 reader.pages.load(Atomic::Relaxed)
7163 };
7164
7165 assert_eq!(read(true), workers, "one page read per stripe and no more");
7166 assert!(read(false) > workers, "a cache that small is read again on every part");
7167 fs::remove_file(path).expect("remove scratch file");
7168 }
7169
7170 #[test]
7175 fn a_damaged_index_page_is_an_error() {
7176 let path = path("damaged-index");
7177 let mut writer =
7178 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
7179 .expect("new file");
7180 writer.append(&sample_ids()).expect("first part");
7181 writer.append(&sample_ids()).expect("second part");
7182 writer.finish().expect("commit");
7183
7184 let reader = Reader::open(&path).expect("valid directory");
7185 let index = reader.table.stripes[0].index;
7186 let mut byte = [0; 1];
7187 read_at(&reader.file, index.offset, &mut byte).expect("the first part length");
7188 let mut file = OpenOptions::new().write(true).open(&path).expect("open index page");
7189 file.seek(SeekFrom::Start(index.offset)).expect("index start");
7190 file.write_all(&[!byte[0]]).expect("damage the first part length");
7191 let error = reader.read(1, &[0]).expect_err("a damaged index must not be used");
7192 assert!(error.message().contains("index page section checksum differs"), "{error}");
7193 fs::remove_file(path).expect("remove scratch file");
7194 }
7195
7196 #[test]
7203 fn every_integer_width_round_trips_through_a_page() {
7204 let path = path("integer-widths");
7205 let columns = [
7206 (LogicalType::TinyInt, vec![Value::TinyInt(i8::MIN), Value::TinyInt(i8::MAX)]),
7207 (LogicalType::UTinyInt, vec![Value::UTinyInt(0), Value::UTinyInt(u8::MAX)]),
7208 (LogicalType::SmallInt, vec![Value::SmallInt(i16::MIN), Value::SmallInt(i16::MAX)]),
7209 (LogicalType::USmallInt, vec![Value::USmallInt(0), Value::USmallInt(u16::MAX)]),
7210 (LogicalType::Integer, vec![Value::Integer(i32::MIN), Value::Integer(i32::MAX)]),
7211 (LogicalType::UInteger, vec![Value::UInteger(0), Value::UInteger(u32::MAX)]),
7212 (LogicalType::BigInt, vec![Value::BigInt(i64::MIN), Value::BigInt(i64::MAX)]),
7213 (LogicalType::UBigInt, vec![Value::UBigInt(0), Value::UBigInt(u64::MAX)]),
7214 ];
7215 let fields = columns
7216 .iter()
7217 .enumerate()
7218 .map(|(at, (ty, _))| Field::required(format!("c{at}"), ty.clone()))
7219 .collect::<Vec<_>>();
7220 let vectors = columns
7221 .iter()
7222 .map(|(ty, values)| Vector::from_values(ty.clone(), values).expect("a vector"))
7223 .collect::<Vec<_>>();
7224 let mut writer = Writer::create(&path, "widths", fields).expect("new file");
7225 writer.append(&Chunk::new(vectors).expect("matching rows")).expect("one stripe");
7226 writer.finish().expect("commit");
7227
7228 let reader = Reader::open(&path).expect("reopen from disk");
7229 let wanted = (0..columns.len()).collect::<Vec<_>>();
7230 let read = reader.read(0, &wanted).expect("every column");
7231 assert_eq!(read.len(), 2);
7232 for (at, (ty, values)) in columns.iter().enumerate() {
7234 assert_eq!(read.value_at(0, at), values[0], "the low end of {ty}");
7235 assert_eq!(read.value_at(1, at), values[1], "the high end of {ty}");
7236 }
7237 fs::remove_file(path).expect("remove scratch file");
7238 }
7239
7240 #[test]
7241 fn numeric_frequency_candidates_keep_bounded_row_ordinals() {
7242 let path = path("frequency-ordinals");
7243 let mut writer =
7244 Writer::create(&path, "items", vec![Field::required("id", LogicalType::BigInt)])
7245 .expect("new file");
7246 let mut values = Vec::new();
7247 for leader in 0..10_i64 {
7248 values.extend(std::iter::repeat_n(leader, 100));
7249 }
7250 values.extend(1_000_i64..41_000);
7251 for part in values.chunks(1_024) {
7252 let vector = Vector::flat(LogicalType::BigInt, Data::Int64(part.to_vec().into()))
7253 .expect("big integers");
7254 writer.append(&Chunk::new(vec![vector]).expect("one column")).expect("one stripe");
7255 }
7256 writer.finish().expect("commit");
7257
7258 let reader = Reader::open(&path).expect("reopen from disk");
7259 let occurrences =
7260 reader.frequency_occurrences(0).expect("valid metadata").expect("bounded ordinals");
7261 assert!(occurrences.omitted_max < 100);
7262 assert!(occurrences.ordinals.len() <= FREQUENCY_ORDINALS);
7263 assert!(occurrences.ordinals.windows(2).all(|pair| pair[0] < pair[1]));
7264 assert_eq!(&occurrences.ordinals[..1_000], &(0_u64..1_000).collect::<Vec<_>>());
7265 fs::remove_file(path).expect("remove scratch file");
7266 }
7267
7268 #[test]
7274 fn a_file_from_another_format_says_which_format_it_is() {
7275 let older = path("older-format");
7276 let mut writer =
7277 Writer::create(&older, "items", vec![Field::new("id", LogicalType::Integer)])
7278 .expect("new file");
7279 let chunk = Chunk::new(vec![
7280 Vector::flat(LogicalType::Integer, Data::Int32(vec![1, 2, 3].into()))
7281 .expect("integers"),
7282 ])
7283 .expect("chunk");
7284 writer.append(&chunk).expect("page written");
7285 writer.finish().expect("commit");
7286
7287 let mut file = OpenOptions::new().write(true).open(&older).expect("open for the header");
7288 file.seek(SeekFrom::Start(8)).expect("the version follows the magic");
7289 file.write_all(&(FORMAT - 1).to_le_bytes()).expect("write an older version");
7290 drop(file);
7291 let complaint = Reader::open(&older).expect_err("an older format is refused").to_string();
7292 assert!(complaint.contains(&format!("format {}", FORMAT - 1)), "{complaint}");
7293 assert!(complaint.contains(&format!("format {FORMAT}")), "{complaint}");
7294
7295 let mut file = OpenOptions::new().write(true).open(&older).expect("open for the header");
7296 file.seek(SeekFrom::Start(0)).expect("the magic is first");
7297 file.write_all(b"NOTRUDB!").expect("write another engine's magic");
7298 drop(file);
7299 let complaint = Reader::open(&older).expect_err("a foreign file is refused").to_string();
7300 assert!(complaint.contains("magic"), "{complaint}");
7301 assert!(!complaint.contains("format"), "a version has nothing to do with it: {complaint}");
7302 fs::remove_file(older).expect("remove scratch file");
7303 }
7304
7305 #[test]
7306 fn an_unfinished_or_damaged_file_does_not_answer_with_partial_rows() {
7307 let unfinished = path("unfinished");
7308 let mut writer =
7309 Writer::create(&unfinished, "items", vec![Field::new("id", LogicalType::Integer)])
7310 .expect("new file");
7311 let chunk = Chunk::new(vec![
7312 Vector::flat(LogicalType::Integer, Data::Int32(vec![1, 2, 3].into()))
7313 .expect("integers"),
7314 ])
7315 .expect("chunk");
7316 writer.append(&chunk).expect("page written");
7317 drop(writer);
7318 assert!(Reader::open(&unfinished).is_err(), "no directory was committed");
7319 fs::remove_file(unfinished).expect("remove scratch file");
7320
7321 let damaged = path("damaged");
7322 let mut writer =
7323 Writer::create(&damaged, "items", vec![Field::new("id", LogicalType::Integer)])
7324 .expect("new file");
7325 writer.append(&chunk).expect("page written");
7326 writer.finish().expect("commit");
7327 let reader = Reader::open(&damaged).expect("valid directory");
7328 let mut file =
7329 OpenOptions::new().write(true).open(&damaged).expect("open for a damaged page");
7330 file.seek(SeekFrom::Start(HEADER + 1)).expect("inside first page");
7331 file.write_all(&[255]).expect("damage one byte");
7332 assert!(reader.read(0, &[0]).is_err(), "page checksum rejects corruption");
7333 fs::remove_file(damaged).expect("remove scratch file");
7334 }
7335
7336 #[test]
7337 fn damaged_lazy_dictionary_payload_is_an_error() {
7338 let path = path("damaged-dictionary");
7339 let mut writer = Writer::create(
7340 &path,
7341 "items",
7342 vec![
7343 Field::required("id", LogicalType::Integer),
7344 Field::new("text", LogicalType::Varchar),
7345 ],
7346 )
7347 .expect("new file");
7348 writer.append(&sample()).expect("stripe written");
7349 writer.finish().expect("commit");
7350
7351 let reader = Reader::open(&path).expect("valid directory");
7352 let dictionary = reader.table.dictionaries[1].expect("string dictionary page");
7353 let mut header = [0; DICTIONARY_HEADER];
7356 read_at(&reader.file, dictionary.offset, &mut header).expect("dictionary header");
7357 let index_len = dictionary_index_len(&header);
7358 let rank_len = last_rank_end(&reader.file, dictionary.offset, &header);
7359 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
7360 file.seek(SeekFrom::Start(dictionary.offset + index_len + rank_len))
7361 .expect("inside dictionary payload");
7362 file.write_all(&[255]).expect("damage dictionary payload");
7363
7364 let chunk = reader.read(0, &[1]).expect("code page and dictionary index remain valid");
7365 let error =
7366 chunk.validate_external().expect_err("payload corruption must reach the caller");
7367 assert!(error.message().contains("payload checksum differs"), "{error}");
7368 fs::remove_file(path).expect("remove scratch file");
7369 }
7370
7371 #[test]
7381 fn a_column_of_all_different_values_is_written_without_a_dictionary() {
7382 let path = path("dictionary-decide");
7383 let rows = 20_000;
7384 let unique =
7386 |row: usize| format!("{row:09} a value that appears exactly once in the table");
7387 let repeated = |row: usize| unique(row / 40);
7389 let mut writer = Writer::create(
7390 &path,
7391 "items",
7392 vec![
7393 Field::required("unique", LogicalType::Varchar),
7394 Field::required("repeated", LogicalType::Varchar),
7395 ],
7396 )
7397 .expect("new file");
7398 for part in (0..rows).step_by(1_000) {
7399 let span = part..(part + 1_000).min(rows);
7400 let left = span.clone().map(|row| Value::Varchar(unique(row))).collect::<Vec<_>>();
7401 let right = span.map(|row| Value::Varchar(repeated(row))).collect::<Vec<_>>();
7402 writer
7403 .append(
7404 &Chunk::new(vec![
7405 Vector::from_values(LogicalType::Varchar, &left).expect("strings"),
7406 Vector::from_values(LogicalType::Varchar, &right).expect("strings"),
7407 ])
7408 .expect("two columns"),
7409 )
7410 .expect("a part");
7411 }
7412 writer.finish().expect("commit");
7413
7414 let reader = Reader::open(&path).expect("reopen from disk");
7415 assert!(
7416 reader.table.dictionaries[0].is_none(),
7417 "a column with no repeats has nothing to say twice"
7418 );
7419 assert!(
7420 reader.table.dictionaries[1].is_some(),
7421 "a column whose values come round again keeps its dictionary"
7422 );
7423 let mut first = 0;
7424 for part in 0..reader.parts() {
7425 let chunk = reader.read(part, &[0, 1]).expect("a part");
7426 for row in 0..chunk.len() {
7427 assert_eq!(chunk.value_at(row, 0), Value::Varchar(unique(first + row)));
7428 assert_eq!(chunk.value_at(row, 1), Value::Varchar(repeated(first + row)));
7429 }
7430 first += chunk.len();
7431 }
7432 assert_eq!(first, rows, "every row was read back");
7433 let raw = (0..rows).map(|row| unique(row).len()).sum::<usize>();
7434 let size = fs::metadata(&path).expect("the file is there").len() as usize;
7435 assert!(size < raw, "a column without a dictionary is still encoded: {size} against {raw}");
7436 fs::remove_file(path).expect("remove scratch file");
7437 }
7438
7439 #[test]
7452 fn a_dictionary_over_many_blocks_checks_every_block_of_it() {
7453 let path = path("dictionary-blocks");
7454 let value = |row: usize| {
7455 let row = row.saturating_sub(8_000);
7456 format!("{row:07} a value long enough to be worth a payload block")
7457 };
7458 let parts = 40;
7459 let per_part = 1000;
7460 let mut writer =
7461 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
7462 .expect("new file");
7463 for part in 0..parts {
7464 let values = (0..per_part)
7465 .map(|row| Value::Varchar(value(part * per_part + row)))
7466 .collect::<Vec<_>>();
7467 let chunk = Chunk::new(vec![
7468 Vector::from_values(LogicalType::Varchar, &values).expect("strings"),
7469 ])
7470 .expect("matching rows");
7471 writer.append(&chunk).expect("a part");
7472 }
7473 writer.finish().expect("commit");
7474
7475 let reader = Reader::open(&path).expect("reopen from disk");
7476 let dictionary = reader.table.dictionaries[0].expect("string dictionary page");
7477 assert!(
7478 parts * per_part > TEXT_PAYLOAD_VALUES * 4,
7479 "the dictionary has to be several blocks for this to be testing anything"
7480 );
7481 for part in [0, parts - 1] {
7482 let chunk = reader.read(part, &[0]).expect("a part");
7483 chunk.validate_external().expect("every payload block checks out");
7484 assert_eq!(chunk.value_at(0, 0), Value::Varchar(value(part * per_part)));
7485 }
7486
7487 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
7488 file.seek(SeekFrom::Start(dictionary.offset + u64::from(dictionary.length) - 4))
7489 .expect("the last bytes of the page are payload");
7490 file.write_all(&[255]).expect("damage the last payload block");
7491 let reader = Reader::open(&path).expect("the directory and the index are untouched");
7492 let chunk = reader.read(parts - 1, &[0]).expect("the code page remains valid");
7493 let error = chunk.validate_external().expect_err("the damage must reach the caller");
7494 assert!(error.message().contains("payload checksum differs"), "{error}");
7495 fs::remove_file(path).expect("remove scratch file");
7496 }
7497
7498 #[test]
7512 fn values_of_different_lengths_read_back_out_of_packed_offsets() {
7513 let path = path("dictionary-offsets");
7514 let value = |row: usize| {
7515 let row = row % 5_000;
7516 if row % 511 == 3 { String::new() } else { "x".repeat(row % 97) + &format!("{row:05}") }
7517 };
7518 let rows = 6_000;
7519 let mut writer =
7520 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
7521 .expect("new file");
7522 let values = (0..rows).map(|row| Value::Varchar(value(row))).collect::<Vec<_>>();
7523 for part in values.chunks(1_000) {
7524 let chunk =
7525 Chunk::new(vec![Vector::from_values(LogicalType::Varchar, part).expect("strings")])
7526 .expect("matching rows");
7527 writer.append(&chunk).expect("a part");
7528 }
7529 writer.finish().expect("commit");
7530
7531 let reader = Reader::open(&path).expect("reopen from disk");
7532 assert!(
7533 rows > TEXT_PAYLOAD_VALUES * 4,
7534 "the dictionary has to be several blocks for this to be testing anything"
7535 );
7536 for part in 0..rows / 1_000 {
7537 let chunk = reader.read(part, &[0]).expect("a part");
7538 for row in 0..1_000 {
7539 let row = part * 1_000 + row;
7540 assert_eq!(
7541 chunk.value_at(row % 1_000, 0),
7542 Value::Varchar(value(row)),
7543 "value {row}"
7544 );
7545 }
7546 }
7547 fs::remove_file(path).expect("remove scratch file");
7548 }
7549
7550 #[test]
7562 fn a_global_dictionary_is_opened_once_however_many_workers_ask_at_once() {
7563 let path = path("dictionary-once");
7564 let parts = 8;
7565 let per_part = 500;
7566 let value =
7567 |row: usize| format!("{row:07} a value long enough to be worth a payload block");
7568 let mut writer =
7569 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
7570 .expect("new file");
7571 for part in 0..parts {
7572 let values = (0..per_part)
7573 .map(|row| Value::Varchar(value(part * per_part + row)))
7574 .collect::<Vec<_>>();
7575 let chunk = Chunk::new(vec![
7576 Vector::from_values(LogicalType::Varchar, &values).expect("strings"),
7577 ])
7578 .expect("matching rows");
7579 writer.append(&chunk).expect("a part");
7580 }
7581 writer.finish().expect("commit");
7582
7583 let reader = Reader::open(&path).expect("reopen from disk");
7584 assert!(reader.table.dictionaries[0].is_some(), "the column has to have one to share");
7585 assert_eq!(reader.reads().dictionaries, 0, "opening the file does not open a dictionary");
7586
7587 let workers = 16;
7588 let gate = std::sync::Barrier::new(workers);
7589 std::thread::scope(|scope| {
7590 for worker in 0..workers {
7591 let reader = reader.clone();
7592 let gate = &gate;
7593 scope.spawn(move || {
7594 gate.wait();
7595 let chunk = reader.read(worker % parts, &[0]).expect("a part");
7596 assert_eq!(
7597 chunk.value_at(0, 0),
7598 Value::Varchar(value((worker % parts) * per_part))
7599 );
7600 });
7601 }
7602 });
7603
7604 assert_eq!(reader.reads().dictionaries, 1, "sixteen workers, one dictionary, one open");
7605 fs::remove_file(path).expect("remove scratch file");
7606 }
7607
7608 #[test]
7613 fn a_damaged_sorted_order_is_an_error() {
7614 let path = path("damaged-order");
7615 let mut writer = Writer::create(
7616 &path,
7617 "items",
7618 vec![
7619 Field::required("id", LogicalType::Integer),
7620 Field::new("text", LogicalType::Varchar),
7621 ],
7622 )
7623 .expect("new file");
7624 writer.append(&sample()).expect("stripe written");
7625 writer.finish().expect("commit");
7626
7627 let reader = Reader::open(&path).expect("valid directory");
7628 let page = reader.table.dictionaries[1].expect("string dictionary page");
7629 let mut header = [0; DICTIONARY_HEADER];
7630 read_at(&reader.file, page.offset, &mut header).expect("dictionary header");
7631 let index_len = dictionary_index_len(&header);
7632 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
7633 file.seek(SeekFrom::Start(page.offset + index_len)).expect("the first head");
7634 file.write_all(&[255]).expect("damage the order");
7635
7636 let dictionary = reader.dictionary(1).expect("read").expect("a string column has one");
7637 let error = dictionary.compare_rank(0, b"anything").expect_err("a damaged order is caught");
7638 assert!(error.message().contains("rank checksum differs"), "{error}");
7639 fs::remove_file(path).expect("remove scratch file");
7640 }
7641
7642 #[test]
7646 fn a_global_dictionary_carries_the_sorted_order_of_its_values() {
7647 let spellings = ["overlong1z", "b", "", "overlong1a", "overlong", "ab", "a", "overlong1"];
7650 let path = path("dictionary-order");
7651 let mut writer =
7652 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7653 .expect("new file");
7654 writer
7655 .append(
7656 &Chunk::new(vec![
7657 Vector::from_values(
7658 LogicalType::Varchar,
7659 &spellings.map(|text| Value::Varchar(text.into())),
7660 )
7661 .expect("strings"),
7662 ])
7663 .expect("one column"),
7664 )
7665 .expect("stripe written");
7666 writer.finish().expect("commit");
7667
7668 let reader = Reader::open(&path).expect("valid directory");
7669 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
7670 let count = dictionary.ranks().expect("a v10 file stores one");
7671 assert_eq!(count, spellings.len(), "every distinct value has a rank");
7672 let order = (0..count)
7673 .map(|rank| dictionary.code_at_rank(rank).expect("a code"))
7674 .collect::<Vec<_>>();
7675 let mut seen = order.clone();
7676 seen.sort_unstable();
7677 assert_eq!(seen, (0..spellings.len() as u32).collect::<Vec<_>>(), "a permutation of codes");
7678
7679 let ranked = order
7680 .iter()
7681 .map(|&code| {
7682 dictionary.try_bytes_at(code as usize).expect("read").expect("a value").to_vec()
7683 })
7684 .collect::<Vec<_>>();
7685 let mut expected = spellings.map(|text| text.as_bytes().to_vec()).to_vec();
7686 expected.sort();
7687 assert_eq!(ranked, expected, "rank order is value order");
7688
7689 for (rank, value) in expected.iter().enumerate() {
7692 assert_eq!(
7693 dictionary.compare_rank(rank, value).expect("compare"),
7694 Ordering::Equal,
7695 "rank {rank} is its own value"
7696 );
7697 if rank > 0 {
7698 assert_eq!(
7699 dictionary.compare_rank(rank - 1, value).expect("compare"),
7700 Ordering::Less,
7701 "rank {rank} follows the one before it"
7702 );
7703 }
7704 }
7705 fs::remove_file(path).expect("remove scratch file");
7706 }
7707
7708 #[test]
7716 fn a_dictionary_sweep_reads_every_value_and_keeps_it_under_the_budget() {
7717 let path = path("dictionary-sweep");
7718 let spellings = (0..2_500)
7721 .map(|index| Value::Varchar(format!("value {index:08} {}", "x".repeat(index % 40))))
7722 .collect::<Vec<_>>();
7723 let mut writer =
7724 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7725 .expect("new file");
7726 for part in spellings.chunks(1_024) {
7729 writer
7730 .append(
7731 &Chunk::new(vec![
7732 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
7733 ])
7734 .expect("one column"),
7735 )
7736 .expect("stripe written");
7737 }
7738 writer.finish().expect("commit");
7739
7740 let reader = Reader::open(&path).expect("valid directory");
7741 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
7742 assert_eq!(dictionary.len(), spellings.len(), "every value is distinct");
7743
7744 let resting = dictionary.footprint();
7745 let mut swept: Vec<Vec<u8>> = Vec::new();
7746 let mut at = 0;
7747 let mut calls = 0;
7748 while at < dictionary.len() {
7749 let stopped = dictionary
7750 .sweep_text(at, dictionary.len(), &mut |index: usize, text: &[u8]| {
7751 assert_eq!(index, swept.len(), "a sweep hands its values over in order");
7752 swept.push(text.to_vec());
7753 Ok(())
7754 })
7755 .expect("a sweep reads");
7756 assert!(stopped > at, "a sweep moves");
7757 at = stopped;
7758 calls += 1;
7759 }
7760 assert_eq!(calls, 3, "a sweep hands over one block at a time");
7761 let after = dictionary.footprint();
7762 assert!(after > resting, "a sweep under the budget keeps what it decoded");
7763
7764 let read = (0..dictionary.len())
7765 .map(|code| dictionary.try_bytes_at(code).expect("read").expect("a value").to_vec())
7766 .collect::<Vec<_>>();
7767 assert_eq!(swept, read, "a sweep answers what a point read answers");
7768 assert_eq!(dictionary.footprint(), after, "a point read of a kept block decodes nothing");
7769 fs::remove_file(path).expect("remove scratch file");
7770 }
7771
7772 #[test]
7783 fn a_sweep_over_a_block_with_a_short_second_run_reads_what_a_point_read_reads() {
7784 let path = path("dictionary-sweep-short-run");
7785 let spellings = (0..2_800)
7786 .map(|index| Value::Varchar(format!("value {index:08} {}", "x".repeat(index % 40))))
7787 .collect::<Vec<_>>();
7788 let mut writer =
7789 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7790 .expect("new file");
7791 for part in spellings.chunks(1_024) {
7792 writer
7793 .append(
7794 &Chunk::new(vec![
7795 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
7796 ])
7797 .expect("one column"),
7798 )
7799 .expect("stripe written");
7800 }
7801 writer.finish().expect("commit");
7802
7803 let reader = Reader::open(&path).expect("valid directory");
7804 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
7805 assert_eq!(dictionary.len(), spellings.len(), "every value is distinct");
7806 let last = dictionary.len() % TEXT_PAYLOAD_VALUES;
7807 assert!(last > TEXT_OFFSET_RUN, "the last block has to reach into a second run of offsets");
7808 assert!(last < TEXT_PAYLOAD_VALUES, "and that second run has to be short of a whole one");
7809
7810 let mut swept: Vec<Vec<u8>> = Vec::new();
7811 let mut at = 0;
7812 while at < dictionary.len() {
7813 let stopped = dictionary
7814 .sweep_text(at, dictionary.len(), &mut |index: usize, text: &[u8]| {
7815 assert_eq!(index, swept.len(), "a sweep hands its values over in order");
7816 swept.push(text.to_vec());
7817 Ok(())
7818 })
7819 .expect("a sweep reads");
7820 assert!(stopped > at, "a sweep moves");
7821 at = stopped;
7822 }
7823 let read = (0..dictionary.len())
7824 .map(|code| dictionary.try_bytes_at(code).expect("read").expect("a value").to_vec())
7825 .collect::<Vec<_>>();
7826 assert_eq!(swept, read, "a sweep answers what a point read answers");
7827 fs::remove_file(path).expect("remove scratch file");
7828 }
7829
7830 #[test]
7840 fn narrowing_a_page_takes_what_fits_and_refuses_what_does_not() {
7841 assert_eq!(fit::<i8>(&[]).expect("an empty page fits anything"), Vec::<i8>::new());
7842 assert_eq!(fit::<i8>(&[-128, 0, 127]).expect("the edges fit"), vec![-128_i8, 0, 127]);
7843 fit::<i8>(&[128]).expect_err("one past the top does not fit");
7844 fit::<i8>(&[-129]).expect_err("one past the bottom does not fit");
7845 assert_eq!(fit::<u8>(&[0, 255]).expect("the edges fit"), vec![0_u8, 255]);
7846 fit::<u8>(&[256]).expect_err("one past the top does not fit");
7847 fit::<u8>(&[-1]).expect_err("a negative does not fit an unsigned page");
7848 assert_eq!(
7849 fit::<i16>(&[-32_768, 0, 32_767]).expect("the edges fit"),
7850 vec![-32_768_i16, 0, 32_767]
7851 );
7852 fit::<i16>(&[32_768]).expect_err("one past the top does not fit");
7853 fit::<i16>(&[-32_769]).expect_err("one past the bottom does not fit");
7854 assert_eq!(fit::<u16>(&[0, 65_535]).expect("the edges fit"), vec![0_u16, 65_535]);
7855 fit::<u16>(&[65_536]).expect_err("one past the top does not fit");
7856 fit::<u16>(&[-1]).expect_err("a negative does not fit an unsigned page");
7857 assert_eq!(
7858 fit::<i32>(&[i64::from(i32::MIN), 0, i64::from(i32::MAX)]).expect("the edges fit"),
7859 vec![i32::MIN, 0, i32::MAX]
7860 );
7861 fit::<i32>(&[i64::from(i32::MAX) + 1]).expect_err("one past the top does not fit");
7862 fit::<i32>(&[i64::from(i32::MIN) - 1]).expect_err("one past the bottom does not fit");
7863 assert_eq!(
7864 fit::<u32>(&[0, 4_294_967_295]).expect("the edges fit"),
7865 vec![0_u32, 4_294_967_295]
7866 );
7867 fit::<u32>(&[4_294_967_296]).expect_err("one past the top does not fit");
7868 fit::<u32>(&[-1]).expect_err("a negative does not fit an unsigned page");
7869
7870 fit::<i8>(&[0, 1, 2, 128, 3]).expect_err("one bad value spoils the page");
7873 }
7874
7875 #[test]
7882 fn the_residue_agrees_with_a_checked_conversion_everywhere() {
7883 for value in -70_000_i64..70_000 {
7884 assert_eq!(fit::<i8>(&[value]).is_ok(), i8::try_from(value).is_ok(), "{value} as i8");
7885 assert_eq!(fit::<u8>(&[value]).is_ok(), u8::try_from(value).is_ok(), "{value} as u8");
7886 assert_eq!(fit::<i16>(&[value]).is_ok(), i16::try_from(value).is_ok(), "{value} i16");
7887 assert_eq!(fit::<u16>(&[value]).is_ok(), u16::try_from(value).is_ok(), "{value} u16");
7888 }
7889 let wide = [i64::MIN, i64::MIN + 1, i64::from(i32::MIN), 0, i64::from(u32::MAX), i64::MAX];
7890 for edge in wide {
7891 for step in -2_i64..=2 {
7892 let value = edge.saturating_add(step);
7893 assert_eq!(
7894 fit::<i32>(&[value]).is_ok(),
7895 i32::try_from(value).is_ok(),
7896 "{value} as i32"
7897 );
7898 assert_eq!(
7899 fit::<u32>(&[value]).is_ok(),
7900 u32::try_from(value).is_ok(),
7901 "{value} as u32"
7902 );
7903 }
7904 }
7905 }
7906
7907 #[test]
7915 fn a_dictionary_at_its_budget_sweeps_without_keeping() {
7916 let path = path("dictionary-budget");
7917 let spellings = (0..2_500)
7918 .map(|index| Value::Varchar(format!("value {index:08} {}", "y".repeat(index % 40))))
7919 .collect::<Vec<_>>();
7920 let mut writer =
7921 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7922 .expect("new file");
7923 for part in spellings.chunks(1_024) {
7924 writer
7925 .append(
7926 &Chunk::new(vec![
7927 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
7928 ])
7929 .expect("one column"),
7930 )
7931 .expect("stripe written");
7932 }
7933 writer.finish().expect("commit");
7934
7935 let reader = Reader::open(&path).expect("valid directory");
7936 let page = reader.table.dictionaries[0].expect("a string column has one");
7937 let file = Arc::clone(&reader.file);
7938 let starved = open_global_dictionary(file, page, &LogicalType::Varchar, 0)
7939 .expect("a dictionary opens whatever it may keep");
7940
7941 let resting = starved.footprint();
7942 let mut swept: Vec<Vec<u8>> = Vec::new();
7943 let mut at = 0;
7944 while at < starved.len() {
7945 at = starved
7946 .sweep_text(at, starved.len(), &mut |_index: usize, text: &[u8]| {
7947 swept.push(text.to_vec());
7948 Ok(())
7949 })
7950 .expect("a sweep reads");
7951 }
7952 assert_eq!(swept.len(), spellings.len(), "a starved sweep still reads every value");
7953 assert_eq!(starved.footprint(), resting, "and keeps no block it decoded");
7954
7955 let generous = reader.dictionary(0).expect("read").expect("a string column has one");
7956 let read = (0..generous.len())
7957 .map(|code| generous.try_bytes_at(code).expect("read").expect("a value").to_vec())
7958 .collect::<Vec<_>>();
7959 assert_eq!(swept, read, "a starved sweep answers what a point read answers");
7960 fs::remove_file(path).expect("remove scratch file");
7961 }
7962
7963 #[test]
7964 fn damaged_membership_cannot_skip_a_string_page() {
7965 let path = path("damaged-membership");
7966 let mut writer = Writer::create(
7967 &path,
7968 "items",
7969 vec![
7970 Field::required("id", LogicalType::Integer),
7971 Field::new("text", LogicalType::Varchar),
7972 ],
7973 )
7974 .expect("new file");
7975 writer.append(&sample()).expect("stripe written");
7976 writer.finish().expect("commit");
7977
7978 let reader = Reader::open(&path).expect("valid directory");
7979 let membership = reader.table.stripes[0].memberships[1].expect("string membership");
7980 let mut file = OpenOptions::new().write(true).open(&path).expect("open membership page");
7981 file.seek(SeekFrom::Start(membership.offset)).expect("membership start");
7982 file.write_all(&[255]).expect("damage membership");
7983 let error = reader.skips_codes(0, 1, &[3]).expect_err("corruption must not skip rows");
7984 assert!(error.message().contains("membership page checksum differs"), "{error}");
7985 fs::remove_file(path).expect("remove scratch file");
7986 }
7987
7988 #[test]
7989 fn membership_delta_stream_is_sorted_exact_and_bounded() {
7990 let unique = unique_codes(&[900, 4, 4, 72, 9, u32::MAX]);
7991 assert_eq!(unique, [4, 9, 72, 900, u32::MAX]);
7992 let encoded = encode_membership(&unique);
7993 assert_eq!(
7994 decode_membership(&encoded).expect("valid membership"),
7995 [4, 9, 72, 900, u32::MAX]
7996 );
7997 let merged = merged_codes(vec![vec![4, 900], vec![9, 900, u32::MAX], vec![72]]);
8000 assert_eq!(merged, [4, 9, 72, 900, u32::MAX]);
8001 assert_eq!(
8002 decode_membership(&encode_membership(&merged)).expect("valid membership"),
8003 unique
8004 );
8005 assert!(decode_membership(&[1, 0x80]).is_err(), "a truncated varint is invalid");
8006 assert!(
8007 decode_membership(&[1, 0xff, 0xff, 0xff, 0xff, 0x10]).is_err(),
8008 "a value past u32 is invalid"
8009 );
8010 }
8011
8012 #[test]
8013 fn a_global_dictionary_may_be_larger_than_one_column_page() {
8014 let dictionary = Page {
8015 offset: HEADER,
8016 length: u32::try_from(MAX_PAGE + 1).expect("the page bound fits on disk"),
8017 hash: 0,
8018 };
8019 let table = Table {
8020 name: "items".to_owned(),
8021 fields: vec![Field::new("text", LogicalType::Varchar)],
8022 stripes: Vec::new(),
8023 rows: 0,
8024 dictionaries: vec![Some(dictionary)],
8025 distincts: vec![None],
8026 frequencies: vec![None],
8027 clustering: None,
8028 };
8029 let directory = encode_directory(&table).expect("directory");
8030 let file_size = dictionary.offset + u64::from(dictionary.length) + 1;
8031
8032 let decoded = decode_directory(&directory, file_size).expect("large lazy dictionary");
8033 assert_eq!(decoded.dictionaries[0].expect("dictionary").length, dictionary.length);
8034 }
8035
8036 #[test]
8037 fn a_column_with_one_value_everywhere_costs_almost_nothing_a_row() {
8038 let path = path("constant-codes");
8039 let mut writer =
8040 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
8041 .expect("new file");
8042 let empty = vec![Value::Varchar(String::new()); 1024];
8043 for _ in 0..4 {
8044 let column = Vector::from_values(LogicalType::Varchar, &empty).expect("strings");
8045 writer.append(&Chunk::new(vec![column]).expect("one column")).expect("a part");
8046 }
8047 writer.finish().expect("commit");
8048
8049 let reader = Reader::open(&path).expect("valid directory");
8050 let pages = reader.layout().columns.first().expect("one column").pages;
8051 assert!(pages < 256, "{pages} bytes of pages for 4,096 rows of one value");
8055 let read = reader.read(3, &[0]).expect("the last part back");
8056 assert_eq!(read.value_at(0, 0), Value::Varchar(String::new()));
8057 assert_eq!(read.value_at(1023, 0), Value::Varchar(String::new()));
8058 fs::remove_file(path).expect("remove scratch file");
8059 }
8060
8061 #[test]
8062 fn a_cascade_value_too_wide_for_its_column_is_refused_rather_than_cut() {
8063 let over = vec![i64::from(i32::MAX) + 1];
8066 let error = narrowed(&LogicalType::Integer, over).expect_err("a page that disagrees");
8067 assert!(format!("{error}").contains("not of its type"), "{error}");
8068 assert!(narrowed(&LogicalType::BigInt, vec![i64::MIN]).is_ok(), "bigint holds all of i64");
8069 assert!(narrowed(&LogicalType::Varchar, vec![0]).is_err(), "strings are not integers");
8070 }
8071
8072 #[test]
8073 fn a_code_stream_the_cascade_cannot_shrink_is_left_alone() {
8074 let mut state: u32 = 0x9e37_79b9;
8078 let spread: Vec<u32> = (0..1024)
8079 .map(|_| {
8080 state ^= state << 13;
8081 state ^= state >> 17;
8082 state ^= state << 5;
8083 state
8084 })
8085 .collect();
8086 assert_eq!(encoded_codes(&spread).expect("no failure"), None);
8087 let near: Vec<u32> = (0..1024).collect();
8088 let coded = encoded_codes(&near).expect("no failure").expect("counting up is packable");
8089 assert!(coded.len() < near.len() * 4, "{} bytes for a run of 1,024", coded.len());
8090 }
8091
8092 #[test]
8098 fn two_writes_of_the_same_rows_give_the_same_bytes() {
8099 fn written(path: &PathBuf) {
8100 let fields = (0..40)
8101 .map(|column| {
8102 let ty =
8103 if column % 4 == 0 { LogicalType::Varchar } else { LogicalType::BigInt };
8104 Field::new(format!("c{column}"), ty)
8105 })
8106 .collect::<Vec<_>>();
8107 let mut writer = Writer::create(path, "wide", fields).expect("new file");
8108 for part in 0..70_u64 {
8109 let columns = (0..40)
8110 .map(|column| {
8111 let values = (0..64_u64)
8112 .map(|row| {
8113 let seed = part.wrapping_mul(31).wrapping_add(row);
8114 if column % 4 == 0 {
8115 Value::Varchar(format!("v{}", seed % 17))
8116 } else {
8117 Value::BigInt(i64::try_from(seed % 97).expect("small"))
8118 }
8119 })
8120 .collect::<Vec<_>>();
8121 let ty = if column % 4 == 0 {
8122 LogicalType::Varchar
8123 } else {
8124 LogicalType::BigInt
8125 };
8126 Vector::from_values(ty, &values).expect("a column")
8127 })
8128 .collect::<Vec<_>>();
8129 writer.append(&Chunk::new(columns).expect("forty columns")).expect("a part");
8130 }
8131 writer.finish().expect("commit");
8132 }
8133
8134 let first = path("repeatable-one");
8135 let second = path("repeatable-two");
8136 written(&first);
8137 written(&second);
8138 let left = fs::read(&first).expect("the first file");
8139 let right = fs::read(&second).expect("the second file");
8140 assert_eq!(left.len(), right.len(), "two writes of the same rows differ in length");
8141 assert!(left == right, "two writes of the same rows differ in their bytes");
8142
8143 let reader = Reader::open(&first).expect("valid directory");
8146 assert_eq!(reader.table().rows(), 70 * 64);
8147 let read = reader.read(0, &[0, 1]).expect("the first part back");
8148 assert_eq!(read.value_at(0, 0), Value::Varchar("v0".to_owned()));
8149 assert_eq!(read.value_at(0, 1), Value::BigInt(0));
8150 fs::remove_file(first).expect("remove scratch file");
8151 fs::remove_file(second).expect("remove scratch file");
8152 }
8153
8154 fn three_tables(path: &PathBuf) {
8156 let writer = Writer::create(
8157 path,
8158 "region",
8159 vec![
8160 Field::new("r_key", LogicalType::Integer),
8161 Field::new("r_name", LogicalType::Varchar),
8162 ],
8163 )
8164 .expect("new file");
8165 let mut writer = writer;
8166 writer
8167 .append(
8168 &Chunk::new(vec![
8169 Vector::from_values(
8170 LogicalType::Integer,
8171 &[Value::Integer(0), Value::Integer(1)],
8172 )
8173 .expect("keys"),
8174 Vector::from_values(
8175 LogicalType::Varchar,
8176 &[Value::Varchar("AFRICA".to_owned()), Value::Varchar("ASIA".to_owned())],
8177 )
8178 .expect("names"),
8179 ])
8180 .expect("two columns"),
8181 )
8182 .expect("a part");
8183 let mut writer = writer
8184 .next("empty", vec![Field::new("nothing", LogicalType::BigInt)])
8185 .expect("a second table");
8186 writer
8187 .append(
8188 &Chunk::new(vec![
8189 Vector::from_values(LogicalType::BigInt, &[Value::BigInt(7)]).expect("a row"),
8190 ])
8191 .expect("one column"),
8192 )
8193 .expect("a part");
8194 let mut writer =
8195 writer.next("wide", vec![Field::new("n", LogicalType::BigInt)]).expect("a third table");
8196 for part in 0..70_i64 {
8197 let values = (0..64).map(|row| Value::BigInt(part * 64 + row)).collect::<Vec<_>>();
8198 writer
8199 .append(
8200 &Chunk::new(vec![
8201 Vector::from_values(LogicalType::BigInt, &values).expect("a column"),
8202 ])
8203 .expect("one column"),
8204 )
8205 .expect("a part");
8206 }
8207 writer.finish().expect("commit");
8208 }
8209
8210 #[test]
8211 fn three_tables_in_one_file_read_back_by_name() {
8212 let file = path("three-tables");
8213 three_tables(&file);
8214 let catalog = Catalog::open(&file).expect("a committed catalog");
8215 assert_eq!(catalog.names().collect::<Vec<_>>(), ["region", "empty", "wide"]);
8216
8217 let region = catalog.table("region").expect("the first table");
8218 assert_eq!(region.table().rows(), 2);
8219 assert_eq!(
8220 region.read(0, &[1]).expect("names").value_at(1, 0),
8221 Value::Varchar("ASIA".to_owned())
8222 );
8223
8224 let wide = catalog.table("wide").expect("the third table");
8225 assert_eq!(wide.table().rows(), 70 * 64);
8226 assert_eq!(wide.read(0, &[0]).expect("the first part").value_at(0, 0), Value::BigInt(0));
8227
8228 let empty = catalog.table("empty").expect("the second table");
8231 assert_eq!(empty.table().rows(), 1);
8232 assert_eq!(empty.read(0, &[0]).expect("the row").value_at(0, 0), Value::BigInt(7));
8233
8234 fs::remove_file(file).expect("remove scratch file");
8235 }
8236
8237 #[test]
8238 fn a_name_the_file_does_not_hold_is_an_error_rather_than_the_first_table() {
8239 let file = path("three-tables-missing");
8240 three_tables(&file);
8241 let catalog = Catalog::open(&file).expect("a committed catalog");
8242 let error = catalog.table("nation").expect_err("no such table");
8243 assert!(error.message().contains("nation"), "{}", error.message());
8244 fs::remove_file(file).expect("remove scratch file");
8245 }
8246
8247 #[test]
8248 fn a_file_of_three_tables_will_not_open_as_one() {
8249 let file = path("three-tables-unnamed");
8250 three_tables(&file);
8251 let error = Reader::open(&file).expect_err("more than one table");
8252 assert!(error.message().contains("more than one table"), "{}", error.message());
8253 fs::remove_file(file).expect("remove scratch file");
8254 }
8255
8256 #[test]
8258 fn decimals_of_every_storage_width_round_trip() {
8259 let file = path("decimals");
8260 let widths = [(4_u8, 2_u8), (9, 2), (18, 4), (38, 6)];
8261 let fields = widths
8262 .iter()
8263 .enumerate()
8264 .map(|(index, (width, scale))| {
8265 Field::new(
8266 format!("d{index}"),
8267 LogicalType::decimal(*width, *scale).expect("a decimal type"),
8268 )
8269 })
8270 .collect::<Vec<_>>();
8271 let mut writer = Writer::create(&file, "money", fields).expect("new file");
8272 let rows: [i128; 3] = [-1234, 0, 999];
8273 let columns = widths
8274 .iter()
8275 .map(|(width, scale)| {
8276 let values = rows
8277 .iter()
8278 .map(|unscaled| Value::Decimal {
8279 unscaled: *unscaled,
8280 width: *width,
8281 scale: *scale,
8282 })
8283 .collect::<Vec<_>>();
8284 Vector::from_values(
8285 LogicalType::decimal(*width, *scale).expect("a decimal type"),
8286 &values,
8287 )
8288 .expect("a decimal column")
8289 })
8290 .collect::<Vec<_>>();
8291 writer.append(&Chunk::new(columns).expect("four columns")).expect("a part");
8292 writer.finish().expect("commit");
8293
8294 let reader = Reader::open(&file).expect("a committed file");
8295 for (index, (width, scale)) in widths.iter().enumerate() {
8296 assert_eq!(
8297 reader.table().fields()[index].ty,
8298 LogicalType::decimal(*width, *scale).expect("a decimal type"),
8299 "column {index} came back as another type"
8300 );
8301 let column = reader.read(0, &[index]).expect("the column");
8302 for (row, unscaled) in rows.iter().enumerate() {
8303 assert_eq!(
8304 column.value_at(row, 0),
8305 Value::Decimal { unscaled: *unscaled, width: *width, scale: *scale },
8306 "column {index} row {row}"
8307 );
8308 }
8309 }
8310 fs::remove_file(file).expect("remove scratch file");
8311 }
8312
8313 #[test]
8314 fn two_tables_of_one_name_are_refused_before_anything_is_committed() {
8315 let file = path("two-of-a-name");
8316 let writer = Writer::create(&file, "t", vec![Field::new("a", LogicalType::BigInt)])
8317 .expect("new file");
8318 let error = writer
8319 .next("t", vec![Field::new("a", LogicalType::BigInt)])
8320 .expect_err("the same name twice");
8321 assert!(error.message().contains("same name"), "{}", error.message());
8322 fs::remove_file(file).expect("remove scratch file");
8323 }
8324
8325 #[test]
8326 fn opening_the_catalog_reads_no_table_directory() {
8327 let file = path("catalog-only");
8328 three_tables(&file);
8329 let catalog = Catalog::open(&file).expect("a committed catalog");
8330 assert_eq!(catalog.opening.reads, 2, "opening the catalog read more than the slot");
8333 assert_eq!(catalog.names().len(), 3);
8334 fs::remove_file(file).expect("remove scratch file");
8335 }
8336
8337 #[test]
8348 fn the_checksum_answers_what_it_has_always_answered() {
8349 let bytes: Vec<u8> =
8350 (0..1000_u32).map(|at| (at.wrapping_mul(31).wrapping_add(7) % 251) as u8).collect();
8351 for (length, expected) in [
8352 (0, 0xef46_db37_51d8_e999),
8353 (1, 0xa96c_7f0c_e858_bbb7),
8354 (3, 0x56e6_9576_32a4_87f9),
8355 (4, 0xc60d_15b1_e3ff_8f04),
8356 (5, 0x8088_1585_8624_dd4e),
8357 (7, 0xafbe_fc3d_6c6f_9a8e),
8358 (8, 0x3da5_c7aa_2696_83e0),
8359 (9, 0x465e_c429_b13c_3892),
8360 (15, 0xdee8_9d8a_065a_6233),
8361 (16, 0x1330_489a_7767_9c80),
8362 (31, 0x3391_303d_485e_846e),
8363 (32, 0x40b7_aff7_5d45_bbc8),
8364 (33, 0x4997_cae4_951c_17a5),
8365 (39, 0x5807_28fd_5c14_5739),
8366 (40, 0xf95c_f6f5_c08a_3d3b),
8367 (63, 0x2944_b4da_fc69_b206),
8368 (64, 0xbb76_f6ef_19bd_5a1b),
8369 (65, 0x814e_0c65_4a9f_d640),
8370 (127, 0x00de_aab1_31cf_f89b),
8371 (1000, 0x9e33_00c1_cde3_c58d),
8372 ] {
8373 assert_eq!(checksum(&bytes[..length]), expected, "the checksum of {length} bytes");
8374 }
8375 assert_eq!(checksum(b"the quick brown fox jumps over the lazy dog"), 0xed71_4233_c5a9_a792);
8376 }
8377 #[test]
8384 fn a_declared_order_comes_back_out_of_the_file() {
8385 let path = path("clustered");
8386 let shipped = vec![
8387 Field::new("key", LogicalType::BigInt),
8388 Field::new("line", LogicalType::Integer),
8389 Field::new("shipdate", LogicalType::Date),
8390 ];
8391 let plain = vec![Field::new("a", LogicalType::Integer)];
8392 let stage_zero = Clustering::new(vec![2, 0, 1], Width::Month, &shipped).expect("valid");
8393
8394 let mut writer = Writer::create(&path, "lineitem", shipped)
8395 .expect("new file")
8396 .declare(stage_zero.clone())
8397 .expect("the columns are the table's");
8398 let column = |ty: LogicalType, values: &[Value]| {
8399 Vector::from_values(ty, values).expect("the values match the type")
8400 };
8401 writer
8402 .append(
8403 &Chunk::new(vec![
8404 column(
8405 LogicalType::BigInt,
8406 &[Value::BigInt(0), Value::BigInt(1), Value::BigInt(2), Value::BigInt(3)],
8407 ),
8408 column(
8409 LogicalType::Integer,
8410 &[
8411 Value::Integer(1),
8412 Value::Integer(1),
8413 Value::Integer(1),
8414 Value::Integer(1),
8415 ],
8416 ),
8417 column(
8418 LogicalType::Date,
8419 &[Value::Date(0), Value::Date(1), Value::Date(2), Value::Date(3)],
8420 ),
8421 ])
8422 .expect("three columns"),
8423 )
8424 .expect("four rows");
8425 let mut writer = writer.next("nation", plain).expect("a second table");
8426 writer
8427 .append(
8428 &Chunk::new(vec![column(LogicalType::Integer, &[Value::Integer(7)])])
8429 .expect("one column"),
8430 )
8431 .expect("one row");
8432 writer.finish().expect("commit");
8433
8434 let catalog = Catalog::open(&path).expect("reopen");
8435 let lineitem = catalog.table("lineitem").expect("the clustered table");
8436 assert_eq!(lineitem.table().clustering(), Some(&stage_zero));
8437 let nation = catalog.table("nation").expect("the plain table");
8438 assert_eq!(nation.table().clustering(), None, "nobody declared one here");
8439
8440 assert_eq!(lineitem.table().rows(), 4);
8443 assert_eq!(nation.table().rows(), 1);
8444 fs::remove_file(&path).ok();
8445 }
8446
8447 #[test]
8449 fn a_declaration_off_the_end_of_the_table_never_reaches_the_file() {
8450 let path = path("clustered-bad");
8451 let writer = Writer::create(&path, "items", vec![Field::new("a", LogicalType::Integer)])
8452 .expect("new file");
8453 let four =
8454 (0..4).map(|at| Field::new(format!("c{at}"), LogicalType::Integer)).collect::<Vec<_>>();
8455 let wrong = Clustering::new(vec![3], Width::Exact, &four).expect("valid against four");
8456 assert!(writer.declare(wrong).is_err(), "the table has one column, not four");
8457 fs::remove_file(&path).ok();
8458 }
8459}