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 = 23;
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, Clone)]
535pub struct StoredPart {
536 pub stripe: usize,
538 pub part: usize,
540 pub row: usize,
542 pub rows: usize,
544 pub encoding: String,
546 pub bytes: u64,
548 pub page: u64,
550 pub offset: u64,
552 pub low: Option<Value>,
554 pub high: Option<Value>,
556 pub nulls: Option<usize>,
558}
559
560#[derive(Debug)]
562struct GlobalDictionary {
563 primary: HashMap<u64, u32>,
564 collisions: HashMap<u64, Vec<u32>>,
565 offsets: Vec<u32>,
566 payload: Vec<u8>,
567 counts: Vec<u64>,
568 nulls: u64,
569}
570
571impl GlobalDictionary {
572 fn new() -> Self {
573 Self {
574 primary: HashMap::new(),
575 collisions: HashMap::new(),
576 offsets: vec![0],
577 payload: Vec::new(),
578 counts: Vec::new(),
579 nulls: 0,
580 }
581 }
582
583 fn bytes(&self, code: u32) -> Option<&[u8]> {
584 let start = *self.offsets.get(code as usize)? as usize;
585 let end = *self.offsets.get(code as usize + 1)? as usize;
586 self.payload.get(start..end)
587 }
588
589 fn code(&mut self, text: &str) -> Result<u32> {
590 let hash = checksum(text.as_bytes());
591 if let Some(&code) = self.primary.get(&hash) {
592 if self.bytes(code) == Some(text.as_bytes()) {
593 return Ok(code);
594 }
595 if let Some(codes) = self.collisions.get(&hash) {
596 if let Some(code) =
597 codes.iter().copied().find(|&code| self.bytes(code) == Some(text.as_bytes()))
598 {
599 return Ok(code);
600 }
601 }
602 let code = self.insert(text)?;
603 self.collisions.entry(hash).or_default().push(code);
604 return Ok(code);
605 }
606 let code = self.insert(text)?;
607 self.primary.insert(hash, code);
608 Ok(code)
609 }
610
611 fn insert(&mut self, text: &str) -> Result<u32> {
612 let code = u32::try_from(self.offsets.len() - 1)
613 .map_err(|_| invalid("global dictionary has too many values"))?;
614 self.payload.extend_from_slice(text.as_bytes());
615 self.offsets.push(
616 u32::try_from(self.payload.len())
617 .map_err(|_| invalid("global dictionary payload exceeds 4 GiB"))?,
618 );
619 self.counts.push(0);
620 Ok(code)
621 }
622
623 fn ranked(&self) -> Vec<(u64, u32)> {
643 let count = self.offsets.len() - 1;
644 let mut ranked = (0..count)
645 .map(|code| {
646 let code = code as u32;
647 (head(self.bytes(code).unwrap_or_default()), code)
648 })
649 .collect::<Vec<_>>();
650 ranked.sort_unstable_by(|left, right| {
651 left.0.cmp(&right.0).then_with(|| self.bytes(left.1).cmp(&self.bytes(right.1)))
652 });
653 ranked
654 }
655
656 fn observe(&mut self, code: u32, null: bool) -> Result<()> {
657 if null {
658 self.nulls = self.nulls.saturating_add(1);
659 return Ok(());
660 }
661 let count = self
662 .counts
663 .get_mut(code as usize)
664 .ok_or_else(|| invalid("global dictionary count code is out of range"))?;
665 *count = count.saturating_add(1);
666 Ok(())
667 }
668}
669
670#[derive(Debug)]
678pub struct Writer {
679 file: File,
680 at: u64,
688 table: Table,
689 generation: u64,
690 order: Vec<((u64, u64), (u64, u64))>,
693 next_order: u64,
694 dictionaries: Vec<Option<GlobalDictionary>>,
695 pending: Vec<PendingChunk>,
696 closed: Vec<Entry>,
698}
699
700#[derive(Debug)]
708struct PendingChunk {
709 order: (u64, u64),
710 chunk: Chunk,
711}
712
713#[derive(Debug)]
719struct ColumnStripe {
720 pages: Vec<Vec<u8>>,
721 codes: Vec<Option<Vec<u32>>>,
722 sieves: Vec<Option<Sieve>>,
723 ranges: Vec<Range>,
724}
725
726fn weight(ty: &LogicalType) -> usize {
734 match ty {
735 LogicalType::Varchar | LogicalType::Blob | LogicalType::Bit => 64,
736 LogicalType::HugeInt
737 | LogicalType::UHugeInt
738 | LogicalType::Uuid
739 | LogicalType::Interval => 16,
740 LogicalType::BigInt
741 | LogicalType::UBigInt
742 | LogicalType::Timestamp
743 | LogicalType::Time
744 | LogicalType::TimeTz
745 | LogicalType::TimestampTz
746 | LogicalType::TimestampS
747 | LogicalType::TimestampMs
748 | LogicalType::TimestampNs
749 | LogicalType::Double
750 | LogicalType::Decimal { .. } => 8,
751 LogicalType::Integer | LogicalType::UInteger | LogicalType::Date | LogicalType::Float => 4,
752 LogicalType::SmallInt | LogicalType::USmallInt => 2,
753 _ => 1,
754 }
755}
756
757pub const STRIPE_PARTS: usize = 64;
764
765const DICTIONARY_DECIDE_ROWS: usize = 4_096;
773
774const DICTIONARY_DISTINCT_IN_TEN: usize = 9;
790
791const INDEX_ENTRY: usize = size_of::<u32>() + size_of::<u64>();
793
794fn index_section(parts: usize) -> Result<usize> {
796 parts
797 .checked_mul(INDEX_ENTRY)
798 .and_then(|bytes| bytes.checked_add(size_of::<u64>()))
799 .ok_or_else(|| invalid("index page length overflow"))
800}
801
802impl Writer {
803 pub fn open(
821 path: impl AsRef<Path>,
822 name: impl Into<String>,
823 fields: Vec<Field>,
824 ) -> Result<Self> {
825 for field in &fields {
826 type_tag(&field.ty)?;
827 }
828 let name = name.into();
829 let path = path.as_ref();
830 let (_, size, slot, bytes, _) = slot_bytes(path)?;
831 let closed = decode_catalog(&bytes, size)?;
832 if closed.iter().any(|held| held.name == name) {
833 return Err(invalid("two tables in one native file have the same name"));
834 }
835 let generation = slot
840 .generation
841 .checked_add(1)
842 .ok_or_else(|| invalid("native file generation overflow"))?;
843 let file = OpenOptions::new().write(true).read(true).open(path).map_err(io)?;
844 Ok(Self {
845 file,
846 at: size,
849 dictionaries: fields
850 .iter()
851 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
852 .collect(),
853 table: Table {
854 name,
855 dictionaries: vec![None; fields.len()],
856 distincts: vec![None; fields.len()],
857 fields,
858 stripes: Vec::new(),
859 rows: 0,
860 frequencies: Vec::new(),
861 clustering: None,
862 },
863 generation,
864 order: Vec::new(),
865 next_order: 0,
866 pending: Vec::with_capacity(STRIPE_PARTS),
867 closed,
868 })
869 }
870
871 pub fn create(
877 path: impl AsRef<Path>,
878 name: impl Into<String>,
879 fields: Vec<Field>,
880 ) -> Result<Self> {
881 for field in &fields {
882 type_tag(&field.ty)?;
883 }
884 let file =
885 OpenOptions::new().write(true).read(true).create_new(true).open(path).map_err(io)?;
886 let mut header = [0; HEADER as usize];
887 header[..8].copy_from_slice(MAGIC);
888 header[8..12].copy_from_slice(&FORMAT.to_le_bytes());
889 write_at(&file, 0, &header)?;
890 Ok(Self {
891 file,
892 at: HEADER,
893 dictionaries: fields
894 .iter()
895 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
896 .collect(),
897 table: Table {
898 name: name.into(),
899 dictionaries: vec![None; fields.len()],
900 distincts: vec![None; fields.len()],
901 fields,
902 stripes: Vec::new(),
903 rows: 0,
904 frequencies: Vec::new(),
905 clustering: None,
906 },
907 generation: 1,
908 order: Vec::new(),
909 next_order: 0,
910 pending: Vec::with_capacity(STRIPE_PARTS),
911 closed: Vec::new(),
912 })
913 }
914
915 pub fn empty(path: impl AsRef<Path>) -> Result<()> {
933 let file =
934 OpenOptions::new().write(true).read(true).create_new(true).open(path).map_err(io)?;
935 let mut header = [0; HEADER as usize];
936 header[..8].copy_from_slice(MAGIC);
937 header[8..12].copy_from_slice(&FORMAT.to_le_bytes());
938 write_at(&file, 0, &header)?;
939 let catalog = encode_catalog(&[])?;
940 write_at(&file, HEADER, &catalog)?;
941 file.sync_all().map_err(io)?;
945 let slot = Slot {
946 offset: HEADER,
947 length: u32::try_from(catalog.len()).map_err(|_| invalid("catalog length overflow"))?,
948 generation: 1,
949 hash: checksum(&catalog),
950 };
951 write_at(&file, slot_offset(1), &slot.bytes())?;
952 file.sync_all().map_err(io)?;
953 Ok(())
954 }
955
956 pub fn next(mut self, name: impl Into<String>, fields: Vec<Field>) -> Result<Self> {
967 for field in &fields {
968 type_tag(&field.ty)?;
969 }
970 let name = name.into();
971 let entry = self.close()?;
972 if self.closed.iter().chain(std::iter::once(&entry)).any(|held| held.name == name) {
973 return Err(invalid("two tables in one native file have the same name"));
974 }
975 let Self { file, at, generation, mut closed, .. } = self;
976 closed.push(entry);
977 Ok(Self {
978 file,
979 at,
980 generation,
981 closed,
982 dictionaries: fields
983 .iter()
984 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
985 .collect(),
986 table: Table {
987 name,
988 dictionaries: vec![None; fields.len()],
989 distincts: vec![None; fields.len()],
990 fields,
991 stripes: Vec::new(),
992 rows: 0,
993 frequencies: Vec::new(),
994 clustering: None,
995 },
996 order: Vec::new(),
997 next_order: 0,
998 pending: Vec::with_capacity(STRIPE_PARTS),
999 })
1000 }
1001
1002 pub fn declare(mut self, clustering: Clustering) -> Result<Self> {
1017 self.table.clustering = Some(Clustering::new(
1020 clustering.columns().to_vec(),
1021 clustering.width(),
1022 &self.table.fields,
1023 )?);
1024 Ok(self)
1025 }
1026
1027 fn put(&mut self, bytes: &[u8]) -> Result<()> {
1032 write_at(&self.file, self.at, bytes)?;
1033 self.at = self
1034 .at
1035 .checked_add(bytes.len() as u64)
1036 .ok_or_else(|| invalid("native file length overflow"))?;
1037 Ok(())
1038 }
1039
1040 pub fn append(&mut self, chunk: &Chunk) -> Result<()> {
1046 let order = (self.next_order, 0);
1047 self.next_order = self.next_order.saturating_add(1);
1048 self.append_at(order, chunk)
1049 }
1050
1051 pub fn append_at(&mut self, order: (u64, u64), chunk: &Chunk) -> Result<()> {
1062 if chunk.is_empty() {
1063 return Ok(());
1064 }
1065 self.admit(chunk)?;
1066 if self.pending.last().is_some_and(|last| last.order > order) {
1067 self.flush_pending()?;
1068 }
1069 self.pending.push(PendingChunk { order, chunk: chunk.clone() });
1074 if self.pending.len() == STRIPE_PARTS {
1075 self.flush_pending()?;
1076 }
1077 Ok(())
1078 }
1079
1080 pub fn append_stripe(&mut self, parts: Vec<((u64, u64), Chunk)>) -> Result<()> {
1096 if parts.len() > STRIPE_PARTS {
1097 return Err(invalid("a stripe was handed more parts than it holds"));
1098 }
1099 self.flush_pending()?;
1102 for (order, chunk) in parts {
1103 if chunk.is_empty() {
1104 continue;
1105 }
1106 self.admit(&chunk)?;
1107 self.pending.push(PendingChunk { order, chunk });
1108 }
1109 self.flush_pending()
1110 }
1111
1112 fn admit(&mut self, chunk: &Chunk) -> Result<()> {
1114 if chunk.width() != self.table.fields.len() {
1115 return Err(invalid("chunk width differs from table schema"));
1116 }
1117 for (index, field) in self.table.fields.iter().enumerate() {
1118 if chunk.column(index)?.logical_type() != &field.ty {
1119 return Err(invalid("chunk type differs from table schema"));
1120 }
1121 }
1122 self.table.rows = self
1123 .table
1124 .rows
1125 .checked_add(chunk.len())
1126 .ok_or_else(|| invalid("row count overflow"))?;
1127 Ok(())
1128 }
1129
1130 fn encode_column(
1161 index: usize,
1162 held: &[PendingChunk],
1163 dictionary: &mut Option<GlobalDictionary>,
1164 ) -> Result<ColumnStripe> {
1165 let deciding = dictionary.as_ref().is_some_and(|held| held.offsets.len() == 1);
1168 let stripe = Self::encode_pages(index, held, dictionary.as_mut())?;
1169 if !deciding {
1170 return Ok(stripe);
1171 }
1172 let rows: usize = held.iter().map(|pending| pending.chunk.len()).sum();
1173 let distinct = dictionary.as_ref().map_or(0, |held| held.offsets.len() - 1);
1174 if rows < DICTIONARY_DECIDE_ROWS
1175 || distinct.saturating_mul(10) <= rows.saturating_mul(DICTIONARY_DISTINCT_IN_TEN)
1176 {
1177 return Ok(stripe);
1178 }
1179 *dictionary = None;
1180 Self::encode_pages(index, held, None)
1181 }
1182
1183 fn encode_pages(
1185 index: usize,
1186 held: &[PendingChunk],
1187 mut dictionary: Option<&mut GlobalDictionary>,
1188 ) -> Result<ColumnStripe> {
1189 let mut stripe = ColumnStripe {
1190 pages: Vec::with_capacity(held.len()),
1191 codes: Vec::with_capacity(held.len()),
1192 sieves: Vec::with_capacity(held.len()),
1193 ranges: Vec::with_capacity(held.len()),
1194 };
1195 for pending in held {
1196 let column = pending.chunk.column(index)?;
1197 let (bytes, unique) = encode(column, dictionary.as_deref_mut())?;
1198 if bytes.len() > MAX_PAGE {
1199 return Err(invalid("column page exceeds the configured bound"));
1200 }
1201 let range = Range::of(column);
1204 let sieve = match dictionary {
1217 Some(_) => None,
1218 None => Sieve::of(column, &range, SIEVE_BUDGET)
1219 .filter(|sieve| sieve.len() < bytes.len()),
1220 };
1221 stripe.pages.push(bytes);
1222 stripe.codes.push(unique);
1223 stripe.sieves.push(sieve);
1224 stripe.ranges.push(range);
1225 }
1226 Ok(stripe)
1227 }
1228
1229 fn encode_columns(&mut self, held: &[PendingChunk]) -> Result<Vec<ColumnStripe>> {
1238 let width = self.table.fields.len();
1239 let workers = std::thread::available_parallelism()
1240 .map_or(1, usize::from)
1241 .min(MAX_ENCODE_WORKERS)
1242 .min(width);
1243 if workers <= 1 || held.len() <= 1 {
1244 return self
1245 .dictionaries
1246 .iter_mut()
1247 .enumerate()
1248 .map(|(index, dictionary)| Self::encode_column(index, held, dictionary))
1249 .collect();
1250 }
1251 let mut jobs: Vec<(usize, Option<GlobalDictionary>)> =
1254 std::mem::take(&mut self.dictionaries).into_iter().enumerate().collect();
1255 jobs.sort_by_key(|(index, _)| weight(&self.table.fields[*index].ty));
1257 let queue = Mutex::new(jobs);
1258 let pieces = std::thread::scope(|scope| {
1259 (0..workers)
1260 .map(|_| {
1261 scope.spawn(|| {
1262 let mut mine = Vec::new();
1263 loop {
1264 let taken = queue
1265 .lock()
1266 .map_err(|_| Error::internal("a native encode worker panicked"))?
1267 .pop();
1268 let Some((index, mut dictionary)) = taken else { break };
1269 let encoded = Self::encode_column(index, held, &mut dictionary)?;
1270 mine.push((index, dictionary, encoded));
1271 }
1272 Ok(mine)
1273 })
1274 })
1275 .collect::<Vec<_>>()
1276 .into_iter()
1277 .map(|handle| {
1278 handle.join().map_err(|_| Error::internal("a native encode worker panicked"))?
1279 })
1280 .collect::<Result<Vec<_>>>()
1281 })?;
1282 let mut dictionaries: Vec<Option<GlobalDictionary>> = (0..width).map(|_| None).collect();
1283 let mut encoded: Vec<Option<ColumnStripe>> = (0..width).map(|_| None).collect();
1284 for piece in pieces {
1285 for (index, dictionary, stripe) in piece {
1286 dictionaries[index] = dictionary;
1287 encoded[index] = Some(stripe);
1288 }
1289 }
1290 self.dictionaries = dictionaries;
1291 encoded
1292 .into_iter()
1293 .map(|stripe| stripe.ok_or_else(|| Error::internal("a column was never encoded")))
1294 .collect()
1295 }
1296
1297 fn flush_pending(&mut self) -> Result<()> {
1299 if self.pending.is_empty() {
1300 return Ok(());
1301 }
1302 let width = self.table.fields.len();
1303 let mut held = std::mem::take(&mut self.pending);
1306 let parts = held.len();
1307 let encoded = self.encode_columns(&held)?;
1308 let mut pages = Vec::with_capacity(width);
1309 let mut memberships = vec![None; width];
1310 let mut ranges = Vec::with_capacity(width);
1311 let mut index = Vec::with_capacity(width.saturating_mul(index_section(parts)?));
1312 for stripe in &encoded {
1313 let offset = self.at;
1314 let section = index.len();
1315 let mut length = 0_usize;
1316 for bytes in &stripe.pages {
1317 write_at(&self.file, self.at + length as u64, bytes)?;
1318 put_u32(
1319 &mut index,
1320 u32::try_from(bytes.len()).map_err(|_| invalid("part length overflow"))?,
1321 );
1322 put_u64(&mut index, checksum(bytes));
1323 length = length
1324 .checked_add(bytes.len())
1325 .ok_or_else(|| invalid("column page length overflow"))?;
1326 }
1327 let hash = checksum(&index[section..]);
1328 put_u64(&mut index, hash);
1329 if length > MAX_PAGE {
1330 return Err(invalid("column page exceeds the configured bound"));
1331 }
1332 self.at = self
1333 .at
1334 .checked_add(length as u64)
1335 .ok_or_else(|| invalid("native file length overflow"))?;
1336 pages.push(Span {
1337 offset,
1338 length: u32::try_from(length).map_err(|_| invalid("page length overflow"))?,
1339 });
1340 ranges.push(merged_range(stripe.ranges.iter().cloned()));
1341 }
1342 for (membership, stripe) in memberships.iter_mut().zip(&encoded) {
1343 if stripe.codes.iter().all(Option::is_none) {
1344 continue;
1345 }
1346 let lists = stripe
1347 .codes
1348 .iter()
1349 .map(|codes| codes.clone().unwrap_or_default())
1350 .collect::<Vec<_>>();
1351 let bytes = encode_membership(&merged_codes(lists));
1352 let offset = self.at;
1353 self.put(&bytes)?;
1354 *membership = Some(Page {
1355 offset,
1356 length: u32::try_from(bytes.len())
1357 .map_err(|_| invalid("membership page length overflow"))?,
1358 hash: checksum(&bytes),
1359 });
1360 }
1361 let mut sieves = vec![None; width];
1362 for (page, stripe) in sieves.iter_mut().zip(&encoded) {
1363 if stripe.sieves.iter().all(Option::is_none) {
1364 continue;
1365 }
1366 let bytes = encode_sieves(stripe.sieves.iter())?;
1367 let offset = self.at;
1368 self.put(&bytes)?;
1369 *page = Some(Page {
1370 offset,
1371 length: u32::try_from(bytes.len())
1372 .map_err(|_| invalid("sieve page length overflow"))?,
1373 hash: checksum(&bytes),
1374 });
1375 }
1376 let mut part_ranges = vec![None; width];
1382 if parts > 1 {
1383 for ((page, stripe), span) in part_ranges.iter_mut().zip(&encoded).zip(&pages) {
1384 let bytes = encode_part_ranges(&stripe.ranges)?;
1385 if bytes.len() >= span.length as usize {
1386 continue;
1387 }
1388 let offset = self.at;
1389 self.put(&bytes)?;
1390 *page = Some(Page {
1391 offset,
1392 length: u32::try_from(bytes.len())
1393 .map_err(|_| invalid("part range page length overflow"))?,
1394 hash: checksum(&bytes),
1395 });
1396 }
1397 }
1398 let offset = self.at;
1399 self.put(&index)?;
1400 let index = Span {
1401 offset,
1402 length: u32::try_from(index.len())
1403 .map_err(|_| invalid("index page length overflow"))?,
1404 };
1405 let mut rows = 0_usize;
1406 let mut lengths = Vec::with_capacity(parts);
1407 let mut span = None;
1408 for pending in held.drain(..) {
1409 let part = pending.chunk.len();
1410 rows = rows.checked_add(part).ok_or_else(|| invalid("row count overflow"))?;
1411 lengths.push(u32::try_from(part).map_err(|_| invalid("part row count overflow"))?);
1412 span = Some(
1413 span.map_or((pending.order, pending.order), |(first, _)| (first, pending.order)),
1414 );
1415 }
1416 self.order.push(span.ok_or_else(|| invalid("a stripe was flushed with no parts"))?);
1417 self.table.stripes.push(Stripe {
1418 rows,
1419 parts: lengths,
1420 index,
1421 pages,
1422 memberships,
1423 sieves,
1424 part_ranges,
1425 zone: Zone::from_ranges(ranges),
1426 });
1427 self.pending = held;
1429 Ok(())
1430 }
1431
1432 fn numeric_frequency(&self, column: usize) -> Result<Option<FrequencySummary>> {
1436 let ty = &self.table.fields[column].ty;
1437 if !matches!(
1438 ty,
1439 LogicalType::TinyInt
1440 | LogicalType::SmallInt
1441 | LogicalType::Integer
1442 | LogicalType::BigInt
1443 | LogicalType::UTinyInt
1444 | LogicalType::USmallInt
1445 | LogicalType::UInteger
1446 | LogicalType::UBigInt
1447 | LogicalType::Date
1448 | LogicalType::Timestamp
1449 ) {
1450 return Ok(None);
1451 }
1452 let mut candidates: HashMap<FrequencyValue, u32> = HashMap::new();
1453 let mut decrements = 0_u64;
1454 self.visit_numeric(column, |_, value| {
1455 if let Some(count) = candidates.get_mut(&value) {
1456 *count = count.saturating_add(1);
1457 } else if candidates.len() < FREQUENCY_CANDIDATES {
1458 candidates.insert(value, 1);
1459 } else {
1460 candidates.retain(|_, count| {
1461 *count -= 1;
1462 *count != 0
1463 });
1464 decrements = decrements.saturating_add(1);
1465 }
1466 })?;
1467 let (exact, ordinals) = if decrements == 0 {
1468 (
1469 candidates
1470 .into_iter()
1471 .map(|(value, count)| (value, u64::from(count)))
1472 .collect::<HashMap<_, _>>(),
1473 Vec::new(),
1474 )
1475 } else {
1476 let mut lower = candidates.values().copied().collect::<Vec<_>>();
1477 lower.sort_unstable_by(|left, right| right.cmp(left));
1478 if lower.len() < FREQUENCY_BUILD_RANK
1479 || u64::from(lower[FREQUENCY_BUILD_RANK - 1]) <= decrements
1480 {
1481 return Ok(None);
1482 }
1483 let mut exact =
1484 candidates.into_keys().map(|value| (value, 0_u64)).collect::<HashMap<_, _>>();
1485 let mut ordinals = Vec::new();
1486 let mut exceeded = false;
1487 self.visit_numeric(column, |ordinal, value| {
1488 if let Some(count) = exact.get_mut(&value) {
1489 *count = count.saturating_add(1);
1490 if !exceeded {
1491 if ordinals.len() < FREQUENCY_ORDINALS {
1492 ordinals.push(ordinal);
1493 } else {
1494 ordinals.clear();
1495 exceeded = true;
1496 }
1497 }
1498 }
1499 })?;
1500 (exact, ordinals)
1501 };
1502 let mut entries = exact
1503 .into_iter()
1504 .map(|(value, count)| FrequencyEntry { value, count })
1505 .collect::<Vec<_>>();
1506 entries.sort_unstable_by(|left, right| {
1507 right.count.cmp(&left.count).then_with(|| frequency_order(left.value, right.value))
1508 });
1509 let omitted_max =
1510 entries.get(FREQUENCY_ENTRIES).map_or(decrements, |entry| decrements.max(entry.count));
1511 entries.truncate(FREQUENCY_ENTRIES);
1512 Ok(Some(FrequencySummary { entries, omitted_max, ordinals }))
1513 }
1514
1515 fn visit_numeric(
1516 &self,
1517 column: usize,
1518 mut visit: impl FnMut(u64, FrequencyValue),
1519 ) -> Result<()> {
1520 let ty = &self.table.fields[column].ty;
1521 let mut start = 0_u64;
1522 for stripe in &self.table.stripes {
1523 let spans = read_index(&self.file, stripe, column)?;
1524 let page = stripe.pages[column];
1525 let mut bytes = vec![0; page.length as usize];
1526 read_at(&self.file, page.offset, &mut bytes)?;
1527 for (span, &rows) in spans.iter().zip(&stripe.parts) {
1528 let part = part_bytes(&bytes, *span)?;
1529 if checksum(part) != span.hash {
1530 return Err(invalid("column page checksum differs while building frequencies"));
1531 }
1532 let rows = rows as usize;
1533 let vector = decode(ty, rows, part, None)?;
1534 for row in 0..rows {
1536 let value = if vector.is_null_at(row) {
1537 FrequencyValue::Null
1538 } else {
1539 let widened = match vector.signed_at(row) {
1543 Some(value) => Some(value),
1544 None => match vector.value_at(row) {
1545 Value::UTinyInt(value) => Some(i128::from(value)),
1546 Value::USmallInt(value) => Some(i128::from(value)),
1547 Value::UInteger(value) => Some(i128::from(value)),
1548 Value::UBigInt(value) => Some(i128::from(value)),
1549 _ => None,
1550 },
1551 };
1552 FrequencyValue::Integer(widened.ok_or_else(|| {
1553 invalid("numeric frequency page did not contain an integer value")
1554 })?)
1555 };
1556 visit(start.saturating_add(row as u64), value);
1557 }
1558 start = start.saturating_add(rows as u64);
1559 }
1560 }
1561 Ok(())
1562 }
1563
1564 fn numeric_frequencies(&self) -> Result<Vec<Option<FrequencySummary>>> {
1572 let mut columns = self
1573 .table
1574 .fields
1575 .iter()
1576 .enumerate()
1577 .filter_map(|(column, field)| {
1578 matches!(
1579 field.ty,
1580 LogicalType::TinyInt
1581 | LogicalType::SmallInt
1582 | LogicalType::Integer
1583 | LogicalType::BigInt
1584 | LogicalType::UTinyInt
1585 | LogicalType::USmallInt
1586 | LogicalType::UInteger
1587 | LogicalType::UBigInt
1588 | LogicalType::Date
1589 | LogicalType::Timestamp
1590 )
1591 .then_some(column)
1592 })
1593 .collect::<Vec<_>>();
1594 let workers = std::thread::available_parallelism()
1595 .map_or(1, usize::from)
1596 .min(MAX_FREQUENCY_WORKERS)
1597 .min(columns.len());
1598 if workers <= 1 {
1599 let mut frequencies = vec![None; self.table.fields.len()];
1600 for column in columns {
1601 frequencies[column] = self.numeric_frequency(column)?;
1602 }
1603 return Ok(frequencies);
1604 }
1605 columns.sort_by_key(|&column| weight(&self.table.fields[column].ty));
1608 let queue = Mutex::new(columns);
1609 let pieces = std::thread::scope(|scope| {
1610 (0..workers)
1611 .map(|_| {
1612 scope.spawn(|| {
1613 let mut mine = Vec::new();
1614 loop {
1615 let taken = queue
1616 .lock()
1617 .map_err(|_| Error::internal("a native frequency worker panicked"))?
1618 .pop();
1619 let Some(column) = taken else { break };
1620 mine.push((column, self.numeric_frequency(column)?));
1621 }
1622 Ok(mine)
1623 })
1624 })
1625 .collect::<Vec<_>>()
1626 .into_iter()
1627 .map(|handle| {
1628 handle
1629 .join()
1630 .map_err(|_| Error::internal("a native frequency worker panicked"))?
1631 })
1632 .collect::<Result<Vec<_>>>()
1633 })?;
1634 let mut frequencies = vec![None; self.table.fields.len()];
1635 for piece in pieces {
1636 for (column, summary) in piece {
1637 frequencies[column] = summary;
1638 }
1639 }
1640 Ok(frequencies)
1641 }
1642
1643 fn close(&mut self) -> Result<Entry> {
1654 self.flush_pending()?;
1655 let mut stripes = std::mem::take(&mut self.order)
1656 .into_iter()
1657 .zip(std::mem::take(&mut self.table.stripes))
1658 .collect::<Vec<_>>();
1659 stripes.sort_by_key(|(order, _)| order.0);
1660 let mut previous: Option<(u64, u64)> = None;
1661 for ((first, last), _) in &stripes {
1662 if previous.is_some_and(|previous| previous >= *first) {
1663 return Err(invalid("chunks did not arrive in source order"));
1664 }
1665 previous = Some(*last);
1666 }
1667 self.table.stripes = stripes.into_iter().map(|(_, stripe)| stripe).collect();
1668 self.table.frequencies = self.numeric_frequencies()?;
1669 let dictionaries = std::mem::take(&mut self.dictionaries);
1670 let orders = rankings(&dictionaries)?;
1671 for (index, (dictionary, order)) in dictionaries.into_iter().zip(orders).enumerate() {
1672 let Some(dictionary) = dictionary else { continue };
1673 self.table.distincts[index] =
1677 Some(dictionary.counts.iter().filter(|count| **count != 0).count() as u64);
1678 self.table.frequencies[index] = Some(code_frequency(&dictionary));
1679 let encoded = encode_global_dictionary(dictionary, &order)?;
1680 let offset = self.at;
1681 self.put(&encoded.index)?;
1682 self.put(&encoded.ranks)?;
1683 for block in &encoded.payload {
1684 self.put(block)?;
1685 }
1686 let payload_len =
1687 encoded.payload.iter().try_fold(0_usize, |len, block| len.checked_add(block.len()));
1688 let length = payload_len
1689 .and_then(|len| len.checked_add(encoded.index.len()))
1690 .and_then(|len| len.checked_add(encoded.ranks.len()))
1691 .ok_or_else(|| invalid("dictionary page length overflow"))?;
1692 self.table.dictionaries[index] = Some(Page {
1693 offset,
1694 length: u32::try_from(length)
1695 .map_err(|_| invalid("dictionary page length overflow"))?,
1696 hash: checksum(&encoded.index),
1697 });
1698 }
1699 let directory = encode_directory(&self.table)?;
1700 if directory.len() > MAX_DIRECTORY {
1701 return Err(invalid("directory exceeds the configured bound"));
1702 }
1703 let offset = self.at;
1704 self.put(&directory)?;
1705 Ok(Entry {
1706 name: self.table.name.clone(),
1707 fields: self.table.fields.clone(),
1708 rows: self.table.rows,
1709 directory: Page {
1710 offset,
1711 length: u32::try_from(directory.len())
1712 .map_err(|_| invalid("directory length overflow"))?,
1713 hash: checksum(&directory),
1714 },
1715 })
1716 }
1717
1718 pub fn finish(mut self) -> Result<Table> {
1728 let entry = self.close()?;
1729 let mut tables = std::mem::take(&mut self.closed);
1730 tables.push(entry);
1731 let catalog = encode_catalog(&tables)?;
1732 if catalog.len() > MAX_DIRECTORY {
1733 return Err(invalid("catalog exceeds the configured bound"));
1734 }
1735 let offset = self.at;
1736 self.put(&catalog)?;
1737 self.file.sync_all().map_err(io)?;
1741 let slot = Slot {
1742 offset,
1743 length: u32::try_from(catalog.len()).map_err(|_| invalid("catalog length overflow"))?,
1744 generation: self.generation,
1745 hash: checksum(&catalog),
1746 };
1747 write_at(&self.file, slot_offset(self.generation), &slot.bytes())?;
1752 self.file.sync_all().map_err(io)?;
1753 Ok(self.table)
1754 }
1755}
1756
1757#[derive(Debug, Clone)]
1759pub struct Reader {
1760 file: Arc<File>,
1761 table: Arc<Table>,
1762 dictionaries: Arc<Vec<OnceLock<Arc<Vector>>>>,
1763 loading: Arc<Vec<Mutex<()>>>,
1772 opened: Arc<AtomicUsize>,
1776 sieves: Arc<Vec<Vec<SieveSlot>>>,
1780 part_ranges: Arc<Vec<Vec<RangeSlot>>>,
1783 places: Arc<Vec<Place>>,
1785 cache: Arc<Vec<Mutex<Cached>>>,
1786 pages: Arc<AtomicUsize>,
1789 indexes: Arc<AtomicUsize>,
1792 kept: Arc<AtomicUsize>,
1795 size: u64,
1797 directory: u64,
1799 opening: Opening,
1801}
1802
1803#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1815pub struct Opening {
1816 pub reads: u32,
1819 pub bytes: u64,
1821}
1822
1823#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1825pub struct Reads {
1826 pub opening: Opening,
1828 pub pages: usize,
1830 pub indexes: usize,
1832 pub dictionaries: usize,
1835}
1836
1837#[derive(Debug, Clone, Copy)]
1839struct Place {
1840 stripe: u32,
1841 part: u32,
1842 rows: u32,
1843}
1844
1845#[derive(Debug, Clone, Copy)]
1847struct PartSpan {
1848 start: usize,
1849 length: usize,
1850 hash: u64,
1851}
1852
1853#[derive(Debug, Clone)]
1859struct CachedColumn {
1860 stripe: usize,
1861 index: Arc<Vec<PartSpan>>,
1862 page: Option<Arc<Vec<u8>>>,
1863}
1864
1865#[derive(Debug, Default)]
1885struct Cached {
1886 pages: Vec<Option<Arc<Vec<u8>>>>,
1887 order: VecDeque<usize>,
1888 loading: Vec<usize>,
1889 index: Vec<Option<Arc<Vec<PartSpan>>>>,
1890}
1891
1892const CACHED_STRIPES_PER_COLUMN: usize = 4;
1904
1905type SieveSlot = OnceLock<Arc<Vec<Option<Sieve>>>>;
1907
1908type RangeSlot = OnceLock<Arc<Vec<Range>>>;
1909
1910#[derive(Debug)]
1911struct NativeText {
1912 file: Arc<File>,
1913 values: usize,
1915 offsets: Vec<u8>,
1924 offset_bits: usize,
1927 ranks: usize,
1929 rank_at: u64,
1933 rank_ends: Vec<u64>,
1937 rank_hashes: Vec<u64>,
1938 rank_blocks: Vec<OnceLock<Result<Vec<u8>>>>,
1939 code_bits: usize,
1942 code_ranks: OnceLock<Option<Vec<u32>>>,
1949 payload: u64,
1950 ends: Vec<u64>,
1953 hashes: Vec<u64>,
1954 blocks: Vec<OnceLock<Result<Vec<u8>>>>,
1956 keep_budget: usize,
1959 payload_kept: AtomicUsize,
1967 searched: Mutex<HashMap<Vec<u8>, (usize, bool)>>,
1984}
1985
1986const TEXT_SEARCH_MEMO: usize = 64;
1991
1992const TEXT_PAYLOAD_VALUES: usize = 1024;
2008
2009const TEXT_KEEP_BUDGET: usize = 256 * 1024 * 1024;
2030
2031const TEXT_OFFSET_RUN: usize = 512;
2038
2039const DICTIONARY_HEADER: usize = 16;
2042
2043const TEXT_RANK_BLOCK: usize = 512;
2054
2055const RANK_BLOCK_HEADER: usize = size_of::<u64>() + 1;
2069
2070impl NativeText {
2071 fn payload_block(&self, block: usize) -> Result<Option<&[u8]>> {
2078 let Some(slot) = self.blocks.get(block) else { return Ok(None) };
2079 let bytes = slot.get_or_init(|| self.decode_block(block)).as_ref().map_err(Clone::clone)?;
2080 Ok(Some(bytes.as_slice()))
2081 }
2082
2083 fn decode_block(&self, block: usize) -> Result<Vec<u8>> {
2088 let start = if block == 0 { 0 } else { self.ends[block - 1] };
2089 let end = self.ends[block];
2090 let len = end
2091 .checked_sub(start)
2092 .ok_or_else(|| invalid("global dictionary block ends before it starts"))?;
2093 let mut stored = vec![
2094 0;
2095 usize::try_from(len).map_err(|_| invalid(
2096 "global dictionary block does not fit in memory"
2097 ))?
2098 ];
2099 read_at(&self.file, self.payload + start, &mut stored)?;
2100 if checksum(&stored) != self.hashes[block] {
2101 return Err(invalid("global dictionary payload checksum differs"));
2102 }
2103 let first = block * TEXT_PAYLOAD_VALUES;
2104 let last = (first + TEXT_PAYLOAD_VALUES).min(self.values);
2105 let want = self.end_within(last - 1)? as usize;
2106 let values = string::decode_flat(&stored)?;
2107 if values.len() != last - first {
2108 return Err(invalid("global dictionary block holds the wrong value count"));
2109 }
2110 let bytes = values.into_bytes();
2111 if bytes.len() != want {
2112 return Err(invalid("global dictionary block decodes to the wrong length"));
2113 }
2114 Ok(bytes)
2115 }
2116
2117 fn end_within(&self, index: usize) -> Result<u32> {
2119 let run = index / TEXT_OFFSET_RUN;
2120 let bytes = self
2121 .offsets
2122 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
2123 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
2124 let end = bitpack::tail_at(bytes, self.offset_bits, index % TEXT_OFFSET_RUN)
2125 .map_err(|_| invalid("global dictionary offsets are short"))?;
2126 u32::try_from(end).map_err(|_| invalid("global dictionary offset is past the payload"))
2127 }
2128
2129 fn ends_within(&self, first: usize, last: usize) -> Result<Vec<u64>> {
2142 let mut ends = Vec::with_capacity(last.saturating_sub(first));
2143 let mut at = first;
2144 while at < last {
2145 let run = at / TEXT_OFFSET_RUN;
2146 let stop = ((run + 1) * TEXT_OFFSET_RUN).min(last);
2147 let held = self.values.saturating_sub(run * TEXT_OFFSET_RUN).min(TEXT_OFFSET_RUN);
2148 let bytes = self
2149 .offsets
2150 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
2151 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
2152 let run_ends = bitpack::unpack_tail(bytes, self.offset_bits, held)
2153 .map_err(|_| invalid("global dictionary offsets are short"))?;
2154 let within = run_ends
2155 .get(at % TEXT_OFFSET_RUN..stop - run * TEXT_OFFSET_RUN)
2156 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
2157 ends.extend_from_slice(within);
2158 at = stop;
2159 }
2160 Ok(ends)
2161 }
2162
2163 fn start_within(&self, index: usize) -> Result<u32> {
2166 if index % TEXT_PAYLOAD_VALUES == 0 { Ok(0) } else { self.end_within(index - 1) }
2167 }
2168
2169 fn span_within(&self, index: usize) -> Result<(u32, u32)> {
2177 let within = index % TEXT_OFFSET_RUN;
2178 let (start, end) = if within == 0 {
2179 (self.start_within(index)?, self.end_within(index)?)
2180 } else {
2181 let run = index / TEXT_OFFSET_RUN;
2182 let bytes = self
2183 .offsets
2184 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
2185 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
2186 let (start, end) = bitpack::tail_pair(bytes, self.offset_bits, within)
2187 .map_err(|_| invalid("global dictionary offsets are short"))?;
2188 let ends = u32::try_from(end)
2189 .map_err(|_| invalid("global dictionary offset is past the payload"))?;
2190 let starts = u32::try_from(start)
2191 .map_err(|_| invalid("global dictionary offset is past the payload"))?;
2192 (starts, ends)
2193 };
2194 if start > end {
2195 return Err(invalid("global dictionary value ends before it starts"));
2196 }
2197 Ok((start, end))
2198 }
2199
2200 fn rank_parts(&self, rank: usize) -> Result<(&[u8], usize)> {
2207 let slot = self
2208 .rank_blocks
2209 .get(rank / TEXT_RANK_BLOCK)
2210 .ok_or_else(|| invalid("global dictionary rank is past the order"))?;
2211 let block = slot
2212 .get_or_init(|| {
2213 let which = rank / TEXT_RANK_BLOCK;
2214 let start = if which == 0 { 0 } else { self.rank_ends[which - 1] };
2215 let end = self.rank_ends[which];
2216 let mut bytes = vec![0; (end - start) as usize];
2217 read_at(&self.file, self.rank_at + start, &mut bytes)?;
2218 if checksum(&bytes)
2219 != *self
2220 .rank_hashes
2221 .get(rank / TEXT_RANK_BLOCK)
2222 .ok_or_else(|| invalid("global dictionary rank block has no checksum"))?
2223 {
2224 return Err(invalid("global dictionary rank checksum differs"));
2225 }
2226 Ok(bytes)
2227 })
2228 .as_ref()
2229 .map_err(Clone::clone)?;
2230 Ok((block.as_slice(), rank % TEXT_RANK_BLOCK))
2231 }
2232
2233 fn head_at(&self, rank: usize) -> Result<u64> {
2235 let (block, within) = self.rank_parts(rank)?;
2236 let (base, width, packed) = rank_heads(block)?;
2237 let above = bitpack::tail_at(packed, width, within)
2238 .map_err(|_| invalid("global dictionary rank block is short of heads"))?;
2239 Ok(base.wrapping_add(above))
2240 }
2241
2242 fn rank_codes<'block>(&self, block: &'block [u8], count: usize) -> Result<&'block [u8]> {
2244 let (_, width, packed) = rank_heads(block)?;
2245 packed
2246 .get(bitpack::tail_len(count, width)..)
2247 .ok_or_else(|| invalid("global dictionary rank block is short of codes"))
2248 }
2249
2250 fn rank_block_len(&self, rank: usize) -> usize {
2252 let first = rank / TEXT_RANK_BLOCK * TEXT_RANK_BLOCK;
2253 TEXT_RANK_BLOCK.min(self.ranks - first)
2254 }
2255}
2256
2257fn rank_heads(block: &[u8]) -> Result<(u64, usize, &[u8])> {
2259 let header = block
2260 .get(..RANK_BLOCK_HEADER)
2261 .ok_or_else(|| invalid("global dictionary rank block is short"))?;
2262 let base = u64::from_le_bytes(header[..8].try_into().expect("eight bytes"));
2263 let width = header[8] as usize;
2264 if width > 64 {
2265 return Err(invalid("global dictionary rank block packs heads past a word"));
2266 }
2267 Ok((base, width, &block[RANK_BLOCK_HEADER..]))
2268}
2269
2270fn offset_width(offsets: &[u32]) -> usize {
2277 let values = offsets.len() - 1;
2278 let mut span = 0;
2279 for first in (0..values).step_by(TEXT_PAYLOAD_VALUES) {
2280 let last = (first + TEXT_PAYLOAD_VALUES).min(values);
2281 span = span.max(offsets[last] - offsets[first]);
2282 }
2283 (u32::BITS - span.leading_zeros()) as usize
2284}
2285
2286fn offset_bytes(values: usize, bits: usize) -> usize {
2289 let full = values / TEXT_OFFSET_RUN;
2290 let rest = values % TEXT_OFFSET_RUN;
2291 full * TEXT_OFFSET_RUN / 8 * bits + bitpack::tail_len(rest, bits)
2292}
2293
2294fn encode_offsets(offsets: &[u32], bits: usize, out: &mut Vec<u8>) -> Result<()> {
2296 let values = offsets.len() - 1;
2297 let mut run = Vec::with_capacity(TEXT_OFFSET_RUN);
2298 for first in (0..values).step_by(TEXT_OFFSET_RUN) {
2299 let last = (first + TEXT_OFFSET_RUN).min(values);
2300 let base = offsets[first / TEXT_PAYLOAD_VALUES * TEXT_PAYLOAD_VALUES];
2301 run.clear();
2302 run.extend((first..last).map(|value| u64::from(offsets[value + 1] - base)));
2303 bitpack::pack_tail(&run, bits, out)
2304 .map_err(|_| invalid("global dictionary offsets do not pack"))?;
2305 }
2306 Ok(())
2307}
2308
2309fn code_width(values: usize) -> usize {
2311 match u64::try_from(values).unwrap_or(u64::MAX) {
2312 0 | 1 => 0,
2313 last => (u64::BITS - (last - 1).leading_zeros()) as usize,
2314 }
2315}
2316
2317impl TextSource for NativeText {
2318 fn len(&self) -> usize {
2319 self.values
2320 }
2321
2322 fn bytes_at(&self, index: usize) -> Result<Option<&[u8]>> {
2323 if index >= self.values {
2324 return Ok(None);
2325 }
2326 let (start, end) = self.span_within(index)?;
2327 if start == end {
2328 return Ok(Some(&[]));
2329 }
2330 let block = index / TEXT_PAYLOAD_VALUES;
2333 let Some(bytes) = self.payload_block(block)? else { return Ok(None) };
2334 Ok(bytes.get(start as usize..end as usize))
2335 }
2336
2337 fn bytes_len_at(&self, index: usize) -> Result<Option<usize>> {
2338 if index >= self.values {
2339 return Ok(None);
2340 }
2341 let (start, end) = self.span_within(index)?;
2342 Ok(Some((end - start) as usize))
2343 }
2344
2345 fn sweep(
2358 &self,
2359 first: usize,
2360 limit: usize,
2361 body: &mut dyn FnMut(usize, &[u8]) -> Result<()>,
2362 ) -> Result<usize> {
2363 let limit = limit.min(self.values);
2364 if first >= limit {
2365 return Ok(first);
2366 }
2367 let block = first / TEXT_PAYLOAD_VALUES;
2368 let last = ((block + 1) * TEXT_PAYLOAD_VALUES).min(limit);
2369 let decoded;
2370 let bytes: &[u8] = match self.blocks.get(block).and_then(OnceLock::get) {
2371 Some(Ok(kept)) => kept,
2372 _ if self.payload_kept.load(Atomic::Relaxed) < self.keep_budget => {
2373 let kept = self
2374 .payload_block(block)?
2375 .ok_or_else(|| invalid("global dictionary block is past the payload"))?;
2376 self.payload_kept.fetch_add(kept.len(), Atomic::Relaxed);
2377 kept
2378 }
2379 _ => {
2380 decoded = self.decode_block(block)?;
2381 &decoded
2382 }
2383 };
2384 let ends = self.ends_within(first, last)?;
2385 if ends.len() != last - first {
2386 return Err(invalid("global dictionary offsets are short"));
2387 }
2388 let mut start = u64::from(self.start_within(first)?);
2389 for (index, &end) in (first..last).zip(&ends) {
2392 let value = usize::try_from(start)
2393 .ok()
2394 .zip(usize::try_from(end).ok())
2395 .and_then(|(from, to)| bytes.get(from..to))
2396 .ok_or_else(|| invalid("global dictionary value is past its block"))?;
2397 body(index, value)?;
2398 start = end;
2399 }
2400 Ok(last)
2401 }
2402
2403 fn ranks(&self) -> Option<usize> {
2404 (self.ranks > 0).then_some(self.ranks)
2405 }
2406
2407 fn below(&self, ranks: usize, wanted: &[u8]) -> Result<(usize, bool)> {
2415 let mut memo = self.searched.lock().map_err(|_| invalid("a poisoned dictionary search"))?;
2416 if let Some(&answer) = memo.get(wanted) {
2417 return Ok(answer);
2418 }
2419 let answer = search_below(self, ranks, wanted)?;
2420 if memo.len() >= TEXT_SEARCH_MEMO {
2421 memo.clear();
2422 }
2423 memo.insert(wanted.to_vec(), answer);
2424 Ok(answer)
2425 }
2426
2427 fn compare_rank(&self, rank: usize, wanted: &[u8]) -> Result<Ordering> {
2428 let settled = self.head_at(rank)?.cmp(&head(wanted));
2432 if settled != Ordering::Equal {
2433 return Ok(settled);
2434 }
2435 let code = self.code_at_rank(rank)?;
2436 let bytes = self
2437 .bytes_at(code as usize)?
2438 .ok_or_else(|| invalid("global dictionary order names a code it does not have"))?;
2439 Ok(bytes.cmp(wanted))
2440 }
2441
2442 fn code_at_rank(&self, rank: usize) -> Result<u32> {
2443 let (block, within) = self.rank_parts(rank)?;
2444 let codes = self.rank_codes(block, self.rank_block_len(rank))?;
2445 let code = bitpack::tail_at(codes, self.code_bits, within)
2446 .map_err(|_| invalid("global dictionary rank block is short of codes"))?;
2447 let code = u32::try_from(code)
2448 .map_err(|_| invalid("global dictionary order names a code it does not have"))?;
2449 if code as usize >= self.len() {
2450 return Err(invalid("global dictionary order names a code it does not have"));
2451 }
2452 Ok(code)
2453 }
2454
2455 fn code_ranks(&self) -> Option<&[u32]> {
2456 if self.ranks == 0 || self.ranks != self.len() {
2460 return None;
2461 }
2462 self.code_ranks
2463 .get_or_init(|| {
2464 let mut ranks = vec![u32::MAX; self.ranks];
2465 for first in (0..self.ranks).step_by(TEXT_RANK_BLOCK) {
2468 let (block, _) = self.rank_parts(first).ok()?;
2469 let count = self.rank_block_len(first);
2470 let codes = self.rank_codes(block, count).ok()?;
2471 for (within, code) in bitpack::unpack_tail(codes, self.code_bits, count)
2472 .ok()?
2473 .into_iter()
2474 .enumerate()
2475 {
2476 let code = usize::try_from(code).ok()?;
2477 *ranks.get_mut(code)? = u32::try_from(first + within).ok()?;
2478 }
2479 }
2480 if ranks.contains(&u32::MAX) {
2481 return None;
2482 }
2483 Some(ranks)
2484 })
2485 .as_deref()
2486 }
2487
2488 fn footprint(&self) -> usize {
2489 self.offsets.capacity()
2490 + self
2491 .code_ranks
2492 .get()
2493 .and_then(Option::as_ref)
2494 .map_or(0, |ranks| ranks.capacity() * size_of::<u32>())
2495 + self.rank_hashes.capacity() * size_of::<u64>()
2496 + self.rank_ends.capacity() * size_of::<u64>()
2497 + self.rank_blocks.capacity() * size_of::<OnceLock<Result<Vec<u8>>>>()
2498 + self
2499 .rank_blocks
2500 .iter()
2501 .filter_map(OnceLock::get)
2502 .filter_map(|result| result.as_ref().ok())
2503 .map(Vec::capacity)
2504 .sum::<usize>()
2505 + self.blocks.capacity() * size_of::<OnceLock<Result<Vec<u8>>>>()
2506 + self.hashes.capacity() * size_of::<u64>()
2507 + self.ends.capacity() * size_of::<u64>()
2508 + self
2509 .blocks
2510 .iter()
2511 .filter_map(OnceLock::get)
2512 .filter_map(|result| result.as_ref().ok())
2513 .map(Vec::capacity)
2514 .sum::<usize>()
2515 }
2516}
2517
2518fn places(table: &Table) -> Result<Vec<Place>> {
2520 let mut places = Vec::with_capacity(table.stripes.len().saturating_mul(STRIPE_PARTS));
2521 for (at, stripe) in table.stripes.iter().enumerate() {
2522 let index = u32::try_from(at).map_err(|_| invalid("too many stripes"))?;
2523 for (part, &rows) in stripe.parts.iter().enumerate() {
2524 places.push(Place {
2525 stripe: index,
2526 part: u32::try_from(part).map_err(|_| invalid("too many parts in a stripe"))?,
2527 rows,
2528 });
2529 }
2530 }
2531 Ok(places)
2532}
2533
2534fn read_index(file: &File, stripe: &Stripe, column: usize) -> Result<Vec<PartSpan>> {
2539 let parts = stripe.parts.len();
2540 let section = index_section(parts)?;
2541 let at = column.checked_mul(section).ok_or_else(|| invalid("index page offset overflow"))?;
2542 let end = at.checked_add(section).ok_or_else(|| invalid("index page offset overflow"))?;
2543 if end > stripe.index.length as usize {
2544 return Err(invalid("index page is shorter than its columns"));
2545 }
2546 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
2547 let mut bytes = vec![0; section];
2548 let offset = stripe
2549 .index
2550 .offset
2551 .checked_add(at as u64)
2552 .ok_or_else(|| invalid("index page offset overflow"))?;
2553 read_at(file, offset, &mut bytes)?;
2554 let entries = section - size_of::<u64>();
2555 let stored = u64::from_le_bytes(bytes[entries..].try_into().expect("eight bytes"));
2556 if checksum(&bytes[..entries]) != stored {
2557 return Err(invalid(&format!(
2560 "index page section checksum differs, column {column} of {parts} parts at {offset}, \
2561 wanted {stored:016x} and got {:016x}",
2562 checksum(&bytes[..entries]),
2563 )));
2564 }
2565 let mut spans = Vec::with_capacity(parts);
2566 let mut start = 0_usize;
2567 for part in 0..parts {
2568 let at = part * INDEX_ENTRY;
2569 let length = u32::from_le_bytes(bytes[at..at + 4].try_into().expect("four bytes")) as usize;
2570 let hash = u64::from_le_bytes(bytes[at + 4..at + 12].try_into().expect("eight bytes"));
2571 spans.push(PartSpan { start, length, hash });
2572 start = start.checked_add(length).ok_or_else(|| invalid("column page length overflow"))?;
2573 }
2574 if start != page.length as usize {
2575 return Err(invalid("column page length differs from its index"));
2576 }
2577 Ok(spans)
2578}
2579
2580fn part_bytes(page: &[u8], span: PartSpan) -> Result<&[u8]> {
2582 let end = span.start.checked_add(span.length).ok_or_else(|| invalid("part range overflow"))?;
2583 page.get(span.start..end).ok_or_else(|| invalid("part exceeds its column page"))
2584}
2585
2586fn remember(cached: &mut Cached, held: &CachedColumn, kept: usize) {
2591 if let Some(slot) = cached.index.get_mut(held.stripe) {
2592 if slot.is_none() {
2593 *slot = Some(Arc::clone(&held.index));
2594 }
2595 }
2596 let Some(page) = held.page.clone() else { return };
2597 let Some(slot) = cached.pages.get_mut(held.stripe) else { return };
2598 if slot.is_none() {
2599 cached.order.push_back(held.stripe);
2600 }
2601 *slot = Some(page);
2602 while cached.order.len() > kept.max(1) {
2603 let Some(oldest) = cached.order.pop_front() else { break };
2604 if let Some(slot) = cached.pages.get_mut(oldest) {
2605 *slot = None;
2606 }
2607 }
2608}
2609
2610#[derive(Debug, Clone)]
2619pub struct Catalog {
2620 file: Arc<File>,
2621 size: u64,
2622 entries: Arc<Vec<Entry>>,
2623 opening: Opening,
2624}
2625
2626impl Catalog {
2627 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
2633 let (file, size, _, bytes, opening) = slot_bytes(path)?;
2634 let entries = decode_catalog(&bytes, size)?;
2635 Ok(Self { file: Arc::new(file), size, entries: Arc::new(entries), opening })
2636 }
2637
2638 pub fn names(&self) -> impl ExactSizeIterator<Item = &str> {
2640 self.entries.iter().map(|entry| entry.name.as_str())
2641 }
2642
2643 #[must_use]
2645 pub fn len(&self) -> usize {
2646 self.entries.len()
2647 }
2648
2649 #[must_use]
2652 pub fn is_empty(&self) -> bool {
2653 self.entries.is_empty()
2654 }
2655
2656 pub fn table(&self, name: &str) -> Result<Reader> {
2662 let entry = self
2663 .entries
2664 .iter()
2665 .find(|entry| entry.name == name)
2666 .ok_or_else(|| invalid(&format!("the file holds no table called {name}")))?;
2667 let mut bytes = vec![0; entry.directory.length as usize];
2668 read_at(&self.file, entry.directory.offset, &mut bytes)?;
2669 if checksum(&bytes) != entry.directory.hash {
2670 return Err(invalid(&format!("the directory of table {name} does not checksum")));
2671 }
2672 let mut opening = self.opening;
2673 opening.reads += 1;
2674 opening.bytes += u64::from(entry.directory.length);
2675 Reader::build(
2676 Arc::clone(&self.file),
2677 self.size,
2678 decode_directory(&bytes, self.size)?,
2679 u64::from(entry.directory.length),
2680 opening,
2681 )
2682 }
2683}
2684
2685fn slot_offset(generation: u64) -> u64 {
2690 16 + (generation - 1) % 2 * SLOT_BYTES as u64
2691}
2692
2693fn slot_bytes(path: impl AsRef<Path>) -> Result<(File, u64, Slot, Vec<u8>, Opening)> {
2698 let mut file = File::open(path).map_err(io)?;
2699 let size = file.metadata().map_err(io)?.len();
2700 if size < HEADER {
2701 return Err(invalid("file is shorter than its header"));
2702 }
2703 let mut header = [0; HEADER as usize];
2704 file.read_exact(&mut header).map_err(io)?;
2705 let mut opening = Opening { reads: 1, bytes: HEADER };
2706 let version = u32::from_le_bytes([header[8], header[9], header[10], header[11]]);
2707 if &header[..8] != MAGIC {
2712 return Err(invalid("the header does not begin with a rudb native magic"));
2713 }
2714 if version != FORMAT {
2715 return Err(invalid(&format!(
2716 "the file is format {version} and this build reads format {FORMAT}, so it has to \
2717 be written again"
2718 )));
2719 }
2720 let mut selected = None;
2721 for start in [16, 16 + SLOT_BYTES] {
2722 let slot = Slot::read(&header[start..start + SLOT_BYTES]);
2723 if slot.generation == 0 || slot.length == 0 || slot.length as usize > MAX_DIRECTORY {
2724 continue;
2725 }
2726 let Some(end) = slot.offset.checked_add(u64::from(slot.length)) else { continue };
2727 if slot.offset < HEADER || end > size {
2728 continue;
2729 }
2730 let mut bytes = vec![0; slot.length as usize];
2731 file.seek(SeekFrom::Start(slot.offset)).map_err(io)?;
2732 file.read_exact(&mut bytes).map_err(io)?;
2733 opening.reads += 1;
2734 opening.bytes += u64::from(slot.length);
2735 if checksum(&bytes) == slot.hash
2736 && selected
2737 .as_ref()
2738 .is_none_or(|(old, _): &(Slot, Vec<u8>)| old.generation < slot.generation)
2739 {
2740 selected = Some((slot, bytes));
2741 }
2742 }
2743 let (slot, bytes) = selected.ok_or_else(|| invalid("no committed directory slot is valid"))?;
2744 Ok((file, size, slot, bytes, opening))
2745}
2746
2747impl Reader {
2748 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
2755 let catalog = Catalog::open(path)?;
2756 let mut names = catalog.names();
2757 let name = names.next().ok_or_else(|| invalid("the file holds no table"))?.to_string();
2758 if names.next().is_some() {
2759 return Err(invalid(
2760 "the file holds more than one table, so it has to be opened by name",
2761 ));
2762 }
2763 catalog.table(&name)
2764 }
2765
2766 fn build(
2768 file: Arc<File>,
2769 size: u64,
2770 table: Table,
2771 directory: u64,
2772 opening: Opening,
2773 ) -> Result<Self> {
2774 let places = places(&table)?;
2775 let dictionaries = (0..table.fields.len()).map(|_| OnceLock::new()).collect();
2776 let table_fields = table.fields.len();
2777 let stripes = table.stripes.len();
2778 let cache = (0..table.fields.len())
2779 .map(|_| {
2780 Mutex::new(Cached {
2781 pages: (0..stripes).map(|_| None).collect(),
2782 index: (0..stripes).map(|_| None).collect(),
2783 ..Cached::default()
2784 })
2785 })
2786 .collect::<Vec<_>>();
2787 let sieves: Vec<Vec<SieveSlot>> = (0..table.fields.len())
2788 .map(|_| table.stripes.iter().map(|_| OnceLock::new()).collect())
2789 .collect();
2790 let part_ranges: Vec<Vec<RangeSlot>> = (0..table.fields.len())
2791 .map(|_| table.stripes.iter().map(|_| OnceLock::new()).collect())
2792 .collect();
2793 Ok(Self {
2794 file,
2795 table: Arc::new(table),
2796 dictionaries: Arc::new(dictionaries),
2797 loading: Arc::new((0..table_fields).map(|_| Mutex::new(())).collect()),
2798 opened: Arc::new(AtomicUsize::new(0)),
2799 sieves: Arc::new(sieves),
2800 part_ranges: Arc::new(part_ranges),
2801 places: Arc::new(places),
2802 cache: Arc::new(cache),
2803 pages: Arc::new(AtomicUsize::new(0)),
2804 indexes: Arc::new(AtomicUsize::new(0)),
2805 kept: Arc::new(AtomicUsize::new(CACHED_STRIPES_PER_COLUMN)),
2806 size,
2807 directory,
2808 opening,
2809 })
2810 }
2811
2812 #[must_use]
2819 pub fn reads(&self) -> Reads {
2820 Reads {
2821 opening: self.opening,
2822 pages: self.pages.load(Atomic::Relaxed),
2823 indexes: self.indexes.load(Atomic::Relaxed),
2824 dictionaries: self.opened.load(Atomic::Relaxed),
2825 }
2826 }
2827
2828 #[must_use]
2833 pub fn layout(&self) -> Layout {
2834 let table = &self.table;
2835 let stripes = table.stripes.as_slice();
2836 let columns = table
2837 .fields
2838 .iter()
2839 .enumerate()
2840 .map(|(at, field)| ColumnLayout {
2841 name: field.name.clone(),
2842 kind: field.ty.to_string(),
2843 pages: sum(stripes.iter().map(|stripe| span_bytes(&stripe.pages, at))),
2844 memberships: sum(stripes.iter().map(|stripe| page_bytes(&stripe.memberships, at))),
2845 sieves: sum(stripes.iter().map(|stripe| page_bytes(&stripe.sieves, at))),
2846 part_ranges: sum(stripes.iter().map(|stripe| page_bytes(&stripe.part_ranges, at))),
2847 dictionary: page_bytes(&table.dictionaries, at),
2848 })
2849 .collect();
2850 Layout {
2851 file: self.size,
2852 rows: table.rows,
2853 stripes: stripes.len(),
2854 parts: self.places.len(),
2855 columns,
2856 indexes: sum(stripes.iter().map(|stripe| u64::from(stripe.index.length))),
2857 directory: self.directory,
2858 header: HEADER,
2859 }
2860 }
2861
2862 pub fn stored(&self, column: usize) -> Result<Vec<StoredPart>> {
2879 let field = self
2880 .table
2881 .fields
2882 .get(column)
2883 .ok_or_else(|| invalid("stored column index out of range"))?;
2884 let mut stored = Vec::with_capacity(self.places.len());
2885 let mut row = 0;
2886 for (at, stripe) in self.table.stripes.iter().enumerate() {
2887 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
2888 let index = read_index(&self.file, stripe, column)?;
2889 let mut bytes = vec![0; page.length as usize];
2890 read_at(&self.file, page.offset, &mut bytes)?;
2891 let ranges = self.stripe_part_ranges(at, column);
2892 for (part, &rows) in stripe.parts.iter().enumerate() {
2893 let span = *index.get(part).ok_or_else(|| invalid("part index out of range"))?;
2894 let held = part_bytes(&bytes, span)?;
2895 let range = ranges.and_then(|held| held.get(part));
2896 stored.push(StoredPart {
2897 stripe: at,
2898 part,
2899 row,
2900 rows: rows as usize,
2901 encoding: page_encoding(&field.ty, rows as usize, held),
2902 bytes: span.length as u64,
2903 page: page.offset,
2904 offset: span.start as u64,
2905 low: range
2906 .and_then(|range| range.low.clone())
2907 .and_then(|bound| bound.into_value(&field.ty)),
2908 high: range
2909 .and_then(|range| range.high.clone())
2910 .and_then(|bound| bound.into_value(&field.ty)),
2911 nulls: range.map(|range| range.nulls),
2912 });
2913 row += rows as usize;
2914 }
2915 }
2916 Ok(stored)
2917 }
2918
2919 #[must_use]
2921 pub fn parts(&self) -> usize {
2922 self.places.len()
2923 }
2924
2925 #[must_use]
2932 pub fn stripe_parts(&self) -> Vec<std::ops::Range<usize>> {
2933 let mut runs = Vec::with_capacity(self.table.stripes.len());
2934 let mut start = 0;
2935 for stripe in &self.table.stripes {
2936 let end = start + stripe.parts.len();
2937 runs.push(start..end);
2938 start = end;
2939 }
2940 runs
2941 }
2942
2943 #[must_use]
2948 pub fn stripe_rows(&self, stripe: usize) -> usize {
2949 self.table.stripes.get(stripe).map_or(0, |held| held.rows)
2950 }
2951
2952 pub fn keep_stripes(&self, stripes: usize) {
2959 self.kept.fetch_max(stripes, Atomic::Relaxed);
2960 }
2961
2962 #[must_use]
2964 pub fn part_rows(&self, at: usize) -> usize {
2965 self.places.get(at).map_or(0, |place| place.rows as usize)
2966 }
2967
2968 #[must_use]
2970 pub fn table(&self) -> &Table {
2971 &self.table
2972 }
2973
2974 pub fn top_frequencies(&self, column: usize, top: usize) -> Result<Option<Vec<(Value, u64)>>> {
2983 let field = self
2984 .table
2985 .fields
2986 .get(column)
2987 .ok_or_else(|| invalid("frequency column index out of range"))?;
2988 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2989 return Ok(None);
2990 };
2991 if top == 0 || summary.entries.len() < top {
2992 return Ok(None);
2993 }
2994 let boundary = summary.entries[top - 1].count;
2995 if boundary <= summary.omitted_max {
2996 return Ok(None);
2997 }
2998 self.decode_frequencies(column, &field.ty, &summary.entries).map(Some)
2999 }
3000
3001 pub fn exact_frequencies(&self, column: usize) -> Result<Option<Vec<(Value, u64)>>> {
3021 let Some(prefix) = self.frequency_prefix(column)? else {
3022 return Ok(None);
3023 };
3024 Ok((prefix.omitted_max == 0).then_some(prefix.entries))
3025 }
3026
3027 pub fn frequency_prefix(&self, column: usize) -> Result<Option<FrequencyPrefix>> {
3050 let field = self
3051 .table
3052 .fields
3053 .get(column)
3054 .ok_or_else(|| invalid("frequency column index out of range"))?;
3055 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
3056 return Ok(None);
3057 };
3058 let entries = self.decode_frequencies(column, &field.ty, &summary.entries)?;
3059 Ok(Some(FrequencyPrefix { entries, omitted_max: summary.omitted_max }))
3060 }
3061
3062 fn decode_frequencies(
3064 &self,
3065 column: usize,
3066 ty: &LogicalType,
3067 entries: &[FrequencyEntry],
3068 ) -> Result<Vec<(Value, u64)>> {
3069 let dictionary = if *ty == LogicalType::Varchar { self.dictionary(column)? } else { None };
3070 let mut out = Vec::with_capacity(entries.len());
3071 for entry in entries {
3072 let value = match entry.value {
3073 FrequencyValue::Null => Value::Null,
3074 FrequencyValue::Integer(value) => match *ty {
3075 LogicalType::TinyInt => Value::TinyInt(
3076 i8::try_from(value)
3077 .map_err(|_| invalid("frequency TINYINT is out of range"))?,
3078 ),
3079 LogicalType::UTinyInt => Value::UTinyInt(
3080 u8::try_from(value)
3081 .map_err(|_| invalid("frequency UTINYINT is out of range"))?,
3082 ),
3083 LogicalType::USmallInt => Value::USmallInt(
3084 u16::try_from(value)
3085 .map_err(|_| invalid("frequency USMALLINT is out of range"))?,
3086 ),
3087 LogicalType::UInteger => Value::UInteger(
3088 u32::try_from(value)
3089 .map_err(|_| invalid("frequency UINTEGER is out of range"))?,
3090 ),
3091 LogicalType::UBigInt => Value::UBigInt(
3092 u64::try_from(value)
3093 .map_err(|_| invalid("frequency UBIGINT is out of range"))?,
3094 ),
3095 LogicalType::SmallInt => Value::SmallInt(
3096 i16::try_from(value)
3097 .map_err(|_| invalid("frequency SMALLINT is out of range"))?,
3098 ),
3099 LogicalType::Integer => Value::Integer(
3100 i32::try_from(value)
3101 .map_err(|_| invalid("frequency INTEGER is out of range"))?,
3102 ),
3103 LogicalType::BigInt => Value::BigInt(
3104 i64::try_from(value)
3105 .map_err(|_| invalid("frequency BIGINT is out of range"))?,
3106 ),
3107 LogicalType::Date => Value::Date(
3108 i32::try_from(value)
3109 .map_err(|_| invalid("frequency DATE is out of range"))?,
3110 ),
3111 LogicalType::Timestamp => Value::Timestamp(
3112 i64::try_from(value)
3113 .map_err(|_| invalid("frequency TIMESTAMP is out of range"))?,
3114 ),
3115 _ => return Err(invalid("integer frequency belongs to another type")),
3116 },
3117 FrequencyValue::Code(code) => dictionary
3118 .as_ref()
3119 .ok_or_else(|| invalid("frequency code has no dictionary"))?
3120 .try_value_at(code as usize)?,
3121 };
3122 out.push((value, entry.count));
3123 }
3124 Ok(out)
3125 }
3126
3127 pub fn frequency_occurrences(&self, column: usize) -> Result<Option<FrequencyOccurrences>> {
3137 self.table
3138 .fields
3139 .get(column)
3140 .ok_or_else(|| invalid("frequency column index out of range"))?;
3141 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
3142 return Ok(None);
3143 };
3144 if summary.ordinals.is_empty() {
3145 return Ok(None);
3146 }
3147 Ok(Some(FrequencyOccurrences {
3148 omitted_max: summary.omitted_max,
3149 ordinals: summary.ordinals.clone(),
3150 }))
3151 }
3152
3153 pub fn distinct_values(&self, column: usize) -> Result<Option<u64>> {
3177 self.table
3178 .distincts
3179 .get(column)
3180 .copied()
3181 .ok_or_else(|| invalid("distinct column index out of range"))
3182 }
3183
3184 pub fn null_count(&self, column: usize) -> Result<u64> {
3195 if column >= self.table.fields.len() {
3196 return Err(invalid("null count column index out of range"));
3197 }
3198 let mut nulls = 0_u64;
3199 for stripe in &self.table.stripes {
3200 let range = stripe
3201 .zone
3202 .column(column)
3203 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
3204 nulls = nulls
3205 .checked_add(range.nulls as u64)
3206 .ok_or_else(|| invalid("null count overflow"))?;
3207 }
3208 Ok(nulls)
3209 }
3210
3211 pub fn text_extremes(&self, column: usize) -> Result<Option<(Value, Value)>> {
3226 if self.null_count(column)? > 0 {
3227 return Ok(None);
3228 }
3229 let Some(dictionary) = self.dictionary(column)? else { return Ok(None) };
3230 let Some(ranks) = dictionary.ranks() else { return Ok(None) };
3231 if ranks == 0 {
3232 return Ok(None);
3233 }
3234 let low = text_at_rank(&dictionary, 0)?;
3235 let high = text_at_rank(&dictionary, ranks - 1)?;
3236 Ok(Some((low, high)))
3237 }
3238
3239 pub fn exact_extremes(&self, column: usize) -> Result<Option<(Bound, Bound)>> {
3262 if column >= self.table.fields.len() {
3263 return Err(invalid("extremes column index out of range"));
3264 }
3265 let mut low: Option<Bound> = None;
3266 let mut high: Option<Bound> = None;
3267 for stripe in &self.table.stripes {
3268 let range = stripe
3269 .zone
3270 .column(column)
3271 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
3272 if !range.exact {
3273 return Ok(None);
3274 }
3275 let (Some(small), Some(large)) = (range.low.as_ref(), range.high.as_ref()) else {
3280 if stripe.rows > range.nulls {
3281 return Ok(None);
3282 }
3283 continue;
3284 };
3285 low = Some(low.map_or_else(|| small.clone(), |held| held.smaller(small.clone())));
3286 high = Some(high.map_or_else(|| large.clone(), |held| held.larger(large.clone())));
3287 }
3288 Ok(low.zip(high))
3289 }
3290
3291 pub fn exact_sum(&self, column: usize) -> Result<Option<(i128, u64)>> {
3304 if column >= self.table.fields.len() {
3305 return Err(invalid("sum column index out of range"));
3306 }
3307 let mut total = 0_i128;
3308 let mut rows = 0_u64;
3309 for stripe in &self.table.stripes {
3310 let range = stripe
3311 .zone
3312 .column(column)
3313 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
3314 let Some(part) = range.sum else { return Ok(None) };
3315 let Some(sum) = total.checked_add(part) else { return Ok(None) };
3316 total = sum;
3317 rows = rows.saturating_add(stripe.rows as u64 - range.nulls as u64);
3318 }
3319 Ok(Some((total, rows)))
3320 }
3321
3322 fn dictionary(&self, column: usize) -> Result<Option<Arc<Vector>>> {
3331 let Some(page) = self.table.dictionaries[column] else { return Ok(None) };
3332 if let Some(dictionary) = self.dictionaries[column].get() {
3333 return Ok(Some(Arc::clone(dictionary)));
3334 }
3335 let _queued = self.loading[column].lock().map_err(|_| invalid("a poisoned dictionary"))?;
3336 if let Some(dictionary) = self.dictionaries[column].get() {
3337 return Ok(Some(Arc::clone(dictionary)));
3338 }
3339 self.opened.fetch_add(1, Atomic::Relaxed);
3340 let dictionary = Arc::new(open_global_dictionary(
3341 Arc::clone(&self.file),
3342 page,
3343 &self.table.fields[column].ty,
3344 TEXT_KEEP_BUDGET,
3345 )?);
3346 let _ = self.dictionaries[column].set(Arc::clone(&dictionary));
3347 Ok(Some(dictionary))
3348 }
3349
3350 pub fn read(&self, part: usize, columns: &[usize]) -> Result<Chunk> {
3359 self.read_impl(part, columns, true)
3360 }
3361
3362 pub fn read_sparse(&self, part: usize, columns: &[usize]) -> Result<Chunk> {
3372 self.read_impl(part, columns, false)
3373 }
3374
3375 pub fn skips_codes(&self, part: usize, column: usize, candidates: &[u32]) -> Result<bool> {
3382 if candidates.is_empty() {
3383 return Ok(true);
3384 }
3385 if candidates.windows(2).any(|pair| pair[0] >= pair[1]) {
3386 return Err(Error::internal("native code candidates are not sorted and unique"));
3387 }
3388 let stripe = self.stripe_of(part)?;
3389 let Some(page) = stripe.memberships.get(column).copied().flatten() else {
3390 return Ok(false);
3391 };
3392 let mut bytes = vec![0; page.length as usize];
3393 read_at(&self.file, page.offset, &mut bytes)?;
3394 if checksum(&bytes) != page.hash {
3395 return Err(invalid("membership page checksum differs"));
3396 }
3397 let codes = decode_membership(&bytes)?;
3398 let mut left = 0;
3399 let mut right = 0;
3400 while left < codes.len() && right < candidates.len() {
3401 match codes[left].cmp(&candidates[right]) {
3402 Ordering::Less => left += 1,
3403 Ordering::Greater => right += 1,
3404 Ordering::Equal => return Ok(false),
3405 }
3406 }
3407 Ok(true)
3408 }
3409
3410 fn stripe_of(&self, part: usize) -> Result<&Stripe> {
3411 let place = self.places.get(part).ok_or_else(|| invalid("part index out of range"))?;
3412 self.table
3413 .stripes
3414 .get(place.stripe as usize)
3415 .ok_or_else(|| invalid("stripe index out of range"))
3416 }
3417
3418 fn held(&self, at: usize, stripe: &Stripe, column: usize, whole: bool) -> Result<CachedColumn> {
3435 let cache = self.cache.get(column).ok_or_else(|| invalid("column index out of range"))?;
3436 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3437 let known = cached.index.get(at).and_then(Clone::clone);
3438 let page = cached.pages.get(at).and_then(Clone::clone);
3439 if let Some(index) = known.clone() {
3440 if !whole || page.is_some() {
3441 return Ok(CachedColumn { stripe: at, index, page });
3442 }
3443 }
3444 if cached.loading.contains(&at) {
3445 drop(cached);
3446 if let Some(index) = known {
3450 return Ok(CachedColumn { stripe: at, index, page: None });
3451 }
3452 let held = self.page_of(stripe, column, at, false, None)?;
3453 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3454 remember(&mut cached, &held, self.kept.load(Atomic::Relaxed));
3455 return Ok(held);
3456 }
3457 cached.loading.push(at);
3458 drop(cached);
3459
3460 let read = self.page_of(stripe, column, at, whole, known);
3461
3462 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3466 if let Some(position) = cached.loading.iter().position(|loading| *loading == at) {
3467 cached.loading.remove(position);
3468 }
3469 let held = read?;
3470 remember(&mut cached, &held, self.kept.load(Atomic::Relaxed));
3471 Ok(held)
3472 }
3473
3474 fn page_of(
3480 &self,
3481 stripe: &Stripe,
3482 column: usize,
3483 at: usize,
3484 whole: bool,
3485 known: Option<Arc<Vec<PartSpan>>>,
3486 ) -> Result<CachedColumn> {
3487 let index = match known {
3488 Some(index) => index,
3489 None => {
3490 self.indexes.fetch_add(1, Atomic::Relaxed);
3491 Arc::new(read_index(&self.file, stripe, column)?)
3492 }
3493 };
3494 let page = if whole {
3495 self.pages.fetch_add(1, Atomic::Relaxed);
3496 let span = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
3497 let mut bytes = vec![0; span.length as usize];
3498 read_at(&self.file, span.offset, &mut bytes)?;
3499 Some(Arc::new(bytes))
3500 } else {
3501 None
3502 };
3503 Ok(CachedColumn { stripe: at, index, page })
3504 }
3505
3506 fn read_impl(&self, at: usize, columns: &[usize], whole: bool) -> Result<Chunk> {
3507 let place = *self.places.get(at).ok_or_else(|| invalid("part index out of range"))?;
3508 let index = place.stripe as usize;
3509 let stripe =
3510 self.table.stripes.get(index).ok_or_else(|| invalid("stripe index out of range"))?;
3511 let rows = place.rows as usize;
3512 let mut picked = Vec::with_capacity(columns.len());
3513 for &column in columns {
3514 let field = self
3515 .table
3516 .fields
3517 .get(column)
3518 .ok_or_else(|| invalid("column index out of range"))?;
3519 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
3520 let held = self.held(index, stripe, column, whole)?;
3521 let span = *held
3522 .index
3523 .get(place.part as usize)
3524 .ok_or_else(|| invalid("part index out of range"))?;
3525 let owned;
3526 let bytes = match &held.page {
3527 Some(held) => part_bytes(held, span)?,
3528 None => {
3529 let offset = page
3530 .offset
3531 .checked_add(span.start as u64)
3532 .ok_or_else(|| invalid("part range overflow"))?;
3533 let mut bytes = vec![0; span.length];
3534 read_at(&self.file, offset, &mut bytes)?;
3535 owned = bytes;
3536 &owned
3537 }
3538 };
3539 if checksum(bytes) != span.hash {
3540 return Err(invalid(&format!(
3541 "column page checksum differs, column {column} part {} at {}+{} of {} bytes, \
3542 wanted {:016x} and got {:016x}",
3543 place.part,
3544 page.offset,
3545 span.start,
3546 span.length,
3547 span.hash,
3548 checksum(bytes),
3549 )));
3550 }
3551 let dictionary = self.dictionary(column)?;
3552 picked.push(decode(&field.ty, rows, bytes, dictionary)?.into_pages());
3558 }
3559 Chunk::with_rows(picked, rows)
3560 }
3561
3562 #[must_use]
3578 pub fn skips(&self, part: usize, probes: &[Probe]) -> bool {
3579 let Some(place) = self.places.get(part).copied() else { return false };
3580 let Some(stripe) = self.table.stripes.get(place.stripe as usize) else { return false };
3581 if stripe.zone.skips(probes) {
3582 return true;
3583 }
3584 probes.iter().any(|probe| self.outside(place, probe) || self.sifted(place, probe))
3585 }
3586
3587 fn outside(&self, place: Place, probe: &Probe) -> bool {
3593 match self.stripe_part_ranges(place.stripe as usize, probe.column) {
3594 Some(ranges) => ranges
3595 .get(place.part as usize)
3596 .is_some_and(|range| range.excludes(probe.op, &probe.value)),
3597 None => false,
3598 }
3599 }
3600
3601 fn stripe_part_ranges(&self, stripe: usize, column: usize) -> Option<&[Range]> {
3607 let slot = self.part_ranges.get(column)?.get(stripe)?;
3608 if let Some(held) = slot.get() {
3609 return Some(held);
3610 }
3611 let page = self.table.stripes.get(stripe)?.part_ranges.get(column).copied().flatten()?;
3612 let mut bytes = vec![0; page.length as usize];
3613 read_at(&self.file, page.offset, &mut bytes).ok()?;
3614 if checksum(&bytes) != page.hash {
3615 return None;
3616 }
3617 let ranges = Arc::new(decode_part_ranges(&bytes).ok()?);
3618 let _ = slot.set(ranges);
3619 slot.get().map(|held| held.as_slice())
3620 }
3621
3622 #[must_use]
3639 pub fn certain(&self, part: usize, probes: &[Probe]) -> bool {
3640 let Some(place) = self.places.get(part).copied() else { return false };
3641 let Some(stripe) = self.table.stripes.get(place.stripe as usize) else { return false };
3642 if stripe.zone.certain(probes) {
3643 return true;
3644 }
3645 probes
3646 .iter()
3647 .all(|probe| stripe.zone.certain(slice::from_ref(probe)) || self.inside(place, probe))
3648 }
3649
3650 fn inside(&self, place: Place, probe: &Probe) -> bool {
3656 match self.stripe_part_ranges(place.stripe as usize, probe.column) {
3657 Some(ranges) => ranges
3658 .get(place.part as usize)
3659 .is_some_and(|range| range.certain(probe.op, &probe.value)),
3660 None => false,
3661 }
3662 }
3663
3664 #[must_use]
3675 pub fn stripe_skips(&self, stripe: usize, probes: &[Probe]) -> bool {
3676 self.table.stripes.get(stripe).is_some_and(|held| held.zone.skips(probes))
3677 }
3678
3679 fn sifted(&self, place: Place, probe: &Probe) -> bool {
3685 if probe.op != Op::Equal {
3686 return false;
3687 }
3688 match self.stripe_sieves(place.stripe as usize, probe.column) {
3689 Some(sieves) => sieves
3690 .get(place.part as usize)
3691 .and_then(Option::as_ref)
3692 .is_some_and(|sieve| sieve.excludes(&probe.value)),
3693 None => false,
3694 }
3695 }
3696
3697 fn stripe_sieves(&self, stripe: usize, column: usize) -> Option<&[Option<Sieve>]> {
3704 let slot = self.sieves.get(column)?.get(stripe)?;
3705 if let Some(held) = slot.get() {
3706 return Some(held);
3707 }
3708 let page = self.table.stripes.get(stripe)?.sieves.get(column).copied().flatten()?;
3709 let mut bytes = vec![0; page.length as usize];
3710 read_at(&self.file, page.offset, &mut bytes).ok()?;
3711 if checksum(&bytes) != page.hash {
3712 return None;
3713 }
3714 let sieves = Arc::new(decode_sieves(&bytes).ok()?);
3715 let _ = slot.set(sieves);
3716 slot.get().map(|held| held.as_slice())
3717 }
3718}
3719
3720fn text_at_rank(dictionary: &Vector, rank: usize) -> Result<Value> {
3722 let code = dictionary.code_at_rank(rank)? as usize;
3723 let text = dictionary
3724 .try_text_at(code)?
3725 .ok_or_else(|| invalid("global dictionary order names a code it does not have"))?;
3726 Ok(Value::Varchar(text.into()))
3727}
3728
3729#[cfg(unix)]
3734fn write_at(file: &File, mut offset: u64, mut bytes: &[u8]) -> Result<()> {
3735 use std::os::unix::fs::FileExt;
3736 while !bytes.is_empty() {
3737 let written = file.write_at(bytes, offset).map_err(io)?;
3738 if written == 0 {
3739 return Err(invalid("a write to the native file wrote nothing"));
3740 }
3741 offset += written as u64;
3742 bytes = &bytes[written..];
3743 }
3744 Ok(())
3745}
3746
3747#[cfg(windows)]
3749fn write_at(file: &File, mut offset: u64, mut bytes: &[u8]) -> Result<()> {
3750 use std::os::windows::fs::FileExt;
3751 while !bytes.is_empty() {
3752 let written = file.seek_write(bytes, offset).map_err(io)?;
3753 if written == 0 {
3754 return Err(invalid("a write to the native file wrote nothing"));
3755 }
3756 offset += written as u64;
3757 bytes = &bytes[written..];
3758 }
3759 Ok(())
3760}
3761
3762#[cfg(not(any(unix, windows)))]
3764fn write_at(file: &File, offset: u64, bytes: &[u8]) -> Result<()> {
3765 use std::io::Write;
3766 let mut file = file.try_clone().map_err(io)?;
3767 file.seek(SeekFrom::Start(offset)).map_err(io)?;
3768 file.write_all(bytes).map_err(io)
3769}
3770
3771#[cfg(unix)]
3781fn read_at(file: &File, mut offset: u64, mut bytes: &mut [u8]) -> Result<()> {
3782 use std::os::unix::fs::FileExt;
3783 while !bytes.is_empty() {
3784 let read = file.read_at(bytes, offset).map_err(io)?;
3785 if read == 0 {
3786 return Err(invalid("column page ends before its declared length"));
3787 }
3788 offset += read as u64;
3789 bytes = &mut bytes[read..];
3790 }
3791 Ok(())
3792}
3793
3794#[cfg(windows)]
3800fn read_at(file: &File, mut offset: u64, mut bytes: &mut [u8]) -> Result<()> {
3801 use std::os::windows::fs::FileExt;
3802 while !bytes.is_empty() {
3803 let read = file.seek_read(bytes, offset).map_err(io)?;
3804 if read == 0 {
3805 return Err(invalid("column page ends before its declared length"));
3806 }
3807 offset += read as u64;
3808 bytes = &mut bytes[read..];
3809 }
3810 Ok(())
3811}
3812
3813#[cfg(not(any(unix, windows)))]
3818fn read_at(file: &File, offset: u64, bytes: &mut [u8]) -> Result<()> {
3819 let mut file = file.try_clone().map_err(io)?;
3820 file.seek(SeekFrom::Start(offset)).map_err(io)?;
3821 file.read_exact(bytes).map_err(io)
3822}
3823
3824fn type_tag(ty: &LogicalType) -> Result<u8> {
3831 match ty {
3832 LogicalType::SmallInt => Ok(1),
3833 LogicalType::Integer => Ok(2),
3834 LogicalType::BigInt => Ok(3),
3835 LogicalType::Varchar => Ok(4),
3836 LogicalType::Date => Ok(5),
3837 LogicalType::Timestamp => Ok(6),
3838 LogicalType::Boolean => Ok(7),
3839 LogicalType::TinyInt => Ok(8),
3840 LogicalType::UTinyInt => Ok(9),
3841 LogicalType::USmallInt => Ok(10),
3842 LogicalType::UInteger => Ok(11),
3843 LogicalType::UBigInt => Ok(12),
3844 LogicalType::Decimal { .. } => Ok(13),
3845 LogicalType::Float => Ok(14),
3846 LogicalType::Double => Ok(15),
3847 LogicalType::HugeInt => Ok(16),
3848 LogicalType::UHugeInt => Ok(17),
3849 LogicalType::Time => Ok(18),
3850 LogicalType::TimeTz => Ok(19),
3851 LogicalType::TimestampTz => Ok(20),
3852 LogicalType::Interval => Ok(21),
3853 LogicalType::Uuid => Ok(22),
3854 LogicalType::Blob => Ok(23),
3855 LogicalType::Bit => Ok(24),
3856 LogicalType::TimestampS => Ok(25),
3857 LogicalType::TimestampMs => Ok(26),
3858 LogicalType::TimestampNs => Ok(27),
3859 _ => Err(Error::not_implemented(format!("native storage for {ty}"))),
3860 }
3861}
3862
3863fn put_type(out: &mut Vec<u8>, ty: &LogicalType) -> Result<()> {
3869 out.push(type_tag(ty)?);
3870 if let LogicalType::Decimal { width, scale } = ty {
3871 out.push(*width);
3872 out.push(*scale);
3873 }
3874 Ok(())
3875}
3876
3877fn read_type(cur: &mut Cursor<'_>) -> Result<LogicalType> {
3879 let tag = cur.u8()?;
3880 if tag == 13 {
3881 let width = cur.u8()?;
3882 let scale = cur.u8()?;
3883 return LogicalType::decimal(width, scale)
3884 .map_err(|_| invalid("decimal column width and scale are not a decimal"));
3885 }
3886 tag_type(tag)
3887}
3888
3889fn tag_type(tag: u8) -> Result<LogicalType> {
3890 match tag {
3891 1 => Ok(LogicalType::SmallInt),
3892 2 => Ok(LogicalType::Integer),
3893 3 => Ok(LogicalType::BigInt),
3894 4 => Ok(LogicalType::Varchar),
3895 5 => Ok(LogicalType::Date),
3896 6 => Ok(LogicalType::Timestamp),
3897 7 => Ok(LogicalType::Boolean),
3898 8 => Ok(LogicalType::TinyInt),
3899 9 => Ok(LogicalType::UTinyInt),
3900 10 => Ok(LogicalType::USmallInt),
3901 11 => Ok(LogicalType::UInteger),
3902 12 => Ok(LogicalType::UBigInt),
3903 14 => Ok(LogicalType::Float),
3904 15 => Ok(LogicalType::Double),
3905 16 => Ok(LogicalType::HugeInt),
3906 17 => Ok(LogicalType::UHugeInt),
3907 18 => Ok(LogicalType::Time),
3908 19 => Ok(LogicalType::TimeTz),
3909 20 => Ok(LogicalType::TimestampTz),
3910 21 => Ok(LogicalType::Interval),
3911 22 => Ok(LogicalType::Uuid),
3912 23 => Ok(LogicalType::Blob),
3913 24 => Ok(LogicalType::Bit),
3914 25 => Ok(LogicalType::TimestampS),
3915 26 => Ok(LogicalType::TimestampMs),
3916 27 => Ok(LogicalType::TimestampNs),
3917 _ => Err(invalid("column type tag is unknown")),
3918 }
3919}
3920
3921fn put_u16(out: &mut Vec<u8>, value: u16) {
3922 out.extend_from_slice(&value.to_le_bytes());
3923}
3924fn put_u32(out: &mut Vec<u8>, value: u32) {
3925 out.extend_from_slice(&value.to_le_bytes());
3926}
3927fn put_u64(out: &mut Vec<u8>, value: u64) {
3928 out.extend_from_slice(&value.to_le_bytes());
3929}
3930fn put_var_u64(out: &mut Vec<u8>, mut value: u64) {
3931 while value >= 0x80 {
3932 out.push((value as u8 & 0x7f) | 0x80);
3933 value >>= 7;
3934 }
3935 out.push(value as u8);
3936}
3937
3938fn frequency_order(left: FrequencyValue, right: FrequencyValue) -> Ordering {
3939 match (left, right) {
3940 (FrequencyValue::Null, FrequencyValue::Null) => Ordering::Equal,
3941 (FrequencyValue::Null, _) => Ordering::Less,
3942 (_, FrequencyValue::Null) => Ordering::Greater,
3943 (FrequencyValue::Integer(left), FrequencyValue::Integer(right)) => left.cmp(&right),
3944 (FrequencyValue::Code(left), FrequencyValue::Code(right)) => left.cmp(&right),
3945 (FrequencyValue::Integer(_), FrequencyValue::Code(_)) => Ordering::Less,
3946 (FrequencyValue::Code(_), FrequencyValue::Integer(_)) => Ordering::Greater,
3947 }
3948}
3949
3950fn code_frequency(dictionary: &GlobalDictionary) -> FrequencySummary {
3951 let mut entries = dictionary
3952 .counts
3953 .iter()
3954 .enumerate()
3955 .filter(|(_, count)| **count != 0)
3956 .map(|(code, &count)| FrequencyEntry { value: FrequencyValue::Code(code as u32), count })
3957 .collect::<Vec<_>>();
3958 if dictionary.nulls != 0 {
3959 entries.push(FrequencyEntry { value: FrequencyValue::Null, count: dictionary.nulls });
3960 }
3961 entries.sort_unstable_by(|left, right| {
3962 right.count.cmp(&left.count).then_with(|| frequency_order(left.value, right.value))
3963 });
3964 let omitted_max = entries.get(FREQUENCY_ENTRIES).map_or(0, |entry| entry.count);
3965 entries.truncate(FREQUENCY_ENTRIES);
3966 FrequencySummary { entries, omitted_max, ordinals: Vec::new() }
3967}
3968
3969fn encode_directory(table: &Table) -> Result<Vec<u8>> {
3970 let mut out = DIRECTORY.to_vec();
3971 let name = table.name.as_bytes();
3972 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("table name too long"))?);
3973 out.extend_from_slice(name);
3974 put_u16(&mut out, u16::try_from(table.fields.len()).map_err(|_| invalid("too many columns"))?);
3975 for field in &table.fields {
3976 let name = field.name.as_bytes();
3977 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("column name too long"))?);
3978 out.extend_from_slice(name);
3979 put_type(&mut out, &field.ty)?;
3980 out.push(u8::from(field.not_null));
3981 }
3982 for dictionary in &table.dictionaries {
3983 match dictionary {
3984 None => out.push(0),
3985 Some(page) => {
3986 out.push(1);
3987 put_u64(&mut out, page.offset);
3988 put_u32(&mut out, page.length);
3989 put_u64(&mut out, page.hash);
3990 }
3991 }
3992 }
3993 for distinct in &table.distincts {
3994 match distinct {
3995 None => out.push(0),
3996 Some(count) => {
3997 out.push(1);
3998 put_u64(&mut out, *count);
3999 }
4000 }
4001 }
4002 put_u64(&mut out, u64::try_from(table.rows).map_err(|_| invalid("row count overflow"))?);
4003 put_u32(&mut out, u32::try_from(table.stripes.len()).map_err(|_| invalid("too many stripes"))?);
4004 for stripe in &table.stripes {
4005 put_u32(
4006 &mut out,
4007 u32::try_from(stripe.parts.len()).map_err(|_| invalid("too many parts in a stripe"))?,
4008 );
4009 for &rows in &stripe.parts {
4010 put_u32(&mut out, rows);
4011 }
4012 put_u64(&mut out, stripe.index.offset);
4013 put_u32(&mut out, stripe.index.length);
4014 for page in &stripe.pages {
4015 put_u64(&mut out, page.offset);
4016 put_u32(&mut out, page.length);
4017 }
4018 for ((field, dictionary), membership) in
4023 table.fields.iter().zip(&table.dictionaries).zip(&stripe.memberships)
4024 {
4025 if field.ty != LogicalType::Varchar || dictionary.is_none() {
4026 continue;
4027 }
4028 let page =
4029 membership.ok_or_else(|| invalid("string page has no code membership index"))?;
4030 put_u64(&mut out, page.offset);
4031 put_u32(&mut out, page.length);
4032 put_u64(&mut out, page.hash);
4033 }
4034 for sieve in &stripe.sieves {
4035 match sieve {
4036 None => out.push(0),
4037 Some(page) => {
4038 out.push(1);
4039 put_u64(&mut out, page.offset);
4040 put_u32(&mut out, page.length);
4041 put_u64(&mut out, page.hash);
4042 }
4043 }
4044 }
4045 for held in &stripe.part_ranges {
4046 match held {
4047 None => out.push(0),
4048 Some(page) => {
4049 out.push(1);
4050 put_u64(&mut out, page.offset);
4051 put_u32(&mut out, page.length);
4052 put_u64(&mut out, page.hash);
4053 }
4054 }
4055 }
4056 for range in stripe.zone.columns() {
4057 put_bound(&mut out, range.low.as_ref())?;
4058 put_bound(&mut out, range.high.as_ref())?;
4059 put_u32(
4060 &mut out,
4061 u32::try_from(range.nulls).map_err(|_| invalid("null count overflow"))?,
4062 );
4063 out.push(u8::from(range.exact));
4064 match range.sum {
4065 None => out.push(0),
4066 Some(total) => {
4067 out.push(1);
4068 out.extend_from_slice(&total.to_le_bytes());
4069 }
4070 }
4071 }
4072 }
4073 out.extend_from_slice(FREQUENCIES);
4074 put_u16(
4075 &mut out,
4076 u16::try_from(table.frequencies.len())
4077 .map_err(|_| invalid("too many frequency columns"))?,
4078 );
4079 for summary in &table.frequencies {
4080 let Some(summary) = summary else {
4081 out.push(0);
4082 continue;
4083 };
4084 out.push(1);
4085 put_u64(&mut out, summary.omitted_max);
4086 put_u32(
4087 &mut out,
4088 u32::try_from(summary.entries.len())
4089 .map_err(|_| invalid("too many frequency entries"))?,
4090 );
4091 for entry in &summary.entries {
4092 match entry.value {
4093 FrequencyValue::Null => out.push(0),
4094 FrequencyValue::Integer(value) => {
4095 out.push(1);
4096 out.extend_from_slice(&value.to_le_bytes());
4097 }
4098 FrequencyValue::Code(value) => {
4099 out.push(2);
4100 put_u32(&mut out, value);
4101 }
4102 }
4103 put_u64(&mut out, entry.count);
4104 }
4105 put_u32(
4106 &mut out,
4107 u32::try_from(summary.ordinals.len())
4108 .map_err(|_| invalid("too many frequency ordinals"))?,
4109 );
4110 let mut previous = 0_u64;
4111 for (at, &ordinal) in summary.ordinals.iter().enumerate() {
4112 let delta = if at == 0 {
4113 ordinal
4114 } else {
4115 ordinal
4116 .checked_sub(previous)
4117 .ok_or_else(|| invalid("frequency ordinals are not ordered"))?
4118 };
4119 if at != 0 && delta == 0 {
4120 return Err(invalid("frequency ordinals are not unique"));
4121 }
4122 put_var_u64(&mut out, delta);
4123 previous = ordinal;
4124 }
4125 }
4126 if let Some(clustering) = &table.clustering {
4129 out.extend_from_slice(CLUSTERING);
4130 out.push(clustering.width().tag());
4131 put_u16(
4132 &mut out,
4133 u16::try_from(clustering.columns().len())
4134 .map_err(|_| invalid("too many clustering columns"))?,
4135 );
4136 for &column in clustering.columns() {
4137 put_u16(
4138 &mut out,
4139 u16::try_from(column).map_err(|_| invalid("clustering column index overflow"))?,
4140 );
4141 }
4142 }
4143 Ok(out)
4144}
4145
4146fn encode_catalog(entries: &[Entry]) -> Result<Vec<u8>> {
4152 let mut out = CATALOG.to_vec();
4153 put_u32(&mut out, u32::try_from(entries.len()).map_err(|_| invalid("too many tables"))?);
4154 for entry in entries {
4155 let name = entry.name.as_bytes();
4156 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("table name too long"))?);
4157 out.extend_from_slice(name);
4158 put_u64(&mut out, u64::try_from(entry.rows).map_err(|_| invalid("row count overflow"))?);
4159 put_u16(
4160 &mut out,
4161 u16::try_from(entry.fields.len()).map_err(|_| invalid("too many columns"))?,
4162 );
4163 for field in &entry.fields {
4164 let name = field.name.as_bytes();
4165 put_u16(
4166 &mut out,
4167 u16::try_from(name.len()).map_err(|_| invalid("column name too long"))?,
4168 );
4169 out.extend_from_slice(name);
4170 put_type(&mut out, &field.ty)?;
4171 out.push(u8::from(field.not_null));
4172 }
4173 put_u64(&mut out, entry.directory.offset);
4174 put_u32(&mut out, entry.directory.length);
4175 put_u64(&mut out, entry.directory.hash);
4176 }
4177 Ok(out)
4178}
4179
4180fn decode_catalog(bytes: &[u8], size: u64) -> Result<Vec<Entry>> {
4183 let mut cur = Cursor { bytes, at: 0 };
4184 if cur.take(8)? != CATALOG {
4185 return Err(invalid("catalog magic differs"));
4186 }
4187 let count = cur.u32()? as usize;
4188 let mut entries: Vec<Entry> = Vec::with_capacity(count.min(1024));
4189 for _ in 0..count {
4190 let name = cur.text()?;
4191 let rows = usize::try_from(cur.u64()?).map_err(|_| invalid("row count does not fit"))?;
4192 let width = cur.u16()? as usize;
4193 let mut fields = Vec::with_capacity(width);
4194 for _ in 0..width {
4195 let name = cur.text()?;
4196 let ty = read_type(&mut cur)?;
4197 let not_null = match cur.u8()? {
4198 0 => false,
4199 1 => true,
4200 _ => return Err(invalid("nullability flag differs")),
4201 };
4202 fields.push(Field { name, ty, not_null });
4203 }
4204 let directory = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4205 let end = directory
4206 .offset
4207 .checked_add(u64::from(directory.length))
4208 .ok_or_else(|| invalid("table directory offset overflow"))?;
4209 if directory.offset < HEADER
4210 || end > size
4211 || directory.length as usize > MAX_DIRECTORY
4212 || directory.length == 0
4213 {
4214 return Err(invalid("table directory range is outside the file"));
4215 }
4216 if entries.iter().any(|held| held.name == name) {
4217 return Err(invalid("two tables in the catalog have the same name"));
4218 }
4219 entries.push(Entry { name, fields, rows, directory });
4220 }
4221 Ok(entries)
4222}
4223
4224struct Cursor<'a> {
4225 bytes: &'a [u8],
4226 at: usize,
4227}
4228impl<'a> Cursor<'a> {
4229 fn take(&mut self, len: usize) -> Result<&'a [u8]> {
4230 let end = self.at.checked_add(len).ok_or_else(|| invalid("directory offset overflow"))?;
4231 let bytes =
4232 self.bytes.get(self.at..end).ok_or_else(|| invalid("directory is truncated"))?;
4233 self.at = end;
4234 Ok(bytes)
4235 }
4236 fn u8(&mut self) -> Result<u8> {
4237 Ok(self.take(1)?[0])
4238 }
4239 fn u16(&mut self) -> Result<u16> {
4240 Ok(u16::from_le_bytes(self.take(2)?.try_into().expect("two bytes")))
4241 }
4242 fn u32(&mut self) -> Result<u32> {
4243 Ok(u32::from_le_bytes(self.take(4)?.try_into().expect("four bytes")))
4244 }
4245 fn u64(&mut self) -> Result<u64> {
4246 Ok(u64::from_le_bytes(self.take(8)?.try_into().expect("eight bytes")))
4247 }
4248 fn var_u64(&mut self) -> Result<u64> {
4249 let mut value = 0_u64;
4250 for shift in (0..=63).step_by(7) {
4251 let byte = self.u8()?;
4252 let part = u64::from(byte & 0x7f);
4253 if shift == 63 && part > 1 {
4254 return Err(invalid("frequency ordinal varint overflows"));
4255 }
4256 value |= part << shift;
4257 if byte & 0x80 == 0 {
4258 return Ok(value);
4259 }
4260 }
4261 Err(invalid("frequency ordinal varint is too long"))
4262 }
4263 fn bound(&mut self) -> Result<Option<Bound>> {
4264 Ok(match self.u8()? {
4265 0 => None,
4266 1 => Some(Bound::Int(i128::from_le_bytes(
4267 self.take(16)?.try_into().expect("sixteen bytes"),
4268 ))),
4269 2 => Some(Bound::Real(f64::from_le_bytes(
4270 self.take(8)?.try_into().expect("eight bytes"),
4271 ))),
4272 3 => {
4273 let length = self.u32()? as usize;
4274 Some(Bound::Bytes(self.take(length)?.to_vec()))
4275 }
4276 4 => {
4277 let unscaled =
4278 i128::from_le_bytes(self.take(16)?.try_into().expect("sixteen bytes"));
4279 Some(Bound::Scaled { unscaled, scale: self.u8()? })
4280 }
4281 _ => return Err(invalid("bound tag differs")),
4282 })
4283 }
4284 fn text(&mut self) -> Result<String> {
4285 let len = self.u16()? as usize;
4286 String::from_utf8(self.take(len)?.to_vec()).map_err(|_| invalid("name is not UTF-8"))
4287 }
4288}
4289
4290fn decode_directory(bytes: &[u8], size: u64) -> Result<Table> {
4291 let mut cur = Cursor { bytes, at: 0 };
4292 if cur.take(8)? != DIRECTORY {
4293 return Err(invalid("directory magic differs"));
4294 }
4295 let name = cur.text()?;
4296 let width = cur.u16()? as usize;
4297 let mut fields = Vec::with_capacity(width);
4298 for _ in 0..width {
4299 let name = cur.text()?;
4300 let ty = read_type(&mut cur)?;
4301 let not_null = match cur.u8()? {
4302 0 => false,
4303 1 => true,
4304 _ => return Err(invalid("nullability flag differs")),
4305 };
4306 fields.push(Field { name, ty, not_null });
4307 }
4308 let mut dictionaries = Vec::with_capacity(width);
4309 for _ in 0..width {
4310 dictionaries.push(match cur.u8()? {
4311 0 => None,
4312 1 => {
4313 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4314 let end = page
4315 .offset
4316 .checked_add(u64::from(page.length))
4317 .ok_or_else(|| invalid("dictionary page offset overflow"))?;
4318 if page.offset < HEADER || end > size {
4323 return Err(invalid("dictionary page range is outside the file"));
4324 }
4325 Some(page)
4326 }
4327 _ => return Err(invalid("dictionary page tag differs")),
4328 });
4329 }
4330 let mut distincts = Vec::with_capacity(width);
4331 for _ in 0..width {
4332 distincts.push(match cur.u8()? {
4333 0 => None,
4334 1 => Some(cur.u64()?),
4335 _ => return Err(invalid("distinct count tag differs")),
4336 });
4337 }
4338 let rows = usize::try_from(cur.u64()?).map_err(|_| invalid("row count does not fit"))?;
4339 let count = cur.u32()? as usize;
4340 let mut stripes = Vec::with_capacity(count);
4341 let mut total = 0_usize;
4342 for _ in 0..count {
4343 let count = cur.u32()? as usize;
4344 if count == 0 || count > STRIPE_PARTS {
4345 return Err(invalid("stripe part count is outside its bound"));
4346 }
4347 let mut parts = Vec::with_capacity(count);
4348 let mut stripe_rows = 0_usize;
4349 for _ in 0..count {
4350 let rows = cur.u32()?;
4351 if rows == 0 {
4352 return Err(invalid("empty part"));
4353 }
4354 parts.push(rows);
4355 stripe_rows = stripe_rows
4356 .checked_add(rows as usize)
4357 .ok_or_else(|| invalid("stripe row count overflow"))?;
4358 }
4359 total =
4360 total.checked_add(stripe_rows).ok_or_else(|| invalid("stripe row count overflow"))?;
4361 let index = Span { offset: cur.u64()?, length: cur.u32()? };
4362 let section = index_section(count)?;
4363 let wanted = section
4364 .checked_mul(width)
4365 .and_then(|bytes| u32::try_from(bytes).ok())
4366 .ok_or_else(|| invalid("index page length overflow"))?;
4367 let end = index
4368 .offset
4369 .checked_add(u64::from(index.length))
4370 .ok_or_else(|| invalid("index page offset overflow"))?;
4371 if index.offset < HEADER || end > size || index.length != wanted {
4372 return Err(invalid("index page range is outside the file"));
4373 }
4374 let mut pages = Vec::with_capacity(width);
4375 for _ in 0..width {
4376 let offset = cur.u64()?;
4377 let length = cur.u32()?;
4378 let end = offset
4379 .checked_add(u64::from(length))
4380 .ok_or_else(|| invalid("page offset overflow"))?;
4381 if offset < HEADER || end > size || length as usize > MAX_PAGE {
4382 return Err(invalid("page range is outside the file"));
4383 }
4384 pages.push(Span { offset, length });
4385 }
4386 let mut memberships = vec![None; width];
4387 for (column, field) in fields.iter().enumerate() {
4388 if field.ty != LogicalType::Varchar || dictionaries[column].is_none() {
4389 continue;
4390 }
4391 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4392 let end = page
4393 .offset
4394 .checked_add(u64::from(page.length))
4395 .ok_or_else(|| invalid("membership page offset overflow"))?;
4396 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
4397 return Err(invalid("membership page range is outside the file"));
4398 }
4399 memberships[column] = Some(page);
4400 }
4401 let mut sieves = vec![None; width];
4402 for sieve in sieves.iter_mut().take(width) {
4403 match cur.u8()? {
4404 0 => continue,
4405 1 => {}
4406 _ => return Err(invalid("a sieve page has an unknown tag")),
4407 }
4408 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4409 let end = page
4410 .offset
4411 .checked_add(u64::from(page.length))
4412 .ok_or_else(|| invalid("sieve page offset overflow"))?;
4413 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
4414 return Err(invalid("sieve page range is outside the file"));
4415 }
4416 *sieve = Some(page);
4417 }
4418 let mut part_ranges = vec![None; width];
4419 for held in part_ranges.iter_mut().take(width) {
4420 match cur.u8()? {
4421 0 => continue,
4422 1 => {}
4423 _ => return Err(invalid("a part range page has an unknown tag")),
4424 }
4425 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
4426 let end = page
4427 .offset
4428 .checked_add(u64::from(page.length))
4429 .ok_or_else(|| invalid("part range page offset overflow"))?;
4430 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
4431 return Err(invalid("part range page range is outside the file"));
4432 }
4433 *held = Some(page);
4434 }
4435 let mut ranges = Vec::with_capacity(width);
4436 for column in 0..width {
4437 let low = cur.bound()?;
4438 let high = cur.bound()?;
4439 let nulls = cur.u32()? as usize;
4440 if nulls > stripe_rows {
4441 return Err(invalid("null count exceeds stripe rows"));
4442 }
4443 let exact = cur.u8()? != 0;
4444 let sum = match cur.u8()? {
4445 0 => None,
4446 1 => Some(i128::from_le_bytes(
4447 cur.take(16)?.try_into().map_err(|_| invalid("a stripe sum is truncated"))?,
4448 )),
4449 _ => return Err(invalid("a stripe sum has an unknown tag")),
4450 };
4451 let ty = &fields.get(column).ok_or_else(|| invalid("a stripe range has no column"))?.ty;
4457 let low = low.map(|bound| scaled_as(bound, ty));
4458 let high = high.map(|bound| scaled_as(bound, ty));
4459 ranges.push(Range { low, high, nulls, exact, sum });
4460 }
4461 stripes.push(Stripe {
4462 rows: stripe_rows,
4463 parts,
4464 index,
4465 pages,
4466 memberships,
4467 sieves,
4468 part_ranges,
4469 zone: Zone::from_ranges(ranges),
4470 });
4471 }
4472 if total != rows {
4473 return Err(invalid("table row count differs from stripes"));
4474 }
4475 let frequencies = if cur.at == bytes.len() {
4476 vec![None; width]
4477 } else {
4478 if cur.take(8)? != FREQUENCIES {
4479 return Err(invalid("directory extension magic differs"));
4480 }
4481 if cur.u16()? as usize != width {
4482 return Err(invalid("frequency column count differs"));
4483 }
4484 let mut frequencies = Vec::with_capacity(width);
4485 for field in &fields {
4486 let summary = match cur.u8()? {
4487 0 => None,
4488 1 => {
4489 let omitted_max = cur.u64()?;
4490 let count = cur.u32()? as usize;
4491 if count > FREQUENCY_ENTRIES {
4492 return Err(invalid("frequency entry count exceeds its bound"));
4493 }
4494 let mut entries = Vec::with_capacity(count);
4495 for _ in 0..count {
4497 let value = match cur.u8()? {
4498 0 => FrequencyValue::Null,
4499 1 => FrequencyValue::Integer(i128::from_le_bytes(
4500 cur.take(16)?.try_into().expect("sixteen bytes"),
4501 )),
4502 2 => FrequencyValue::Code(cur.u32()?),
4503 _ => return Err(invalid("frequency value tag differs")),
4504 };
4505 let valid = matches!(
4506 (&field.ty, value),
4507 (_, FrequencyValue::Null)
4508 | (LogicalType::Varchar, FrequencyValue::Code(_))
4509 | (
4510 LogicalType::TinyInt
4511 | LogicalType::SmallInt
4512 | LogicalType::Integer
4513 | LogicalType::BigInt
4514 | LogicalType::UTinyInt
4515 | LogicalType::USmallInt
4516 | LogicalType::UInteger
4517 | LogicalType::UBigInt
4518 | LogicalType::Date
4519 | LogicalType::Timestamp,
4520 FrequencyValue::Integer(_),
4521 )
4522 );
4523 if !valid {
4524 return Err(invalid("frequency value does not match its column"));
4525 }
4526 let count = cur.u64()?;
4527 if count == 0 || count > rows as u64 {
4528 return Err(invalid("frequency count is outside the table"));
4529 }
4530 entries.push(FrequencyEntry { value, count });
4531 }
4532 if entries.windows(2).any(|pair| pair[0].count < pair[1].count) {
4533 return Err(invalid("frequency entries are not descending"));
4534 }
4535 let ordinals = {
4536 let ordinal_count = cur.u32()? as usize;
4537 if ordinal_count > FREQUENCY_ORDINALS || ordinal_count > rows {
4538 return Err(invalid("frequency ordinal count exceeds its bound"));
4539 }
4540 let mut ordinals = Vec::with_capacity(ordinal_count);
4541 let mut previous = 0_u64;
4542 for at in 0..ordinal_count {
4543 let delta = cur.var_u64()?;
4544 if at != 0 && delta == 0 {
4545 return Err(invalid("frequency ordinals are not increasing"));
4546 }
4547 let ordinal = if at == 0 {
4548 delta
4549 } else {
4550 previous
4551 .checked_add(delta)
4552 .ok_or_else(|| invalid("frequency ordinal overflows"))?
4553 };
4554 if ordinal >= rows as u64 {
4555 return Err(invalid("frequency ordinal is outside the table"));
4556 }
4557 ordinals.push(ordinal);
4558 previous = ordinal;
4559 }
4560 ordinals
4561 };
4562 Some(FrequencySummary { entries, omitted_max, ordinals })
4563 }
4564 _ => return Err(invalid("frequency summary tag differs")),
4565 };
4566 frequencies.push(summary);
4567 }
4568 frequencies
4569 };
4570 let clustering = if cur.at == bytes.len() {
4571 None
4572 } else {
4573 if cur.take(8)? != CLUSTERING {
4574 return Err(invalid("directory extension magic differs"));
4575 }
4576 let bucket =
4577 Width::from_tag(cur.u8()?).ok_or_else(|| invalid("clustering width tag differs"))?;
4578 let count = cur.u16()? as usize;
4579 let mut columns = Vec::with_capacity(count.min(fields.len()));
4580 for _ in 0..count {
4581 columns.push(u32::from(cur.u16()?));
4582 }
4583 Some(Clustering::new(columns, bucket, &fields).map_err(|_| {
4586 invalid("stored clustering declaration does not match the table it is on")
4587 })?)
4588 };
4589 if cur.at != bytes.len() {
4590 return Err(invalid("directory has trailing bytes"));
4591 }
4592 Ok(Table { name, fields, stripes, rows, dictionaries, distincts, frequencies, clustering })
4593}
4594
4595fn put_bound(out: &mut Vec<u8>, bound: Option<&Bound>) -> Result<()> {
4596 match bound {
4597 None => out.push(0),
4598 Some(Bound::Int(value)) => {
4599 out.push(1);
4600 out.extend_from_slice(&value.to_le_bytes());
4601 }
4602 Some(Bound::Real(value)) => {
4603 out.push(2);
4604 out.extend_from_slice(&value.to_le_bytes());
4605 }
4606 Some(Bound::Bytes(value)) => {
4607 out.push(3);
4608 put_u32(out, u32::try_from(value.len()).map_err(|_| invalid("bound length overflow"))?);
4609 out.extend_from_slice(value);
4610 }
4611 Some(Bound::Scaled { unscaled, scale }) => {
4612 out.push(4);
4613 out.extend_from_slice(&unscaled.to_le_bytes());
4614 out.push(*scale);
4615 }
4616 }
4617 Ok(())
4618}
4619
4620#[derive(Debug)]
4637struct Codes;
4638
4639impl chooser::Chooser for Codes {
4640 fn name(&self) -> &'static str {
4641 "codes"
4642 }
4643
4644 fn narrow_strings(
4645 &self,
4646 _values: &[&[u8]],
4647 offered: &[string::Kind],
4648 _depth: u8,
4649 ) -> Vec<string::Kind> {
4650 offered.to_vec()
4653 }
4654
4655 fn narrow_integers(
4656 &self,
4657 _values: &[i64],
4658 offered: &[integer::Kind],
4659 depth: u8,
4660 ) -> Vec<integer::Kind> {
4661 let keep: &[integer::Kind] = if depth == 0 {
4662 &[integer::Kind::Constant, integer::Kind::Packed, integer::Kind::Rle]
4663 } else {
4664 &[integer::Kind::Constant, integer::Kind::Packed]
4665 };
4666 let narrowed: Vec<integer::Kind> =
4667 offered.iter().copied().filter(|kind| keep.contains(kind)).collect();
4668 if narrowed.is_empty() { offered.to_vec() } else { narrowed }
4671 }
4672}
4673
4674#[derive(Debug)]
4686struct Fixed;
4687
4688impl chooser::Chooser for Fixed {
4689 fn name(&self) -> &'static str {
4690 "fixed"
4691 }
4692
4693 fn narrow_strings(
4694 &self,
4695 _values: &[&[u8]],
4696 offered: &[string::Kind],
4697 _depth: u8,
4698 ) -> Vec<string::Kind> {
4699 offered.to_vec()
4700 }
4701
4702 fn narrow_integers(
4703 &self,
4704 _values: &[i64],
4705 offered: &[integer::Kind],
4706 depth: u8,
4707 ) -> Vec<integer::Kind> {
4708 let keep: &[integer::Kind] = if depth == 0 {
4709 &[
4710 integer::Kind::Constant,
4711 integer::Kind::Packed,
4712 integer::Kind::Delta,
4713 integer::Kind::Rle,
4714 integer::Kind::Sparse,
4715 integer::Kind::Strided,
4716 ]
4717 } else {
4718 &[integer::Kind::Constant, integer::Kind::Packed, integer::Kind::Delta]
4719 };
4720 let narrowed: Vec<integer::Kind> =
4721 offered.iter().copied().filter(|kind| keep.contains(kind)).collect();
4722 if narrowed.is_empty() { offered.to_vec() } else { narrowed }
4723 }
4724}
4725
4726fn widened(data: &Data) -> Option<Vec<i64>> {
4733 match data {
4734 Data::Int8(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4735 Data::UInt8(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4736 Data::Int16(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4737 Data::UInt16(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4738 Data::Int32(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4739 Data::UInt32(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4740 Data::Int64(values) => Some(values.to_vec()),
4741 _ => None,
4742 }
4743}
4744
4745trait Narrow: Copy {
4752 const BIASED: (u32, u64);
4757
4758 fn narrow(value: i64) -> Self;
4760}
4761
4762#[allow(clippy::cast_sign_loss, reason = "a residue is a bit pattern and not a number")]
4779fn residue<T: Narrow>(value: i64) -> u64 {
4780 let (bits, bias) = T::BIASED;
4781 (value as u64).wrapping_add(bias) >> bits
4782}
4783
4784macro_rules! narrows {
4789 ($($ty:ty => $bias:expr),* $(,)?) => {$(
4790 impl Narrow for $ty {
4791 const BIASED: (u32, u64) = (<$ty>::BITS, $bias);
4792
4793 #[allow(
4794 clippy::cast_possible_truncation,
4795 clippy::cast_sign_loss,
4796 reason = "the caller has checked the bits this truncates away"
4797 )]
4798 fn narrow(value: i64) -> Self {
4799 value as Self
4800 }
4801 }
4802 )*};
4803}
4804
4805narrows! {
4806 i8 => 1 << 7,
4807 u8 => 0,
4808 i16 => 1 << 15,
4809 u16 => 0,
4810 i32 => 1 << 31,
4811 u32 => 0,
4812}
4813
4814fn fit<T: Narrow>(values: &[i64]) -> Result<Vec<T>> {
4827 let mut spilled = 0u64;
4828 for value in values {
4829 spilled |= residue::<T>(*value);
4830 }
4831 if spilled != 0 {
4832 return Err(invalid("page value is not of its type"));
4833 }
4834 Ok(values.iter().map(|value| T::narrow(*value)).collect())
4835}
4836
4837fn narrowed(ty: &LogicalType, values: Vec<i64>) -> Result<Data> {
4842 Ok(match ty {
4843 LogicalType::TinyInt => Data::Int8(fit::<i8>(&values)?.into()),
4844 LogicalType::UTinyInt => Data::UInt8(fit::<u8>(&values)?.into()),
4845 LogicalType::SmallInt => Data::Int16(fit::<i16>(&values)?.into()),
4846 LogicalType::USmallInt => Data::UInt16(fit::<u16>(&values)?.into()),
4847 LogicalType::Integer | LogicalType::Date => Data::Int32(fit::<i32>(&values)?.into()),
4848 LogicalType::UInteger => Data::UInt32(fit::<u32>(&values)?.into()),
4849 LogicalType::BigInt
4850 | LogicalType::Timestamp
4851 | LogicalType::Time
4852 | LogicalType::TimeTz
4853 | LogicalType::TimestampTz
4854 | LogicalType::TimestampS
4855 | LogicalType::TimestampMs
4856 | LogicalType::TimestampNs => Data::Int64(values.into()),
4857 LogicalType::Decimal { .. } => match ty.physical() {
4860 PhysicalType::Int16 => Data::Int16(fit::<i16>(&values)?.into()),
4861 PhysicalType::Int32 => Data::Int32(fit::<i32>(&values)?.into()),
4862 PhysicalType::Int64 => Data::Int64(values.into()),
4863 _ => return Err(invalid("cascade codec belongs to a decimal that is not an integer")),
4864 },
4865 _ => return Err(invalid("cascade codec belongs to a page that is not integers")),
4866 })
4867}
4868
4869fn plain_width(ty: &LogicalType) -> Option<usize> {
4872 Some(match ty {
4873 LogicalType::TinyInt | LogicalType::UTinyInt => 1,
4874 LogicalType::SmallInt | LogicalType::USmallInt => 2,
4875 LogicalType::Integer | LogicalType::UInteger | LogicalType::Date => 4,
4876 LogicalType::BigInt
4877 | LogicalType::Timestamp
4878 | LogicalType::Time
4879 | LogicalType::TimeTz
4880 | LogicalType::TimestampTz
4881 | LogicalType::TimestampS
4882 | LogicalType::TimestampMs
4883 | LogicalType::TimestampNs => 8,
4884 LogicalType::Decimal { .. } => match ty.physical() {
4885 PhysicalType::Int16 => 2,
4886 PhysicalType::Int32 => 4,
4887 PhysicalType::Int64 => 8,
4888 _ => return None,
4891 },
4892 _ => return None,
4893 })
4894}
4895
4896fn cascaded(
4902 flat: &Vector,
4903 ty: &LogicalType,
4904 packed: Option<&Packed<'_>>,
4905) -> Result<Option<Vec<u8>>> {
4906 let (Some(width), Some(data)) = (plain_width(ty), flat.data()) else { return Ok(None) };
4907 let Some(values) = widened(data) else { return Ok(None) };
4908 let plain = values.len().saturating_mul(width);
4909 let best = match packed {
4910 Some(packed) => plain.min(21 + size_of_val(packed.words())),
4912 None => plain,
4913 };
4914 let out = integer::encode_with(&values, &Fixed)?;
4915 Ok((out.len() < best).then_some(out))
4916}
4917
4918fn text_compressed(flat: &Vector) -> Result<Option<Vec<u8>>> {
4956 let mut values: Vec<&[u8]> = Vec::with_capacity(flat.len());
4957 let mut payload = 0_usize;
4958 for row in 0..flat.len() {
4959 let text = flat.text_at(row).unwrap_or("").as_bytes();
4960 payload = payload.saturating_add(text.len());
4961 values.push(text);
4962 }
4963 let plain = (flat.len() + 1).saturating_mul(4).saturating_add(payload);
4965 let Some(out) = string::encode_only(string::Kind::Fsst, &values)? else {
4966 return Ok(None);
4967 };
4968 Ok((out.len() < plain).then_some(out))
4969}
4970
4971fn encoded_codes(codes: &[u32]) -> Result<Option<Vec<u8>>> {
4972 let wide: Vec<i64> = codes.iter().map(|code| i64::from(*code)).collect();
4973 let coded = integer::encode_with(&wide, &Codes)?;
4974 let plain = codes.len().saturating_mul(size_of::<u32>());
4975 Ok((coded.len() < plain).then_some(coded))
4976}
4977
4978fn encode(
4979 vector: &Vector,
4980 global: Option<&mut GlobalDictionary>,
4981) -> Result<(Vec<u8>, Option<Vec<u32>>)> {
4982 let ty = vector.logical_type();
4983 let flat = vector.flatten()?;
4985 let mut out = Vec::new();
4986 let mut global_codes = None;
4987 if let Some(global) = global {
4988 let mut codes = Vec::with_capacity(flat.len());
4989 for row in 0..flat.len() {
4990 let text = flat.text_at(row).unwrap_or("");
4991 let code = global.code(text)?;
4992 global.observe(code, flat.is_null_at(row))?;
4993 codes.push(code);
4994 }
4995 global_codes = Some(codes);
4996 }
4997 let membership = global_codes.as_deref().map(unique_codes);
4998 let dictionary = if global_codes.is_none() && ty == &LogicalType::Varchar {
4999 string_dictionary(&flat)?
5000 } else {
5001 None
5002 };
5003 let compressed_text =
5004 if global_codes.is_none() && dictionary.is_none() && ty == &LogicalType::Varchar {
5005 text_compressed(&flat)?
5006 } else {
5007 None
5008 };
5009 let packed_vector = if dictionary.is_none() && global_codes.is_none() {
5010 Some(flat.bit_packed()?)
5011 } else {
5012 None
5013 };
5014 let packed = packed_vector.as_ref().and_then(Vector::packed_parts);
5015 let coded = match global_codes.as_deref() {
5016 Some(codes) => encoded_codes(codes)?,
5017 None => None,
5018 };
5019 let cascade = if dictionary.is_none() && global_codes.is_none() {
5023 cascaded(&flat, ty, packed.as_ref())?
5024 } else {
5025 None
5026 };
5027 out.push(if coded.is_some() {
5028 4
5029 } else if cascade.is_some() {
5030 5
5031 } else if global_codes.is_some() {
5032 3
5033 } else if dictionary.is_some() {
5034 1
5035 } else if compressed_text.is_some() {
5036 6
5037 } else if packed.is_some() {
5038 2
5039 } else {
5040 0
5041 });
5042 let nulls = flat.validity();
5043 let flag = match nulls {
5044 Validity::AllValid => 0,
5045 Validity::AllInvalid => 1,
5046 Validity::Mask(_) => 2,
5047 };
5048 out.push(flag);
5049 if flag == 2 {
5050 for group in (0..vector.len()).step_by(8) {
5051 let mut bits = 0_u8;
5052 for bit in 0..8 {
5053 if group + bit < vector.len() && !flat.is_null_at(group + bit) {
5054 bits |= 1 << bit;
5055 }
5056 }
5057 out.push(bits);
5058 }
5059 }
5060 if let Some(coded) = coded {
5061 out.extend_from_slice(&coded);
5062 return Ok((out, membership));
5063 }
5064 if let Some(cascade) = cascade {
5065 out.extend_from_slice(&cascade);
5066 return Ok((out, membership));
5067 }
5068 if let Some(codes) = global_codes {
5069 for code in codes {
5070 put_u32(&mut out, code);
5071 }
5072 return Ok((out, membership));
5073 }
5074 if let Some(dictionary) = dictionary {
5075 out.extend_from_slice(&dictionary);
5076 return Ok((out, membership));
5077 }
5078 if let Some(compressed_text) = compressed_text {
5079 out.extend_from_slice(&compressed_text);
5080 return Ok((out, membership));
5081 }
5082 if let Some(packed) = packed {
5083 if packed.offset() != 0 {
5084 return Err(invalid("writer received a sliced packed vector"));
5085 }
5086 out.push(u8::try_from(packed.width()).map_err(|_| invalid("packed width overflow"))?);
5087 out.extend_from_slice(&packed.base().to_le_bytes());
5088 put_u32(
5089 &mut out,
5090 u32::try_from(packed.words().len()).map_err(|_| invalid("too many packed words"))?,
5091 );
5092 for word in packed.words() {
5093 put_u64(&mut out, *word);
5094 }
5095 return Ok((out, membership));
5096 }
5097 let data = flat.data().ok_or_else(|| invalid("scalar column did not flatten"))?;
5098 match (ty, data) {
5099 (LogicalType::TinyInt, Data::Int8(values)) => {
5100 for value in &**values {
5101 out.extend_from_slice(&value.to_le_bytes());
5102 }
5103 }
5104 (LogicalType::UTinyInt, Data::UInt8(values)) => {
5105 for value in &**values {
5106 out.extend_from_slice(&value.to_le_bytes());
5107 }
5108 }
5109 (LogicalType::SmallInt, Data::Int16(values)) => {
5110 for value in &**values {
5111 out.extend_from_slice(&value.to_le_bytes());
5112 }
5113 }
5114 (LogicalType::USmallInt, Data::UInt16(values)) => {
5115 for value in &**values {
5116 out.extend_from_slice(&value.to_le_bytes());
5117 }
5118 }
5119 (LogicalType::UInteger, Data::UInt32(values)) => {
5120 for value in &**values {
5121 out.extend_from_slice(&value.to_le_bytes());
5122 }
5123 }
5124 (LogicalType::UBigInt, Data::UInt64(values)) => {
5125 for value in &**values {
5126 out.extend_from_slice(&value.to_le_bytes());
5127 }
5128 }
5129 (LogicalType::Integer | LogicalType::Date, Data::Int32(values)) => {
5130 for value in &**values {
5131 out.extend_from_slice(&value.to_le_bytes());
5132 }
5133 }
5134 (
5135 LogicalType::BigInt
5136 | LogicalType::Timestamp
5137 | LogicalType::Time
5138 | LogicalType::TimeTz
5139 | LogicalType::TimestampTz
5140 | LogicalType::TimestampS
5141 | LogicalType::TimestampMs
5142 | LogicalType::TimestampNs,
5143 Data::Int64(values),
5144 ) => {
5145 for value in &**values {
5146 out.extend_from_slice(&value.to_le_bytes());
5147 }
5148 }
5149 (LogicalType::HugeInt | LogicalType::Uuid, Data::Int128(values)) => {
5152 for value in &**values {
5153 out.extend_from_slice(&value.to_le_bytes());
5154 }
5155 }
5156 (LogicalType::UHugeInt, Data::UInt128(values)) => {
5157 for value in &**values {
5158 out.extend_from_slice(&value.to_le_bytes());
5159 }
5160 }
5161 (LogicalType::Float, Data::Float32(values)) => {
5164 for value in &**values {
5165 out.extend_from_slice(&value.to_le_bytes());
5166 }
5167 }
5168 (LogicalType::Double, Data::Float64(values)) => {
5169 for value in &**values {
5170 out.extend_from_slice(&value.to_le_bytes());
5171 }
5172 }
5173 (LogicalType::Interval, Data::Interval(values)) => {
5177 for (months, days, micros) in &**values {
5178 out.extend_from_slice(&months.to_le_bytes());
5179 out.extend_from_slice(&days.to_le_bytes());
5180 out.extend_from_slice(µs.to_le_bytes());
5181 }
5182 }
5183 (LogicalType::Boolean, Data::Bool(values)) => {
5184 for value in &**values {
5185 out.push(u8::from(*value));
5186 }
5187 }
5188 (LogicalType::Decimal { .. }, Data::Int16(values)) => {
5191 for value in &**values {
5192 out.extend_from_slice(&value.to_le_bytes());
5193 }
5194 }
5195 (LogicalType::Decimal { .. }, Data::Int32(values)) => {
5196 for value in &**values {
5197 out.extend_from_slice(&value.to_le_bytes());
5198 }
5199 }
5200 (LogicalType::Decimal { .. }, Data::Int64(values)) => {
5201 for value in &**values {
5202 out.extend_from_slice(&value.to_le_bytes());
5203 }
5204 }
5205 (LogicalType::Decimal { .. }, Data::Int128(values)) => {
5206 for value in &**values {
5207 out.extend_from_slice(&value.to_le_bytes());
5208 }
5209 }
5210 (LogicalType::Varchar | LogicalType::Blob | LogicalType::Bit, Data::Varlen(values)) => {
5215 let mut bytes = Vec::new();
5216 put_u32(&mut out, 0);
5217 for row in 0..vector.len() {
5218 let value = values.bytes(row).ok_or_else(|| invalid("string view is invalid"))?;
5219 bytes.extend_from_slice(value);
5220 put_u32(
5221 &mut out,
5222 u32::try_from(bytes.len())
5223 .map_err(|_| invalid("string payload exceeds 4GiB"))?,
5224 );
5225 }
5226 out.extend_from_slice(&bytes);
5227 }
5228 _ => return Err(Error::not_implemented(format!("native page for {ty}"))),
5229 }
5230 Ok((out, membership))
5231}
5232
5233fn put_varint(out: &mut Vec<u8>, mut value: u32) {
5234 while value >= 0x80 {
5235 out.push((value as u8 & 0x7f) | 0x80);
5236 value >>= 7;
5237 }
5238 out.push(value as u8);
5239}
5240
5241fn unique_codes(codes: &[u32]) -> Vec<u32> {
5243 let mut unique = codes.to_vec();
5244 unique.sort_unstable();
5245 unique.dedup();
5246 unique
5247}
5248
5249fn merged_codes(lists: Vec<Vec<u32>>) -> Vec<u32> {
5255 let mut lists = lists;
5256 while lists.len() > 1 {
5257 let mut next = Vec::with_capacity(lists.len().div_ceil(2));
5258 for pair in lists.chunks(2) {
5259 match pair {
5260 [left, right] => next.push(merged_pair(left, right)),
5261 [only] => next.push(only.clone()),
5262 _ => {}
5263 }
5264 }
5265 lists = next;
5266 }
5267 lists.pop().unwrap_or_default()
5268}
5269
5270fn merged_pair(left: &[u32], right: &[u32]) -> Vec<u32> {
5271 let mut out = Vec::with_capacity(left.len().saturating_add(right.len()));
5272 let mut at = 0;
5273 let mut to = 0;
5274 while at < left.len() && to < right.len() {
5275 match left[at].cmp(&right[to]) {
5276 Ordering::Less => {
5277 out.push(left[at]);
5278 at += 1;
5279 }
5280 Ordering::Greater => {
5281 out.push(right[to]);
5282 to += 1;
5283 }
5284 Ordering::Equal => {
5285 out.push(left[at]);
5286 at += 1;
5287 to += 1;
5288 }
5289 }
5290 }
5291 out.extend_from_slice(&left[at..]);
5292 out.extend_from_slice(&right[to..]);
5293 out
5294}
5295
5296fn merged_range(ranges: impl Iterator<Item = Range>) -> Range {
5301 let mut merged = Range::default();
5302 let mut first = true;
5303 for range in ranges {
5304 merged.nulls = merged.nulls.saturating_add(range.nulls);
5305 merged.sum = match (merged.sum.take(), range.sum) {
5309 (Some(held), Some(next)) if !first => held.checked_add(next),
5310 (_, next) if first => next,
5311 _ => None,
5312 };
5313 merged.exact = if first { range.exact } else { merged.exact && range.exact };
5314 if first {
5315 merged.low = range.low;
5316 merged.high = range.high;
5317 first = false;
5318 continue;
5319 }
5320 merged.low = match (merged.low.take(), range.low) {
5321 (Some(held), Some(next)) => Some(held.smaller(next)),
5322 _ => None,
5323 };
5324 merged.high = match (merged.high.take(), range.high) {
5325 (Some(held), Some(next)) => Some(held.larger(next)),
5326 _ => None,
5327 };
5328 }
5329 merged
5330}
5331
5332fn shortened(bound: Option<Bound>, high: bool) -> Option<Bound> {
5345 match bound {
5346 Some(Bound::Bytes(mut value)) if value.len() > PART_BOUND_BYTES => {
5347 value.truncate(PART_BOUND_BYTES);
5348 if !high {
5349 return Some(Bound::Bytes(value));
5350 }
5351 while let Some(last) = value.pop() {
5352 if last < u8::MAX {
5353 value.push(last + 1);
5354 return Some(Bound::Bytes(value));
5355 }
5356 }
5357 None
5358 }
5359 other => other,
5360 }
5361}
5362
5363fn encode_part_ranges(ranges: &[Range]) -> Result<Vec<u8>> {
5371 let mut out = Vec::new();
5372 put_u32(
5373 &mut out,
5374 u32::try_from(ranges.len()).map_err(|_| invalid("too many parts in a stripe"))?,
5375 );
5376 for range in ranges {
5377 put_bound(&mut out, shortened(range.low.clone(), false).as_ref())?;
5378 put_bound(&mut out, shortened(range.high.clone(), true).as_ref())?;
5379 put_u32(&mut out, u32::try_from(range.nulls).map_err(|_| invalid("null count overflow"))?);
5380 }
5381 Ok(out)
5382}
5383
5384fn decode_part_ranges(bytes: &[u8]) -> Result<Vec<Range>> {
5386 let mut cur = Cursor { bytes, at: 0 };
5387 let parts = cur.u32()? as usize;
5388 let mut out = Vec::new();
5389 for _ in 0..parts {
5390 let low = cur.bound()?;
5391 let high = cur.bound()?;
5392 let nulls = cur.u32()? as usize;
5393 out.push(Range { low, high, nulls, exact: false, sum: None });
5394 }
5395 Ok(out)
5396}
5397
5398fn encode_sieves<'a>(sieves: impl Iterator<Item = &'a Option<Sieve>>) -> Result<Vec<u8>> {
5399 let held: Vec<&Option<Sieve>> = sieves.collect();
5400 let mut out = Vec::new();
5401 put_u32(
5402 &mut out,
5403 u32::try_from(held.len()).map_err(|_| invalid("too many parts in a stripe"))?,
5404 );
5405 for sieve in &held {
5406 let length = sieve.as_ref().map_or(0, Sieve::len);
5407 put_u32(&mut out, u32::try_from(length).map_err(|_| invalid("sieve length overflow"))?);
5408 }
5409 for sieve in held.into_iter().flatten() {
5411 out.extend_from_slice(&sieve.to_bytes());
5412 }
5413 Ok(out)
5414}
5415
5416fn decode_sieves(bytes: &[u8]) -> Result<Vec<Option<Sieve>>> {
5422 let parts = u32::from_le_bytes(
5423 bytes
5424 .get(..4)
5425 .ok_or_else(|| invalid("sieve page is truncated"))?
5426 .try_into()
5427 .map_err(|_| invalid("sieve page is truncated"))?,
5428 ) as usize;
5429 let mut lengths = Vec::with_capacity(parts);
5430 for part in 0..parts {
5431 let at = 4 + part * 4;
5432 let field = bytes.get(at..at + 4).ok_or_else(|| invalid("sieve page is truncated"))?;
5433 lengths.push(u32::from_le_bytes(
5434 field.try_into().map_err(|_| invalid("sieve page is truncated"))?,
5435 ) as usize);
5436 }
5437 let mut at = 4 + parts * 4;
5438 let mut out = Vec::with_capacity(parts);
5439 for length in lengths {
5440 if length == 0 {
5441 out.push(None);
5442 continue;
5443 }
5444 let end = at.checked_add(length).ok_or_else(|| invalid("sieve page is truncated"))?;
5445 let field = bytes.get(at..end).ok_or_else(|| invalid("sieve page is truncated"))?;
5446 out.push(Sieve::from_bytes(field));
5447 at = end;
5448 }
5449 if at != bytes.len() {
5450 return Err(invalid("sieve page has trailing bytes"));
5451 }
5452 Ok(out)
5453}
5454
5455fn encode_membership(unique: &[u32]) -> Vec<u8> {
5461 let mut out = Vec::with_capacity(unique.len().saturating_mul(2).saturating_add(5));
5462 put_varint(&mut out, u32::try_from(unique.len()).unwrap_or(u32::MAX));
5463 let mut previous = 0;
5464 for (at, &code) in unique.iter().enumerate() {
5465 put_varint(&mut out, if at == 0 { code } else { code - previous });
5466 previous = code;
5467 }
5468 out
5469}
5470
5471fn take_varint(bytes: &[u8], at: &mut usize) -> Result<u32> {
5472 let mut value = 0_u32;
5473 for shift in (0..35).step_by(7) {
5474 let byte = *bytes.get(*at).ok_or_else(|| invalid("membership varint is truncated"))?;
5475 *at += 1;
5476 let part = u32::from(byte & 0x7f);
5477 if shift == 28 && part > 0x0f {
5478 return Err(invalid("membership varint overflow"));
5479 }
5480 value = value
5481 .checked_add(
5482 part.checked_shl(shift).ok_or_else(|| invalid("membership varint overflow"))?,
5483 )
5484 .ok_or_else(|| invalid("membership varint overflow"))?;
5485 if byte & 0x80 == 0 {
5486 return Ok(value);
5487 }
5488 }
5489 Err(invalid("membership varint is too long"))
5490}
5491
5492fn decode_membership(bytes: &[u8]) -> Result<Vec<u32>> {
5493 let mut at = 0;
5494 let count = take_varint(bytes, &mut at)? as usize;
5495 let mut codes = Vec::with_capacity(count);
5496 let mut previous = 0_u32;
5497 for index in 0..count {
5498 let delta = take_varint(bytes, &mut at)?;
5499 let code = if index == 0 {
5500 delta
5501 } else {
5502 previous.checked_add(delta).ok_or_else(|| invalid("membership code overflow"))?
5503 };
5504 if index > 0 && code <= previous {
5505 return Err(invalid("membership codes are not increasing"));
5506 }
5507 codes.push(code);
5508 previous = code;
5509 }
5510 if at != bytes.len() {
5511 return Err(invalid("membership page has trailing bytes"));
5512 }
5513 Ok(codes)
5514}
5515
5516fn string_dictionary(vector: &Vector) -> Result<Option<Vec<u8>>> {
5517 let mut by_text = HashMap::new();
5518 let mut values = Vec::new();
5519 let mut codes = Vec::with_capacity(vector.len());
5520 let mut plain_bytes = 0_usize;
5521 for row in 0..vector.len() {
5522 let text = vector.text_at(row).unwrap_or("");
5523 plain_bytes = plain_bytes.saturating_add(text.len());
5524 let code = match by_text.get(text) {
5525 Some(&code) => code,
5526 None => {
5527 let code = u32::try_from(values.len())
5528 .map_err(|_| invalid("too many dictionary values"))?;
5529 by_text.insert(text, code);
5530 values.push(text);
5531 code
5532 }
5533 };
5534 codes.push(code);
5535 }
5536 let dictionary_bytes = values.iter().map(|value| value.len()).sum::<usize>();
5537 let encoded = 8_usize
5538 .saturating_add((values.len() + 1).saturating_mul(4))
5539 .saturating_add(dictionary_bytes)
5540 .saturating_add(codes.len().saturating_mul(4));
5541 let plain = (vector.len() + 1).saturating_mul(4).saturating_add(plain_bytes);
5542 if encoded >= plain {
5543 return Ok(None);
5544 }
5545 let mut out = Vec::with_capacity(encoded);
5546 put_u32(
5547 &mut out,
5548 u32::try_from(values.len()).map_err(|_| invalid("too many dictionary values"))?,
5549 );
5550 put_u32(
5551 &mut out,
5552 u32::try_from(dictionary_bytes).map_err(|_| invalid("dictionary payload exceeds 4GiB"))?,
5553 );
5554 let mut offset = 0_u32;
5555 put_u32(&mut out, offset);
5556 for value in &values {
5557 offset = offset
5558 .checked_add(
5559 u32::try_from(value.len()).map_err(|_| invalid("dictionary value is too long"))?,
5560 )
5561 .ok_or_else(|| invalid("dictionary payload exceeds 4GiB"))?;
5562 put_u32(&mut out, offset);
5563 }
5564 for value in values {
5565 out.extend_from_slice(value.as_bytes());
5566 }
5567 for code in codes {
5568 put_u32(&mut out, code);
5569 }
5570 Ok(Some(out))
5571}
5572
5573struct EncodedDictionary {
5574 index: Vec<u8>,
5575 ranks: Vec<u8>,
5576 payload: Vec<Vec<u8>>,
5579}
5580
5581fn head(bytes: &[u8]) -> u64 {
5583 let mut word = [0; 8];
5584 let take = bytes.len().min(8);
5585 word[..take].copy_from_slice(&bytes[..take]);
5586 u64::from_be_bytes(word)
5587}
5588
5589fn rankings(dictionaries: &[Option<GlobalDictionary>]) -> Result<Vec<Vec<(u64, u32)>>> {
5597 let present =
5598 dictionaries.iter().enumerate().filter(|(_, held)| held.is_some()).map(|(at, _)| at);
5599 let present = present.collect::<Vec<_>>();
5600 let mut orders = vec![Vec::new(); dictionaries.len()];
5601 let workers = std::thread::available_parallelism()
5602 .map_or(1, usize::from)
5603 .min(MAX_FREQUENCY_WORKERS)
5604 .min(present.len());
5605 if workers <= 1 {
5606 for at in present {
5607 if let Some(dictionary) = &dictionaries[at] {
5608 orders[at] = dictionary.ranked();
5609 }
5610 }
5611 return Ok(orders);
5612 }
5613 let width = present.len().div_ceil(workers);
5614 let pieces = std::thread::scope(|scope| {
5615 present
5616 .chunks(width)
5617 .map(|columns| {
5618 scope.spawn(|| {
5619 columns
5620 .iter()
5621 .filter_map(|&at| dictionaries[at].as_ref().map(|held| (at, held.ranked())))
5622 .collect::<Vec<_>>()
5623 })
5624 })
5625 .collect::<Vec<_>>()
5626 .into_iter()
5627 .map(|handle| {
5628 handle.join().map_err(|_| Error::internal("a dictionary sort worker panicked"))
5629 })
5630 .collect::<Result<Vec<_>>>()
5631 })?;
5632 for piece in pieces {
5633 for (at, order) in piece {
5634 orders[at] = order;
5635 }
5636 }
5637 Ok(orders)
5638}
5639
5640fn encode_global_dictionary(
5641 dictionary: GlobalDictionary,
5642 order: &[(u64, u32)],
5643) -> Result<EncodedDictionary> {
5644 let values = dictionary.offsets.len() - 1;
5645 if order.len() != values {
5646 return Err(invalid("global dictionary order does not cover its values"));
5647 }
5648 let blocks = values.div_ceil(TEXT_PAYLOAD_VALUES);
5649 let payload = encode_payload(&dictionary)?;
5650 if payload.len() != blocks {
5651 return Err(invalid("global dictionary payload is not the blocks it says it is"));
5652 }
5653 let (ranks, rank_ends) = encode_ranks(order, code_width(values))?;
5654 let rank_blocks = values.div_ceil(TEXT_RANK_BLOCK);
5655 let offset_bits = offset_width(&dictionary.offsets);
5656 let mut index = Vec::with_capacity(
5657 DICTIONARY_HEADER + offset_bytes(values, offset_bits) + (blocks + rank_blocks) * 16,
5658 );
5659 put_u32(
5660 &mut index,
5661 u32::try_from(values).map_err(|_| invalid("global dictionary has too many values"))?,
5662 );
5663 put_u32(&mut index, TEXT_PAYLOAD_VALUES as u32);
5664 put_u32(
5665 &mut index,
5666 u32::try_from(blocks).map_err(|_| invalid("global dictionary has too many blocks"))?,
5667 );
5668 put_u32(&mut index, offset_bits as u32);
5669 encode_offsets(&dictionary.offsets, offset_bits, &mut index)?;
5670 let mut at = 0_u64;
5674 for block in &payload {
5675 at = at
5676 .checked_add(block.len() as u64)
5677 .ok_or_else(|| invalid("global dictionary payload overflow"))?;
5678 put_u64(&mut index, at);
5679 }
5680 for block in &payload {
5681 put_u64(&mut index, checksum(block));
5682 }
5683 if rank_ends.len() != rank_blocks {
5686 return Err(invalid("global dictionary order is not the blocks it says it is"));
5687 }
5688 for end in &rank_ends {
5689 put_u64(&mut index, *end);
5690 }
5691 let mut at = 0_usize;
5692 for end in &rank_ends {
5693 let end = usize::try_from(*end).map_err(|_| invalid("global dictionary order overflow"))?;
5694 put_u64(&mut index, checksum(&ranks[at..end]));
5695 at = end;
5696 }
5697 Ok(EncodedDictionary { index, ranks, payload })
5698}
5699
5700const PAYLOAD_SAMPLE_BLOCKS: usize = 8;
5707
5708fn payload_shapes() -> Vec<chooser::Settled> {
5734 let integers = vec![integer::Kind::Packed];
5735 [
5736 vec![string::Kind::Front, string::Kind::Lz],
5737 vec![string::Kind::Lz, string::Kind::Fsst],
5738 vec![string::Kind::Lz, string::Kind::Plain],
5739 vec![string::Kind::Fsst],
5740 vec![string::Kind::Plain],
5741 ]
5742 .into_iter()
5743 .map(|strings| chooser::Settled::new(strings, integers.clone()))
5744 .collect()
5745}
5746
5747fn encode_payload(dictionary: &GlobalDictionary) -> Result<Vec<Vec<u8>>> {
5753 let values = dictionary.offsets.len() - 1;
5754 let blocks = values.div_ceil(TEXT_PAYLOAD_VALUES);
5755 let run = |block: usize| {
5756 let first = block * TEXT_PAYLOAD_VALUES;
5757 let last = (first + TEXT_PAYLOAD_VALUES).min(values);
5758 (first..last)
5759 .map(|value| {
5760 let from = dictionary.offsets[value] as usize;
5761 let to = dictionary.offsets[value + 1] as usize;
5762 &dictionary.payload[from..to]
5763 })
5764 .collect::<Vec<_>>()
5765 };
5766 let shape = (blocks > PAYLOAD_SAMPLE_BLOCKS).then(|| settle_shape(&run, blocks)).transpose()?;
5769 let one = |block: usize| match &shape {
5770 Some(shape) => string::encode_with(&run(block), shape),
5771 None => string::encode(&run(block)),
5772 };
5773 let workers = std::thread::available_parallelism()
5774 .map_or(1, usize::from)
5775 .min(MAX_FREQUENCY_WORKERS)
5776 .min(blocks);
5777 if workers <= 1 {
5778 return (0..blocks).map(one).collect();
5779 }
5780 let next = AtomicUsize::new(0);
5781 let pieces = std::thread::scope(|scope| {
5782 (0..workers)
5783 .map(|_| {
5784 scope.spawn(|| {
5785 let mut mine = Vec::new();
5786 loop {
5787 let block = next.fetch_add(1, Atomic::Relaxed);
5788 if block >= blocks {
5789 break;
5790 }
5791 mine.push((block, one(block)?));
5792 }
5793 Ok(mine)
5794 })
5795 })
5796 .collect::<Vec<_>>()
5797 .into_iter()
5798 .map(|handle| {
5799 handle.join().map_err(|_| Error::internal("a dictionary encode worker panicked"))?
5800 })
5801 .collect::<Result<Vec<_>>>()
5802 })?;
5803 let mut payload = vec![Vec::new(); blocks];
5804 for piece in pieces {
5805 for (block, bytes) in piece {
5806 payload[block] = bytes;
5807 }
5808 }
5809 Ok(payload)
5810}
5811
5812fn settle_shape<'a>(
5820 run: &dyn Fn(usize) -> Vec<&'a [u8]>,
5821 blocks: usize,
5822) -> Result<chooser::Settled> {
5823 let last = blocks - 1;
5824 let sample = (0..PAYLOAD_SAMPLE_BLOCKS)
5825 .map(|region| run(region * last / (PAYLOAD_SAMPLE_BLOCKS - 1)))
5826 .collect::<Vec<_>>();
5827 let mut best: Option<(chooser::Settled, usize)> = None;
5828 for shape in payload_shapes() {
5829 let mut size = 0;
5830 for block in &sample {
5831 size += string::encode_with(block, &shape)?.len();
5832 }
5833 if best.as_ref().is_none_or(|(_, smallest)| size < *smallest) {
5834 best = Some((shape, size));
5835 }
5836 }
5837 best.map(|(shape, _)| shape)
5838 .ok_or_else(|| invalid("no shape applies to a global dictionary payload"))
5839}
5840
5841fn encode_ranks(order: &[(u64, u32)], code_bits: usize) -> Result<(Vec<u8>, Vec<u64>)> {
5848 let mut out = Vec::with_capacity(order.len() * 4);
5849 let mut ends = Vec::with_capacity(order.len().div_ceil(TEXT_RANK_BLOCK));
5850 let mut heads = Vec::with_capacity(TEXT_RANK_BLOCK);
5851 let mut codes = Vec::with_capacity(TEXT_RANK_BLOCK);
5852 for block in order.chunks(TEXT_RANK_BLOCK) {
5853 let base = block.first().map_or(0, |&(head, _)| head);
5856 let span = block.last().map_or(0, |&(head, _)| head.wrapping_sub(base));
5857 let width = (u64::BITS - span.leading_zeros()) as usize;
5858 heads.clear();
5859 codes.clear();
5860 for &(head, code) in block {
5861 heads.push(head.wrapping_sub(base));
5862 codes.push(u64::from(code));
5863 }
5864 put_u64(&mut out, base);
5865 out.push(width as u8);
5866 bitpack::pack_tail(&heads, width, &mut out)
5867 .map_err(|_| invalid("global dictionary heads do not pack"))?;
5868 bitpack::pack_tail(&codes, code_bits, &mut out)
5869 .map_err(|_| invalid("global dictionary codes do not pack"))?;
5870 ends.push(out.len() as u64);
5871 }
5872 Ok((out, ends))
5873}
5874
5875fn open_global_dictionary(
5882 file: Arc<File>,
5883 page: Page,
5884 ty: &LogicalType,
5885 keep_budget: usize,
5886) -> Result<Vector> {
5887 if ty != &LogicalType::Varchar {
5888 return Err(invalid("global dictionary belongs to a non-string column"));
5889 }
5890 let mut header = [0; DICTIONARY_HEADER];
5891 read_at(&file, page.offset, &mut header)?;
5892 let count = u32::from_le_bytes(header[0..4].try_into().expect("four bytes")) as usize;
5893 let per_block = u32::from_le_bytes(header[4..8].try_into().expect("four bytes")) as usize;
5894 let blocks = u32::from_le_bytes(header[8..12].try_into().expect("four bytes")) as usize;
5895 let offset_bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
5896 if per_block != TEXT_PAYLOAD_VALUES {
5897 return Err(invalid("global dictionary block width differs"));
5898 }
5899 if blocks != count.div_ceil(TEXT_PAYLOAD_VALUES) {
5900 return Err(invalid("global dictionary block count differs from its value count"));
5901 }
5902 if offset_bits > u32::BITS as usize {
5903 return Err(invalid("global dictionary packs offsets past a payload"));
5904 }
5905 let offset_len = offset_bytes(count, offset_bits);
5906 let ranks = count;
5911 let rank_blocks = ranks.div_ceil(TEXT_RANK_BLOCK);
5912 let hash_len = blocks
5915 .checked_add(rank_blocks)
5916 .and_then(|words| words.checked_mul(16))
5917 .ok_or_else(|| invalid("global dictionary block count overflow"))?;
5918 let index_len = DICTIONARY_HEADER
5919 .checked_add(offset_len)
5920 .and_then(|len| len.checked_add(hash_len))
5921 .ok_or_else(|| invalid("global dictionary header overflow"))?;
5922 if index_len > page.length as usize {
5923 return Err(invalid("global dictionary offset index exceeds its page"));
5924 }
5925 let mut index = vec![0; index_len];
5926 index[..DICTIONARY_HEADER].copy_from_slice(&header);
5927 read_at(&file, page.offset + DICTIONARY_HEADER as u64, &mut index[DICTIONARY_HEADER..])?;
5928 if checksum(&index) != page.hash {
5929 return Err(invalid("global dictionary index checksum differs"));
5930 }
5931 let offsets = index[DICTIONARY_HEADER..DICTIONARY_HEADER + offset_len].to_vec();
5932 let mut words = index[DICTIONARY_HEADER + offset_len..]
5933 .chunks_exact(8)
5934 .map(|part| u64::from_le_bytes(part.try_into().expect("eight bytes")))
5935 .collect::<Vec<_>>();
5936 let mut hashes = words.split_off(blocks);
5937 let mut rank_ends = hashes.split_off(blocks);
5938 let rank_hashes = rank_ends.split_off(rank_blocks);
5939 let ends = words;
5940 if rank_ends.windows(2).any(|pair| pair[0] >= pair[1]) {
5943 return Err(invalid("global dictionary order blocks do not rise"));
5944 }
5945 let rank_len = usize::try_from(rank_ends.last().copied().unwrap_or_default())
5946 .map_err(|_| invalid("global dictionary rank overflow"))?;
5947 let body_len = index_len
5948 .checked_add(rank_len)
5949 .ok_or_else(|| invalid("global dictionary header overflow"))?;
5950 if body_len > page.length as usize {
5951 return Err(invalid("global dictionary order exceeds its page"));
5952 }
5953 let stored_len = page.length as usize - body_len;
5956 if ends.last().copied().unwrap_or_default() as usize != stored_len
5957 || ends.windows(2).any(|pair| pair[0] > pair[1])
5958 {
5959 return Err(invalid("global dictionary blocks do not bound the payload"));
5960 }
5961 Vector::external_text(
5962 LogicalType::Varchar,
5963 Arc::new(NativeText {
5964 file,
5965 values: count,
5966 offsets,
5967 offset_bits,
5968 ranks,
5969 rank_at: page.offset + index_len as u64,
5970 rank_ends,
5971 rank_hashes,
5972 rank_blocks: (0..rank_blocks).map(|_| OnceLock::new()).collect(),
5973 code_bits: code_width(count),
5974 code_ranks: OnceLock::new(),
5975 payload: page.offset + body_len as u64,
5976 ends,
5977 hashes,
5978 blocks: (0..blocks).map(|_| OnceLock::new()).collect(),
5979 keep_budget,
5980 payload_kept: AtomicUsize::new(0),
5981 searched: Mutex::new(HashMap::new()),
5982 }),
5983 )
5984}
5985
5986fn page_encoding(ty: &LogicalType, rows: usize, bytes: &[u8]) -> String {
5999 fn cascade_at(rows: usize, bytes: &[u8]) -> Result<(u8, usize)> {
6001 let mut cur = Cursor { bytes, at: 0 };
6002 let codec = cur.u8()?;
6003 if cur.u8()? == 2 {
6004 cur.take(rows.div_ceil(8))?;
6005 }
6006 Ok((codec, cur.at))
6007 }
6008 let Ok((codec, at)) = cascade_at(rows, bytes) else {
6009 return "UNREADABLE".to_string();
6010 };
6011 let tail = &bytes[at..];
6012 let described = |described: Result<String>| described.unwrap_or_else(|_| "UNREADABLE".into());
6013 match codec {
6014 0 => match ty {
6015 LogicalType::Varchar | LogicalType::Blob => "PLAIN".to_string(),
6016 _ => "FIXED".to_string(),
6017 },
6018 1 => "DICT(PLAIN)".to_string(),
6019 2 => "FOR+BITPACK".to_string(),
6020 3 => "TABLE DICT".to_string(),
6021 4 => format!("TABLE DICT({})", described(integer::describe(tail))),
6022 5 => described(integer::describe(tail)),
6023 6 => described(string::describe(tail)),
6024 other => format!("CODEC {other}"),
6025 }
6026}
6027
6028fn decode(
6029 ty: &LogicalType,
6030 rows: usize,
6031 bytes: &[u8],
6032 global: Option<Arc<Vector>>,
6033) -> Result<Vector> {
6034 let mut cur = Cursor { bytes, at: 0 };
6035 let codec = cur.u8()?;
6036 let flag = cur.u8()?;
6037 let validity = match flag {
6038 0 => Validity::AllValid,
6039 1 => Validity::AllInvalid,
6040 2 => {
6041 let mask = cur.take(rows.div_ceil(8))?;
6042 Validity::from_iter(rows, |row| mask[row / 8] >> (row % 8) & 1 == 1)
6043 }
6044 _ => return Err(invalid("page validity tag differs")),
6045 };
6046 if codec == 1 {
6047 if ty != &LogicalType::Varchar {
6048 return Err(invalid("dictionary codec belongs to a non-string page"));
6049 }
6050 let count = cur.u32()? as usize;
6051 let payload_len = cur.u32()? as usize;
6052 let offset_bytes = cur.take(
6053 (count + 1)
6054 .checked_mul(4)
6055 .ok_or_else(|| invalid("dictionary offset count overflow"))?,
6056 )?;
6057 let offsets = offset_bytes
6058 .chunks_exact(4)
6059 .map(|part| u32::from_le_bytes(part.try_into().expect("four bytes")))
6060 .collect::<Vec<_>>();
6061 let payload = cur.take(payload_len)?.to_vec();
6062 if offsets.first() != Some(&0)
6063 || offsets.last().copied().map(|last| last as usize) != Some(payload.len())
6064 || offsets.windows(2).any(|pair| pair[0] > pair[1])
6065 {
6066 return Err(invalid("dictionary offsets do not bound the payload"));
6067 }
6068 let mut strings = StringColumn::over(Buffer::from_vec(payload).into_page());
6071 for pair in offsets.windows(2) {
6072 strings.push_in_place(pair[0] as usize, (pair[1] - pair[0]) as usize)?;
6073 }
6074 let mut codes = Vec::with_capacity(rows);
6075 for _ in 0..rows {
6076 codes.push(cur.u32()?);
6077 }
6078 if codes.iter().any(|code| *code as usize >= count) {
6079 return Err(invalid("dictionary code is out of range"));
6080 }
6081 if cur.at != bytes.len() {
6082 return Err(invalid("dictionary page has trailing bytes"));
6083 }
6084 let dictionary = Vector::flat(LogicalType::Varchar, Data::Varlen(strings))?;
6085 return Ok(Vector::dictionary(codes, dictionary)?.with_validity(validity));
6086 }
6087 if codec == 3 || codec == 4 {
6088 let dictionary = global.ok_or_else(|| invalid("global code page has no dictionary"))?;
6089 let codes = if codec == 4 {
6090 let wide = integer::decode(&bytes[cur.at..])?;
6093 if wide.len() != rows {
6094 return Err(invalid("encoded code page holds the wrong number of rows"));
6095 }
6096 let mut codes = Vec::with_capacity(wide.len());
6103 let mut seen = 0_i64;
6104 for &code in &wide {
6105 seen |= code;
6106 codes.push(code as u32);
6107 }
6108 if seen < 0 || seen > i64::from(u32::MAX) {
6109 return Err(invalid("code is not a code"));
6110 }
6111 codes
6112 } else {
6113 let mut codes = Vec::with_capacity(rows);
6114 for _ in 0..rows {
6115 codes.push(cur.u32()?);
6116 }
6117 if cur.at != bytes.len() {
6118 return Err(invalid("global code page has trailing bytes"));
6119 }
6120 codes
6121 };
6122 let highest = codes.iter().copied().max();
6123 return Ok(Vector::stable_dictionary_validated(codes, dictionary, highest)?
6124 .with_validity(validity));
6125 }
6126 if codec == 6 {
6127 if ty != &LogicalType::Varchar {
6128 return Err(invalid("compressed text codec belongs to a non-string page"));
6129 }
6130 let (payload, ends) = string::decode_flat(&bytes[cur.at..])?.into_parts();
6134 if ends.len() != rows {
6135 return Err(invalid("compressed text page holds the wrong number of rows"));
6136 }
6137 let mut values = StringColumn::over(Buffer::from_vec(payload).into_page());
6140 let mut start = 0;
6141 for end in ends {
6142 let len = end
6143 .checked_sub(start)
6144 .ok_or_else(|| invalid("compressed text value ends before it starts"))?;
6145 values.push_in_place(start, len)?;
6146 start = end;
6147 }
6148 return Ok(Vector::flat(ty.clone(), Data::Varlen(values))?.with_validity(validity));
6149 }
6150 if codec == 5 {
6151 let values = integer::decode(&bytes[cur.at..])?;
6153 if values.len() != rows {
6154 return Err(invalid("cascade page holds the wrong number of rows"));
6155 }
6156 let data = narrowed(ty, values)?;
6157 return Ok(Vector::flat(ty.clone(), data)?.with_validity(validity));
6158 }
6159 if codec == 2 {
6160 let width = u32::from(cur.u8()?);
6161 let base = i128::from_le_bytes(cur.take(16)?.try_into().expect("sixteen bytes"));
6162 let count = cur.u32()? as usize;
6163 let mut words = Vec::with_capacity(count);
6164 for _ in 0..count {
6165 words.push(cur.u64()?);
6166 }
6167 if cur.at != bytes.len() {
6168 return Err(invalid("packed page has trailing bytes"));
6169 }
6170 return Ok(Vector::packed(ty.clone(), words, width, base, rows)?.with_validity(validity));
6171 }
6172 if codec != 0 {
6173 return Err(invalid("page codec is unknown"));
6174 }
6175 let data = match ty {
6176 LogicalType::TinyInt => {
6177 let values = cur.take(rows)?;
6178 Data::Int8(values.iter().map(|item| *item as i8).collect::<Vec<_>>().into())
6179 }
6180 LogicalType::UTinyInt => Data::UInt8(cur.take(rows)?.to_vec().into()),
6181 LogicalType::SmallInt => {
6182 let values =
6183 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
6184 Data::Int16(
6185 values
6186 .chunks_exact(2)
6187 .map(|item| i16::from_le_bytes(item.try_into().expect("two bytes")))
6188 .collect::<Vec<_>>()
6189 .into(),
6190 )
6191 }
6192 LogicalType::USmallInt => {
6193 let values =
6194 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
6195 Data::UInt16(
6196 values
6197 .chunks_exact(2)
6198 .map(|item| u16::from_le_bytes(item.try_into().expect("two bytes")))
6199 .collect::<Vec<_>>()
6200 .into(),
6201 )
6202 }
6203 LogicalType::UInteger => {
6204 let values =
6205 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
6206 Data::UInt32(
6207 values
6208 .chunks_exact(4)
6209 .map(|item| u32::from_le_bytes(item.try_into().expect("four bytes")))
6210 .collect::<Vec<_>>()
6211 .into(),
6212 )
6213 }
6214 LogicalType::UBigInt => {
6215 let values =
6216 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
6217 Data::UInt64(
6218 values
6219 .chunks_exact(8)
6220 .map(|item| u64::from_le_bytes(item.try_into().expect("eight bytes")))
6221 .collect::<Vec<_>>()
6222 .into(),
6223 )
6224 }
6225 LogicalType::Integer | LogicalType::Date => {
6226 let values =
6227 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
6228 Data::Int32(
6229 values
6230 .chunks_exact(4)
6231 .map(|item| i32::from_le_bytes(item.try_into().expect("four bytes")))
6232 .collect::<Vec<_>>()
6233 .into(),
6234 )
6235 }
6236 LogicalType::BigInt
6237 | LogicalType::Timestamp
6238 | LogicalType::Time
6239 | LogicalType::TimeTz
6240 | LogicalType::TimestampTz
6241 | LogicalType::TimestampS
6242 | LogicalType::TimestampMs
6243 | LogicalType::TimestampNs => {
6244 let values =
6245 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
6246 Data::Int64(
6247 values
6248 .chunks_exact(8)
6249 .map(|item| i64::from_le_bytes(item.try_into().expect("eight bytes")))
6250 .collect::<Vec<_>>()
6251 .into(),
6252 )
6253 }
6254 LogicalType::HugeInt | LogicalType::Uuid => {
6255 let values =
6256 cur.take(rows.checked_mul(16).ok_or_else(|| invalid("page size overflow"))?)?;
6257 Data::Int128(
6258 values
6259 .chunks_exact(16)
6260 .map(|item| i128::from_le_bytes(item.try_into().expect("sixteen bytes")))
6261 .collect::<Vec<_>>()
6262 .into(),
6263 )
6264 }
6265 LogicalType::UHugeInt => {
6266 let values =
6267 cur.take(rows.checked_mul(16).ok_or_else(|| invalid("page size overflow"))?)?;
6268 Data::UInt128(
6269 values
6270 .chunks_exact(16)
6271 .map(|item| u128::from_le_bytes(item.try_into().expect("sixteen bytes")))
6272 .collect::<Vec<_>>()
6273 .into(),
6274 )
6275 }
6276 LogicalType::Float => {
6277 let values =
6278 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
6279 Data::Float32(
6280 values
6281 .chunks_exact(4)
6282 .map(|item| f32::from_le_bytes(item.try_into().expect("four bytes")))
6283 .collect::<Vec<_>>()
6284 .into(),
6285 )
6286 }
6287 LogicalType::Double => {
6288 let values =
6289 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
6290 Data::Float64(
6291 values
6292 .chunks_exact(8)
6293 .map(|item| f64::from_le_bytes(item.try_into().expect("eight bytes")))
6294 .collect::<Vec<_>>()
6295 .into(),
6296 )
6297 }
6298 LogicalType::Interval => {
6299 let values =
6300 cur.take(rows.checked_mul(16).ok_or_else(|| invalid("page size overflow"))?)?;
6301 Data::Interval(
6302 values
6303 .chunks_exact(16)
6304 .map(|item| {
6305 (
6306 i32::from_le_bytes(item[..4].try_into().expect("four bytes")),
6307 i32::from_le_bytes(item[4..8].try_into().expect("four bytes")),
6308 i64::from_le_bytes(item[8..].try_into().expect("eight bytes")),
6309 )
6310 })
6311 .collect::<Vec<_>>()
6312 .into(),
6313 )
6314 }
6315 LogicalType::Boolean => {
6316 let values = cur.take(rows)?;
6317 if values.iter().any(|value| *value > 1) {
6318 return Err(invalid("boolean page has another value"));
6319 }
6320 Data::Bool(values.iter().map(|value| *value == 1).collect::<Vec<_>>().into())
6321 }
6322 LogicalType::Decimal { .. } => match ty.physical() {
6325 PhysicalType::Int16 => {
6326 let values =
6327 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
6328 Data::Int16(
6329 values
6330 .chunks_exact(2)
6331 .map(|item| i16::from_le_bytes(item.try_into().expect("two bytes")))
6332 .collect::<Vec<_>>()
6333 .into(),
6334 )
6335 }
6336 PhysicalType::Int32 => {
6337 let values =
6338 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
6339 Data::Int32(
6340 values
6341 .chunks_exact(4)
6342 .map(|item| i32::from_le_bytes(item.try_into().expect("four bytes")))
6343 .collect::<Vec<_>>()
6344 .into(),
6345 )
6346 }
6347 PhysicalType::Int64 => {
6348 let values =
6349 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
6350 Data::Int64(
6351 values
6352 .chunks_exact(8)
6353 .map(|item| i64::from_le_bytes(item.try_into().expect("eight bytes")))
6354 .collect::<Vec<_>>()
6355 .into(),
6356 )
6357 }
6358 _ => {
6359 let values =
6360 cur.take(rows.checked_mul(16).ok_or_else(|| invalid("page size overflow"))?)?;
6361 Data::Int128(
6362 values
6363 .chunks_exact(16)
6364 .map(|item| i128::from_le_bytes(item.try_into().expect("sixteen bytes")))
6365 .collect::<Vec<_>>()
6366 .into(),
6367 )
6368 }
6369 },
6370 LogicalType::Varchar | LogicalType::Blob | LogicalType::Bit => {
6371 let offset_bytes = cur
6372 .take((rows + 1).checked_mul(4).ok_or_else(|| invalid("offset count overflow"))?)?;
6373 let offsets = offset_bytes
6374 .chunks_exact(4)
6375 .map(|part| u32::from_le_bytes(part.try_into().expect("four bytes")))
6376 .collect::<Vec<_>>();
6377 let payload = cur.take(bytes.len() - cur.at)?.to_vec();
6378 if offsets.first() != Some(&0)
6379 || offsets.last().copied().map(|last| last as usize) != Some(payload.len())
6380 || offsets.windows(2).any(|pair| pair[0] > pair[1])
6381 {
6382 return Err(invalid("string offsets do not bound the payload"));
6383 }
6384 let mut values = StringColumn::over(Buffer::from_vec(payload).into_page());
6392 let text = ty == &LogicalType::Varchar;
6393 for pair in offsets.windows(2) {
6394 let (at, len) = (pair[0] as usize, (pair[1] - pair[0]) as usize);
6395 if text {
6396 values.push_in_place(at, len)?;
6397 } else {
6398 values.push_bytes_in_place(at, len)?;
6399 }
6400 }
6401 Data::Varlen(values)
6402 }
6403 _ => return Err(Error::not_implemented(format!("native page for {ty}"))),
6404 };
6405 if cur.at != bytes.len() {
6406 return Err(invalid("page has trailing bytes"));
6407 }
6408 Ok(Vector::flat(ty.clone(), data)?.with_validity(validity))
6409}
6410
6411#[cfg(test)]
6412mod tests {
6413 use std::fs;
6414 use std::io::{Seek, SeekFrom, Write};
6415 use std::path::PathBuf;
6416 use std::time::{SystemTime, UNIX_EPOCH};
6417
6418 use rudb_common::Stat;
6419 use rudb_common::Value;
6420 use rudb_common::bounds::{Frequencies, Op, Remainder, Zones};
6421 use rudb_common::stat::Provenance;
6422
6423 use super::*;
6424
6425 #[test]
6426 fn checksum_matches_fixed_vectors() {
6427 assert_eq!(checksum(b""), 0xef46_db37_51d8_e999);
6428 assert_eq!(checksum(b"a"), 0xd24e_c4f1_a98c_6e5b);
6429 assert_eq!(checksum(b"abc"), 0x44bc_2cf5_ad77_0999);
6430 }
6431
6432 fn path(label: &str) -> PathBuf {
6433 let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
6434 std::env::temp_dir().join(format!("rudb-native-{label}-{}-{stamp}.rdb", std::process::id()))
6435 }
6436
6437 #[test]
6439 fn a_read_at_an_offset_ignores_where_another_thread_left_the_cursor() {
6440 const SPANS: usize = 64;
6441 const SPAN: usize = 512;
6442 let path = path("positional");
6443 let content: Vec<u8> =
6444 (0..SPANS).flat_map(|span| std::iter::repeat_n(span as u8, SPAN)).collect();
6445 fs::write(&path, &content).expect("the file is written");
6446 let file = Arc::new(File::open(&path).expect("the file opens"));
6447 std::thread::scope(|scope| {
6448 for _ in 0..8 {
6449 let file = Arc::clone(&file);
6450 scope.spawn(move || {
6451 for _ in 0..64 {
6452 for span in 0..SPANS {
6453 let mut bytes = [0_u8; SPAN];
6454 read_at(&file, (span * SPAN) as u64, &mut bytes)
6455 .expect("the span reads");
6456 assert!(
6457 bytes.iter().all(|byte| *byte == span as u8),
6458 "span {span} came back as {}",
6459 bytes[0],
6460 );
6461 }
6462 }
6463 });
6464 }
6465 });
6466 let mut past = [0_u8; SPAN];
6467 let end = (SPANS * SPAN) as u64;
6468 let error = read_at(&file, end, &mut past).expect_err("a read past the end is refused");
6469 assert!(error.message().contains("ends before its declared length"), "{error}");
6470 drop(file);
6471 let _ = fs::remove_file(&path);
6472 }
6473
6474 #[test]
6480 fn a_writer_puts_a_page_where_it_said_it_did_wherever_the_cursor_has_got_to() {
6481 let path = path("cursor");
6482 let mut writer = Writer::create(
6483 &path,
6484 "items",
6485 vec![
6486 Field::required("id", LogicalType::Integer),
6487 Field::new("text", LogicalType::Varchar),
6488 ],
6489 )
6490 .expect("new file");
6491 writer.append(&sample()).expect("first part");
6492 writer.file.seek(SeekFrom::Start(0)).expect("the cursor goes back to the header");
6493 writer.append(&sample()).expect("second part");
6494 writer.file.seek(SeekFrom::Start(1)).expect("and somewhere useless again");
6495 writer.finish().expect("commit");
6496 let reader = Reader::open(&path).expect("reopen from disk");
6497 assert_eq!(reader.table().rows(), 6);
6498 let ids = reader.read(0, &[0]).expect("the integer page reads back");
6499 assert_eq!(ids.value_at(0, 0), Value::Integer(4));
6500 assert_eq!(ids.value_at(2, 0), Value::Integer(-2));
6501 let text = reader.read(1, &[1]).expect("the text page reads back");
6502 assert_eq!(text.value_at(1, 0), Value::Null);
6503 assert_eq!(text.value_at(2, 0), Value::Varchar("long text after a slash".into()));
6504 let end = reader.table().stripes().iter().flat_map(|stripe| {
6507 stripe
6508 .pages
6509 .iter()
6510 .map(|page| page.offset + u64::from(page.length))
6511 .chain(std::iter::once(stripe.index.offset + u64::from(stripe.index.length)))
6512 });
6513 let last = end.fold(HEADER, u64::max);
6514 let directory = fs::metadata(&path).expect("the file is there").len();
6515 assert!(last <= directory, "a page runs to {last} in a file of {directory} bytes");
6516 fs::remove_file(path).expect("remove scratch file");
6517 }
6518
6519 fn dictionary_index_len(header: &[u8; DICTIONARY_HEADER]) -> u64 {
6525 let count = u64::from(u32::from_le_bytes(header[0..4].try_into().expect("four bytes")));
6526 let blocks = u64::from(u32::from_le_bytes(header[8..12].try_into().expect("four bytes")));
6527 let bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
6528 let rank_blocks = count.div_ceil(TEXT_RANK_BLOCK as u64);
6529 DICTIONARY_HEADER as u64
6530 + offset_bytes(count as usize, bits) as u64
6531 + (blocks + rank_blocks) * 16
6532 }
6533
6534 fn last_rank_end(file: &File, offset: u64, header: &[u8; DICTIONARY_HEADER]) -> u64 {
6536 let count = u64::from(u32::from_le_bytes(header[0..4].try_into().expect("four bytes")));
6537 let blocks = u64::from(u32::from_le_bytes(header[8..12].try_into().expect("four bytes")));
6538 let bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
6539 let rank_blocks = count.div_ceil(TEXT_RANK_BLOCK as u64);
6540 let at = offset
6541 + DICTIONARY_HEADER as u64
6542 + offset_bytes(count as usize, bits) as u64
6543 + blocks * 16
6544 + (rank_blocks - 1) * 8;
6545 let mut end = [0; 8];
6546 read_at(file, at, &mut end).expect("the last rank block end");
6547 u64::from_le_bytes(end)
6548 }
6549
6550 fn sample() -> Chunk {
6551 Chunk::new(vec![
6552 Vector::from_values(
6553 LogicalType::Integer,
6554 &[Value::Integer(4), Value::Integer(9), Value::Integer(-2)],
6555 )
6556 .expect("integers"),
6557 Vector::from_values(
6558 LogicalType::Varchar,
6559 &[
6560 Value::Varchar("alpha".into()),
6561 Value::Null,
6562 Value::Varchar("long text after a slash".into()),
6563 ],
6564 )
6565 .expect("strings"),
6566 ])
6567 .expect("matching rows")
6568 }
6569
6570 fn sample_ids() -> Chunk {
6571 Chunk::new(vec![
6572 Vector::flat(LogicalType::Integer, Data::Int32(vec![7, 8, 9].into()))
6573 .expect("integers"),
6574 ])
6575 .expect("one column")
6576 }
6577
6578 #[test]
6579 fn the_planner_gets_the_null_count_off_the_same_directory_the_bounds_are_in() {
6580 let path = path("nulls_for_the_planner");
6583 let mut writer =
6584 Writer::create(&path, "items", vec![Field::new("a", LogicalType::Integer)])
6585 .expect("new file");
6586 let rows = Chunk::new(vec![
6587 Vector::from_values(
6588 LogicalType::Integer,
6589 &[
6590 Value::Integer(4),
6591 Value::Null,
6592 Value::Integer(9),
6593 Value::Null,
6594 Value::Integer(1),
6595 Value::Integer(2),
6596 ],
6597 )
6598 .expect("integers"),
6599 ])
6600 .expect("one column");
6601 writer.append(&rows).expect("the only part");
6602 writer.finish().expect("commit");
6603 let reader = Reader::open(&path).expect("reopen from disk");
6604 let stripes = Stripes::new(reader);
6605 let column = stripes.column("a").expect("the file has that column");
6606 assert_eq!(stripes.nulls(column), Stat::exact(2, Provenance::NullCount));
6607 assert_eq!(stripes.nulls(column + 1), Stat::Unknown);
6610 fs::remove_file(&path).expect("clean up");
6611 }
6612
6613 #[test]
6614 fn the_planner_gets_a_row_count_per_value_off_a_complete_synopsis() {
6615 let path = path("frequencies_for_the_planner");
6620 let mut writer =
6621 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6622 .expect("new file");
6623 let rows = Chunk::new(vec![
6624 Vector::from_values(
6625 LogicalType::Integer,
6626 &[
6627 Value::Integer(4),
6628 Value::Integer(4),
6629 Value::Integer(4),
6630 Value::Integer(9),
6631 Value::Integer(9),
6632 Value::Integer(1),
6633 ],
6634 )
6635 .expect("integers"),
6636 ])
6637 .expect("one column");
6638 writer.append(&rows).expect("the only part");
6639 writer.finish().expect("commit");
6640 let reader = Reader::open(&path).expect("reopen from disk");
6641 let common = Common::new(reader);
6642 assert_eq!(common.rows(), 6);
6643 let column = common.column("id").expect("the file has that column");
6644 assert_eq!(common.column("nothing"), None);
6645 assert_eq!(
6646 common.rows_with(column, &Bound::Int(4)),
6647 Stat::exact(3, Provenance::FrequencySynopsis)
6648 );
6649 assert_eq!(
6651 common.rows_with(column, &Bound::Int(7)),
6652 Stat::exact(0, Provenance::FrequencySynopsis)
6653 );
6654 assert_eq!(common.rows_with(column, &Bound::Bytes(b"four".to_vec())), Stat::Unknown);
6657 assert_eq!(common.remainder(column), None);
6660 fs::remove_file(&path).expect("clean up");
6661 }
6662
6663 #[test]
6664 fn the_planner_gets_an_exact_count_for_a_leading_value_of_an_incomplete_synopsis() {
6665 let path = path("frequency_prefix_for_the_planner");
6672 let mut writer =
6673 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6674 .expect("new file");
6675 let mut values = vec![Value::Integer(1); 10_000];
6676 for _ in 0..10 {
6677 values.extend((0..600).map(|tail| Value::Integer(1_000 + tail)));
6678 }
6679 for part in values.chunks(8_000) {
6682 let rows = Chunk::new(vec![
6683 Vector::from_values(LogicalType::Integer, part).expect("integers"),
6684 ])
6685 .expect("one column");
6686 writer.append(&rows).expect("a part");
6687 }
6688 writer.finish().expect("commit");
6689 let reader = Reader::open(&path).expect("reopen from disk");
6690 let prefix =
6691 reader.frequency_prefix(0).expect("a readable synopsis").expect("the column has one");
6692 assert_eq!(prefix.entries.len(), 512);
6695 assert_eq!(prefix.omitted_max, 10);
6696 let common = Common::new(reader);
6697 assert_eq!(common.rows(), 16_000);
6698 let column = common.column("id").expect("the file has that column");
6699 assert_eq!(
6700 common.rows_with(column, &Bound::Int(1)),
6701 Stat::exact(10_000, Provenance::FrequencySynopsis)
6702 );
6703 assert_eq!(
6705 common.rows_with(column, &Bound::Int(1_100)),
6706 Stat::exact(10, Provenance::FrequencySynopsis)
6707 );
6708 assert_eq!(common.rows_with(column, &Bound::Int(1_550)), Stat::Unknown);
6711 assert_eq!(common.rows_with(column, &Bound::Int(9_999)), Stat::Unknown);
6714 let remainder = common.remainder(column).expect("the list is a prefix");
6718 assert_eq!(remainder, Remainder { rows: 890, listed: 512, most: 10 });
6719 assert_eq!(remainder.rows / (601 - remainder.listed), 10);
6720 fs::remove_file(&path).expect("clean up");
6721 }
6722
6723 #[test]
6725 fn a_file_holding_no_table_commits_and_opens_and_a_table_can_be_added_to_it() {
6726 let path = path("empty");
6727 Writer::empty(&path).expect("a file with nothing in it");
6728 let catalog = Catalog::open(&path).expect("the empty file opens");
6729 assert_eq!(catalog.len(), 0);
6730 assert!(catalog.is_empty());
6731 assert_eq!(catalog.names().count(), 0);
6732 let mut writer =
6735 Writer::open(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6736 .expect("a table goes into the empty file");
6737 writer.append(&sample_ids()).expect("rows");
6738 writer.finish().expect("commit");
6739 let catalog = Catalog::open(&path).expect("the file opens again");
6740 assert_eq!(catalog.names().collect::<Vec<_>>(), vec!["items"]);
6741 fs::remove_file(&path).expect("clean up");
6742 }
6743
6744 #[test]
6745 fn committed_file_reopens_and_reads_only_requested_columns() {
6746 let path = path("reopen");
6747 let mut writer = Writer::create(
6748 &path,
6749 "items",
6750 vec![
6751 Field::required("id", LogicalType::Integer),
6752 Field::new("text", LogicalType::Varchar),
6753 ],
6754 )
6755 .expect("new file");
6756 writer.append(&sample()).expect("first part");
6757 writer.append(&sample()).expect("second part");
6758 writer.finish().expect("commit");
6759 let reader = Reader::open(&path).expect("reopen from disk");
6760 assert_eq!(reader.table().rows(), 6);
6761 assert_eq!(reader.table().stripes().len(), 1);
6764 assert_eq!(reader.parts(), 2);
6765 assert_eq!(reader.part_rows(0), 3);
6766 assert_eq!(reader.part_rows(1), 3);
6767 let text = reader.read(1, &[1]).expect("only text page");
6768 assert_eq!(text.width(), 1);
6769 assert_eq!(text.value_at(1, 0), Value::Null);
6770 assert_eq!(text.value_at(2, 0), Value::Varchar("long text after a slash".into()));
6771 let sparse = reader.read_sparse(1, &[1]).expect("one part without its whole page");
6772 assert_eq!(sparse.width(), 1);
6773 assert_eq!(sparse.value_at(1, 0), Value::Null);
6774 assert_eq!(sparse.value_at(2, 0), Value::Varchar("long text after a slash".into()));
6775 assert!(!reader.skips_codes(0, 1, &[0]).expect("alpha is in the stripe"));
6776 assert!(!reader.skips_codes(0, 1, &[2]).expect("long text is in the stripe"));
6777 assert!(reader.skips_codes(0, 1, &[3]).expect("unknown code is absent"));
6778 let count = reader.read(0, &[]).expect("no page is needed for count");
6779 assert_eq!(count.len(), 3);
6780 assert!(reader.skips(0, &[Probe { column: 0, op: Op::Greater, value: Bound::Int(100) }]));
6781 assert!(!reader.skips(0, &[Probe { column: 0, op: Op::Greater, value: Bound::Int(0) }]));
6782 let integers = reader.top_frequencies(0, 1).expect("valid integer synopsis").expect("kept");
6783 assert_eq!(
6784 integers,
6785 vec![(Value::Integer(-2), 2), (Value::Integer(4), 2), (Value::Integer(9), 2),]
6786 );
6787 let strings = reader.top_frequencies(1, 1).expect("valid string synopsis").expect("kept");
6788 assert_eq!(strings.len(), 3);
6789 assert!(strings.contains(&(Value::Null, 2)));
6790 assert!(strings.contains(&(Value::Varchar("alpha".into()), 2)));
6791 assert!(strings.contains(&(Value::Varchar("long text after a slash".into()), 2)));
6792 fs::remove_file(path).expect("remove scratch file");
6793 }
6794
6795 #[test]
6803 fn runs_handed_over_out_of_order_still_read_back_in_source_order() {
6804 let path = path("interleaved-runs");
6805 let mut writer =
6806 Writer::create(&path, "interleaved", vec![Field::new("v", LogicalType::BigInt)])
6807 .expect("new file");
6808 for morsel in [2_u64, 0, 3, 1] {
6809 let parts = (0..4_u64)
6810 .map(|chunk| {
6811 let first = i64::try_from(morsel * 32 + chunk * 8).expect("small");
6812 let values =
6813 (0..8_i64).map(|row| Value::BigInt(first + row)).collect::<Vec<_>>();
6814 let column =
6815 Vector::from_values(LogicalType::BigInt, &values).expect("a column");
6816 ((morsel, chunk), Chunk::new(vec![column]).expect("one column"))
6817 })
6818 .collect::<Vec<_>>();
6819 writer.append_stripe(parts).expect("a stripe");
6820 }
6821 writer.finish().expect("commit");
6822
6823 let reader = Reader::open(&path).expect("valid directory");
6824 assert_eq!(reader.table().stripes().len(), 4, "a run is a stripe of its own");
6825 assert_eq!(reader.table().rows(), 128);
6826 for part in 0..16_usize {
6827 let read = reader.read(part, &[0]).expect("a part back");
6828 for row in 0..8_usize {
6829 let want = i64::try_from(part * 8 + row).expect("small");
6830 assert_eq!(read.value_at(row, 0), Value::BigInt(want), "part {part} row {row}");
6831 }
6832 }
6833 fs::remove_file(path).expect("remove scratch file");
6834 }
6835
6836 #[test]
6839 fn runs_that_overlap_each_other_are_refused_at_commit() {
6840 let path = path("overlapping-runs");
6841 let mut writer =
6842 Writer::create(&path, "overlapping", vec![Field::new("v", LogicalType::BigInt)])
6843 .expect("new file");
6844 let one = |order: (u64, u64)| {
6845 let column =
6846 Vector::from_values(LogicalType::BigInt, &[Value::BigInt(1)]).expect("a column");
6847 (order, Chunk::new(vec![column]).expect("one column"))
6848 };
6849 writer.append_stripe(vec![one((0, 0)), one((0, 2))]).expect("a stripe");
6852 writer.append_stripe(vec![one((0, 1))]).expect("a stripe");
6853 let error = writer.finish().expect_err("the runs overlap");
6854 assert!(error.message().contains("source order"), "{error}");
6855 fs::remove_file(path).expect("remove scratch file");
6856 }
6857
6858 #[test]
6861 fn a_run_longer_than_a_stripe_is_refused() {
6862 let path = path("overlong-run");
6863 let mut writer =
6864 Writer::create(&path, "overlong", vec![Field::new("v", LogicalType::BigInt)])
6865 .expect("new file");
6866 let parts = (0..=STRIPE_PARTS)
6867 .map(|at| {
6868 let column = Vector::from_values(LogicalType::BigInt, &[Value::BigInt(1)])
6869 .expect("a column");
6870 let chunk = Chunk::new(vec![column]).expect("one column");
6871 ((0, u64::try_from(at).expect("small")), chunk)
6872 })
6873 .collect::<Vec<_>>();
6874 let error = writer.append_stripe(parts).expect_err("one part too many");
6875 assert!(error.message().contains("more parts than it holds"), "{error}");
6876 fs::remove_file(path).expect("remove scratch file");
6877 }
6878
6879 #[test]
6885 fn parts_past_the_stripe_bound_start_a_new_stripe() {
6886 let path = path("stripe-bound");
6887 let mut writer = Writer::create(
6888 &path,
6889 "items",
6890 vec![
6891 Field::required("id", LogicalType::Integer),
6892 Field::new("text", LogicalType::Varchar),
6893 ],
6894 )
6895 .expect("new file");
6896 let parts = STRIPE_PARTS * 2 + 3;
6897 for part in 0..parts {
6898 let id = part as i32;
6899 let chunk = Chunk::new(vec![
6900 Vector::from_values(
6901 LogicalType::Integer,
6902 &[Value::Integer(id), Value::Integer(-id)],
6903 )
6904 .expect("integers"),
6905 Vector::from_values(
6906 LogicalType::Varchar,
6907 &[Value::Varchar(format!("value {part}")), Value::Null],
6908 )
6909 .expect("strings"),
6910 ])
6911 .expect("matching rows");
6912 writer.append(&chunk).expect("one part");
6913 }
6914 writer.finish().expect("commit");
6915
6916 let reader = Reader::open(&path).expect("reopen from disk");
6917 assert_eq!(reader.parts(), parts);
6918 assert_eq!(reader.table().rows(), parts * 2);
6919 assert_eq!(reader.table().stripes().len(), parts.div_ceil(STRIPE_PARTS));
6920 assert_eq!(reader.table().stripes()[0].parts(), STRIPE_PARTS);
6921 assert_eq!(reader.table().stripes()[0].rows(), STRIPE_PARTS * 2);
6922 assert_eq!(reader.table().stripes()[2].parts(), 3);
6923 for part in (0..parts).rev() {
6926 let dense = reader.read(part, &[0, 1]).expect("a whole page read");
6927 let sparse = reader.read_sparse(part, &[0, 1]).expect("one part read");
6928 for chunk in [&dense, &sparse] {
6929 assert_eq!(chunk.len(), 2, "part {part} has its own row count");
6930 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6931 assert_eq!(chunk.value_at(1, 0), Value::Integer(-(part as i32)));
6932 assert_eq!(chunk.value_at(0, 1), Value::Varchar(format!("value {part}")));
6933 assert_eq!(chunk.value_at(1, 1), Value::Null);
6934 }
6935 }
6936 let above = [Probe { column: 0, op: Op::Greater, value: Bound::Int(100) }];
6939 assert!(reader.skips(0, &above), "the first stripe stops at 63");
6940 assert!(!reader.skips(STRIPE_PARTS * 2, &above), "the third stripe reaches 130");
6941 fs::remove_file(path).expect("remove scratch file");
6942 }
6943
6944 fn scattered(n: i64) -> i64 {
6946 n.wrapping_mul(-7_046_029_254_386_353_131)
6947 }
6948
6949 #[test]
6955 fn a_part_is_skipped_when_its_sieve_does_not_hold_the_constant() {
6956 let path = path("sieve-skip");
6957 let mut writer =
6958 Writer::create(&path, "hits", vec![Field::required("id", LogicalType::BigInt)])
6959 .expect("new file");
6960 let parts = STRIPE_PARTS + 3;
6961 let per_part = 128;
6965 for part in 0..parts {
6966 let held: Vec<Value> = (0..per_part)
6967 .map(|row| Value::BigInt(scattered((part * per_part + row) as i64)))
6968 .collect();
6969 let chunk =
6970 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6971 .expect("one column");
6972 writer.append(&chunk).expect("one part");
6973 }
6974 writer.finish().expect("commit");
6975
6976 let reader = Reader::open(&path).expect("reopen from disk");
6977 let probe = |value: i64| Probe {
6978 column: 0,
6979 op: Op::Equal,
6980 value: Bound::Int(i128::from(scattered(value))),
6981 };
6982 for wanted in [0_i64, (per_part + 1) as i64, (parts * per_part - 1) as i64] {
6983 let tests = [probe(wanted)];
6984 let kept: Vec<usize> = (0..parts).filter(|&part| !reader.skips(part, &tests)).collect();
6985 let home = wanted as usize / per_part;
6986 assert!(kept.contains(&home), "the part holding {wanted} is read");
6987 assert!(kept.len() <= 2, "{wanted} keeps {kept:?}, which is more than one stray part");
6991 }
6992 let absent = [probe((parts * per_part) as i64 + 1)];
6993 let kept = (0..parts).filter(|&part| !reader.skips(part, &absent)).count();
6994 assert!(kept <= 1, "{kept} parts of {parts} kept a value no part holds");
6995 let tests = [probe(0)];
6998 assert!(
6999 reader.table().stripes().iter().all(|stripe| !stripe.zone.skips(&tests)),
7000 "the bounds rule out no stripe at all"
7001 );
7002 fs::remove_file(path).expect("remove scratch file");
7003 }
7004
7005 #[test]
7011 fn a_part_is_skipped_when_its_own_bounds_rule_out_a_comparison_the_stripe_keeps() {
7012 let path = path("part-range-skip");
7013 let mut writer =
7014 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
7015 .expect("new file");
7016 let parts = STRIPE_PARTS + 3;
7017 let per_part = 128;
7018 for part in 0..parts {
7019 let held: Vec<Value> = (0..per_part)
7023 .map(|row| {
7024 Value::BigInt((part * 1_000) as i64 + (scattered(row as i64).rem_euclid(900)))
7025 })
7026 .collect();
7027 let chunk =
7028 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
7029 .expect("one column");
7030 writer.append(&chunk).expect("one part");
7031 }
7032 writer.finish().expect("commit");
7033
7034 let reader = Reader::open(&path).expect("reopen from disk");
7035 let under = [Probe { column: 0, op: Op::Less, value: Bound::Int(3_000) }];
7036 let kept: Vec<usize> = (0..parts).filter(|&part| !reader.skips(part, &under)).collect();
7037 assert_eq!(kept, vec![0, 1, 2], "only the three parts that start under three thousand");
7038 assert!(!reader.stripe_skips(0, &under), "the stripe reaches from zero and keeps itself");
7040 fs::remove_file(path).expect("remove scratch file");
7041 }
7042
7043 #[test]
7047 fn a_part_is_waved_through_when_its_own_bounds_pass_a_comparison_the_stripe_cannot() {
7048 let path = path("part-range-certain");
7049 let mut writer =
7050 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
7051 .expect("new file");
7052 let parts = STRIPE_PARTS + 3;
7053 let per_part = 128;
7054 for part in 0..parts {
7055 let held: Vec<Value> = (0..per_part)
7056 .map(|row| {
7057 Value::BigInt((part * 1_000) as i64 + (scattered(row as i64).rem_euclid(900)))
7058 })
7059 .collect();
7060 let chunk =
7061 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
7062 .expect("one column");
7063 writer.append(&chunk).expect("one part");
7064 }
7065 writer.finish().expect("commit");
7066
7067 let reader = Reader::open(&path).expect("reopen from disk");
7068 let under = [Probe { column: 0, op: Op::Less, value: Bound::Int(3_000) }];
7069 let waved: Vec<usize> = (0..parts).filter(|&part| reader.certain(part, &under)).collect();
7070 assert_eq!(waved, vec![0, 1, 2], "the three parts that end under three thousand");
7071 assert!(!reader.stripe_skips(0, &under), "the stripe straddles the comparison");
7074 fs::remove_file(path).expect("remove scratch file");
7075 }
7076
7077 #[test]
7080 fn a_stripe_of_one_part_writes_no_range_page_and_a_stripe_of_many_does() {
7081 for (parts, wanted) in [(1_usize, false), (STRIPE_PARTS, true)] {
7082 let path = path("part-range-page");
7083 let mut writer =
7084 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
7085 .expect("new file");
7086 for part in 0..parts {
7087 let held: Vec<Value> = (0..128)
7088 .map(|row| {
7089 Value::BigInt((part * 1_000) as i64 + scattered(row as i64).rem_euclid(900))
7090 })
7091 .collect();
7092 let chunk = Chunk::new(vec![
7093 Vector::from_values(LogicalType::BigInt, &held).expect("numbers"),
7094 ])
7095 .expect("one column");
7096 writer.append(&chunk).expect("one part");
7097 }
7098 writer.finish().expect("commit");
7099 let reader = Reader::open(&path).expect("reopen from disk");
7100 let bytes = reader.layout().columns[0].part_ranges;
7101 assert_eq!(bytes > 0, wanted, "{parts} parts wrote {bytes} bytes of ranges");
7102 fs::remove_file(path).expect("remove scratch file");
7103 }
7104 }
7105
7106 #[test]
7109 fn a_string_end_that_is_cut_down_still_covers_the_value_it_came_from() {
7110 let long = vec![b'a'; PART_BOUND_BYTES * 2];
7111 let low = shortened(Some(Bound::Bytes(long.clone())), false).expect("a low end");
7112 let high = shortened(Some(Bound::Bytes(long.clone())), true).expect("a high end");
7113 let Bound::Bytes(low) = low else { panic!("a string stays a string") };
7114 let Bound::Bytes(high) = high else { panic!("a string stays a string") };
7115 assert!(low.len() <= PART_BOUND_BYTES && high.len() <= PART_BOUND_BYTES);
7116 assert!(low.as_slice() <= long.as_slice(), "the low end is at or under the value");
7117 assert!(high.as_slice() >= long.as_slice(), "the high end is at or over the value");
7118 }
7119
7120 #[test]
7123 fn a_string_end_with_no_room_to_step_up_gives_up_the_bound() {
7124 let long = vec![u8::MAX; PART_BOUND_BYTES * 2];
7125 assert_eq!(shortened(Some(Bound::Bytes(long.clone())), true), None);
7126 let low = shortened(Some(Bound::Bytes(long)), false).expect("a low end is still a prefix");
7127 assert_eq!(low, Bound::Bytes(vec![u8::MAX; PART_BOUND_BYTES]));
7128 }
7129
7130 #[test]
7142 fn what_a_column_is_stored_as_follows_the_order_the_rows_were_written_in() {
7143 let parts = 4;
7144 let per_part = 1024;
7145 let rows = parts * per_part;
7146 let written = |name: &str, keys: &[i64]| {
7147 let path = path(name);
7148 let fields = vec![Field::required("key", LogicalType::BigInt)];
7149 let mut writer = Writer::create(&path, "keys", fields).expect("new file");
7150 for part in 0..parts {
7151 let values: Vec<Value> = keys[part * per_part..(part + 1) * per_part]
7152 .iter()
7153 .map(|key| Value::BigInt(*key))
7154 .collect();
7155 let chunk = Chunk::new(vec![
7156 Vector::from_values(LogicalType::BigInt, &values).expect("numbers"),
7157 ])
7158 .expect("one column");
7159 writer.append(&chunk).expect("one part");
7160 }
7161 writer.finish().expect("commit");
7162 path
7163 };
7164 let climbing = |step: &dyn Fn(usize) -> i64| {
7167 let mut key = 0;
7168 (0..rows)
7169 .map(|row| {
7170 key += step(row);
7171 key
7172 })
7173 .collect::<Vec<i64>>()
7174 };
7175 let ascending = climbing(&|row| (row % 3) as i64);
7176 let sparse = climbing(&|row| ((row * 2_654_435_761) % 4096) as i64);
7180 let near_path = written("stored-near", &ascending);
7181 let far_path = written("stored-far", &sparse);
7182
7183 let one = Reader::open(&near_path).expect("reopen from disk");
7184 let other = Reader::open(&far_path).expect("reopen from disk");
7185 let near = one.stored(0).expect("the column is stored");
7186 let far = other.stored(0).expect("the column is stored");
7187 assert_eq!(near.len(), parts, "one row per part");
7188 assert_eq!(far.len(), parts);
7189 let total = |stored: &[StoredPart]| stored.iter().map(|part| part.bytes).sum::<u64>();
7192 assert_eq!(total(&near), one.layout().columns[0].pages);
7193 assert_eq!(total(&far), other.layout().columns[0].pages);
7194 assert!(
7195 total(&near) * 2 < total(&far),
7196 "the sparse keys cost more, {} against {}",
7197 total(&far),
7198 total(&near)
7199 );
7200 for (at, part) in near.iter().enumerate() {
7202 assert_eq!(part.part, at);
7203 assert_eq!(part.row, at * per_part);
7204 assert_eq!(part.rows, per_part);
7205 let held = &ascending[at * per_part..(at + 1) * per_part];
7206 assert_eq!(part.low, Some(Value::BigInt(held[0])));
7207 assert_eq!(part.high, Some(Value::BigInt(held[per_part - 1])));
7208 assert_eq!(part.nulls, Some(0));
7209 }
7210 assert!(near[0].encoding.contains("DELTA"), "{}", near[0].encoding);
7213 assert!(far[0].encoding.contains("DELTA"), "{}", far[0].encoding);
7214 assert_ne!(near[0].encoding, far[0].encoding);
7215 fs::remove_file(near_path).expect("remove scratch file");
7216 fs::remove_file(far_path).expect("remove scratch file");
7217 }
7218
7219 #[test]
7229 fn a_sieve_larger_than_the_part_it_indexes_is_not_written() {
7230 let path = path("sieve-pays");
7231 let fields = vec![
7232 Field::required("spread", LogicalType::BigInt),
7233 Field::required("repeated", LogicalType::BigInt),
7234 ];
7235 let mut writer = Writer::create(&path, "hits", fields).expect("new file");
7236 let parts = 3;
7237 let per_part = 1024;
7238 for part in 0..parts {
7239 let base = (part * per_part) as i64;
7240 let spread: Vec<Value> =
7241 (0..per_part).map(|row| Value::BigInt(scattered(base + row as i64))).collect();
7242 let repeated: Vec<Value> =
7243 (0..per_part).map(|row| Value::BigInt(scattered((row / 256) as i64))).collect();
7244 let chunk = Chunk::new(vec![
7245 Vector::from_values(LogicalType::BigInt, &spread).expect("numbers"),
7246 Vector::from_values(LogicalType::BigInt, &repeated).expect("numbers"),
7247 ])
7248 .expect("two columns");
7249 writer.append(&chunk).expect("one part");
7250 }
7251 writer.finish().expect("commit");
7252
7253 let reader = Reader::open(&path).expect("reopen from disk");
7254 let layout = reader.layout();
7255 let spread = &layout.columns[0];
7256 let repeated = &layout.columns[1];
7257 assert!(spread.sieves > 0, "a column whose parts are worth a filter keeps one");
7258 assert_eq!(
7259 repeated.sieves, 0,
7260 "a column whose filter costs more than its parts keeps none"
7261 );
7262 for column in &layout.columns {
7265 assert!(
7266 column.sieves < column.pages,
7267 "{} spends {} on sieves over {} of data",
7268 column.name,
7269 column.sieves,
7270 column.pages
7271 );
7272 }
7273 let absent = [Probe {
7275 column: 0,
7276 op: Op::Equal,
7277 value: Bound::Int(i128::from(scattered((parts * per_part) as i64 + 1))),
7278 }];
7279 assert!((0..parts).all(|part| reader.skips(part, &absent)), "no part holds it");
7280 fs::remove_file(path).expect("remove scratch file");
7281 }
7282
7283 #[test]
7289 fn a_damaged_sieve_page_is_read_through_rather_than_refused() {
7290 let path = path("sieve-damaged");
7291 let mut writer =
7292 Writer::create(&path, "hits", vec![Field::required("id", LogicalType::BigInt)])
7293 .expect("new file");
7294 let rows = 128;
7295 let held: Vec<Value> = (0..rows).map(|row| Value::BigInt(scattered(row))).collect();
7296 let chunk =
7297 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
7298 .expect("one column");
7299 writer.append(&chunk).expect("one part");
7300 writer.finish().expect("commit");
7301
7302 let page =
7303 Reader::open(&path).expect("reopen").table.stripes[0].sieves[0].expect("a sieve page");
7304 let mut file = OpenOptions::new().write(true).open(&path).expect("open the sieve page");
7305 file.seek(SeekFrom::Start(page.offset + u64::from(page.length) - 1)).expect("seek");
7306 file.write_all(&[0xff]).expect("damage one byte");
7307 drop(file);
7308
7309 let reader = Reader::open(&path).expect("reopen the damaged file");
7310 let absent =
7311 [Probe { column: 0, op: Op::Equal, value: Bound::Int(i128::from(scattered(99))) }];
7312 assert!(!reader.skips(0, &absent), "a sieve that cannot be read skips nothing");
7313 assert_eq!(
7314 reader.read(0, &[0]).expect("the rows are untouched").len(),
7315 usize::try_from(rows).expect("a small count")
7316 );
7317 fs::remove_file(path).expect("remove scratch file");
7318 }
7319
7320 #[test]
7331 fn workers_that_want_the_same_stripe_read_it_once() {
7332 let path = path("single-flight");
7333 let mut writer =
7334 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
7335 .expect("new file");
7336 for part in 0..STRIPE_PARTS {
7337 let id = part as i32;
7338 let chunk = Chunk::new(vec![
7339 Vector::from_values(
7340 LogicalType::Integer,
7341 &[Value::Integer(id), Value::Integer(-id)],
7342 )
7343 .expect("integers"),
7344 ])
7345 .expect("matching rows");
7346 writer.append(&chunk).expect("one part");
7347 }
7348 writer.finish().expect("commit");
7349
7350 let reader = Reader::open(&path).expect("reopen from disk");
7351 assert_eq!(reader.table().stripes().len(), 1, "one stripe is the point of the test");
7352 let barrier = std::sync::Barrier::new(8);
7353 std::thread::scope(|scope| {
7354 for worker in 0..8 {
7355 let reader = &reader;
7356 let barrier = &barrier;
7357 scope.spawn(move || {
7358 barrier.wait();
7359 for part in (worker..STRIPE_PARTS).step_by(8) {
7360 let chunk = reader.read(part, &[0]).expect("a whole page read");
7361 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
7362 assert_eq!(chunk.value_at(1, 0), Value::Integer(-(part as i32)));
7363 }
7364 });
7365 }
7366 });
7367 assert_eq!(reader.pages.load(Atomic::Relaxed), 1, "one stripe, one page read, whoever won");
7368 fs::remove_file(path).expect("remove scratch file");
7369 }
7370
7371 #[test]
7384 fn opening_costs_the_same_over_a_thousand_times_the_rows() {
7385 let opened = |label: &str, rows_per_part: i32| {
7386 let path = path(label);
7387 let mut writer =
7388 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
7389 .expect("new file");
7390 for part in 0..STRIPE_PARTS * 3 {
7391 let values = (0..rows_per_part)
7395 .map(|row| {
7396 Value::Integer((part as i32 * rows_per_part + row).wrapping_mul(2_654_435))
7397 })
7398 .collect::<Vec<_>>();
7399 let chunk = Chunk::new(vec![
7400 Vector::from_values(LogicalType::Integer, &values).expect("integers"),
7401 ])
7402 .expect("matching rows");
7403 writer.append(&chunk).expect("one part");
7404 }
7405 writer.finish().expect("commit");
7406 let reader = Reader::open(&path).expect("reopen from disk");
7407 let size = fs::metadata(&path).expect("the file is there").len();
7408 let out = (reader.reads(), reader.table().stripes().len(), size);
7409 fs::remove_file(path).expect("remove scratch file");
7410 out
7411 };
7412
7413 let (thin, thin_stripes, thin_size) = opened("open-thin", 1);
7414 let (fat, fat_stripes, fat_size) = opened("open-fat", 1000);
7415 assert_eq!(
7416 thin_stripes, fat_stripes,
7417 "the same stripe count is what makes this a fair ask"
7418 );
7419 assert!(
7420 fat_size > thin_size * 50,
7421 "the fat file has to actually be larger, and it is {fat_size} against {thin_size}"
7422 );
7423
7424 assert_eq!(thin.opening.reads, fat.opening.reads, "the same reads either way");
7425 assert_eq!(thin.pages, 0, "opening read a page");
7426 assert_eq!(fat.pages, 0, "opening read a page");
7427 assert_eq!(thin.indexes, 0, "opening read an index");
7428 assert_eq!(fat.indexes, 0, "opening read an index");
7429 assert!(
7432 fat.opening.bytes < thin.opening.bytes * 2,
7433 "opening the thin file read {} bytes and the fat one read {}",
7434 thin.opening.bytes,
7435 fat.opening.bytes
7436 );
7437 }
7438
7439 #[test]
7447 fn two_opens_of_one_file_cost_the_same_and_the_second_is_not_cheaper() {
7448 let path = path("open-twice");
7449 let mut writer =
7450 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
7451 .expect("new file");
7452 for part in 0..STRIPE_PARTS * 3 {
7453 let chunk = Chunk::new(vec![
7454 Vector::from_values(LogicalType::Integer, &[Value::Integer(part as i32)])
7455 .expect("integers"),
7456 ])
7457 .expect("matching rows");
7458 writer.append(&chunk).expect("one part");
7459 }
7460 writer.finish().expect("commit");
7461
7462 let first = Reader::open(&path).expect("open");
7463 for part in 0..first.parts() {
7466 first.read(part, &[0]).expect("a part");
7467 }
7468 assert!(first.reads().pages > 0, "the scan has to have read something");
7469 let second = Reader::open(&path).expect("open again");
7470
7471 assert_eq!(first.reads().opening, second.reads().opening);
7472 assert_eq!(
7473 second.reads().pages,
7474 0,
7475 "the second open read a page off the back of the first"
7476 );
7477 assert_eq!(second.reads().indexes, 0, "the second open read an index it inherited");
7478 fs::remove_file(path).expect("remove scratch file");
7479 }
7480
7481 #[test]
7489 fn an_index_is_read_once_per_stripe_however_often_the_page_is_evicted() {
7490 let path = path("index-cache");
7491 let mut writer =
7492 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
7493 .expect("new file");
7494 let parts = STRIPE_PARTS * (CACHED_STRIPES_PER_COLUMN + 2);
7495 for part in 0..parts {
7496 let id = part as i32;
7497 let chunk = Chunk::new(vec![
7498 Vector::from_values(LogicalType::Integer, &[Value::Integer(id)]).expect("integers"),
7499 ])
7500 .expect("matching rows");
7501 writer.append(&chunk).expect("one part");
7502 }
7503 writer.finish().expect("commit");
7504
7505 let reader = Reader::open(&path).expect("reopen from disk");
7506 let stripes = reader.table().stripes().len();
7507 assert!(stripes > CACHED_STRIPES_PER_COLUMN, "the page cache has to be too small for this");
7508 for _ in 0..2 {
7510 for part in 0..parts {
7511 let chunk = reader.read(part, &[0]).expect("a part");
7512 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
7513 }
7514 }
7515 assert_eq!(reader.indexes.load(Atomic::Relaxed), stripes, "one index read per stripe");
7516 assert!(
7517 reader.pages.load(Atomic::Relaxed) > stripes,
7518 "the pages are the ones that get read again, which is what makes the index count mean \
7519 something"
7520 );
7521 fs::remove_file(path).expect("remove scratch file");
7522 }
7523
7524 #[test]
7533 fn a_worker_per_stripe_reads_its_page_once_when_the_cache_was_told_to_expect_it() {
7534 let workers = CACHED_STRIPES_PER_COLUMN + 4;
7535 let path = path("stripe-per-worker");
7536 let mut writer =
7537 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
7538 .expect("new file");
7539 for part in 0..STRIPE_PARTS * workers {
7540 let chunk = Chunk::new(vec![
7541 Vector::from_values(LogicalType::Integer, &[Value::Integer(part as i32)])
7542 .expect("integers"),
7543 ])
7544 .expect("matching rows");
7545 writer.append(&chunk).expect("one part");
7546 }
7547 writer.finish().expect("commit");
7548
7549 let read = |told: bool| {
7550 let reader = Reader::open(&path).expect("reopen from disk");
7551 assert_eq!(reader.table().stripes().len(), workers, "a stripe per worker");
7552 if told {
7553 reader.keep_stripes(workers);
7554 }
7555 let barrier = std::sync::Barrier::new(workers);
7556 std::thread::scope(|scope| {
7557 for (worker, run) in reader.stripe_parts().into_iter().enumerate() {
7558 let reader = &reader;
7559 let barrier = &barrier;
7560 scope.spawn(move || {
7561 for part in run {
7562 barrier.wait();
7563 let chunk = reader.read(part, &[0]).expect("a part of my own stripe");
7564 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
7565 }
7566 assert!(worker < workers);
7567 });
7568 }
7569 });
7570 reader.pages.load(Atomic::Relaxed)
7571 };
7572
7573 assert_eq!(read(true), workers, "one page read per stripe and no more");
7574 assert!(read(false) > workers, "a cache that small is read again on every part");
7575 fs::remove_file(path).expect("remove scratch file");
7576 }
7577
7578 #[test]
7583 fn a_damaged_index_page_is_an_error() {
7584 let path = path("damaged-index");
7585 let mut writer =
7586 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
7587 .expect("new file");
7588 writer.append(&sample_ids()).expect("first part");
7589 writer.append(&sample_ids()).expect("second part");
7590 writer.finish().expect("commit");
7591
7592 let reader = Reader::open(&path).expect("valid directory");
7593 let index = reader.table.stripes[0].index;
7594 let mut byte = [0; 1];
7595 read_at(&reader.file, index.offset, &mut byte).expect("the first part length");
7596 let mut file = OpenOptions::new().write(true).open(&path).expect("open index page");
7597 file.seek(SeekFrom::Start(index.offset)).expect("index start");
7598 file.write_all(&[!byte[0]]).expect("damage the first part length");
7599 let error = reader.read(1, &[0]).expect_err("a damaged index must not be used");
7600 assert!(error.message().contains("index page section checksum differs"), "{error}");
7601 fs::remove_file(path).expect("remove scratch file");
7602 }
7603
7604 #[test]
7611 fn every_integer_width_round_trips_through_a_page() {
7612 let path = path("integer-widths");
7613 let columns = [
7614 (LogicalType::TinyInt, vec![Value::TinyInt(i8::MIN), Value::TinyInt(i8::MAX)]),
7615 (LogicalType::UTinyInt, vec![Value::UTinyInt(0), Value::UTinyInt(u8::MAX)]),
7616 (LogicalType::SmallInt, vec![Value::SmallInt(i16::MIN), Value::SmallInt(i16::MAX)]),
7617 (LogicalType::USmallInt, vec![Value::USmallInt(0), Value::USmallInt(u16::MAX)]),
7618 (LogicalType::Integer, vec![Value::Integer(i32::MIN), Value::Integer(i32::MAX)]),
7619 (LogicalType::UInteger, vec![Value::UInteger(0), Value::UInteger(u32::MAX)]),
7620 (LogicalType::BigInt, vec![Value::BigInt(i64::MIN), Value::BigInt(i64::MAX)]),
7621 (LogicalType::UBigInt, vec![Value::UBigInt(0), Value::UBigInt(u64::MAX)]),
7622 ];
7623 let fields = columns
7624 .iter()
7625 .enumerate()
7626 .map(|(at, (ty, _))| Field::required(format!("c{at}"), ty.clone()))
7627 .collect::<Vec<_>>();
7628 let vectors = columns
7629 .iter()
7630 .map(|(ty, values)| Vector::from_values(ty.clone(), values).expect("a vector"))
7631 .collect::<Vec<_>>();
7632 let mut writer = Writer::create(&path, "widths", fields).expect("new file");
7633 writer.append(&Chunk::new(vectors).expect("matching rows")).expect("one stripe");
7634 writer.finish().expect("commit");
7635
7636 let reader = Reader::open(&path).expect("reopen from disk");
7637 let wanted = (0..columns.len()).collect::<Vec<_>>();
7638 let read = reader.read(0, &wanted).expect("every column");
7639 assert_eq!(read.len(), 2);
7640 for (at, (ty, values)) in columns.iter().enumerate() {
7642 assert_eq!(read.value_at(0, at), values[0], "the low end of {ty}");
7643 assert_eq!(read.value_at(1, at), values[1], "the high end of {ty}");
7644 }
7645 fs::remove_file(path).expect("remove scratch file");
7646 }
7647
7648 #[test]
7659 fn every_other_type_the_format_knows_round_trips_through_a_page() {
7660 let path = path("other-types");
7661 let columns = [
7662 (LogicalType::Float, vec![Value::Float(f32::MIN), Value::Float(-0.0)]),
7663 (LogicalType::Double, vec![Value::Double(f64::MIN), Value::Double(f64::MAX)]),
7664 (LogicalType::HugeInt, vec![Value::HugeInt(i128::MIN), Value::HugeInt(i128::MAX)]),
7665 (LogicalType::UHugeInt, vec![Value::UHugeInt(0), Value::UHugeInt(u128::MAX)]),
7666 (LogicalType::Time, vec![Value::Time(0), Value::Time(86_399_999_999)]),
7667 (LogicalType::TimeTz, vec![Value::TimeTz(-50_400_000_000), Value::TimeTz(0)]),
7668 (
7669 LogicalType::TimestampTz,
7670 vec![Value::TimestampTz(i64::MIN + 1), Value::TimestampTz(i64::MAX)],
7671 ),
7672 (
7673 LogicalType::Interval,
7674 vec![
7675 Value::Interval { months: i32::MIN, days: i32::MAX, micros: i64::MIN },
7676 Value::Interval { months: 13, days: -1, micros: 1 },
7677 ],
7678 ),
7679 (
7680 LogicalType::Blob,
7681 vec![Value::Blob(vec![0, 0xff, 0x80, 0xfe]), Value::Blob(Vec::new())],
7682 ),
7683 ];
7684 let fields = columns
7685 .iter()
7686 .enumerate()
7687 .map(|(at, (ty, _))| Field::required(format!("c{at}"), ty.clone()))
7688 .collect::<Vec<_>>();
7689 let vectors = columns
7690 .iter()
7691 .map(|(ty, values)| Vector::from_values(ty.clone(), values).expect("a vector"))
7692 .collect::<Vec<_>>();
7693 let mut writer = Writer::create(&path, "others", fields).expect("new file");
7694 writer.append(&Chunk::new(vectors).expect("matching rows")).expect("one stripe");
7695 writer.finish().expect("commit");
7696
7697 let reader = Reader::open(&path).expect("reopen from disk");
7698 let wanted = (0..columns.len()).collect::<Vec<_>>();
7699 let read = reader.read(0, &wanted).expect("every column");
7700 assert_eq!(read.len(), 2);
7701 for (at, (ty, values)) in columns.iter().enumerate() {
7702 assert_eq!(read.value_at(0, at), values[0], "the low end of {ty}");
7703 assert_eq!(read.value_at(1, at), values[1], "the high end of {ty}");
7704 }
7705 let Value::Float(zero) = read.value_at(1, 0) else { panic!("a float stays a float") };
7708 assert!(zero.is_sign_negative(), "a negative zero came back as {zero}");
7709
7710 fs::remove_file(path).expect("remove scratch file");
7711 }
7712
7713 #[test]
7719 fn a_nan_survives_being_written_down() {
7720 let path = path("nan");
7721 let nan = Vector::from_values(LogicalType::Double, &[Value::Double(f64::NAN)])
7722 .expect("a NaN vector");
7723 let mut writer =
7724 Writer::create(&path, "nan", vec![Field::required("d", LogicalType::Double)])
7725 .expect("new file");
7726 writer.append(&Chunk::new(vec![nan]).expect("one column")).expect("one stripe");
7727 writer.finish().expect("commit");
7728 let read = Reader::open(&path).expect("reopen").read(0, &[0]).expect("the column");
7729 let Value::Double(back) = read.value_at(0, 0) else { panic!("a double stays a double") };
7730 assert!(back.is_nan(), "a NaN came back as {back}");
7731 fs::remove_file(path).expect("remove scratch file");
7732 }
7733
7734 #[test]
7741 fn a_uuid_and_a_bit_string_come_back_as_the_bits_that_went_in() {
7742 let path = path("uuid-and-bit");
7743 let uuids = vec![0_i128, i128::MIN, -1];
7744 let mut bits = StringColumn::new();
7745 for value in [&b"\x02\xff"[..], &b""[..], &b"\x00\x01\x02\x03\x04\x05"[..]] {
7746 bits.push_bytes(value);
7747 }
7748 let expected = bits.clone();
7749 let fields =
7750 vec![Field::required("u", LogicalType::Uuid), Field::required("b", LogicalType::Bit)];
7751 let vectors = vec![
7752 Vector::flat(LogicalType::Uuid, Data::Int128(uuids.clone().into())).expect("uuids"),
7753 Vector::flat(LogicalType::Bit, Data::Varlen(bits)).expect("bit strings"),
7754 ];
7755 let mut writer = Writer::create(&path, "ids", fields).expect("new file");
7756 writer.append(&Chunk::new(vectors).expect("matching rows")).expect("one stripe");
7757 writer.finish().expect("commit");
7758
7759 let reader = Reader::open(&path).expect("reopen from disk");
7760 let read = reader.read(0, &[0, 1]).expect("both columns").flatten().expect("flat");
7761 let Some(Data::Int128(back)) = read.column(0).expect("the uuids").data() else {
7762 panic!("a uuid column is the 128 bit lane")
7763 };
7764 assert_eq!(back.as_slice(), uuids.as_slice());
7765 let Some(Data::Varlen(back)) = read.column(1).expect("the bits").data() else {
7766 panic!("a bit column is bytes")
7767 };
7768 for row in 0..expected.len() {
7769 assert_eq!(back.bytes(row), expected.bytes(row), "row {row} of the bit column");
7770 }
7771 fs::remove_file(path).expect("remove scratch file");
7772 }
7773
7774 #[test]
7775 fn numeric_frequency_candidates_keep_bounded_row_ordinals() {
7776 let path = path("frequency-ordinals");
7777 let mut writer =
7778 Writer::create(&path, "items", vec![Field::required("id", LogicalType::BigInt)])
7779 .expect("new file");
7780 let mut values = Vec::new();
7781 for leader in 0..10_i64 {
7782 values.extend(std::iter::repeat_n(leader, 100));
7783 }
7784 values.extend(1_000_i64..41_000);
7785 for part in values.chunks(1_024) {
7786 let vector = Vector::flat(LogicalType::BigInt, Data::Int64(part.to_vec().into()))
7787 .expect("big integers");
7788 writer.append(&Chunk::new(vec![vector]).expect("one column")).expect("one stripe");
7789 }
7790 writer.finish().expect("commit");
7791
7792 let reader = Reader::open(&path).expect("reopen from disk");
7793 let occurrences =
7794 reader.frequency_occurrences(0).expect("valid metadata").expect("bounded ordinals");
7795 assert!(occurrences.omitted_max < 100);
7796 assert!(occurrences.ordinals.len() <= FREQUENCY_ORDINALS);
7797 assert!(occurrences.ordinals.windows(2).all(|pair| pair[0] < pair[1]));
7798 assert_eq!(&occurrences.ordinals[..1_000], &(0_u64..1_000).collect::<Vec<_>>());
7799 fs::remove_file(path).expect("remove scratch file");
7800 }
7801
7802 #[test]
7808 fn a_file_from_another_format_says_which_format_it_is() {
7809 let older = path("older-format");
7810 let mut writer =
7811 Writer::create(&older, "items", vec![Field::new("id", LogicalType::Integer)])
7812 .expect("new file");
7813 let chunk = Chunk::new(vec![
7814 Vector::flat(LogicalType::Integer, Data::Int32(vec![1, 2, 3].into()))
7815 .expect("integers"),
7816 ])
7817 .expect("chunk");
7818 writer.append(&chunk).expect("page written");
7819 writer.finish().expect("commit");
7820
7821 let mut file = OpenOptions::new().write(true).open(&older).expect("open for the header");
7822 file.seek(SeekFrom::Start(8)).expect("the version follows the magic");
7823 file.write_all(&(FORMAT - 1).to_le_bytes()).expect("write an older version");
7824 drop(file);
7825 let complaint = Reader::open(&older).expect_err("an older format is refused").to_string();
7826 assert!(complaint.contains(&format!("format {}", FORMAT - 1)), "{complaint}");
7827 assert!(complaint.contains(&format!("format {FORMAT}")), "{complaint}");
7828
7829 let mut file = OpenOptions::new().write(true).open(&older).expect("open for the header");
7830 file.seek(SeekFrom::Start(0)).expect("the magic is first");
7831 file.write_all(b"NOTRUDB!").expect("write another engine's magic");
7832 drop(file);
7833 let complaint = Reader::open(&older).expect_err("a foreign file is refused").to_string();
7834 assert!(complaint.contains("magic"), "{complaint}");
7835 assert!(!complaint.contains("format"), "a version has nothing to do with it: {complaint}");
7836 fs::remove_file(older).expect("remove scratch file");
7837 }
7838
7839 #[test]
7840 fn an_unfinished_or_damaged_file_does_not_answer_with_partial_rows() {
7841 let unfinished = path("unfinished");
7842 let mut writer =
7843 Writer::create(&unfinished, "items", vec![Field::new("id", LogicalType::Integer)])
7844 .expect("new file");
7845 let chunk = Chunk::new(vec![
7846 Vector::flat(LogicalType::Integer, Data::Int32(vec![1, 2, 3].into()))
7847 .expect("integers"),
7848 ])
7849 .expect("chunk");
7850 writer.append(&chunk).expect("page written");
7851 drop(writer);
7852 assert!(Reader::open(&unfinished).is_err(), "no directory was committed");
7853 fs::remove_file(unfinished).expect("remove scratch file");
7854
7855 let damaged = path("damaged");
7856 let mut writer =
7857 Writer::create(&damaged, "items", vec![Field::new("id", LogicalType::Integer)])
7858 .expect("new file");
7859 writer.append(&chunk).expect("page written");
7860 writer.finish().expect("commit");
7861 let reader = Reader::open(&damaged).expect("valid directory");
7862 let mut file =
7863 OpenOptions::new().write(true).open(&damaged).expect("open for a damaged page");
7864 file.seek(SeekFrom::Start(HEADER + 1)).expect("inside first page");
7865 file.write_all(&[255]).expect("damage one byte");
7866 assert!(reader.read(0, &[0]).is_err(), "page checksum rejects corruption");
7867 fs::remove_file(damaged).expect("remove scratch file");
7868 }
7869
7870 #[test]
7871 fn damaged_lazy_dictionary_payload_is_an_error() {
7872 let path = path("damaged-dictionary");
7873 let mut writer = Writer::create(
7874 &path,
7875 "items",
7876 vec![
7877 Field::required("id", LogicalType::Integer),
7878 Field::new("text", LogicalType::Varchar),
7879 ],
7880 )
7881 .expect("new file");
7882 writer.append(&sample()).expect("stripe written");
7883 writer.finish().expect("commit");
7884
7885 let reader = Reader::open(&path).expect("valid directory");
7886 let dictionary = reader.table.dictionaries[1].expect("string dictionary page");
7887 let mut header = [0; DICTIONARY_HEADER];
7890 read_at(&reader.file, dictionary.offset, &mut header).expect("dictionary header");
7891 let index_len = dictionary_index_len(&header);
7892 let rank_len = last_rank_end(&reader.file, dictionary.offset, &header);
7893 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
7894 file.seek(SeekFrom::Start(dictionary.offset + index_len + rank_len))
7895 .expect("inside dictionary payload");
7896 file.write_all(&[255]).expect("damage dictionary payload");
7897
7898 let chunk = reader.read(0, &[1]).expect("code page and dictionary index remain valid");
7899 let error =
7900 chunk.validate_external().expect_err("payload corruption must reach the caller");
7901 assert!(error.message().contains("payload checksum differs"), "{error}");
7902 fs::remove_file(path).expect("remove scratch file");
7903 }
7904
7905 #[test]
7915 fn a_column_of_all_different_values_is_written_without_a_dictionary() {
7916 let path = path("dictionary-decide");
7917 let rows = 20_000;
7918 let unique =
7920 |row: usize| format!("{row:09} a value that appears exactly once in the table");
7921 let repeated = |row: usize| unique(row / 40);
7923 let mut writer = Writer::create(
7924 &path,
7925 "items",
7926 vec![
7927 Field::required("unique", LogicalType::Varchar),
7928 Field::required("repeated", LogicalType::Varchar),
7929 ],
7930 )
7931 .expect("new file");
7932 for part in (0..rows).step_by(1_000) {
7933 let span = part..(part + 1_000).min(rows);
7934 let left = span.clone().map(|row| Value::Varchar(unique(row))).collect::<Vec<_>>();
7935 let right = span.map(|row| Value::Varchar(repeated(row))).collect::<Vec<_>>();
7936 writer
7937 .append(
7938 &Chunk::new(vec![
7939 Vector::from_values(LogicalType::Varchar, &left).expect("strings"),
7940 Vector::from_values(LogicalType::Varchar, &right).expect("strings"),
7941 ])
7942 .expect("two columns"),
7943 )
7944 .expect("a part");
7945 }
7946 writer.finish().expect("commit");
7947
7948 let reader = Reader::open(&path).expect("reopen from disk");
7949 assert!(
7950 reader.table.dictionaries[0].is_none(),
7951 "a column with no repeats has nothing to say twice"
7952 );
7953 assert!(
7954 reader.table.dictionaries[1].is_some(),
7955 "a column whose values come round again keeps its dictionary"
7956 );
7957 let mut first = 0;
7958 for part in 0..reader.parts() {
7959 let chunk = reader.read(part, &[0, 1]).expect("a part");
7960 for row in 0..chunk.len() {
7961 assert_eq!(chunk.value_at(row, 0), Value::Varchar(unique(first + row)));
7962 assert_eq!(chunk.value_at(row, 1), Value::Varchar(repeated(first + row)));
7963 }
7964 first += chunk.len();
7965 }
7966 assert_eq!(first, rows, "every row was read back");
7967 let raw = (0..rows).map(|row| unique(row).len()).sum::<usize>();
7968 let size = fs::metadata(&path).expect("the file is there").len() as usize;
7969 assert!(size < raw, "a column without a dictionary is still encoded: {size} against {raw}");
7970 fs::remove_file(path).expect("remove scratch file");
7971 }
7972
7973 #[test]
7986 fn a_dictionary_over_many_blocks_checks_every_block_of_it() {
7987 let path = path("dictionary-blocks");
7988 let value = |row: usize| {
7989 let row = row.saturating_sub(8_000);
7990 format!("{row:07} a value long enough to be worth a payload block")
7991 };
7992 let parts = 40;
7993 let per_part = 1000;
7994 let mut writer =
7995 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
7996 .expect("new file");
7997 for part in 0..parts {
7998 let values = (0..per_part)
7999 .map(|row| Value::Varchar(value(part * per_part + row)))
8000 .collect::<Vec<_>>();
8001 let chunk = Chunk::new(vec![
8002 Vector::from_values(LogicalType::Varchar, &values).expect("strings"),
8003 ])
8004 .expect("matching rows");
8005 writer.append(&chunk).expect("a part");
8006 }
8007 writer.finish().expect("commit");
8008
8009 let reader = Reader::open(&path).expect("reopen from disk");
8010 let dictionary = reader.table.dictionaries[0].expect("string dictionary page");
8011 assert!(
8012 parts * per_part > TEXT_PAYLOAD_VALUES * 4,
8013 "the dictionary has to be several blocks for this to be testing anything"
8014 );
8015 for part in [0, parts - 1] {
8016 let chunk = reader.read(part, &[0]).expect("a part");
8017 chunk.validate_external().expect("every payload block checks out");
8018 assert_eq!(chunk.value_at(0, 0), Value::Varchar(value(part * per_part)));
8019 }
8020
8021 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
8022 file.seek(SeekFrom::Start(dictionary.offset + u64::from(dictionary.length) - 4))
8023 .expect("the last bytes of the page are payload");
8024 file.write_all(&[255]).expect("damage the last payload block");
8025 let reader = Reader::open(&path).expect("the directory and the index are untouched");
8026 let chunk = reader.read(parts - 1, &[0]).expect("the code page remains valid");
8027 let error = chunk.validate_external().expect_err("the damage must reach the caller");
8028 assert!(error.message().contains("payload checksum differs"), "{error}");
8029 fs::remove_file(path).expect("remove scratch file");
8030 }
8031
8032 #[test]
8046 fn values_of_different_lengths_read_back_out_of_packed_offsets() {
8047 let path = path("dictionary-offsets");
8048 let value = |row: usize| {
8049 let row = row % 5_000;
8050 if row % 511 == 3 { String::new() } else { "x".repeat(row % 97) + &format!("{row:05}") }
8051 };
8052 let rows = 6_000;
8053 let mut writer =
8054 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
8055 .expect("new file");
8056 let values = (0..rows).map(|row| Value::Varchar(value(row))).collect::<Vec<_>>();
8057 for part in values.chunks(1_000) {
8058 let chunk =
8059 Chunk::new(vec![Vector::from_values(LogicalType::Varchar, part).expect("strings")])
8060 .expect("matching rows");
8061 writer.append(&chunk).expect("a part");
8062 }
8063 writer.finish().expect("commit");
8064
8065 let reader = Reader::open(&path).expect("reopen from disk");
8066 assert!(
8067 rows > TEXT_PAYLOAD_VALUES * 4,
8068 "the dictionary has to be several blocks for this to be testing anything"
8069 );
8070 for part in 0..rows / 1_000 {
8071 let chunk = reader.read(part, &[0]).expect("a part");
8072 for row in 0..1_000 {
8073 let row = part * 1_000 + row;
8074 assert_eq!(
8075 chunk.value_at(row % 1_000, 0),
8076 Value::Varchar(value(row)),
8077 "value {row}"
8078 );
8079 }
8080 }
8081 fs::remove_file(path).expect("remove scratch file");
8082 }
8083
8084 #[test]
8096 fn a_global_dictionary_is_opened_once_however_many_workers_ask_at_once() {
8097 let path = path("dictionary-once");
8098 let parts = 8;
8099 let per_part = 500;
8100 let value =
8101 |row: usize| format!("{row:07} a value long enough to be worth a payload block");
8102 let mut writer =
8103 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
8104 .expect("new file");
8105 for part in 0..parts {
8106 let values = (0..per_part)
8107 .map(|row| Value::Varchar(value(part * per_part + row)))
8108 .collect::<Vec<_>>();
8109 let chunk = Chunk::new(vec![
8110 Vector::from_values(LogicalType::Varchar, &values).expect("strings"),
8111 ])
8112 .expect("matching rows");
8113 writer.append(&chunk).expect("a part");
8114 }
8115 writer.finish().expect("commit");
8116
8117 let reader = Reader::open(&path).expect("reopen from disk");
8118 assert!(reader.table.dictionaries[0].is_some(), "the column has to have one to share");
8119 assert_eq!(reader.reads().dictionaries, 0, "opening the file does not open a dictionary");
8120
8121 let workers = 16;
8122 let gate = std::sync::Barrier::new(workers);
8123 std::thread::scope(|scope| {
8124 for worker in 0..workers {
8125 let reader = reader.clone();
8126 let gate = &gate;
8127 scope.spawn(move || {
8128 gate.wait();
8129 let chunk = reader.read(worker % parts, &[0]).expect("a part");
8130 assert_eq!(
8131 chunk.value_at(0, 0),
8132 Value::Varchar(value((worker % parts) * per_part))
8133 );
8134 });
8135 }
8136 });
8137
8138 assert_eq!(reader.reads().dictionaries, 1, "sixteen workers, one dictionary, one open");
8139 fs::remove_file(path).expect("remove scratch file");
8140 }
8141
8142 #[test]
8147 fn a_damaged_sorted_order_is_an_error() {
8148 let path = path("damaged-order");
8149 let mut writer = Writer::create(
8150 &path,
8151 "items",
8152 vec![
8153 Field::required("id", LogicalType::Integer),
8154 Field::new("text", LogicalType::Varchar),
8155 ],
8156 )
8157 .expect("new file");
8158 writer.append(&sample()).expect("stripe written");
8159 writer.finish().expect("commit");
8160
8161 let reader = Reader::open(&path).expect("valid directory");
8162 let page = reader.table.dictionaries[1].expect("string dictionary page");
8163 let mut header = [0; DICTIONARY_HEADER];
8164 read_at(&reader.file, page.offset, &mut header).expect("dictionary header");
8165 let index_len = dictionary_index_len(&header);
8166 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
8167 file.seek(SeekFrom::Start(page.offset + index_len)).expect("the first head");
8168 file.write_all(&[255]).expect("damage the order");
8169
8170 let dictionary = reader.dictionary(1).expect("read").expect("a string column has one");
8171 let error = dictionary.compare_rank(0, b"anything").expect_err("a damaged order is caught");
8172 assert!(error.message().contains("rank checksum differs"), "{error}");
8173 fs::remove_file(path).expect("remove scratch file");
8174 }
8175
8176 #[test]
8180 fn a_global_dictionary_carries_the_sorted_order_of_its_values() {
8181 let spellings = ["overlong1z", "b", "", "overlong1a", "overlong", "ab", "a", "overlong1"];
8184 let path = path("dictionary-order");
8185 let mut writer =
8186 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
8187 .expect("new file");
8188 writer
8189 .append(
8190 &Chunk::new(vec![
8191 Vector::from_values(
8192 LogicalType::Varchar,
8193 &spellings.map(|text| Value::Varchar(text.into())),
8194 )
8195 .expect("strings"),
8196 ])
8197 .expect("one column"),
8198 )
8199 .expect("stripe written");
8200 writer.finish().expect("commit");
8201
8202 let reader = Reader::open(&path).expect("valid directory");
8203 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
8204 let count = dictionary.ranks().expect("a v10 file stores one");
8205 assert_eq!(count, spellings.len(), "every distinct value has a rank");
8206 let order = (0..count)
8207 .map(|rank| dictionary.code_at_rank(rank).expect("a code"))
8208 .collect::<Vec<_>>();
8209 let mut seen = order.clone();
8210 seen.sort_unstable();
8211 assert_eq!(seen, (0..spellings.len() as u32).collect::<Vec<_>>(), "a permutation of codes");
8212
8213 let ranked = order
8214 .iter()
8215 .map(|&code| {
8216 dictionary.try_bytes_at(code as usize).expect("read").expect("a value").to_vec()
8217 })
8218 .collect::<Vec<_>>();
8219 let mut expected = spellings.map(|text| text.as_bytes().to_vec()).to_vec();
8220 expected.sort();
8221 assert_eq!(ranked, expected, "rank order is value order");
8222
8223 for (rank, value) in expected.iter().enumerate() {
8226 assert_eq!(
8227 dictionary.compare_rank(rank, value).expect("compare"),
8228 Ordering::Equal,
8229 "rank {rank} is its own value"
8230 );
8231 if rank > 0 {
8232 assert_eq!(
8233 dictionary.compare_rank(rank - 1, value).expect("compare"),
8234 Ordering::Less,
8235 "rank {rank} follows the one before it"
8236 );
8237 }
8238 }
8239 fs::remove_file(path).expect("remove scratch file");
8240 }
8241
8242 #[test]
8250 fn a_dictionary_sweep_reads_every_value_and_keeps_it_under_the_budget() {
8251 let path = path("dictionary-sweep");
8252 let spellings = (0..2_500)
8255 .map(|index| Value::Varchar(format!("value {index:08} {}", "x".repeat(index % 40))))
8256 .collect::<Vec<_>>();
8257 let mut writer =
8258 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
8259 .expect("new file");
8260 for part in spellings.chunks(1_024) {
8263 writer
8264 .append(
8265 &Chunk::new(vec![
8266 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
8267 ])
8268 .expect("one column"),
8269 )
8270 .expect("stripe written");
8271 }
8272 writer.finish().expect("commit");
8273
8274 let reader = Reader::open(&path).expect("valid directory");
8275 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
8276 assert_eq!(dictionary.len(), spellings.len(), "every value is distinct");
8277
8278 let resting = dictionary.footprint();
8279 let mut swept: Vec<Vec<u8>> = Vec::new();
8280 let mut at = 0;
8281 let mut calls = 0;
8282 while at < dictionary.len() {
8283 let stopped = dictionary
8284 .sweep_text(at, dictionary.len(), &mut |index: usize, text: &[u8]| {
8285 assert_eq!(index, swept.len(), "a sweep hands its values over in order");
8286 swept.push(text.to_vec());
8287 Ok(())
8288 })
8289 .expect("a sweep reads");
8290 assert!(stopped > at, "a sweep moves");
8291 at = stopped;
8292 calls += 1;
8293 }
8294 assert_eq!(calls, 3, "a sweep hands over one block at a time");
8295 let after = dictionary.footprint();
8296 assert!(after > resting, "a sweep under the budget keeps what it decoded");
8297
8298 let read = (0..dictionary.len())
8299 .map(|code| dictionary.try_bytes_at(code).expect("read").expect("a value").to_vec())
8300 .collect::<Vec<_>>();
8301 assert_eq!(swept, read, "a sweep answers what a point read answers");
8302 assert_eq!(dictionary.footprint(), after, "a point read of a kept block decodes nothing");
8303 fs::remove_file(path).expect("remove scratch file");
8304 }
8305
8306 #[test]
8317 fn a_sweep_over_a_block_with_a_short_second_run_reads_what_a_point_read_reads() {
8318 let path = path("dictionary-sweep-short-run");
8319 let spellings = (0..2_800)
8320 .map(|index| Value::Varchar(format!("value {index:08} {}", "x".repeat(index % 40))))
8321 .collect::<Vec<_>>();
8322 let mut writer =
8323 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
8324 .expect("new file");
8325 for part in spellings.chunks(1_024) {
8326 writer
8327 .append(
8328 &Chunk::new(vec![
8329 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
8330 ])
8331 .expect("one column"),
8332 )
8333 .expect("stripe written");
8334 }
8335 writer.finish().expect("commit");
8336
8337 let reader = Reader::open(&path).expect("valid directory");
8338 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
8339 assert_eq!(dictionary.len(), spellings.len(), "every value is distinct");
8340 let last = dictionary.len() % TEXT_PAYLOAD_VALUES;
8341 assert!(last > TEXT_OFFSET_RUN, "the last block has to reach into a second run of offsets");
8342 assert!(last < TEXT_PAYLOAD_VALUES, "and that second run has to be short of a whole one");
8343
8344 let mut swept: Vec<Vec<u8>> = Vec::new();
8345 let mut at = 0;
8346 while at < dictionary.len() {
8347 let stopped = dictionary
8348 .sweep_text(at, dictionary.len(), &mut |index: usize, text: &[u8]| {
8349 assert_eq!(index, swept.len(), "a sweep hands its values over in order");
8350 swept.push(text.to_vec());
8351 Ok(())
8352 })
8353 .expect("a sweep reads");
8354 assert!(stopped > at, "a sweep moves");
8355 at = stopped;
8356 }
8357 let read = (0..dictionary.len())
8358 .map(|code| dictionary.try_bytes_at(code).expect("read").expect("a value").to_vec())
8359 .collect::<Vec<_>>();
8360 assert_eq!(swept, read, "a sweep answers what a point read answers");
8361 fs::remove_file(path).expect("remove scratch file");
8362 }
8363
8364 #[test]
8374 fn narrowing_a_page_takes_what_fits_and_refuses_what_does_not() {
8375 assert_eq!(fit::<i8>(&[]).expect("an empty page fits anything"), Vec::<i8>::new());
8376 assert_eq!(fit::<i8>(&[-128, 0, 127]).expect("the edges fit"), vec![-128_i8, 0, 127]);
8377 fit::<i8>(&[128]).expect_err("one past the top does not fit");
8378 fit::<i8>(&[-129]).expect_err("one past the bottom does not fit");
8379 assert_eq!(fit::<u8>(&[0, 255]).expect("the edges fit"), vec![0_u8, 255]);
8380 fit::<u8>(&[256]).expect_err("one past the top does not fit");
8381 fit::<u8>(&[-1]).expect_err("a negative does not fit an unsigned page");
8382 assert_eq!(
8383 fit::<i16>(&[-32_768, 0, 32_767]).expect("the edges fit"),
8384 vec![-32_768_i16, 0, 32_767]
8385 );
8386 fit::<i16>(&[32_768]).expect_err("one past the top does not fit");
8387 fit::<i16>(&[-32_769]).expect_err("one past the bottom does not fit");
8388 assert_eq!(fit::<u16>(&[0, 65_535]).expect("the edges fit"), vec![0_u16, 65_535]);
8389 fit::<u16>(&[65_536]).expect_err("one past the top does not fit");
8390 fit::<u16>(&[-1]).expect_err("a negative does not fit an unsigned page");
8391 assert_eq!(
8392 fit::<i32>(&[i64::from(i32::MIN), 0, i64::from(i32::MAX)]).expect("the edges fit"),
8393 vec![i32::MIN, 0, i32::MAX]
8394 );
8395 fit::<i32>(&[i64::from(i32::MAX) + 1]).expect_err("one past the top does not fit");
8396 fit::<i32>(&[i64::from(i32::MIN) - 1]).expect_err("one past the bottom does not fit");
8397 assert_eq!(
8398 fit::<u32>(&[0, 4_294_967_295]).expect("the edges fit"),
8399 vec![0_u32, 4_294_967_295]
8400 );
8401 fit::<u32>(&[4_294_967_296]).expect_err("one past the top does not fit");
8402 fit::<u32>(&[-1]).expect_err("a negative does not fit an unsigned page");
8403
8404 fit::<i8>(&[0, 1, 2, 128, 3]).expect_err("one bad value spoils the page");
8407 }
8408
8409 #[test]
8416 fn the_residue_agrees_with_a_checked_conversion_everywhere() {
8417 for value in -70_000_i64..70_000 {
8418 assert_eq!(fit::<i8>(&[value]).is_ok(), i8::try_from(value).is_ok(), "{value} as i8");
8419 assert_eq!(fit::<u8>(&[value]).is_ok(), u8::try_from(value).is_ok(), "{value} as u8");
8420 assert_eq!(fit::<i16>(&[value]).is_ok(), i16::try_from(value).is_ok(), "{value} i16");
8421 assert_eq!(fit::<u16>(&[value]).is_ok(), u16::try_from(value).is_ok(), "{value} u16");
8422 }
8423 let wide = [i64::MIN, i64::MIN + 1, i64::from(i32::MIN), 0, i64::from(u32::MAX), i64::MAX];
8424 for edge in wide {
8425 for step in -2_i64..=2 {
8426 let value = edge.saturating_add(step);
8427 assert_eq!(
8428 fit::<i32>(&[value]).is_ok(),
8429 i32::try_from(value).is_ok(),
8430 "{value} as i32"
8431 );
8432 assert_eq!(
8433 fit::<u32>(&[value]).is_ok(),
8434 u32::try_from(value).is_ok(),
8435 "{value} as u32"
8436 );
8437 }
8438 }
8439 }
8440
8441 #[test]
8449 fn a_dictionary_at_its_budget_sweeps_without_keeping() {
8450 let path = path("dictionary-budget");
8451 let spellings = (0..2_500)
8452 .map(|index| Value::Varchar(format!("value {index:08} {}", "y".repeat(index % 40))))
8453 .collect::<Vec<_>>();
8454 let mut writer =
8455 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
8456 .expect("new file");
8457 for part in spellings.chunks(1_024) {
8458 writer
8459 .append(
8460 &Chunk::new(vec![
8461 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
8462 ])
8463 .expect("one column"),
8464 )
8465 .expect("stripe written");
8466 }
8467 writer.finish().expect("commit");
8468
8469 let reader = Reader::open(&path).expect("valid directory");
8470 let page = reader.table.dictionaries[0].expect("a string column has one");
8471 let file = Arc::clone(&reader.file);
8472 let starved = open_global_dictionary(file, page, &LogicalType::Varchar, 0)
8473 .expect("a dictionary opens whatever it may keep");
8474
8475 let resting = starved.footprint();
8476 let mut swept: Vec<Vec<u8>> = Vec::new();
8477 let mut at = 0;
8478 while at < starved.len() {
8479 at = starved
8480 .sweep_text(at, starved.len(), &mut |_index: usize, text: &[u8]| {
8481 swept.push(text.to_vec());
8482 Ok(())
8483 })
8484 .expect("a sweep reads");
8485 }
8486 assert_eq!(swept.len(), spellings.len(), "a starved sweep still reads every value");
8487 assert_eq!(starved.footprint(), resting, "and keeps no block it decoded");
8488
8489 let generous = reader.dictionary(0).expect("read").expect("a string column has one");
8490 let read = (0..generous.len())
8491 .map(|code| generous.try_bytes_at(code).expect("read").expect("a value").to_vec())
8492 .collect::<Vec<_>>();
8493 assert_eq!(swept, read, "a starved sweep answers what a point read answers");
8494 fs::remove_file(path).expect("remove scratch file");
8495 }
8496
8497 #[test]
8498 fn damaged_membership_cannot_skip_a_string_page() {
8499 let path = path("damaged-membership");
8500 let mut writer = Writer::create(
8501 &path,
8502 "items",
8503 vec![
8504 Field::required("id", LogicalType::Integer),
8505 Field::new("text", LogicalType::Varchar),
8506 ],
8507 )
8508 .expect("new file");
8509 writer.append(&sample()).expect("stripe written");
8510 writer.finish().expect("commit");
8511
8512 let reader = Reader::open(&path).expect("valid directory");
8513 let membership = reader.table.stripes[0].memberships[1].expect("string membership");
8514 let mut file = OpenOptions::new().write(true).open(&path).expect("open membership page");
8515 file.seek(SeekFrom::Start(membership.offset)).expect("membership start");
8516 file.write_all(&[255]).expect("damage membership");
8517 let error = reader.skips_codes(0, 1, &[3]).expect_err("corruption must not skip rows");
8518 assert!(error.message().contains("membership page checksum differs"), "{error}");
8519 fs::remove_file(path).expect("remove scratch file");
8520 }
8521
8522 #[test]
8523 fn membership_delta_stream_is_sorted_exact_and_bounded() {
8524 let unique = unique_codes(&[900, 4, 4, 72, 9, u32::MAX]);
8525 assert_eq!(unique, [4, 9, 72, 900, u32::MAX]);
8526 let encoded = encode_membership(&unique);
8527 assert_eq!(
8528 decode_membership(&encoded).expect("valid membership"),
8529 [4, 9, 72, 900, u32::MAX]
8530 );
8531 let merged = merged_codes(vec![vec![4, 900], vec![9, 900, u32::MAX], vec![72]]);
8534 assert_eq!(merged, [4, 9, 72, 900, u32::MAX]);
8535 assert_eq!(
8536 decode_membership(&encode_membership(&merged)).expect("valid membership"),
8537 unique
8538 );
8539 assert!(decode_membership(&[1, 0x80]).is_err(), "a truncated varint is invalid");
8540 assert!(
8541 decode_membership(&[1, 0xff, 0xff, 0xff, 0xff, 0x10]).is_err(),
8542 "a value past u32 is invalid"
8543 );
8544 }
8545
8546 #[test]
8547 fn a_global_dictionary_may_be_larger_than_one_column_page() {
8548 let dictionary = Page {
8549 offset: HEADER,
8550 length: u32::try_from(MAX_PAGE + 1).expect("the page bound fits on disk"),
8551 hash: 0,
8552 };
8553 let table = Table {
8554 name: "items".to_owned(),
8555 fields: vec![Field::new("text", LogicalType::Varchar)],
8556 stripes: Vec::new(),
8557 rows: 0,
8558 dictionaries: vec![Some(dictionary)],
8559 distincts: vec![None],
8560 frequencies: vec![None],
8561 clustering: None,
8562 };
8563 let directory = encode_directory(&table).expect("directory");
8564 let file_size = dictionary.offset + u64::from(dictionary.length) + 1;
8565
8566 let decoded = decode_directory(&directory, file_size).expect("large lazy dictionary");
8567 assert_eq!(decoded.dictionaries[0].expect("dictionary").length, dictionary.length);
8568 }
8569
8570 #[test]
8571 fn a_column_with_one_value_everywhere_costs_almost_nothing_a_row() {
8572 let path = path("constant-codes");
8573 let mut writer =
8574 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
8575 .expect("new file");
8576 let empty = vec![Value::Varchar(String::new()); 1024];
8577 for _ in 0..4 {
8578 let column = Vector::from_values(LogicalType::Varchar, &empty).expect("strings");
8579 writer.append(&Chunk::new(vec![column]).expect("one column")).expect("a part");
8580 }
8581 writer.finish().expect("commit");
8582
8583 let reader = Reader::open(&path).expect("valid directory");
8584 let pages = reader.layout().columns.first().expect("one column").pages;
8585 assert!(pages < 256, "{pages} bytes of pages for 4,096 rows of one value");
8589 let read = reader.read(3, &[0]).expect("the last part back");
8590 assert_eq!(read.value_at(0, 0), Value::Varchar(String::new()));
8591 assert_eq!(read.value_at(1023, 0), Value::Varchar(String::new()));
8592 fs::remove_file(path).expect("remove scratch file");
8593 }
8594
8595 #[test]
8596 fn a_cascade_value_too_wide_for_its_column_is_refused_rather_than_cut() {
8597 let over = vec![i64::from(i32::MAX) + 1];
8600 let error = narrowed(&LogicalType::Integer, over).expect_err("a page that disagrees");
8601 assert!(format!("{error}").contains("not of its type"), "{error}");
8602 assert!(narrowed(&LogicalType::BigInt, vec![i64::MIN]).is_ok(), "bigint holds all of i64");
8603 assert!(narrowed(&LogicalType::Varchar, vec![0]).is_err(), "strings are not integers");
8604 }
8605
8606 #[test]
8607 fn a_code_stream_the_cascade_cannot_shrink_is_left_alone() {
8608 let mut state: u32 = 0x9e37_79b9;
8612 let spread: Vec<u32> = (0..1024)
8613 .map(|_| {
8614 state ^= state << 13;
8615 state ^= state >> 17;
8616 state ^= state << 5;
8617 state
8618 })
8619 .collect();
8620 assert_eq!(encoded_codes(&spread).expect("no failure"), None);
8621 let near: Vec<u32> = (0..1024).collect();
8622 let coded = encoded_codes(&near).expect("no failure").expect("counting up is packable");
8623 assert!(coded.len() < near.len() * 4, "{} bytes for a run of 1,024", coded.len());
8624 }
8625
8626 #[test]
8632 fn two_writes_of_the_same_rows_give_the_same_bytes() {
8633 fn written(path: &PathBuf) {
8634 let fields = (0..40)
8635 .map(|column| {
8636 let ty =
8637 if column % 4 == 0 { LogicalType::Varchar } else { LogicalType::BigInt };
8638 Field::new(format!("c{column}"), ty)
8639 })
8640 .collect::<Vec<_>>();
8641 let mut writer = Writer::create(path, "wide", fields).expect("new file");
8642 for part in 0..70_u64 {
8643 let columns = (0..40)
8644 .map(|column| {
8645 let values = (0..64_u64)
8646 .map(|row| {
8647 let seed = part.wrapping_mul(31).wrapping_add(row);
8648 if column % 4 == 0 {
8649 Value::Varchar(format!("v{}", seed % 17))
8650 } else {
8651 Value::BigInt(i64::try_from(seed % 97).expect("small"))
8652 }
8653 })
8654 .collect::<Vec<_>>();
8655 let ty = if column % 4 == 0 {
8656 LogicalType::Varchar
8657 } else {
8658 LogicalType::BigInt
8659 };
8660 Vector::from_values(ty, &values).expect("a column")
8661 })
8662 .collect::<Vec<_>>();
8663 writer.append(&Chunk::new(columns).expect("forty columns")).expect("a part");
8664 }
8665 writer.finish().expect("commit");
8666 }
8667
8668 let first = path("repeatable-one");
8669 let second = path("repeatable-two");
8670 written(&first);
8671 written(&second);
8672 let left = fs::read(&first).expect("the first file");
8673 let right = fs::read(&second).expect("the second file");
8674 assert_eq!(left.len(), right.len(), "two writes of the same rows differ in length");
8675 assert!(left == right, "two writes of the same rows differ in their bytes");
8676
8677 let reader = Reader::open(&first).expect("valid directory");
8680 assert_eq!(reader.table().rows(), 70 * 64);
8681 let read = reader.read(0, &[0, 1]).expect("the first part back");
8682 assert_eq!(read.value_at(0, 0), Value::Varchar("v0".to_owned()));
8683 assert_eq!(read.value_at(0, 1), Value::BigInt(0));
8684 fs::remove_file(first).expect("remove scratch file");
8685 fs::remove_file(second).expect("remove scratch file");
8686 }
8687
8688 fn three_tables(path: &PathBuf) {
8690 let writer = Writer::create(
8691 path,
8692 "region",
8693 vec![
8694 Field::new("r_key", LogicalType::Integer),
8695 Field::new("r_name", LogicalType::Varchar),
8696 ],
8697 )
8698 .expect("new file");
8699 let mut writer = writer;
8700 writer
8701 .append(
8702 &Chunk::new(vec![
8703 Vector::from_values(
8704 LogicalType::Integer,
8705 &[Value::Integer(0), Value::Integer(1)],
8706 )
8707 .expect("keys"),
8708 Vector::from_values(
8709 LogicalType::Varchar,
8710 &[Value::Varchar("AFRICA".to_owned()), Value::Varchar("ASIA".to_owned())],
8711 )
8712 .expect("names"),
8713 ])
8714 .expect("two columns"),
8715 )
8716 .expect("a part");
8717 let mut writer = writer
8718 .next("empty", vec![Field::new("nothing", LogicalType::BigInt)])
8719 .expect("a second table");
8720 writer
8721 .append(
8722 &Chunk::new(vec![
8723 Vector::from_values(LogicalType::BigInt, &[Value::BigInt(7)]).expect("a row"),
8724 ])
8725 .expect("one column"),
8726 )
8727 .expect("a part");
8728 let mut writer =
8729 writer.next("wide", vec![Field::new("n", LogicalType::BigInt)]).expect("a third table");
8730 for part in 0..70_i64 {
8731 let values = (0..64).map(|row| Value::BigInt(part * 64 + row)).collect::<Vec<_>>();
8732 writer
8733 .append(
8734 &Chunk::new(vec![
8735 Vector::from_values(LogicalType::BigInt, &values).expect("a column"),
8736 ])
8737 .expect("one column"),
8738 )
8739 .expect("a part");
8740 }
8741 writer.finish().expect("commit");
8742 }
8743
8744 #[test]
8745 fn three_tables_in_one_file_read_back_by_name() {
8746 let file = path("three-tables");
8747 three_tables(&file);
8748 let catalog = Catalog::open(&file).expect("a committed catalog");
8749 assert_eq!(catalog.names().collect::<Vec<_>>(), ["region", "empty", "wide"]);
8750
8751 let region = catalog.table("region").expect("the first table");
8752 assert_eq!(region.table().rows(), 2);
8753 assert_eq!(
8754 region.read(0, &[1]).expect("names").value_at(1, 0),
8755 Value::Varchar("ASIA".to_owned())
8756 );
8757
8758 let wide = catalog.table("wide").expect("the third table");
8759 assert_eq!(wide.table().rows(), 70 * 64);
8760 assert_eq!(wide.read(0, &[0]).expect("the first part").value_at(0, 0), Value::BigInt(0));
8761
8762 let empty = catalog.table("empty").expect("the second table");
8765 assert_eq!(empty.table().rows(), 1);
8766 assert_eq!(empty.read(0, &[0]).expect("the row").value_at(0, 0), Value::BigInt(7));
8767
8768 fs::remove_file(file).expect("remove scratch file");
8769 }
8770
8771 #[test]
8772 fn a_name_the_file_does_not_hold_is_an_error_rather_than_the_first_table() {
8773 let file = path("three-tables-missing");
8774 three_tables(&file);
8775 let catalog = Catalog::open(&file).expect("a committed catalog");
8776 let error = catalog.table("nation").expect_err("no such table");
8777 assert!(error.message().contains("nation"), "{}", error.message());
8778 fs::remove_file(file).expect("remove scratch file");
8779 }
8780
8781 #[test]
8782 fn a_file_of_three_tables_will_not_open_as_one() {
8783 let file = path("three-tables-unnamed");
8784 three_tables(&file);
8785 let error = Reader::open(&file).expect_err("more than one table");
8786 assert!(error.message().contains("more than one table"), "{}", error.message());
8787 fs::remove_file(file).expect("remove scratch file");
8788 }
8789
8790 #[test]
8792 fn decimals_of_every_storage_width_round_trip() {
8793 let file = path("decimals");
8794 let widths = [(4_u8, 2_u8), (9, 2), (18, 4), (38, 6)];
8795 let fields = widths
8796 .iter()
8797 .enumerate()
8798 .map(|(index, (width, scale))| {
8799 Field::new(
8800 format!("d{index}"),
8801 LogicalType::decimal(*width, *scale).expect("a decimal type"),
8802 )
8803 })
8804 .collect::<Vec<_>>();
8805 let mut writer = Writer::create(&file, "money", fields).expect("new file");
8806 let rows: [i128; 3] = [-1234, 0, 999];
8807 let columns = widths
8808 .iter()
8809 .map(|(width, scale)| {
8810 let values = rows
8811 .iter()
8812 .map(|unscaled| Value::Decimal {
8813 unscaled: *unscaled,
8814 width: *width,
8815 scale: *scale,
8816 })
8817 .collect::<Vec<_>>();
8818 Vector::from_values(
8819 LogicalType::decimal(*width, *scale).expect("a decimal type"),
8820 &values,
8821 )
8822 .expect("a decimal column")
8823 })
8824 .collect::<Vec<_>>();
8825 writer.append(&Chunk::new(columns).expect("four columns")).expect("a part");
8826 writer.finish().expect("commit");
8827
8828 let reader = Reader::open(&file).expect("a committed file");
8829 for (index, (width, scale)) in widths.iter().enumerate() {
8830 assert_eq!(
8831 reader.table().fields()[index].ty,
8832 LogicalType::decimal(*width, *scale).expect("a decimal type"),
8833 "column {index} came back as another type"
8834 );
8835 let column = reader.read(0, &[index]).expect("the column");
8836 for (row, unscaled) in rows.iter().enumerate() {
8837 assert_eq!(
8838 column.value_at(row, 0),
8839 Value::Decimal { unscaled: *unscaled, width: *width, scale: *scale },
8840 "column {index} row {row}"
8841 );
8842 }
8843 }
8844 fs::remove_file(file).expect("remove scratch file");
8845 }
8846
8847 #[test]
8848 fn two_tables_of_one_name_are_refused_before_anything_is_committed() {
8849 let file = path("two-of-a-name");
8850 let writer = Writer::create(&file, "t", vec![Field::new("a", LogicalType::BigInt)])
8851 .expect("new file");
8852 let error = writer
8853 .next("t", vec![Field::new("a", LogicalType::BigInt)])
8854 .expect_err("the same name twice");
8855 assert!(error.message().contains("same name"), "{}", error.message());
8856 fs::remove_file(file).expect("remove scratch file");
8857 }
8858
8859 #[test]
8860 fn opening_the_catalog_reads_no_table_directory() {
8861 let file = path("catalog-only");
8862 three_tables(&file);
8863 let catalog = Catalog::open(&file).expect("a committed catalog");
8864 assert_eq!(catalog.opening.reads, 2, "opening the catalog read more than the slot");
8867 assert_eq!(catalog.names().len(), 3);
8868 fs::remove_file(file).expect("remove scratch file");
8869 }
8870
8871 #[test]
8882 fn the_checksum_answers_what_it_has_always_answered() {
8883 let bytes: Vec<u8> =
8884 (0..1000_u32).map(|at| (at.wrapping_mul(31).wrapping_add(7) % 251) as u8).collect();
8885 for (length, expected) in [
8886 (0, 0xef46_db37_51d8_e999),
8887 (1, 0xa96c_7f0c_e858_bbb7),
8888 (3, 0x56e6_9576_32a4_87f9),
8889 (4, 0xc60d_15b1_e3ff_8f04),
8890 (5, 0x8088_1585_8624_dd4e),
8891 (7, 0xafbe_fc3d_6c6f_9a8e),
8892 (8, 0x3da5_c7aa_2696_83e0),
8893 (9, 0x465e_c429_b13c_3892),
8894 (15, 0xdee8_9d8a_065a_6233),
8895 (16, 0x1330_489a_7767_9c80),
8896 (31, 0x3391_303d_485e_846e),
8897 (32, 0x40b7_aff7_5d45_bbc8),
8898 (33, 0x4997_cae4_951c_17a5),
8899 (39, 0x5807_28fd_5c14_5739),
8900 (40, 0xf95c_f6f5_c08a_3d3b),
8901 (63, 0x2944_b4da_fc69_b206),
8902 (64, 0xbb76_f6ef_19bd_5a1b),
8903 (65, 0x814e_0c65_4a9f_d640),
8904 (127, 0x00de_aab1_31cf_f89b),
8905 (1000, 0x9e33_00c1_cde3_c58d),
8906 ] {
8907 assert_eq!(checksum(&bytes[..length]), expected, "the checksum of {length} bytes");
8908 }
8909 assert_eq!(checksum(b"the quick brown fox jumps over the lazy dog"), 0xed71_4233_c5a9_a792);
8910 }
8911 #[test]
8918 fn a_declared_order_comes_back_out_of_the_file() {
8919 let path = path("clustered");
8920 let shipped = vec![
8921 Field::new("key", LogicalType::BigInt),
8922 Field::new("line", LogicalType::Integer),
8923 Field::new("shipdate", LogicalType::Date),
8924 ];
8925 let plain = vec![Field::new("a", LogicalType::Integer)];
8926 let stage_zero = Clustering::new(vec![2, 0, 1], Width::Month, &shipped).expect("valid");
8927
8928 let mut writer = Writer::create(&path, "lineitem", shipped)
8929 .expect("new file")
8930 .declare(stage_zero.clone())
8931 .expect("the columns are the table's");
8932 let column = |ty: LogicalType, values: &[Value]| {
8933 Vector::from_values(ty, values).expect("the values match the type")
8934 };
8935 writer
8936 .append(
8937 &Chunk::new(vec![
8938 column(
8939 LogicalType::BigInt,
8940 &[Value::BigInt(0), Value::BigInt(1), Value::BigInt(2), Value::BigInt(3)],
8941 ),
8942 column(
8943 LogicalType::Integer,
8944 &[
8945 Value::Integer(1),
8946 Value::Integer(1),
8947 Value::Integer(1),
8948 Value::Integer(1),
8949 ],
8950 ),
8951 column(
8952 LogicalType::Date,
8953 &[Value::Date(0), Value::Date(1), Value::Date(2), Value::Date(3)],
8954 ),
8955 ])
8956 .expect("three columns"),
8957 )
8958 .expect("four rows");
8959 let mut writer = writer.next("nation", plain).expect("a second table");
8960 writer
8961 .append(
8962 &Chunk::new(vec![column(LogicalType::Integer, &[Value::Integer(7)])])
8963 .expect("one column"),
8964 )
8965 .expect("one row");
8966 writer.finish().expect("commit");
8967
8968 let catalog = Catalog::open(&path).expect("reopen");
8969 let lineitem = catalog.table("lineitem").expect("the clustered table");
8970 assert_eq!(lineitem.table().clustering(), Some(&stage_zero));
8971 let nation = catalog.table("nation").expect("the plain table");
8972 assert_eq!(nation.table().clustering(), None, "nobody declared one here");
8973
8974 assert_eq!(lineitem.table().rows(), 4);
8977 assert_eq!(nation.table().rows(), 1);
8978 fs::remove_file(&path).ok();
8979 }
8980
8981 #[test]
8983 fn a_declaration_off_the_end_of_the_table_never_reaches_the_file() {
8984 let path = path("clustered-bad");
8985 let writer = Writer::create(&path, "items", vec![Field::new("a", LogicalType::Integer)])
8986 .expect("new file");
8987 let four =
8988 (0..4).map(|at| Field::new(format!("c{at}"), LogicalType::Integer)).collect::<Vec<_>>();
8989 let wrong = Clustering::new(vec![3], Width::Exact, &four).expect("valid against four");
8990 assert!(writer.declare(wrong).is_err(), "the table has one column, not four");
8991 fs::remove_file(&path).ok();
8992 }
8993}