1#![forbid(unsafe_code)]
34
35use std::cmp::Ordering;
36use std::collections::{HashMap, VecDeque};
37use std::fs::{File, OpenOptions};
38use std::io::{Read, Seek, SeekFrom};
39use std::mem::{size_of, size_of_val};
40use std::path::Path;
41use std::slice;
42use std::sync::atomic::{AtomicUsize, Ordering as Atomic};
43use std::sync::{Arc, Mutex, OnceLock};
44
45use rudb_common::bounds::{Bound, Op, scaled_as};
46use rudb_common::{Error, Field, LogicalType, Result, Value};
47use rudb_encoding::{bitpack, chooser, integer, string};
48use rudb_storage::sieve::Sieve;
49use rudb_storage::{Probe, Range, Zone};
50use rudb_vector::string::StringColumn;
51use rudb_vector::validity::Validity;
52use rudb_vector::{Buffer, Chunk, Data, Packed, TextSource, Vector};
53
54const MAGIC: &[u8; 8] = b"RUDBNV10";
55const DIRECTORY: &[u8; 8] = b"RUDBDI10";
56const CATALOG: &[u8; 8] = b"RUDBCA10";
57const FORMAT: u32 = 22;
58const HEADER: u64 = 80;
59const SLOT_BYTES: usize = 28;
60const MAX_PAGE: usize = 256 * 1024 * 1024;
61const MAX_DIRECTORY: usize = 128 * 1024 * 1024;
62const FREQUENCIES: &[u8; 8] = b"RUDBFQ2\0";
63const FREQUENCY_CANDIDATES: usize = 32_768;
64const FREQUENCY_ENTRIES: usize = 512;
65const FREQUENCY_BUILD_RANK: usize = 10;
66const FREQUENCY_ORDINALS: usize = 65_536;
67const MAX_FREQUENCY_WORKERS: usize = 32;
74
75const MAX_ENCODE_WORKERS: usize = 32;
82
83const SIEVE_BUDGET: usize = 8 * 1024;
91
92const PART_BOUND_BYTES: usize = 24;
101
102fn io(error: std::io::Error) -> Error {
103 Error::io(error.to_string())
104}
105
106fn invalid(message: &str) -> Error {
107 Error::invalid_input(format!("invalid rudb native file: {message}"))
108}
109
110fn sum(counts: impl Iterator<Item = u64>) -> u64 {
112 counts.fold(0, u64::saturating_add)
113}
114
115fn span_bytes(spans: &[Span], at: usize) -> u64 {
117 spans.get(at).map_or(0, |span| u64::from(span.length))
118}
119
120fn page_bytes(pages: &[Option<Page>], at: usize) -> u64 {
122 pages.get(at).and_then(Option::as_ref).map_or(0, Page::bytes)
123}
124
125fn checksum(bytes: &[u8]) -> u64 {
126 const P1: u64 = 11_400_714_785_074_694_791;
127 const P2: u64 = 14_029_467_366_897_019_727;
128 const P3: u64 = 1_609_587_929_392_839_161;
129 const P4: u64 = 9_650_029_242_287_828_579;
130 const P5: u64 = 2_870_177_450_012_600_261;
131 let round = |state: u64, word: u64| {
132 state.wrapping_add(word.wrapping_mul(P2)).rotate_left(31).wrapping_mul(P1)
133 };
134 let merge = |state: u64, lane: u64| (state ^ round(0, lane)).wrapping_mul(P1).wrapping_add(P4);
135 let word =
136 |at: usize| u64::from_le_bytes(bytes[at..at + 8].try_into().expect("eight checksum bytes"));
137
138 let mut at = 0;
139 let mut hash = if bytes.len() >= 32 {
140 let mut one = P1.wrapping_add(P2);
141 let mut two = P2;
142 let mut three = 0;
143 let mut four = 0_u64.wrapping_sub(P1);
144 while at + 32 <= bytes.len() {
145 one = round(one, word(at));
146 two = round(two, word(at + 8));
147 three = round(three, word(at + 16));
148 four = round(four, word(at + 24));
149 at += 32;
150 }
151 let combined = one
152 .rotate_left(1)
153 .wrapping_add(two.rotate_left(7))
154 .wrapping_add(three.rotate_left(12))
155 .wrapping_add(four.rotate_left(18));
156 merge(merge(merge(merge(combined, one), two), three), four)
157 } else {
158 P5
159 };
160 hash = hash.wrapping_add(bytes.len() as u64);
161 while at + 8 <= bytes.len() {
162 hash ^= round(0, word(at));
163 hash = hash.rotate_left(27).wrapping_mul(P1).wrapping_add(P4);
164 at += 8;
165 }
166 if at + 4 <= bytes.len() {
167 let tail = u32::from_le_bytes(bytes[at..at + 4].try_into().expect("four checksum bytes"));
168 hash ^= u64::from(tail).wrapping_mul(P1);
169 hash = hash.rotate_left(23).wrapping_mul(P2).wrapping_add(P3);
170 at += 4;
171 }
172 while at < bytes.len() {
173 hash ^= u64::from(bytes[at]).wrapping_mul(P5);
174 hash = hash.rotate_left(11).wrapping_mul(P1);
175 at += 1;
176 }
177 hash ^= hash >> 33;
178 hash = hash.wrapping_mul(P2);
179 hash ^= hash >> 29;
180 hash = hash.wrapping_mul(P3);
181 hash ^ (hash >> 32)
182}
183
184#[derive(Debug, Clone, Copy)]
185struct Slot {
186 offset: u64,
187 length: u32,
188 generation: u64,
189 hash: u64,
190}
191
192impl Slot {
193 fn bytes(self) -> [u8; SLOT_BYTES] {
194 let mut result = [0; SLOT_BYTES];
195 result[..8].copy_from_slice(&self.offset.to_le_bytes());
196 result[8..12].copy_from_slice(&self.length.to_le_bytes());
197 result[12..20].copy_from_slice(&self.generation.to_le_bytes());
198 result[20..28].copy_from_slice(&self.hash.to_le_bytes());
199 result
200 }
201
202 fn read(bytes: &[u8]) -> Self {
203 Self {
204 offset: u64::from_le_bytes(bytes[..8].try_into().expect("eight bytes")),
205 length: u32::from_le_bytes(bytes[8..12].try_into().expect("four bytes")),
206 generation: u64::from_le_bytes(bytes[12..20].try_into().expect("eight bytes")),
207 hash: u64::from_le_bytes(bytes[20..28].try_into().expect("eight bytes")),
208 }
209 }
210}
211
212#[derive(Debug, Clone, Copy)]
213struct Page {
214 offset: u64,
215 length: u32,
216 hash: u64,
217}
218
219impl Page {
220 fn bytes(&self) -> u64 {
222 u64::from(self.length)
223 }
224}
225
226#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
227enum FrequencyValue {
228 Null,
229 Integer(i128),
230 Code(u32),
231}
232
233#[derive(Debug, Clone)]
234struct FrequencyEntry {
235 value: FrequencyValue,
236 count: u64,
237}
238
239#[derive(Debug, Clone)]
244struct FrequencySummary {
245 entries: Vec<FrequencyEntry>,
246 omitted_max: u64,
247 ordinals: Vec<u64>,
248}
249
250#[derive(Debug, Clone, PartialEq, Eq)]
252pub struct FrequencyOccurrences {
253 pub omitted_max: u64,
255 pub ordinals: Vec<u64>,
257}
258
259#[derive(Debug, Clone, Copy, Default)]
266struct Span {
267 offset: u64,
268 length: u32,
269}
270
271#[derive(Debug, Clone)]
273pub struct Stripe {
274 rows: usize,
275 parts: Vec<u32>,
278 index: Span,
282 pages: Vec<Span>,
283 memberships: Vec<Option<Page>>,
284 sieves: Vec<Option<Page>>,
287 part_ranges: Vec<Option<Page>>,
298 zone: Zone,
299}
300
301impl Stripe {
302 #[must_use]
304 pub fn rows(&self) -> usize {
305 self.rows
306 }
307
308 #[must_use]
310 pub fn parts(&self) -> usize {
311 self.parts.len()
312 }
313}
314
315#[derive(Debug, Clone)]
317pub struct Table {
318 name: String,
319 fields: Vec<Field>,
320 stripes: Vec<Stripe>,
321 rows: usize,
322 dictionaries: Vec<Option<Page>>,
323 frequencies: Vec<Option<FrequencySummary>>,
324 distincts: Vec<Option<u64>>,
334}
335
336impl Table {
337 #[must_use]
339 pub fn name(&self) -> &str {
340 &self.name
341 }
342
343 #[must_use]
345 pub fn fields(&self) -> &[Field] {
346 &self.fields
347 }
348
349 #[must_use]
351 pub fn rows(&self) -> usize {
352 self.rows
353 }
354
355 #[must_use]
357 pub fn stripes(&self) -> &[Stripe] {
358 &self.stripes
359 }
360}
361
362#[derive(Debug, Clone)]
374struct Entry {
375 name: String,
376 fields: Vec<Field>,
377 rows: usize,
378 directory: Page,
380}
381
382#[derive(Debug, Clone)]
384pub struct ColumnLayout {
385 pub name: String,
387 pub kind: String,
389 pub pages: u64,
391 pub memberships: u64,
393 pub sieves: u64,
395 pub part_ranges: u64,
397 pub dictionary: u64,
399}
400
401impl ColumnLayout {
402 #[must_use]
404 pub fn total(&self) -> u64 {
405 self.pages
406 .saturating_add(self.memberships)
407 .saturating_add(self.sieves)
408 .saturating_add(self.part_ranges)
409 .saturating_add(self.dictionary)
410 }
411}
412
413#[derive(Debug, Clone)]
424pub struct Layout {
425 pub file: u64,
427 pub rows: usize,
429 pub stripes: usize,
431 pub parts: usize,
433 pub columns: Vec<ColumnLayout>,
435 pub indexes: u64,
438 pub directory: u64,
440 pub header: u64,
442}
443
444impl Layout {
445 #[must_use]
447 pub fn columns_total(&self) -> u64 {
448 self.columns.iter().map(ColumnLayout::total).fold(0, u64::saturating_add)
449 }
450
451 #[must_use]
457 pub fn unaccounted(&self) -> u64 {
458 self.file
459 .saturating_sub(self.columns_total())
460 .saturating_sub(self.indexes)
461 .saturating_sub(self.directory)
462 .saturating_sub(self.header)
463 }
464}
465
466#[derive(Debug)]
468struct GlobalDictionary {
469 primary: HashMap<u64, u32>,
470 collisions: HashMap<u64, Vec<u32>>,
471 offsets: Vec<u32>,
472 payload: Vec<u8>,
473 counts: Vec<u64>,
474 nulls: u64,
475}
476
477impl GlobalDictionary {
478 fn new() -> Self {
479 Self {
480 primary: HashMap::new(),
481 collisions: HashMap::new(),
482 offsets: vec![0],
483 payload: Vec::new(),
484 counts: Vec::new(),
485 nulls: 0,
486 }
487 }
488
489 fn bytes(&self, code: u32) -> Option<&[u8]> {
490 let start = *self.offsets.get(code as usize)? as usize;
491 let end = *self.offsets.get(code as usize + 1)? as usize;
492 self.payload.get(start..end)
493 }
494
495 fn code(&mut self, text: &str) -> Result<u32> {
496 let hash = checksum(text.as_bytes());
497 if let Some(&code) = self.primary.get(&hash) {
498 if self.bytes(code) == Some(text.as_bytes()) {
499 return Ok(code);
500 }
501 if let Some(codes) = self.collisions.get(&hash) {
502 if let Some(code) =
503 codes.iter().copied().find(|&code| self.bytes(code) == Some(text.as_bytes()))
504 {
505 return Ok(code);
506 }
507 }
508 let code = self.insert(text)?;
509 self.collisions.entry(hash).or_default().push(code);
510 return Ok(code);
511 }
512 let code = self.insert(text)?;
513 self.primary.insert(hash, code);
514 Ok(code)
515 }
516
517 fn insert(&mut self, text: &str) -> Result<u32> {
518 let code = u32::try_from(self.offsets.len() - 1)
519 .map_err(|_| invalid("global dictionary has too many values"))?;
520 self.payload.extend_from_slice(text.as_bytes());
521 self.offsets.push(
522 u32::try_from(self.payload.len())
523 .map_err(|_| invalid("global dictionary payload exceeds 4 GiB"))?,
524 );
525 self.counts.push(0);
526 Ok(code)
527 }
528
529 fn ranked(&self) -> Vec<(u64, u32)> {
549 let count = self.offsets.len() - 1;
550 let mut ranked = (0..count)
551 .map(|code| {
552 let code = code as u32;
553 (head(self.bytes(code).unwrap_or_default()), code)
554 })
555 .collect::<Vec<_>>();
556 ranked.sort_unstable_by(|left, right| {
557 left.0.cmp(&right.0).then_with(|| self.bytes(left.1).cmp(&self.bytes(right.1)))
558 });
559 ranked
560 }
561
562 fn observe(&mut self, code: u32, null: bool) -> Result<()> {
563 if null {
564 self.nulls = self.nulls.saturating_add(1);
565 return Ok(());
566 }
567 let count = self
568 .counts
569 .get_mut(code as usize)
570 .ok_or_else(|| invalid("global dictionary count code is out of range"))?;
571 *count = count.saturating_add(1);
572 Ok(())
573 }
574}
575
576#[derive(Debug)]
584pub struct Writer {
585 file: File,
586 at: u64,
594 table: Table,
595 generation: u64,
596 order: Vec<((u64, u64), (u64, u64))>,
599 next_order: u64,
600 dictionaries: Vec<Option<GlobalDictionary>>,
601 pending: Vec<PendingChunk>,
602 closed: Vec<Entry>,
604}
605
606#[derive(Debug)]
614struct PendingChunk {
615 order: (u64, u64),
616 chunk: Chunk,
617}
618
619#[derive(Debug)]
625struct ColumnStripe {
626 pages: Vec<Vec<u8>>,
627 codes: Vec<Option<Vec<u32>>>,
628 sieves: Vec<Option<Sieve>>,
629 ranges: Vec<Range>,
630}
631
632fn weight(ty: &LogicalType) -> usize {
640 match ty {
641 LogicalType::Varchar | LogicalType::Blob => 64,
642 LogicalType::BigInt
643 | LogicalType::UBigInt
644 | LogicalType::Timestamp
645 | LogicalType::Double
646 | LogicalType::Decimal { .. } => 8,
647 LogicalType::Integer | LogicalType::UInteger | LogicalType::Date | LogicalType::Float => 4,
648 LogicalType::SmallInt | LogicalType::USmallInt => 2,
649 _ => 1,
650 }
651}
652
653pub const STRIPE_PARTS: usize = 64;
660
661const INDEX_ENTRY: usize = size_of::<u32>() + size_of::<u64>();
663
664fn index_section(parts: usize) -> Result<usize> {
666 parts
667 .checked_mul(INDEX_ENTRY)
668 .and_then(|bytes| bytes.checked_add(size_of::<u64>()))
669 .ok_or_else(|| invalid("index page length overflow"))
670}
671
672impl Writer {
673 pub fn create(
679 path: impl AsRef<Path>,
680 name: impl Into<String>,
681 fields: Vec<Field>,
682 ) -> Result<Self> {
683 for field in &fields {
684 type_tag(&field.ty)?;
685 }
686 let file =
687 OpenOptions::new().write(true).read(true).create_new(true).open(path).map_err(io)?;
688 let mut header = [0; HEADER as usize];
689 header[..8].copy_from_slice(MAGIC);
690 header[8..12].copy_from_slice(&FORMAT.to_le_bytes());
691 write_at(&file, 0, &header)?;
692 Ok(Self {
693 file,
694 at: HEADER,
695 dictionaries: fields
696 .iter()
697 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
698 .collect(),
699 table: Table {
700 name: name.into(),
701 dictionaries: vec![None; fields.len()],
702 distincts: vec![None; fields.len()],
703 fields,
704 stripes: Vec::new(),
705 rows: 0,
706 frequencies: Vec::new(),
707 },
708 generation: 1,
709 order: Vec::new(),
710 next_order: 0,
711 pending: Vec::with_capacity(STRIPE_PARTS),
712 closed: Vec::new(),
713 })
714 }
715
716 pub fn next(mut self, name: impl Into<String>, fields: Vec<Field>) -> Result<Self> {
727 for field in &fields {
728 type_tag(&field.ty)?;
729 }
730 let name = name.into();
731 let entry = self.close()?;
732 if self.closed.iter().chain(std::iter::once(&entry)).any(|held| held.name == name) {
733 return Err(invalid("two tables in one native file have the same name"));
734 }
735 let Self { file, at, generation, mut closed, .. } = self;
736 closed.push(entry);
737 Ok(Self {
738 file,
739 at,
740 generation,
741 closed,
742 dictionaries: fields
743 .iter()
744 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
745 .collect(),
746 table: Table {
747 name,
748 dictionaries: vec![None; fields.len()],
749 distincts: vec![None; fields.len()],
750 fields,
751 stripes: Vec::new(),
752 rows: 0,
753 frequencies: Vec::new(),
754 },
755 order: Vec::new(),
756 next_order: 0,
757 pending: Vec::with_capacity(STRIPE_PARTS),
758 })
759 }
760
761 fn put(&mut self, bytes: &[u8]) -> Result<()> {
766 write_at(&self.file, self.at, bytes)?;
767 self.at = self
768 .at
769 .checked_add(bytes.len() as u64)
770 .ok_or_else(|| invalid("native file length overflow"))?;
771 Ok(())
772 }
773
774 pub fn append(&mut self, chunk: &Chunk) -> Result<()> {
780 let order = (self.next_order, 0);
781 self.next_order = self.next_order.saturating_add(1);
782 self.append_at(order, chunk)
783 }
784
785 pub fn append_at(&mut self, order: (u64, u64), chunk: &Chunk) -> Result<()> {
796 if chunk.is_empty() {
797 return Ok(());
798 }
799 self.admit(chunk)?;
800 if self.pending.last().is_some_and(|last| last.order > order) {
801 self.flush_pending()?;
802 }
803 self.pending.push(PendingChunk { order, chunk: chunk.clone() });
808 if self.pending.len() == STRIPE_PARTS {
809 self.flush_pending()?;
810 }
811 Ok(())
812 }
813
814 pub fn append_stripe(&mut self, parts: Vec<((u64, u64), Chunk)>) -> Result<()> {
830 if parts.len() > STRIPE_PARTS {
831 return Err(invalid("a stripe was handed more parts than it holds"));
832 }
833 self.flush_pending()?;
836 for (order, chunk) in parts {
837 if chunk.is_empty() {
838 continue;
839 }
840 self.admit(&chunk)?;
841 self.pending.push(PendingChunk { order, chunk });
842 }
843 self.flush_pending()
844 }
845
846 fn admit(&mut self, chunk: &Chunk) -> Result<()> {
848 if chunk.width() != self.table.fields.len() {
849 return Err(invalid("chunk width differs from table schema"));
850 }
851 for (index, field) in self.table.fields.iter().enumerate() {
852 if chunk.column(index)?.logical_type() != &field.ty {
853 return Err(invalid("chunk type differs from table schema"));
854 }
855 }
856 self.table.rows = self
857 .table
858 .rows
859 .checked_add(chunk.len())
860 .ok_or_else(|| invalid("row count overflow"))?;
861 Ok(())
862 }
863
864 fn encode_column(
872 index: usize,
873 held: &[PendingChunk],
874 mut dictionary: Option<&mut GlobalDictionary>,
875 ) -> Result<ColumnStripe> {
876 let mut stripe = ColumnStripe {
877 pages: Vec::with_capacity(held.len()),
878 codes: Vec::with_capacity(held.len()),
879 sieves: Vec::with_capacity(held.len()),
880 ranges: Vec::with_capacity(held.len()),
881 };
882 for pending in held {
883 let column = pending.chunk.column(index)?;
884 let (bytes, unique) = encode(column, dictionary.as_deref_mut())?;
885 if bytes.len() > MAX_PAGE {
886 return Err(invalid("column page exceeds the configured bound"));
887 }
888 let range = Range::of(column);
891 let sieve = match dictionary {
904 Some(_) => None,
905 None => Sieve::of(column, &range, SIEVE_BUDGET)
906 .filter(|sieve| sieve.len() < bytes.len()),
907 };
908 stripe.pages.push(bytes);
909 stripe.codes.push(unique);
910 stripe.sieves.push(sieve);
911 stripe.ranges.push(range);
912 }
913 Ok(stripe)
914 }
915
916 fn encode_columns(&mut self, held: &[PendingChunk]) -> Result<Vec<ColumnStripe>> {
925 let width = self.table.fields.len();
926 let workers = std::thread::available_parallelism()
927 .map_or(1, usize::from)
928 .min(MAX_ENCODE_WORKERS)
929 .min(width);
930 if workers <= 1 || held.len() <= 1 {
931 return self
932 .dictionaries
933 .iter_mut()
934 .enumerate()
935 .map(|(index, dictionary)| Self::encode_column(index, held, dictionary.as_mut()))
936 .collect();
937 }
938 let mut jobs: Vec<(usize, Option<GlobalDictionary>)> =
941 std::mem::take(&mut self.dictionaries).into_iter().enumerate().collect();
942 jobs.sort_by_key(|(index, _)| weight(&self.table.fields[*index].ty));
944 let queue = Mutex::new(jobs);
945 let pieces = std::thread::scope(|scope| {
946 (0..workers)
947 .map(|_| {
948 scope.spawn(|| {
949 let mut mine = Vec::new();
950 loop {
951 let taken = queue
952 .lock()
953 .map_err(|_| Error::internal("a native encode worker panicked"))?
954 .pop();
955 let Some((index, mut dictionary)) = taken else { break };
956 let encoded = Self::encode_column(index, held, dictionary.as_mut())?;
957 mine.push((index, dictionary, encoded));
958 }
959 Ok(mine)
960 })
961 })
962 .collect::<Vec<_>>()
963 .into_iter()
964 .map(|handle| {
965 handle.join().map_err(|_| Error::internal("a native encode worker panicked"))?
966 })
967 .collect::<Result<Vec<_>>>()
968 })?;
969 let mut dictionaries: Vec<Option<GlobalDictionary>> = (0..width).map(|_| None).collect();
970 let mut encoded: Vec<Option<ColumnStripe>> = (0..width).map(|_| None).collect();
971 for piece in pieces {
972 for (index, dictionary, stripe) in piece {
973 dictionaries[index] = dictionary;
974 encoded[index] = Some(stripe);
975 }
976 }
977 self.dictionaries = dictionaries;
978 encoded
979 .into_iter()
980 .map(|stripe| stripe.ok_or_else(|| Error::internal("a column was never encoded")))
981 .collect()
982 }
983
984 fn flush_pending(&mut self) -> Result<()> {
986 if self.pending.is_empty() {
987 return Ok(());
988 }
989 let width = self.table.fields.len();
990 let mut held = std::mem::take(&mut self.pending);
993 let parts = held.len();
994 let encoded = self.encode_columns(&held)?;
995 let mut pages = Vec::with_capacity(width);
996 let mut memberships = vec![None; width];
997 let mut ranges = Vec::with_capacity(width);
998 let mut index = Vec::with_capacity(width.saturating_mul(index_section(parts)?));
999 for stripe in &encoded {
1000 let offset = self.at;
1001 let section = index.len();
1002 let mut length = 0_usize;
1003 for bytes in &stripe.pages {
1004 write_at(&self.file, self.at + length as u64, bytes)?;
1005 put_u32(
1006 &mut index,
1007 u32::try_from(bytes.len()).map_err(|_| invalid("part length overflow"))?,
1008 );
1009 put_u64(&mut index, checksum(bytes));
1010 length = length
1011 .checked_add(bytes.len())
1012 .ok_or_else(|| invalid("column page length overflow"))?;
1013 }
1014 let hash = checksum(&index[section..]);
1015 put_u64(&mut index, hash);
1016 if length > MAX_PAGE {
1017 return Err(invalid("column page exceeds the configured bound"));
1018 }
1019 self.at = self
1020 .at
1021 .checked_add(length as u64)
1022 .ok_or_else(|| invalid("native file length overflow"))?;
1023 pages.push(Span {
1024 offset,
1025 length: u32::try_from(length).map_err(|_| invalid("page length overflow"))?,
1026 });
1027 ranges.push(merged_range(stripe.ranges.iter().cloned()));
1028 }
1029 for (membership, stripe) in memberships.iter_mut().zip(&encoded) {
1030 if stripe.codes.iter().all(Option::is_none) {
1031 continue;
1032 }
1033 let lists = stripe
1034 .codes
1035 .iter()
1036 .map(|codes| codes.clone().unwrap_or_default())
1037 .collect::<Vec<_>>();
1038 let bytes = encode_membership(&merged_codes(lists));
1039 let offset = self.at;
1040 self.put(&bytes)?;
1041 *membership = Some(Page {
1042 offset,
1043 length: u32::try_from(bytes.len())
1044 .map_err(|_| invalid("membership page length overflow"))?,
1045 hash: checksum(&bytes),
1046 });
1047 }
1048 let mut sieves = vec![None; width];
1049 for (page, stripe) in sieves.iter_mut().zip(&encoded) {
1050 if stripe.sieves.iter().all(Option::is_none) {
1051 continue;
1052 }
1053 let bytes = encode_sieves(stripe.sieves.iter())?;
1054 let offset = self.at;
1055 self.put(&bytes)?;
1056 *page = Some(Page {
1057 offset,
1058 length: u32::try_from(bytes.len())
1059 .map_err(|_| invalid("sieve page length overflow"))?,
1060 hash: checksum(&bytes),
1061 });
1062 }
1063 let mut part_ranges = vec![None; width];
1069 if parts > 1 {
1070 for ((page, stripe), span) in part_ranges.iter_mut().zip(&encoded).zip(&pages) {
1071 let bytes = encode_part_ranges(&stripe.ranges)?;
1072 if bytes.len() >= span.length as usize {
1073 continue;
1074 }
1075 let offset = self.at;
1076 self.put(&bytes)?;
1077 *page = Some(Page {
1078 offset,
1079 length: u32::try_from(bytes.len())
1080 .map_err(|_| invalid("part range page length overflow"))?,
1081 hash: checksum(&bytes),
1082 });
1083 }
1084 }
1085 let offset = self.at;
1086 self.put(&index)?;
1087 let index = Span {
1088 offset,
1089 length: u32::try_from(index.len())
1090 .map_err(|_| invalid("index page length overflow"))?,
1091 };
1092 let mut rows = 0_usize;
1093 let mut lengths = Vec::with_capacity(parts);
1094 let mut span = None;
1095 for pending in held.drain(..) {
1096 let part = pending.chunk.len();
1097 rows = rows.checked_add(part).ok_or_else(|| invalid("row count overflow"))?;
1098 lengths.push(u32::try_from(part).map_err(|_| invalid("part row count overflow"))?);
1099 span = Some(
1100 span.map_or((pending.order, pending.order), |(first, _)| (first, pending.order)),
1101 );
1102 }
1103 self.order.push(span.ok_or_else(|| invalid("a stripe was flushed with no parts"))?);
1104 self.table.stripes.push(Stripe {
1105 rows,
1106 parts: lengths,
1107 index,
1108 pages,
1109 memberships,
1110 sieves,
1111 part_ranges,
1112 zone: Zone::from_ranges(ranges),
1113 });
1114 self.pending = held;
1116 Ok(())
1117 }
1118
1119 fn numeric_frequency(&self, column: usize) -> Result<Option<FrequencySummary>> {
1123 let ty = &self.table.fields[column].ty;
1124 if !matches!(
1125 ty,
1126 LogicalType::TinyInt
1127 | LogicalType::SmallInt
1128 | LogicalType::Integer
1129 | LogicalType::BigInt
1130 | LogicalType::UTinyInt
1131 | LogicalType::USmallInt
1132 | LogicalType::UInteger
1133 | LogicalType::UBigInt
1134 | LogicalType::Date
1135 | LogicalType::Timestamp
1136 ) {
1137 return Ok(None);
1138 }
1139 let mut candidates: HashMap<FrequencyValue, u32> = HashMap::new();
1140 let mut decrements = 0_u64;
1141 self.visit_numeric(column, |_, value| {
1142 if let Some(count) = candidates.get_mut(&value) {
1143 *count = count.saturating_add(1);
1144 } else if candidates.len() < FREQUENCY_CANDIDATES {
1145 candidates.insert(value, 1);
1146 } else {
1147 candidates.retain(|_, count| {
1148 *count -= 1;
1149 *count != 0
1150 });
1151 decrements = decrements.saturating_add(1);
1152 }
1153 })?;
1154 let (exact, ordinals) = if decrements == 0 {
1155 (
1156 candidates
1157 .into_iter()
1158 .map(|(value, count)| (value, u64::from(count)))
1159 .collect::<HashMap<_, _>>(),
1160 Vec::new(),
1161 )
1162 } else {
1163 let mut lower = candidates.values().copied().collect::<Vec<_>>();
1164 lower.sort_unstable_by(|left, right| right.cmp(left));
1165 if lower.len() < FREQUENCY_BUILD_RANK
1166 || u64::from(lower[FREQUENCY_BUILD_RANK - 1]) <= decrements
1167 {
1168 return Ok(None);
1169 }
1170 let mut exact =
1171 candidates.into_keys().map(|value| (value, 0_u64)).collect::<HashMap<_, _>>();
1172 let mut ordinals = Vec::new();
1173 let mut exceeded = false;
1174 self.visit_numeric(column, |ordinal, value| {
1175 if let Some(count) = exact.get_mut(&value) {
1176 *count = count.saturating_add(1);
1177 if !exceeded {
1178 if ordinals.len() < FREQUENCY_ORDINALS {
1179 ordinals.push(ordinal);
1180 } else {
1181 ordinals.clear();
1182 exceeded = true;
1183 }
1184 }
1185 }
1186 })?;
1187 (exact, ordinals)
1188 };
1189 let mut entries = exact
1190 .into_iter()
1191 .map(|(value, count)| FrequencyEntry { value, count })
1192 .collect::<Vec<_>>();
1193 entries.sort_unstable_by(|left, right| {
1194 right.count.cmp(&left.count).then_with(|| frequency_order(left.value, right.value))
1195 });
1196 let omitted_max =
1197 entries.get(FREQUENCY_ENTRIES).map_or(decrements, |entry| decrements.max(entry.count));
1198 entries.truncate(FREQUENCY_ENTRIES);
1199 Ok(Some(FrequencySummary { entries, omitted_max, ordinals }))
1200 }
1201
1202 fn visit_numeric(
1203 &self,
1204 column: usize,
1205 mut visit: impl FnMut(u64, FrequencyValue),
1206 ) -> Result<()> {
1207 let ty = &self.table.fields[column].ty;
1208 let mut start = 0_u64;
1209 for stripe in &self.table.stripes {
1210 let spans = read_index(&self.file, stripe, column)?;
1211 let page = stripe.pages[column];
1212 let mut bytes = vec![0; page.length as usize];
1213 read_at(&self.file, page.offset, &mut bytes)?;
1214 for (span, &rows) in spans.iter().zip(&stripe.parts) {
1215 let part = part_bytes(&bytes, *span)?;
1216 if checksum(part) != span.hash {
1217 return Err(invalid("column page checksum differs while building frequencies"));
1218 }
1219 let rows = rows as usize;
1220 let vector = decode(ty, rows, part, None)?;
1221 for row in 0..rows {
1223 let value = if vector.is_null_at(row) {
1224 FrequencyValue::Null
1225 } else {
1226 let widened = match vector.signed_at(row) {
1230 Some(value) => Some(value),
1231 None => match vector.value_at(row) {
1232 Value::UTinyInt(value) => Some(i128::from(value)),
1233 Value::USmallInt(value) => Some(i128::from(value)),
1234 Value::UInteger(value) => Some(i128::from(value)),
1235 Value::UBigInt(value) => Some(i128::from(value)),
1236 _ => None,
1237 },
1238 };
1239 FrequencyValue::Integer(widened.ok_or_else(|| {
1240 invalid("numeric frequency page did not contain an integer value")
1241 })?)
1242 };
1243 visit(start.saturating_add(row as u64), value);
1244 }
1245 start = start.saturating_add(rows as u64);
1246 }
1247 }
1248 Ok(())
1249 }
1250
1251 fn numeric_frequencies(&self) -> Result<Vec<Option<FrequencySummary>>> {
1259 let mut columns = self
1260 .table
1261 .fields
1262 .iter()
1263 .enumerate()
1264 .filter_map(|(column, field)| {
1265 matches!(
1266 field.ty,
1267 LogicalType::TinyInt
1268 | LogicalType::SmallInt
1269 | LogicalType::Integer
1270 | LogicalType::BigInt
1271 | LogicalType::UTinyInt
1272 | LogicalType::USmallInt
1273 | LogicalType::UInteger
1274 | LogicalType::UBigInt
1275 | LogicalType::Date
1276 | LogicalType::Timestamp
1277 )
1278 .then_some(column)
1279 })
1280 .collect::<Vec<_>>();
1281 let workers = std::thread::available_parallelism()
1282 .map_or(1, usize::from)
1283 .min(MAX_FREQUENCY_WORKERS)
1284 .min(columns.len());
1285 if workers <= 1 {
1286 let mut frequencies = vec![None; self.table.fields.len()];
1287 for column in columns {
1288 frequencies[column] = self.numeric_frequency(column)?;
1289 }
1290 return Ok(frequencies);
1291 }
1292 columns.sort_by_key(|&column| weight(&self.table.fields[column].ty));
1295 let queue = Mutex::new(columns);
1296 let pieces = std::thread::scope(|scope| {
1297 (0..workers)
1298 .map(|_| {
1299 scope.spawn(|| {
1300 let mut mine = Vec::new();
1301 loop {
1302 let taken = queue
1303 .lock()
1304 .map_err(|_| Error::internal("a native frequency worker panicked"))?
1305 .pop();
1306 let Some(column) = taken else { break };
1307 mine.push((column, self.numeric_frequency(column)?));
1308 }
1309 Ok(mine)
1310 })
1311 })
1312 .collect::<Vec<_>>()
1313 .into_iter()
1314 .map(|handle| {
1315 handle
1316 .join()
1317 .map_err(|_| Error::internal("a native frequency worker panicked"))?
1318 })
1319 .collect::<Result<Vec<_>>>()
1320 })?;
1321 let mut frequencies = vec![None; self.table.fields.len()];
1322 for piece in pieces {
1323 for (column, summary) in piece {
1324 frequencies[column] = summary;
1325 }
1326 }
1327 Ok(frequencies)
1328 }
1329
1330 fn close(&mut self) -> Result<Entry> {
1341 self.flush_pending()?;
1342 let mut stripes = std::mem::take(&mut self.order)
1343 .into_iter()
1344 .zip(std::mem::take(&mut self.table.stripes))
1345 .collect::<Vec<_>>();
1346 stripes.sort_by_key(|(order, _)| order.0);
1347 let mut previous: Option<(u64, u64)> = None;
1348 for ((first, last), _) in &stripes {
1349 if previous.is_some_and(|previous| previous >= *first) {
1350 return Err(invalid("chunks did not arrive in source order"));
1351 }
1352 previous = Some(*last);
1353 }
1354 self.table.stripes = stripes.into_iter().map(|(_, stripe)| stripe).collect();
1355 self.table.frequencies = self.numeric_frequencies()?;
1356 let dictionaries = std::mem::take(&mut self.dictionaries);
1357 let orders = rankings(&dictionaries)?;
1358 for (index, (dictionary, order)) in dictionaries.into_iter().zip(orders).enumerate() {
1359 let Some(dictionary) = dictionary else { continue };
1360 self.table.distincts[index] =
1364 Some(dictionary.counts.iter().filter(|count| **count != 0).count() as u64);
1365 self.table.frequencies[index] = Some(code_frequency(&dictionary));
1366 let encoded = encode_global_dictionary(dictionary, &order)?;
1367 let offset = self.at;
1368 self.put(&encoded.index)?;
1369 self.put(&encoded.ranks)?;
1370 for block in &encoded.payload {
1371 self.put(block)?;
1372 }
1373 let payload_len =
1374 encoded.payload.iter().try_fold(0_usize, |len, block| len.checked_add(block.len()));
1375 let length = payload_len
1376 .and_then(|len| len.checked_add(encoded.index.len()))
1377 .and_then(|len| len.checked_add(encoded.ranks.len()))
1378 .ok_or_else(|| invalid("dictionary page length overflow"))?;
1379 self.table.dictionaries[index] = Some(Page {
1380 offset,
1381 length: u32::try_from(length)
1382 .map_err(|_| invalid("dictionary page length overflow"))?,
1383 hash: checksum(&encoded.index),
1384 });
1385 }
1386 let directory = encode_directory(&self.table)?;
1387 if directory.len() > MAX_DIRECTORY {
1388 return Err(invalid("directory exceeds the configured bound"));
1389 }
1390 let offset = self.at;
1391 self.put(&directory)?;
1392 Ok(Entry {
1393 name: self.table.name.clone(),
1394 fields: self.table.fields.clone(),
1395 rows: self.table.rows,
1396 directory: Page {
1397 offset,
1398 length: u32::try_from(directory.len())
1399 .map_err(|_| invalid("directory length overflow"))?,
1400 hash: checksum(&directory),
1401 },
1402 })
1403 }
1404
1405 pub fn finish(mut self) -> Result<Table> {
1415 let entry = self.close()?;
1416 let mut tables = std::mem::take(&mut self.closed);
1417 tables.push(entry);
1418 let catalog = encode_catalog(&tables)?;
1419 if catalog.len() > MAX_DIRECTORY {
1420 return Err(invalid("catalog exceeds the configured bound"));
1421 }
1422 let offset = self.at;
1423 self.put(&catalog)?;
1424 self.file.sync_all().map_err(io)?;
1428 let slot = Slot {
1429 offset,
1430 length: u32::try_from(catalog.len()).map_err(|_| invalid("catalog length overflow"))?,
1431 generation: self.generation,
1432 hash: checksum(&catalog),
1433 };
1434 write_at(&self.file, 16, &slot.bytes())?;
1437 self.file.sync_all().map_err(io)?;
1438 Ok(self.table)
1439 }
1440}
1441
1442#[derive(Debug, Clone)]
1444pub struct Reader {
1445 file: Arc<File>,
1446 table: Arc<Table>,
1447 dictionaries: Arc<Vec<OnceLock<Arc<Vector>>>>,
1448 loading: Arc<Vec<Mutex<()>>>,
1457 opened: Arc<AtomicUsize>,
1461 sieves: Arc<Vec<Vec<SieveSlot>>>,
1465 part_ranges: Arc<Vec<Vec<RangeSlot>>>,
1468 places: Arc<Vec<Place>>,
1470 cache: Arc<Vec<Mutex<Cached>>>,
1471 pages: Arc<AtomicUsize>,
1474 indexes: Arc<AtomicUsize>,
1477 kept: Arc<AtomicUsize>,
1480 size: u64,
1482 directory: u64,
1484 opening: Opening,
1486}
1487
1488#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1500pub struct Opening {
1501 pub reads: u32,
1504 pub bytes: u64,
1506}
1507
1508#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1510pub struct Reads {
1511 pub opening: Opening,
1513 pub pages: usize,
1515 pub indexes: usize,
1517 pub dictionaries: usize,
1520}
1521
1522#[derive(Debug, Clone, Copy)]
1524struct Place {
1525 stripe: u32,
1526 part: u32,
1527 rows: u32,
1528}
1529
1530#[derive(Debug, Clone, Copy)]
1532struct PartSpan {
1533 start: usize,
1534 length: usize,
1535 hash: u64,
1536}
1537
1538#[derive(Debug, Clone)]
1544struct CachedColumn {
1545 stripe: usize,
1546 index: Arc<Vec<PartSpan>>,
1547 page: Option<Arc<Vec<u8>>>,
1548}
1549
1550#[derive(Debug, Default)]
1570struct Cached {
1571 pages: Vec<Option<Arc<Vec<u8>>>>,
1572 order: VecDeque<usize>,
1573 loading: Vec<usize>,
1574 index: Vec<Option<Arc<Vec<PartSpan>>>>,
1575}
1576
1577const CACHED_STRIPES_PER_COLUMN: usize = 4;
1589
1590type SieveSlot = OnceLock<Arc<Vec<Option<Sieve>>>>;
1592
1593type RangeSlot = OnceLock<Arc<Vec<Range>>>;
1594
1595#[derive(Debug)]
1596struct NativeText {
1597 file: Arc<File>,
1598 values: usize,
1600 offsets: Vec<u8>,
1609 offset_bits: usize,
1612 ranks: usize,
1614 rank_at: u64,
1618 rank_ends: Vec<u64>,
1622 rank_hashes: Vec<u64>,
1623 rank_blocks: Vec<OnceLock<Result<Vec<u8>>>>,
1624 code_bits: usize,
1627 code_ranks: OnceLock<Option<Vec<u32>>>,
1634 payload: u64,
1635 ends: Vec<u64>,
1638 hashes: Vec<u64>,
1639 blocks: Vec<OnceLock<Result<Vec<u8>>>>,
1641 keep_budget: usize,
1644 payload_kept: AtomicUsize,
1652}
1653
1654const TEXT_PAYLOAD_VALUES: usize = 1024;
1670
1671const TEXT_KEEP_BUDGET: usize = 256 * 1024 * 1024;
1692
1693const TEXT_OFFSET_RUN: usize = 512;
1700
1701const DICTIONARY_HEADER: usize = 16;
1704
1705const TEXT_RANK_BLOCK: usize = 512;
1716
1717const RANK_BLOCK_HEADER: usize = size_of::<u64>() + 1;
1731
1732impl NativeText {
1733 fn payload_block(&self, block: usize) -> Result<Option<&[u8]>> {
1740 let Some(slot) = self.blocks.get(block) else { return Ok(None) };
1741 let bytes = slot.get_or_init(|| self.decode_block(block)).as_ref().map_err(Clone::clone)?;
1742 Ok(Some(bytes.as_slice()))
1743 }
1744
1745 fn decode_block(&self, block: usize) -> Result<Vec<u8>> {
1750 let start = if block == 0 { 0 } else { self.ends[block - 1] };
1751 let end = self.ends[block];
1752 let len = end
1753 .checked_sub(start)
1754 .ok_or_else(|| invalid("global dictionary block ends before it starts"))?;
1755 let mut stored = vec![
1756 0;
1757 usize::try_from(len).map_err(|_| invalid(
1758 "global dictionary block does not fit in memory"
1759 ))?
1760 ];
1761 read_at(&self.file, self.payload + start, &mut stored)?;
1762 if checksum(&stored) != self.hashes[block] {
1763 return Err(invalid("global dictionary payload checksum differs"));
1764 }
1765 let first = block * TEXT_PAYLOAD_VALUES;
1766 let last = (first + TEXT_PAYLOAD_VALUES).min(self.values);
1767 let want = self.end_within(last - 1)? as usize;
1768 let values = string::decode_flat(&stored)?;
1769 if values.len() != last - first {
1770 return Err(invalid("global dictionary block holds the wrong value count"));
1771 }
1772 let bytes = values.into_bytes();
1773 if bytes.len() != want {
1774 return Err(invalid("global dictionary block decodes to the wrong length"));
1775 }
1776 Ok(bytes)
1777 }
1778
1779 fn end_within(&self, index: usize) -> Result<u32> {
1781 let run = index / TEXT_OFFSET_RUN;
1782 let bytes = self
1783 .offsets
1784 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
1785 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1786 let end = bitpack::tail_at(bytes, self.offset_bits, index % TEXT_OFFSET_RUN)
1787 .map_err(|_| invalid("global dictionary offsets are short"))?;
1788 u32::try_from(end).map_err(|_| invalid("global dictionary offset is past the payload"))
1789 }
1790
1791 fn ends_within(&self, first: usize, last: usize) -> Result<Vec<u64>> {
1804 let mut ends = Vec::with_capacity(last.saturating_sub(first));
1805 let mut at = first;
1806 while at < last {
1807 let run = at / TEXT_OFFSET_RUN;
1808 let stop = ((run + 1) * TEXT_OFFSET_RUN).min(last);
1809 let held = self.values.saturating_sub(run * TEXT_OFFSET_RUN).min(TEXT_OFFSET_RUN);
1810 let bytes = self
1811 .offsets
1812 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
1813 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1814 let run_ends = bitpack::unpack_tail(bytes, self.offset_bits, held)
1815 .map_err(|_| invalid("global dictionary offsets are short"))?;
1816 let within = run_ends
1817 .get(at % TEXT_OFFSET_RUN..stop - run * TEXT_OFFSET_RUN)
1818 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1819 ends.extend_from_slice(within);
1820 at = stop;
1821 }
1822 Ok(ends)
1823 }
1824
1825 fn start_within(&self, index: usize) -> Result<u32> {
1828 if index % TEXT_PAYLOAD_VALUES == 0 { Ok(0) } else { self.end_within(index - 1) }
1829 }
1830
1831 fn span_within(&self, index: usize) -> Result<(u32, u32)> {
1839 let within = index % TEXT_OFFSET_RUN;
1840 let (start, end) = if within == 0 {
1841 (self.start_within(index)?, self.end_within(index)?)
1842 } else {
1843 let run = index / TEXT_OFFSET_RUN;
1844 let bytes = self
1845 .offsets
1846 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
1847 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1848 let (start, end) = bitpack::tail_pair(bytes, self.offset_bits, within)
1849 .map_err(|_| invalid("global dictionary offsets are short"))?;
1850 let ends = u32::try_from(end)
1851 .map_err(|_| invalid("global dictionary offset is past the payload"))?;
1852 let starts = u32::try_from(start)
1853 .map_err(|_| invalid("global dictionary offset is past the payload"))?;
1854 (starts, ends)
1855 };
1856 if start > end {
1857 return Err(invalid("global dictionary value ends before it starts"));
1858 }
1859 Ok((start, end))
1860 }
1861
1862 fn rank_parts(&self, rank: usize) -> Result<(&[u8], usize)> {
1869 let slot = self
1870 .rank_blocks
1871 .get(rank / TEXT_RANK_BLOCK)
1872 .ok_or_else(|| invalid("global dictionary rank is past the order"))?;
1873 let block = slot
1874 .get_or_init(|| {
1875 let which = rank / TEXT_RANK_BLOCK;
1876 let start = if which == 0 { 0 } else { self.rank_ends[which - 1] };
1877 let end = self.rank_ends[which];
1878 let mut bytes = vec![0; (end - start) as usize];
1879 read_at(&self.file, self.rank_at + start, &mut bytes)?;
1880 if checksum(&bytes)
1881 != *self
1882 .rank_hashes
1883 .get(rank / TEXT_RANK_BLOCK)
1884 .ok_or_else(|| invalid("global dictionary rank block has no checksum"))?
1885 {
1886 return Err(invalid("global dictionary rank checksum differs"));
1887 }
1888 Ok(bytes)
1889 })
1890 .as_ref()
1891 .map_err(Clone::clone)?;
1892 Ok((block.as_slice(), rank % TEXT_RANK_BLOCK))
1893 }
1894
1895 fn head_at(&self, rank: usize) -> Result<u64> {
1897 let (block, within) = self.rank_parts(rank)?;
1898 let (base, width, packed) = rank_heads(block)?;
1899 let above = bitpack::tail_at(packed, width, within)
1900 .map_err(|_| invalid("global dictionary rank block is short of heads"))?;
1901 Ok(base.wrapping_add(above))
1902 }
1903
1904 fn rank_codes<'block>(&self, block: &'block [u8], count: usize) -> Result<&'block [u8]> {
1906 let (_, width, packed) = rank_heads(block)?;
1907 packed
1908 .get(bitpack::tail_len(count, width)..)
1909 .ok_or_else(|| invalid("global dictionary rank block is short of codes"))
1910 }
1911
1912 fn rank_block_len(&self, rank: usize) -> usize {
1914 let first = rank / TEXT_RANK_BLOCK * TEXT_RANK_BLOCK;
1915 TEXT_RANK_BLOCK.min(self.ranks - first)
1916 }
1917}
1918
1919fn rank_heads(block: &[u8]) -> Result<(u64, usize, &[u8])> {
1921 let header = block
1922 .get(..RANK_BLOCK_HEADER)
1923 .ok_or_else(|| invalid("global dictionary rank block is short"))?;
1924 let base = u64::from_le_bytes(header[..8].try_into().expect("eight bytes"));
1925 let width = header[8] as usize;
1926 if width > 64 {
1927 return Err(invalid("global dictionary rank block packs heads past a word"));
1928 }
1929 Ok((base, width, &block[RANK_BLOCK_HEADER..]))
1930}
1931
1932fn offset_width(offsets: &[u32]) -> usize {
1939 let values = offsets.len() - 1;
1940 let mut span = 0;
1941 for first in (0..values).step_by(TEXT_PAYLOAD_VALUES) {
1942 let last = (first + TEXT_PAYLOAD_VALUES).min(values);
1943 span = span.max(offsets[last] - offsets[first]);
1944 }
1945 (u32::BITS - span.leading_zeros()) as usize
1946}
1947
1948fn offset_bytes(values: usize, bits: usize) -> usize {
1951 let full = values / TEXT_OFFSET_RUN;
1952 let rest = values % TEXT_OFFSET_RUN;
1953 full * TEXT_OFFSET_RUN / 8 * bits + bitpack::tail_len(rest, bits)
1954}
1955
1956fn encode_offsets(offsets: &[u32], bits: usize, out: &mut Vec<u8>) -> Result<()> {
1958 let values = offsets.len() - 1;
1959 let mut run = Vec::with_capacity(TEXT_OFFSET_RUN);
1960 for first in (0..values).step_by(TEXT_OFFSET_RUN) {
1961 let last = (first + TEXT_OFFSET_RUN).min(values);
1962 let base = offsets[first / TEXT_PAYLOAD_VALUES * TEXT_PAYLOAD_VALUES];
1963 run.clear();
1964 run.extend((first..last).map(|value| u64::from(offsets[value + 1] - base)));
1965 bitpack::pack_tail(&run, bits, out)
1966 .map_err(|_| invalid("global dictionary offsets do not pack"))?;
1967 }
1968 Ok(())
1969}
1970
1971fn code_width(values: usize) -> usize {
1973 match u64::try_from(values).unwrap_or(u64::MAX) {
1974 0 | 1 => 0,
1975 last => (u64::BITS - (last - 1).leading_zeros()) as usize,
1976 }
1977}
1978
1979impl TextSource for NativeText {
1980 fn len(&self) -> usize {
1981 self.values
1982 }
1983
1984 fn bytes_at(&self, index: usize) -> Result<Option<&[u8]>> {
1985 if index >= self.values {
1986 return Ok(None);
1987 }
1988 let (start, end) = self.span_within(index)?;
1989 if start == end {
1990 return Ok(Some(&[]));
1991 }
1992 let block = index / TEXT_PAYLOAD_VALUES;
1995 let Some(bytes) = self.payload_block(block)? else { return Ok(None) };
1996 Ok(bytes.get(start as usize..end as usize))
1997 }
1998
1999 fn bytes_len_at(&self, index: usize) -> Result<Option<usize>> {
2000 if index >= self.values {
2001 return Ok(None);
2002 }
2003 let (start, end) = self.span_within(index)?;
2004 Ok(Some((end - start) as usize))
2005 }
2006
2007 fn sweep(
2020 &self,
2021 first: usize,
2022 limit: usize,
2023 body: &mut dyn FnMut(usize, &[u8]) -> Result<()>,
2024 ) -> Result<usize> {
2025 let limit = limit.min(self.values);
2026 if first >= limit {
2027 return Ok(first);
2028 }
2029 let block = first / TEXT_PAYLOAD_VALUES;
2030 let last = ((block + 1) * TEXT_PAYLOAD_VALUES).min(limit);
2031 let decoded;
2032 let bytes: &[u8] = match self.blocks.get(block).and_then(OnceLock::get) {
2033 Some(Ok(kept)) => kept,
2034 _ if self.payload_kept.load(Atomic::Relaxed) < self.keep_budget => {
2035 let kept = self
2036 .payload_block(block)?
2037 .ok_or_else(|| invalid("global dictionary block is past the payload"))?;
2038 self.payload_kept.fetch_add(kept.len(), Atomic::Relaxed);
2039 kept
2040 }
2041 _ => {
2042 decoded = self.decode_block(block)?;
2043 &decoded
2044 }
2045 };
2046 let ends = self.ends_within(first, last)?;
2047 if ends.len() != last - first {
2048 return Err(invalid("global dictionary offsets are short"));
2049 }
2050 let mut start = u64::from(self.start_within(first)?);
2051 for (index, &end) in (first..last).zip(&ends) {
2054 let value = usize::try_from(start)
2055 .ok()
2056 .zip(usize::try_from(end).ok())
2057 .and_then(|(from, to)| bytes.get(from..to))
2058 .ok_or_else(|| invalid("global dictionary value is past its block"))?;
2059 body(index, value)?;
2060 start = end;
2061 }
2062 Ok(last)
2063 }
2064
2065 fn ranks(&self) -> Option<usize> {
2066 (self.ranks > 0).then_some(self.ranks)
2067 }
2068
2069 fn compare_rank(&self, rank: usize, wanted: &[u8]) -> Result<Ordering> {
2070 let settled = self.head_at(rank)?.cmp(&head(wanted));
2074 if settled != Ordering::Equal {
2075 return Ok(settled);
2076 }
2077 let code = self.code_at_rank(rank)?;
2078 let bytes = self
2079 .bytes_at(code as usize)?
2080 .ok_or_else(|| invalid("global dictionary order names a code it does not have"))?;
2081 Ok(bytes.cmp(wanted))
2082 }
2083
2084 fn code_at_rank(&self, rank: usize) -> Result<u32> {
2085 let (block, within) = self.rank_parts(rank)?;
2086 let codes = self.rank_codes(block, self.rank_block_len(rank))?;
2087 let code = bitpack::tail_at(codes, self.code_bits, within)
2088 .map_err(|_| invalid("global dictionary rank block is short of codes"))?;
2089 let code = u32::try_from(code)
2090 .map_err(|_| invalid("global dictionary order names a code it does not have"))?;
2091 if code as usize >= self.len() {
2092 return Err(invalid("global dictionary order names a code it does not have"));
2093 }
2094 Ok(code)
2095 }
2096
2097 fn code_ranks(&self) -> Option<&[u32]> {
2098 if self.ranks == 0 || self.ranks != self.len() {
2102 return None;
2103 }
2104 self.code_ranks
2105 .get_or_init(|| {
2106 let mut ranks = vec![u32::MAX; self.ranks];
2107 for first in (0..self.ranks).step_by(TEXT_RANK_BLOCK) {
2110 let (block, _) = self.rank_parts(first).ok()?;
2111 let count = self.rank_block_len(first);
2112 let codes = self.rank_codes(block, count).ok()?;
2113 for (within, code) in bitpack::unpack_tail(codes, self.code_bits, count)
2114 .ok()?
2115 .into_iter()
2116 .enumerate()
2117 {
2118 let code = usize::try_from(code).ok()?;
2119 *ranks.get_mut(code)? = u32::try_from(first + within).ok()?;
2120 }
2121 }
2122 if ranks.contains(&u32::MAX) {
2123 return None;
2124 }
2125 Some(ranks)
2126 })
2127 .as_deref()
2128 }
2129
2130 fn footprint(&self) -> usize {
2131 self.offsets.capacity()
2132 + self
2133 .code_ranks
2134 .get()
2135 .and_then(Option::as_ref)
2136 .map_or(0, |ranks| ranks.capacity() * size_of::<u32>())
2137 + self.rank_hashes.capacity() * size_of::<u64>()
2138 + self.rank_ends.capacity() * size_of::<u64>()
2139 + self.rank_blocks.capacity() * size_of::<OnceLock<Result<Vec<u8>>>>()
2140 + self
2141 .rank_blocks
2142 .iter()
2143 .filter_map(OnceLock::get)
2144 .filter_map(|result| result.as_ref().ok())
2145 .map(Vec::capacity)
2146 .sum::<usize>()
2147 + self.blocks.capacity() * size_of::<OnceLock<Result<Vec<u8>>>>()
2148 + self.hashes.capacity() * size_of::<u64>()
2149 + self.ends.capacity() * size_of::<u64>()
2150 + self
2151 .blocks
2152 .iter()
2153 .filter_map(OnceLock::get)
2154 .filter_map(|result| result.as_ref().ok())
2155 .map(Vec::capacity)
2156 .sum::<usize>()
2157 }
2158}
2159
2160fn places(table: &Table) -> Result<Vec<Place>> {
2162 let mut places = Vec::with_capacity(table.stripes.len().saturating_mul(STRIPE_PARTS));
2163 for (at, stripe) in table.stripes.iter().enumerate() {
2164 let index = u32::try_from(at).map_err(|_| invalid("too many stripes"))?;
2165 for (part, &rows) in stripe.parts.iter().enumerate() {
2166 places.push(Place {
2167 stripe: index,
2168 part: u32::try_from(part).map_err(|_| invalid("too many parts in a stripe"))?,
2169 rows,
2170 });
2171 }
2172 }
2173 Ok(places)
2174}
2175
2176fn read_index(file: &File, stripe: &Stripe, column: usize) -> Result<Vec<PartSpan>> {
2181 let parts = stripe.parts.len();
2182 let section = index_section(parts)?;
2183 let at = column.checked_mul(section).ok_or_else(|| invalid("index page offset overflow"))?;
2184 let end = at.checked_add(section).ok_or_else(|| invalid("index page offset overflow"))?;
2185 if end > stripe.index.length as usize {
2186 return Err(invalid("index page is shorter than its columns"));
2187 }
2188 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
2189 let mut bytes = vec![0; section];
2190 let offset = stripe
2191 .index
2192 .offset
2193 .checked_add(at as u64)
2194 .ok_or_else(|| invalid("index page offset overflow"))?;
2195 read_at(file, offset, &mut bytes)?;
2196 let entries = section - size_of::<u64>();
2197 let stored = u64::from_le_bytes(bytes[entries..].try_into().expect("eight bytes"));
2198 if checksum(&bytes[..entries]) != stored {
2199 return Err(invalid(&format!(
2202 "index page section checksum differs, column {column} of {parts} parts at {offset}, \
2203 wanted {stored:016x} and got {:016x}",
2204 checksum(&bytes[..entries]),
2205 )));
2206 }
2207 let mut spans = Vec::with_capacity(parts);
2208 let mut start = 0_usize;
2209 for part in 0..parts {
2210 let at = part * INDEX_ENTRY;
2211 let length = u32::from_le_bytes(bytes[at..at + 4].try_into().expect("four bytes")) as usize;
2212 let hash = u64::from_le_bytes(bytes[at + 4..at + 12].try_into().expect("eight bytes"));
2213 spans.push(PartSpan { start, length, hash });
2214 start = start.checked_add(length).ok_or_else(|| invalid("column page length overflow"))?;
2215 }
2216 if start != page.length as usize {
2217 return Err(invalid("column page length differs from its index"));
2218 }
2219 Ok(spans)
2220}
2221
2222fn part_bytes(page: &[u8], span: PartSpan) -> Result<&[u8]> {
2224 let end = span.start.checked_add(span.length).ok_or_else(|| invalid("part range overflow"))?;
2225 page.get(span.start..end).ok_or_else(|| invalid("part exceeds its column page"))
2226}
2227
2228fn remember(cached: &mut Cached, held: &CachedColumn, kept: usize) {
2233 if let Some(slot) = cached.index.get_mut(held.stripe) {
2234 if slot.is_none() {
2235 *slot = Some(Arc::clone(&held.index));
2236 }
2237 }
2238 let Some(page) = held.page.clone() else { return };
2239 let Some(slot) = cached.pages.get_mut(held.stripe) else { return };
2240 if slot.is_none() {
2241 cached.order.push_back(held.stripe);
2242 }
2243 *slot = Some(page);
2244 while cached.order.len() > kept.max(1) {
2245 let Some(oldest) = cached.order.pop_front() else { break };
2246 if let Some(slot) = cached.pages.get_mut(oldest) {
2247 *slot = None;
2248 }
2249 }
2250}
2251
2252#[derive(Debug, Clone)]
2261pub struct Catalog {
2262 file: Arc<File>,
2263 size: u64,
2264 entries: Arc<Vec<Entry>>,
2265 opening: Opening,
2266}
2267
2268impl Catalog {
2269 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
2275 let (file, size, bytes, opening) = slot_bytes(path)?;
2276 let entries = decode_catalog(&bytes, size)?;
2277 Ok(Self { file: Arc::new(file), size, entries: Arc::new(entries), opening })
2278 }
2279
2280 pub fn names(&self) -> impl ExactSizeIterator<Item = &str> {
2282 self.entries.iter().map(|entry| entry.name.as_str())
2283 }
2284
2285 #[must_use]
2287 pub fn len(&self) -> usize {
2288 self.entries.len()
2289 }
2290
2291 #[must_use]
2293 pub fn is_empty(&self) -> bool {
2294 self.entries.is_empty()
2295 }
2296
2297 pub fn table(&self, name: &str) -> Result<Reader> {
2303 let entry = self
2304 .entries
2305 .iter()
2306 .find(|entry| entry.name == name)
2307 .ok_or_else(|| invalid(&format!("the file holds no table called {name}")))?;
2308 let mut bytes = vec![0; entry.directory.length as usize];
2309 read_at(&self.file, entry.directory.offset, &mut bytes)?;
2310 if checksum(&bytes) != entry.directory.hash {
2311 return Err(invalid(&format!("the directory of table {name} does not checksum")));
2312 }
2313 let mut opening = self.opening;
2314 opening.reads += 1;
2315 opening.bytes += u64::from(entry.directory.length);
2316 Reader::build(
2317 Arc::clone(&self.file),
2318 self.size,
2319 decode_directory(&bytes, self.size)?,
2320 u64::from(entry.directory.length),
2321 opening,
2322 )
2323 }
2324}
2325
2326fn slot_bytes(path: impl AsRef<Path>) -> Result<(File, u64, Vec<u8>, Opening)> {
2331 let mut file = File::open(path).map_err(io)?;
2332 let size = file.metadata().map_err(io)?.len();
2333 if size < HEADER {
2334 return Err(invalid("file is shorter than its header"));
2335 }
2336 let mut header = [0; HEADER as usize];
2337 file.read_exact(&mut header).map_err(io)?;
2338 let mut opening = Opening { reads: 1, bytes: HEADER };
2339 let version = u32::from_le_bytes([header[8], header[9], header[10], header[11]]);
2340 if &header[..8] != MAGIC {
2345 return Err(invalid("the header does not begin with a rudb native magic"));
2346 }
2347 if version != FORMAT {
2348 return Err(invalid(&format!(
2349 "the file is format {version} and this build reads format {FORMAT}, so it has to \
2350 be written again"
2351 )));
2352 }
2353 let mut selected = None;
2354 for start in [16, 16 + SLOT_BYTES] {
2355 let slot = Slot::read(&header[start..start + SLOT_BYTES]);
2356 if slot.generation == 0 || slot.length == 0 || slot.length as usize > MAX_DIRECTORY {
2357 continue;
2358 }
2359 let Some(end) = slot.offset.checked_add(u64::from(slot.length)) else { continue };
2360 if slot.offset < HEADER || end > size {
2361 continue;
2362 }
2363 let mut bytes = vec![0; slot.length as usize];
2364 file.seek(SeekFrom::Start(slot.offset)).map_err(io)?;
2365 file.read_exact(&mut bytes).map_err(io)?;
2366 opening.reads += 1;
2367 opening.bytes += u64::from(slot.length);
2368 if checksum(&bytes) == slot.hash
2369 && selected
2370 .as_ref()
2371 .is_none_or(|(old, _): &(Slot, Vec<u8>)| old.generation < slot.generation)
2372 {
2373 selected = Some((slot, bytes));
2374 }
2375 }
2376 let (_, bytes) = selected.ok_or_else(|| invalid("no committed directory slot is valid"))?;
2377 Ok((file, size, bytes, opening))
2378}
2379
2380impl Reader {
2381 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
2388 let catalog = Catalog::open(path)?;
2389 let mut names = catalog.names();
2390 let name = names.next().ok_or_else(|| invalid("the file holds no table"))?.to_string();
2391 if names.next().is_some() {
2392 return Err(invalid(
2393 "the file holds more than one table, so it has to be opened by name",
2394 ));
2395 }
2396 catalog.table(&name)
2397 }
2398
2399 fn build(
2401 file: Arc<File>,
2402 size: u64,
2403 table: Table,
2404 directory: u64,
2405 opening: Opening,
2406 ) -> Result<Self> {
2407 let places = places(&table)?;
2408 let dictionaries = (0..table.fields.len()).map(|_| OnceLock::new()).collect();
2409 let table_fields = table.fields.len();
2410 let stripes = table.stripes.len();
2411 let cache = (0..table.fields.len())
2412 .map(|_| {
2413 Mutex::new(Cached {
2414 pages: (0..stripes).map(|_| None).collect(),
2415 index: (0..stripes).map(|_| None).collect(),
2416 ..Cached::default()
2417 })
2418 })
2419 .collect::<Vec<_>>();
2420 let sieves: Vec<Vec<SieveSlot>> = (0..table.fields.len())
2421 .map(|_| table.stripes.iter().map(|_| OnceLock::new()).collect())
2422 .collect();
2423 let part_ranges: Vec<Vec<RangeSlot>> = (0..table.fields.len())
2424 .map(|_| table.stripes.iter().map(|_| OnceLock::new()).collect())
2425 .collect();
2426 Ok(Self {
2427 file,
2428 table: Arc::new(table),
2429 dictionaries: Arc::new(dictionaries),
2430 loading: Arc::new((0..table_fields).map(|_| Mutex::new(())).collect()),
2431 opened: Arc::new(AtomicUsize::new(0)),
2432 sieves: Arc::new(sieves),
2433 part_ranges: Arc::new(part_ranges),
2434 places: Arc::new(places),
2435 cache: Arc::new(cache),
2436 pages: Arc::new(AtomicUsize::new(0)),
2437 indexes: Arc::new(AtomicUsize::new(0)),
2438 kept: Arc::new(AtomicUsize::new(CACHED_STRIPES_PER_COLUMN)),
2439 size,
2440 directory,
2441 opening,
2442 })
2443 }
2444
2445 #[must_use]
2452 pub fn reads(&self) -> Reads {
2453 Reads {
2454 opening: self.opening,
2455 pages: self.pages.load(Atomic::Relaxed),
2456 indexes: self.indexes.load(Atomic::Relaxed),
2457 dictionaries: self.opened.load(Atomic::Relaxed),
2458 }
2459 }
2460
2461 #[must_use]
2466 pub fn layout(&self) -> Layout {
2467 let table = &self.table;
2468 let stripes = table.stripes.as_slice();
2469 let columns = table
2470 .fields
2471 .iter()
2472 .enumerate()
2473 .map(|(at, field)| ColumnLayout {
2474 name: field.name.clone(),
2475 kind: field.ty.to_string(),
2476 pages: sum(stripes.iter().map(|stripe| span_bytes(&stripe.pages, at))),
2477 memberships: sum(stripes.iter().map(|stripe| page_bytes(&stripe.memberships, at))),
2478 sieves: sum(stripes.iter().map(|stripe| page_bytes(&stripe.sieves, at))),
2479 part_ranges: sum(stripes.iter().map(|stripe| page_bytes(&stripe.part_ranges, at))),
2480 dictionary: page_bytes(&table.dictionaries, at),
2481 })
2482 .collect();
2483 Layout {
2484 file: self.size,
2485 rows: table.rows,
2486 stripes: stripes.len(),
2487 parts: self.places.len(),
2488 columns,
2489 indexes: sum(stripes.iter().map(|stripe| u64::from(stripe.index.length))),
2490 directory: self.directory,
2491 header: HEADER,
2492 }
2493 }
2494
2495 #[must_use]
2497 pub fn parts(&self) -> usize {
2498 self.places.len()
2499 }
2500
2501 #[must_use]
2508 pub fn stripe_parts(&self) -> Vec<std::ops::Range<usize>> {
2509 let mut runs = Vec::with_capacity(self.table.stripes.len());
2510 let mut start = 0;
2511 for stripe in &self.table.stripes {
2512 let end = start + stripe.parts.len();
2513 runs.push(start..end);
2514 start = end;
2515 }
2516 runs
2517 }
2518
2519 #[must_use]
2524 pub fn stripe_rows(&self, stripe: usize) -> usize {
2525 self.table.stripes.get(stripe).map_or(0, |held| held.rows)
2526 }
2527
2528 pub fn keep_stripes(&self, stripes: usize) {
2535 self.kept.fetch_max(stripes, Atomic::Relaxed);
2536 }
2537
2538 #[must_use]
2540 pub fn part_rows(&self, at: usize) -> usize {
2541 self.places.get(at).map_or(0, |place| place.rows as usize)
2542 }
2543
2544 #[must_use]
2546 pub fn table(&self) -> &Table {
2547 &self.table
2548 }
2549
2550 pub fn top_frequencies(&self, column: usize, top: usize) -> Result<Option<Vec<(Value, u64)>>> {
2559 let field = self
2560 .table
2561 .fields
2562 .get(column)
2563 .ok_or_else(|| invalid("frequency column index out of range"))?;
2564 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2565 return Ok(None);
2566 };
2567 if top == 0 || summary.entries.len() < top {
2568 return Ok(None);
2569 }
2570 let boundary = summary.entries[top - 1].count;
2571 if boundary <= summary.omitted_max {
2572 return Ok(None);
2573 }
2574 self.decode_frequencies(column, &field.ty, &summary.entries).map(Some)
2575 }
2576
2577 pub fn exact_frequencies(&self, column: usize) -> Result<Option<Vec<(Value, u64)>>> {
2597 let field = self
2598 .table
2599 .fields
2600 .get(column)
2601 .ok_or_else(|| invalid("frequency column index out of range"))?;
2602 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2603 return Ok(None);
2604 };
2605 if summary.omitted_max > 0 {
2606 return Ok(None);
2607 }
2608 self.decode_frequencies(column, &field.ty, &summary.entries).map(Some)
2609 }
2610
2611 fn decode_frequencies(
2613 &self,
2614 column: usize,
2615 ty: &LogicalType,
2616 entries: &[FrequencyEntry],
2617 ) -> Result<Vec<(Value, u64)>> {
2618 let dictionary = if *ty == LogicalType::Varchar { self.dictionary(column)? } else { None };
2619 let mut out = Vec::with_capacity(entries.len());
2620 for entry in entries {
2621 let value = match entry.value {
2622 FrequencyValue::Null => Value::Null,
2623 FrequencyValue::Integer(value) => match *ty {
2624 LogicalType::TinyInt => Value::TinyInt(
2625 i8::try_from(value)
2626 .map_err(|_| invalid("frequency TINYINT is out of range"))?,
2627 ),
2628 LogicalType::UTinyInt => Value::UTinyInt(
2629 u8::try_from(value)
2630 .map_err(|_| invalid("frequency UTINYINT is out of range"))?,
2631 ),
2632 LogicalType::USmallInt => Value::USmallInt(
2633 u16::try_from(value)
2634 .map_err(|_| invalid("frequency USMALLINT is out of range"))?,
2635 ),
2636 LogicalType::UInteger => Value::UInteger(
2637 u32::try_from(value)
2638 .map_err(|_| invalid("frequency UINTEGER is out of range"))?,
2639 ),
2640 LogicalType::UBigInt => Value::UBigInt(
2641 u64::try_from(value)
2642 .map_err(|_| invalid("frequency UBIGINT is out of range"))?,
2643 ),
2644 LogicalType::SmallInt => Value::SmallInt(
2645 i16::try_from(value)
2646 .map_err(|_| invalid("frequency SMALLINT is out of range"))?,
2647 ),
2648 LogicalType::Integer => Value::Integer(
2649 i32::try_from(value)
2650 .map_err(|_| invalid("frequency INTEGER is out of range"))?,
2651 ),
2652 LogicalType::BigInt => Value::BigInt(
2653 i64::try_from(value)
2654 .map_err(|_| invalid("frequency BIGINT is out of range"))?,
2655 ),
2656 LogicalType::Date => Value::Date(
2657 i32::try_from(value)
2658 .map_err(|_| invalid("frequency DATE is out of range"))?,
2659 ),
2660 LogicalType::Timestamp => Value::Timestamp(
2661 i64::try_from(value)
2662 .map_err(|_| invalid("frequency TIMESTAMP is out of range"))?,
2663 ),
2664 _ => return Err(invalid("integer frequency belongs to another type")),
2665 },
2666 FrequencyValue::Code(code) => dictionary
2667 .as_ref()
2668 .ok_or_else(|| invalid("frequency code has no dictionary"))?
2669 .try_value_at(code as usize)?,
2670 };
2671 out.push((value, entry.count));
2672 }
2673 Ok(out)
2674 }
2675
2676 pub fn frequency_occurrences(&self, column: usize) -> Result<Option<FrequencyOccurrences>> {
2686 self.table
2687 .fields
2688 .get(column)
2689 .ok_or_else(|| invalid("frequency column index out of range"))?;
2690 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2691 return Ok(None);
2692 };
2693 if summary.ordinals.is_empty() {
2694 return Ok(None);
2695 }
2696 Ok(Some(FrequencyOccurrences {
2697 omitted_max: summary.omitted_max,
2698 ordinals: summary.ordinals.clone(),
2699 }))
2700 }
2701
2702 pub fn distinct_values(&self, column: usize) -> Result<Option<u64>> {
2726 self.table
2727 .distincts
2728 .get(column)
2729 .copied()
2730 .ok_or_else(|| invalid("distinct column index out of range"))
2731 }
2732
2733 pub fn null_count(&self, column: usize) -> Result<u64> {
2744 if column >= self.table.fields.len() {
2745 return Err(invalid("null count column index out of range"));
2746 }
2747 let mut nulls = 0_u64;
2748 for stripe in &self.table.stripes {
2749 let range = stripe
2750 .zone
2751 .column(column)
2752 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
2753 nulls = nulls
2754 .checked_add(range.nulls as u64)
2755 .ok_or_else(|| invalid("null count overflow"))?;
2756 }
2757 Ok(nulls)
2758 }
2759
2760 pub fn text_extremes(&self, column: usize) -> Result<Option<(Value, Value)>> {
2775 if self.null_count(column)? > 0 {
2776 return Ok(None);
2777 }
2778 let Some(dictionary) = self.dictionary(column)? else { return Ok(None) };
2779 let Some(ranks) = dictionary.ranks() else { return Ok(None) };
2780 if ranks == 0 {
2781 return Ok(None);
2782 }
2783 let low = text_at_rank(&dictionary, 0)?;
2784 let high = text_at_rank(&dictionary, ranks - 1)?;
2785 Ok(Some((low, high)))
2786 }
2787
2788 pub fn exact_extremes(&self, column: usize) -> Result<Option<(Bound, Bound)>> {
2811 if column >= self.table.fields.len() {
2812 return Err(invalid("extremes column index out of range"));
2813 }
2814 let mut low: Option<Bound> = None;
2815 let mut high: Option<Bound> = None;
2816 for stripe in &self.table.stripes {
2817 let range = stripe
2818 .zone
2819 .column(column)
2820 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
2821 if !range.exact {
2822 return Ok(None);
2823 }
2824 let (Some(small), Some(large)) = (range.low.as_ref(), range.high.as_ref()) else {
2829 if stripe.rows > range.nulls {
2830 return Ok(None);
2831 }
2832 continue;
2833 };
2834 low = Some(low.map_or_else(|| small.clone(), |held| held.smaller(small.clone())));
2835 high = Some(high.map_or_else(|| large.clone(), |held| held.larger(large.clone())));
2836 }
2837 Ok(low.zip(high))
2838 }
2839
2840 pub fn exact_sum(&self, column: usize) -> Result<Option<(i128, u64)>> {
2853 if column >= self.table.fields.len() {
2854 return Err(invalid("sum column index out of range"));
2855 }
2856 let mut total = 0_i128;
2857 let mut rows = 0_u64;
2858 for stripe in &self.table.stripes {
2859 let range = stripe
2860 .zone
2861 .column(column)
2862 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
2863 let Some(part) = range.sum else { return Ok(None) };
2864 let Some(sum) = total.checked_add(part) else { return Ok(None) };
2865 total = sum;
2866 rows = rows.saturating_add(stripe.rows as u64 - range.nulls as u64);
2867 }
2868 Ok(Some((total, rows)))
2869 }
2870
2871 fn dictionary(&self, column: usize) -> Result<Option<Arc<Vector>>> {
2880 let Some(page) = self.table.dictionaries[column] else { return Ok(None) };
2881 if let Some(dictionary) = self.dictionaries[column].get() {
2882 return Ok(Some(Arc::clone(dictionary)));
2883 }
2884 let _queued = self.loading[column].lock().map_err(|_| invalid("a poisoned dictionary"))?;
2885 if let Some(dictionary) = self.dictionaries[column].get() {
2886 return Ok(Some(Arc::clone(dictionary)));
2887 }
2888 self.opened.fetch_add(1, Atomic::Relaxed);
2889 let dictionary = Arc::new(open_global_dictionary(
2890 Arc::clone(&self.file),
2891 page,
2892 &self.table.fields[column].ty,
2893 TEXT_KEEP_BUDGET,
2894 )?);
2895 let _ = self.dictionaries[column].set(Arc::clone(&dictionary));
2896 Ok(Some(dictionary))
2897 }
2898
2899 pub fn read(&self, part: usize, columns: &[usize]) -> Result<Chunk> {
2908 self.read_impl(part, columns, true)
2909 }
2910
2911 pub fn read_sparse(&self, part: usize, columns: &[usize]) -> Result<Chunk> {
2921 self.read_impl(part, columns, false)
2922 }
2923
2924 pub fn skips_codes(&self, part: usize, column: usize, candidates: &[u32]) -> Result<bool> {
2931 if candidates.is_empty() {
2932 return Ok(true);
2933 }
2934 if candidates.windows(2).any(|pair| pair[0] >= pair[1]) {
2935 return Err(Error::internal("native code candidates are not sorted and unique"));
2936 }
2937 let stripe = self.stripe_of(part)?;
2938 let Some(page) = stripe.memberships.get(column).copied().flatten() else {
2939 return Ok(false);
2940 };
2941 let mut bytes = vec![0; page.length as usize];
2942 read_at(&self.file, page.offset, &mut bytes)?;
2943 if checksum(&bytes) != page.hash {
2944 return Err(invalid("membership page checksum differs"));
2945 }
2946 let codes = decode_membership(&bytes)?;
2947 let mut left = 0;
2948 let mut right = 0;
2949 while left < codes.len() && right < candidates.len() {
2950 match codes[left].cmp(&candidates[right]) {
2951 Ordering::Less => left += 1,
2952 Ordering::Greater => right += 1,
2953 Ordering::Equal => return Ok(false),
2954 }
2955 }
2956 Ok(true)
2957 }
2958
2959 fn stripe_of(&self, part: usize) -> Result<&Stripe> {
2960 let place = self.places.get(part).ok_or_else(|| invalid("part index out of range"))?;
2961 self.table
2962 .stripes
2963 .get(place.stripe as usize)
2964 .ok_or_else(|| invalid("stripe index out of range"))
2965 }
2966
2967 fn held(&self, at: usize, stripe: &Stripe, column: usize, whole: bool) -> Result<CachedColumn> {
2984 let cache = self.cache.get(column).ok_or_else(|| invalid("column index out of range"))?;
2985 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
2986 let known = cached.index.get(at).and_then(Clone::clone);
2987 let page = cached.pages.get(at).and_then(Clone::clone);
2988 if let Some(index) = known.clone() {
2989 if !whole || page.is_some() {
2990 return Ok(CachedColumn { stripe: at, index, page });
2991 }
2992 }
2993 if cached.loading.contains(&at) {
2994 drop(cached);
2995 if let Some(index) = known {
2999 return Ok(CachedColumn { stripe: at, index, page: None });
3000 }
3001 let held = self.page_of(stripe, column, at, false, None)?;
3002 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3003 remember(&mut cached, &held, self.kept.load(Atomic::Relaxed));
3004 return Ok(held);
3005 }
3006 cached.loading.push(at);
3007 drop(cached);
3008
3009 let read = self.page_of(stripe, column, at, whole, known);
3010
3011 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3015 if let Some(position) = cached.loading.iter().position(|loading| *loading == at) {
3016 cached.loading.remove(position);
3017 }
3018 let held = read?;
3019 remember(&mut cached, &held, self.kept.load(Atomic::Relaxed));
3020 Ok(held)
3021 }
3022
3023 fn page_of(
3029 &self,
3030 stripe: &Stripe,
3031 column: usize,
3032 at: usize,
3033 whole: bool,
3034 known: Option<Arc<Vec<PartSpan>>>,
3035 ) -> Result<CachedColumn> {
3036 let index = match known {
3037 Some(index) => index,
3038 None => {
3039 self.indexes.fetch_add(1, Atomic::Relaxed);
3040 Arc::new(read_index(&self.file, stripe, column)?)
3041 }
3042 };
3043 let page = if whole {
3044 self.pages.fetch_add(1, Atomic::Relaxed);
3045 let span = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
3046 let mut bytes = vec![0; span.length as usize];
3047 read_at(&self.file, span.offset, &mut bytes)?;
3048 Some(Arc::new(bytes))
3049 } else {
3050 None
3051 };
3052 Ok(CachedColumn { stripe: at, index, page })
3053 }
3054
3055 fn read_impl(&self, at: usize, columns: &[usize], whole: bool) -> Result<Chunk> {
3056 let place = *self.places.get(at).ok_or_else(|| invalid("part index out of range"))?;
3057 let index = place.stripe as usize;
3058 let stripe =
3059 self.table.stripes.get(index).ok_or_else(|| invalid("stripe index out of range"))?;
3060 let rows = place.rows as usize;
3061 let mut picked = Vec::with_capacity(columns.len());
3062 for &column in columns {
3063 let field = self
3064 .table
3065 .fields
3066 .get(column)
3067 .ok_or_else(|| invalid("column index out of range"))?;
3068 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
3069 let held = self.held(index, stripe, column, whole)?;
3070 let span = *held
3071 .index
3072 .get(place.part as usize)
3073 .ok_or_else(|| invalid("part index out of range"))?;
3074 let owned;
3075 let bytes = match &held.page {
3076 Some(held) => part_bytes(held, span)?,
3077 None => {
3078 let offset = page
3079 .offset
3080 .checked_add(span.start as u64)
3081 .ok_or_else(|| invalid("part range overflow"))?;
3082 let mut bytes = vec![0; span.length];
3083 read_at(&self.file, offset, &mut bytes)?;
3084 owned = bytes;
3085 &owned
3086 }
3087 };
3088 if checksum(bytes) != span.hash {
3089 return Err(invalid(&format!(
3090 "column page checksum differs, column {column} part {} at {}+{} of {} bytes, \
3091 wanted {:016x} and got {:016x}",
3092 place.part,
3093 page.offset,
3094 span.start,
3095 span.length,
3096 span.hash,
3097 checksum(bytes),
3098 )));
3099 }
3100 let dictionary = self.dictionary(column)?;
3101 picked.push(decode(&field.ty, rows, bytes, dictionary)?);
3102 }
3103 Chunk::with_rows(picked, rows)
3104 }
3105
3106 #[must_use]
3122 pub fn skips(&self, part: usize, probes: &[Probe]) -> bool {
3123 let Some(place) = self.places.get(part).copied() else { return false };
3124 let Some(stripe) = self.table.stripes.get(place.stripe as usize) else { return false };
3125 if stripe.zone.skips(probes) {
3126 return true;
3127 }
3128 probes.iter().any(|probe| self.outside(place, probe) || self.sifted(place, probe))
3129 }
3130
3131 fn outside(&self, place: Place, probe: &Probe) -> bool {
3137 match self.stripe_part_ranges(place.stripe as usize, probe.column) {
3138 Some(ranges) => ranges
3139 .get(place.part as usize)
3140 .is_some_and(|range| range.excludes(probe.op, &probe.value)),
3141 None => false,
3142 }
3143 }
3144
3145 fn stripe_part_ranges(&self, stripe: usize, column: usize) -> Option<&[Range]> {
3151 let slot = self.part_ranges.get(column)?.get(stripe)?;
3152 if let Some(held) = slot.get() {
3153 return Some(held);
3154 }
3155 let page = self.table.stripes.get(stripe)?.part_ranges.get(column).copied().flatten()?;
3156 let mut bytes = vec![0; page.length as usize];
3157 read_at(&self.file, page.offset, &mut bytes).ok()?;
3158 if checksum(&bytes) != page.hash {
3159 return None;
3160 }
3161 let ranges = Arc::new(decode_part_ranges(&bytes).ok()?);
3162 let _ = slot.set(ranges);
3163 slot.get().map(|held| held.as_slice())
3164 }
3165
3166 #[must_use]
3183 pub fn certain(&self, part: usize, probes: &[Probe]) -> bool {
3184 let Some(place) = self.places.get(part).copied() else { return false };
3185 let Some(stripe) = self.table.stripes.get(place.stripe as usize) else { return false };
3186 if stripe.zone.certain(probes) {
3187 return true;
3188 }
3189 probes
3190 .iter()
3191 .all(|probe| stripe.zone.certain(slice::from_ref(probe)) || self.inside(place, probe))
3192 }
3193
3194 fn inside(&self, place: Place, probe: &Probe) -> bool {
3200 match self.stripe_part_ranges(place.stripe as usize, probe.column) {
3201 Some(ranges) => ranges
3202 .get(place.part as usize)
3203 .is_some_and(|range| range.certain(probe.op, &probe.value)),
3204 None => false,
3205 }
3206 }
3207
3208 #[must_use]
3219 pub fn stripe_skips(&self, stripe: usize, probes: &[Probe]) -> bool {
3220 self.table.stripes.get(stripe).is_some_and(|held| held.zone.skips(probes))
3221 }
3222
3223 fn sifted(&self, place: Place, probe: &Probe) -> bool {
3229 if probe.op != Op::Equal {
3230 return false;
3231 }
3232 match self.stripe_sieves(place.stripe as usize, probe.column) {
3233 Some(sieves) => sieves
3234 .get(place.part as usize)
3235 .and_then(Option::as_ref)
3236 .is_some_and(|sieve| sieve.excludes(&probe.value)),
3237 None => false,
3238 }
3239 }
3240
3241 fn stripe_sieves(&self, stripe: usize, column: usize) -> Option<&[Option<Sieve>]> {
3248 let slot = self.sieves.get(column)?.get(stripe)?;
3249 if let Some(held) = slot.get() {
3250 return Some(held);
3251 }
3252 let page = self.table.stripes.get(stripe)?.sieves.get(column).copied().flatten()?;
3253 let mut bytes = vec![0; page.length as usize];
3254 read_at(&self.file, page.offset, &mut bytes).ok()?;
3255 if checksum(&bytes) != page.hash {
3256 return None;
3257 }
3258 let sieves = Arc::new(decode_sieves(&bytes).ok()?);
3259 let _ = slot.set(sieves);
3260 slot.get().map(|held| held.as_slice())
3261 }
3262}
3263
3264fn text_at_rank(dictionary: &Vector, rank: usize) -> Result<Value> {
3266 let code = dictionary.code_at_rank(rank)? as usize;
3267 let text = dictionary
3268 .try_text_at(code)?
3269 .ok_or_else(|| invalid("global dictionary order names a code it does not have"))?;
3270 Ok(Value::Varchar(text.into()))
3271}
3272
3273#[cfg(unix)]
3278fn write_at(file: &File, mut offset: u64, mut bytes: &[u8]) -> Result<()> {
3279 use std::os::unix::fs::FileExt;
3280 while !bytes.is_empty() {
3281 let written = file.write_at(bytes, offset).map_err(io)?;
3282 if written == 0 {
3283 return Err(invalid("a write to the native file wrote nothing"));
3284 }
3285 offset += written as u64;
3286 bytes = &bytes[written..];
3287 }
3288 Ok(())
3289}
3290
3291#[cfg(windows)]
3293fn write_at(file: &File, mut offset: u64, mut bytes: &[u8]) -> Result<()> {
3294 use std::os::windows::fs::FileExt;
3295 while !bytes.is_empty() {
3296 let written = file.seek_write(bytes, offset).map_err(io)?;
3297 if written == 0 {
3298 return Err(invalid("a write to the native file wrote nothing"));
3299 }
3300 offset += written as u64;
3301 bytes = &bytes[written..];
3302 }
3303 Ok(())
3304}
3305
3306#[cfg(not(any(unix, windows)))]
3308fn write_at(file: &File, offset: u64, bytes: &[u8]) -> Result<()> {
3309 use std::io::Write;
3310 let mut file = file.try_clone().map_err(io)?;
3311 file.seek(SeekFrom::Start(offset)).map_err(io)?;
3312 file.write_all(bytes).map_err(io)
3313}
3314
3315#[cfg(unix)]
3325fn read_at(file: &File, mut offset: u64, mut bytes: &mut [u8]) -> Result<()> {
3326 use std::os::unix::fs::FileExt;
3327 while !bytes.is_empty() {
3328 let read = file.read_at(bytes, offset).map_err(io)?;
3329 if read == 0 {
3330 return Err(invalid("column page ends before its declared length"));
3331 }
3332 offset += read as u64;
3333 bytes = &mut bytes[read..];
3334 }
3335 Ok(())
3336}
3337
3338#[cfg(windows)]
3344fn read_at(file: &File, mut offset: u64, mut bytes: &mut [u8]) -> Result<()> {
3345 use std::os::windows::fs::FileExt;
3346 while !bytes.is_empty() {
3347 let read = file.seek_read(bytes, offset).map_err(io)?;
3348 if read == 0 {
3349 return Err(invalid("column page ends before its declared length"));
3350 }
3351 offset += read as u64;
3352 bytes = &mut bytes[read..];
3353 }
3354 Ok(())
3355}
3356
3357#[cfg(not(any(unix, windows)))]
3362fn read_at(file: &File, offset: u64, bytes: &mut [u8]) -> Result<()> {
3363 let mut file = file.try_clone().map_err(io)?;
3364 file.seek(SeekFrom::Start(offset)).map_err(io)?;
3365 file.read_exact(bytes).map_err(io)
3366}
3367
3368fn type_tag(ty: &LogicalType) -> Result<u8> {
3369 match ty {
3370 LogicalType::SmallInt => Ok(1),
3371 LogicalType::Integer => Ok(2),
3372 LogicalType::BigInt => Ok(3),
3373 LogicalType::Varchar => Ok(4),
3374 LogicalType::Date => Ok(5),
3375 LogicalType::Timestamp => Ok(6),
3376 LogicalType::Boolean => Ok(7),
3377 LogicalType::TinyInt => Ok(8),
3378 LogicalType::UTinyInt => Ok(9),
3379 LogicalType::USmallInt => Ok(10),
3380 LogicalType::UInteger => Ok(11),
3381 LogicalType::UBigInt => Ok(12),
3382 _ => Err(Error::not_implemented(format!("native storage for {ty}"))),
3383 }
3384}
3385
3386fn tag_type(tag: u8) -> Result<LogicalType> {
3387 match tag {
3388 1 => Ok(LogicalType::SmallInt),
3389 2 => Ok(LogicalType::Integer),
3390 3 => Ok(LogicalType::BigInt),
3391 4 => Ok(LogicalType::Varchar),
3392 5 => Ok(LogicalType::Date),
3393 6 => Ok(LogicalType::Timestamp),
3394 7 => Ok(LogicalType::Boolean),
3395 8 => Ok(LogicalType::TinyInt),
3396 9 => Ok(LogicalType::UTinyInt),
3397 10 => Ok(LogicalType::USmallInt),
3398 11 => Ok(LogicalType::UInteger),
3399 12 => Ok(LogicalType::UBigInt),
3400 _ => Err(invalid("column type tag is unknown")),
3401 }
3402}
3403
3404fn put_u16(out: &mut Vec<u8>, value: u16) {
3405 out.extend_from_slice(&value.to_le_bytes());
3406}
3407fn put_u32(out: &mut Vec<u8>, value: u32) {
3408 out.extend_from_slice(&value.to_le_bytes());
3409}
3410fn put_u64(out: &mut Vec<u8>, value: u64) {
3411 out.extend_from_slice(&value.to_le_bytes());
3412}
3413fn put_var_u64(out: &mut Vec<u8>, mut value: u64) {
3414 while value >= 0x80 {
3415 out.push((value as u8 & 0x7f) | 0x80);
3416 value >>= 7;
3417 }
3418 out.push(value as u8);
3419}
3420
3421fn frequency_order(left: FrequencyValue, right: FrequencyValue) -> Ordering {
3422 match (left, right) {
3423 (FrequencyValue::Null, FrequencyValue::Null) => Ordering::Equal,
3424 (FrequencyValue::Null, _) => Ordering::Less,
3425 (_, FrequencyValue::Null) => Ordering::Greater,
3426 (FrequencyValue::Integer(left), FrequencyValue::Integer(right)) => left.cmp(&right),
3427 (FrequencyValue::Code(left), FrequencyValue::Code(right)) => left.cmp(&right),
3428 (FrequencyValue::Integer(_), FrequencyValue::Code(_)) => Ordering::Less,
3429 (FrequencyValue::Code(_), FrequencyValue::Integer(_)) => Ordering::Greater,
3430 }
3431}
3432
3433fn code_frequency(dictionary: &GlobalDictionary) -> FrequencySummary {
3434 let mut entries = dictionary
3435 .counts
3436 .iter()
3437 .enumerate()
3438 .filter(|(_, count)| **count != 0)
3439 .map(|(code, &count)| FrequencyEntry { value: FrequencyValue::Code(code as u32), count })
3440 .collect::<Vec<_>>();
3441 if dictionary.nulls != 0 {
3442 entries.push(FrequencyEntry { value: FrequencyValue::Null, count: dictionary.nulls });
3443 }
3444 entries.sort_unstable_by(|left, right| {
3445 right.count.cmp(&left.count).then_with(|| frequency_order(left.value, right.value))
3446 });
3447 let omitted_max = entries.get(FREQUENCY_ENTRIES).map_or(0, |entry| entry.count);
3448 entries.truncate(FREQUENCY_ENTRIES);
3449 FrequencySummary { entries, omitted_max, ordinals: Vec::new() }
3450}
3451
3452fn encode_directory(table: &Table) -> Result<Vec<u8>> {
3453 let mut out = DIRECTORY.to_vec();
3454 let name = table.name.as_bytes();
3455 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("table name too long"))?);
3456 out.extend_from_slice(name);
3457 put_u16(&mut out, u16::try_from(table.fields.len()).map_err(|_| invalid("too many columns"))?);
3458 for field in &table.fields {
3459 let name = field.name.as_bytes();
3460 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("column name too long"))?);
3461 out.extend_from_slice(name);
3462 out.push(type_tag(&field.ty)?);
3463 out.push(u8::from(field.not_null));
3464 }
3465 for dictionary in &table.dictionaries {
3466 match dictionary {
3467 None => out.push(0),
3468 Some(page) => {
3469 out.push(1);
3470 put_u64(&mut out, page.offset);
3471 put_u32(&mut out, page.length);
3472 put_u64(&mut out, page.hash);
3473 }
3474 }
3475 }
3476 for distinct in &table.distincts {
3477 match distinct {
3478 None => out.push(0),
3479 Some(count) => {
3480 out.push(1);
3481 put_u64(&mut out, *count);
3482 }
3483 }
3484 }
3485 put_u64(&mut out, u64::try_from(table.rows).map_err(|_| invalid("row count overflow"))?);
3486 put_u32(&mut out, u32::try_from(table.stripes.len()).map_err(|_| invalid("too many stripes"))?);
3487 for stripe in &table.stripes {
3488 put_u32(
3489 &mut out,
3490 u32::try_from(stripe.parts.len()).map_err(|_| invalid("too many parts in a stripe"))?,
3491 );
3492 for &rows in &stripe.parts {
3493 put_u32(&mut out, rows);
3494 }
3495 put_u64(&mut out, stripe.index.offset);
3496 put_u32(&mut out, stripe.index.length);
3497 for page in &stripe.pages {
3498 put_u64(&mut out, page.offset);
3499 put_u32(&mut out, page.length);
3500 }
3501 for (field, membership) in table.fields.iter().zip(&stripe.memberships) {
3502 if field.ty != LogicalType::Varchar {
3503 continue;
3504 }
3505 let page =
3506 membership.ok_or_else(|| invalid("string page has no code membership index"))?;
3507 put_u64(&mut out, page.offset);
3508 put_u32(&mut out, page.length);
3509 put_u64(&mut out, page.hash);
3510 }
3511 for sieve in &stripe.sieves {
3512 match sieve {
3513 None => out.push(0),
3514 Some(page) => {
3515 out.push(1);
3516 put_u64(&mut out, page.offset);
3517 put_u32(&mut out, page.length);
3518 put_u64(&mut out, page.hash);
3519 }
3520 }
3521 }
3522 for held in &stripe.part_ranges {
3523 match held {
3524 None => out.push(0),
3525 Some(page) => {
3526 out.push(1);
3527 put_u64(&mut out, page.offset);
3528 put_u32(&mut out, page.length);
3529 put_u64(&mut out, page.hash);
3530 }
3531 }
3532 }
3533 for range in stripe.zone.columns() {
3534 put_bound(&mut out, range.low.as_ref())?;
3535 put_bound(&mut out, range.high.as_ref())?;
3536 put_u32(
3537 &mut out,
3538 u32::try_from(range.nulls).map_err(|_| invalid("null count overflow"))?,
3539 );
3540 out.push(u8::from(range.exact));
3541 match range.sum {
3542 None => out.push(0),
3543 Some(total) => {
3544 out.push(1);
3545 out.extend_from_slice(&total.to_le_bytes());
3546 }
3547 }
3548 }
3549 }
3550 out.extend_from_slice(FREQUENCIES);
3551 put_u16(
3552 &mut out,
3553 u16::try_from(table.frequencies.len())
3554 .map_err(|_| invalid("too many frequency columns"))?,
3555 );
3556 for summary in &table.frequencies {
3557 let Some(summary) = summary else {
3558 out.push(0);
3559 continue;
3560 };
3561 out.push(1);
3562 put_u64(&mut out, summary.omitted_max);
3563 put_u32(
3564 &mut out,
3565 u32::try_from(summary.entries.len())
3566 .map_err(|_| invalid("too many frequency entries"))?,
3567 );
3568 for entry in &summary.entries {
3569 match entry.value {
3570 FrequencyValue::Null => out.push(0),
3571 FrequencyValue::Integer(value) => {
3572 out.push(1);
3573 out.extend_from_slice(&value.to_le_bytes());
3574 }
3575 FrequencyValue::Code(value) => {
3576 out.push(2);
3577 put_u32(&mut out, value);
3578 }
3579 }
3580 put_u64(&mut out, entry.count);
3581 }
3582 put_u32(
3583 &mut out,
3584 u32::try_from(summary.ordinals.len())
3585 .map_err(|_| invalid("too many frequency ordinals"))?,
3586 );
3587 let mut previous = 0_u64;
3588 for (at, &ordinal) in summary.ordinals.iter().enumerate() {
3589 let delta = if at == 0 {
3590 ordinal
3591 } else {
3592 ordinal
3593 .checked_sub(previous)
3594 .ok_or_else(|| invalid("frequency ordinals are not ordered"))?
3595 };
3596 if at != 0 && delta == 0 {
3597 return Err(invalid("frequency ordinals are not unique"));
3598 }
3599 put_var_u64(&mut out, delta);
3600 previous = ordinal;
3601 }
3602 }
3603 Ok(out)
3604}
3605
3606fn encode_catalog(entries: &[Entry]) -> Result<Vec<u8>> {
3612 let mut out = CATALOG.to_vec();
3613 put_u32(&mut out, u32::try_from(entries.len()).map_err(|_| invalid("too many tables"))?);
3614 for entry in entries {
3615 let name = entry.name.as_bytes();
3616 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("table name too long"))?);
3617 out.extend_from_slice(name);
3618 put_u64(&mut out, u64::try_from(entry.rows).map_err(|_| invalid("row count overflow"))?);
3619 put_u16(
3620 &mut out,
3621 u16::try_from(entry.fields.len()).map_err(|_| invalid("too many columns"))?,
3622 );
3623 for field in &entry.fields {
3624 let name = field.name.as_bytes();
3625 put_u16(
3626 &mut out,
3627 u16::try_from(name.len()).map_err(|_| invalid("column name too long"))?,
3628 );
3629 out.extend_from_slice(name);
3630 out.push(type_tag(&field.ty)?);
3631 out.push(u8::from(field.not_null));
3632 }
3633 put_u64(&mut out, entry.directory.offset);
3634 put_u32(&mut out, entry.directory.length);
3635 put_u64(&mut out, entry.directory.hash);
3636 }
3637 Ok(out)
3638}
3639
3640fn decode_catalog(bytes: &[u8], size: u64) -> Result<Vec<Entry>> {
3643 let mut cur = Cursor { bytes, at: 0 };
3644 if cur.take(8)? != CATALOG {
3645 return Err(invalid("catalog magic differs"));
3646 }
3647 let count = cur.u32()? as usize;
3648 let mut entries: Vec<Entry> = Vec::with_capacity(count.min(1024));
3649 for _ in 0..count {
3650 let name = cur.text()?;
3651 let rows = usize::try_from(cur.u64()?).map_err(|_| invalid("row count does not fit"))?;
3652 let width = cur.u16()? as usize;
3653 let mut fields = Vec::with_capacity(width);
3654 for _ in 0..width {
3655 let name = cur.text()?;
3656 let ty = tag_type(cur.u8()?)?;
3657 let not_null = match cur.u8()? {
3658 0 => false,
3659 1 => true,
3660 _ => return Err(invalid("nullability flag differs")),
3661 };
3662 fields.push(Field { name, ty, not_null });
3663 }
3664 let directory = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3665 let end = directory
3666 .offset
3667 .checked_add(u64::from(directory.length))
3668 .ok_or_else(|| invalid("table directory offset overflow"))?;
3669 if directory.offset < HEADER
3670 || end > size
3671 || directory.length as usize > MAX_DIRECTORY
3672 || directory.length == 0
3673 {
3674 return Err(invalid("table directory range is outside the file"));
3675 }
3676 if entries.iter().any(|held| held.name == name) {
3677 return Err(invalid("two tables in the catalog have the same name"));
3678 }
3679 entries.push(Entry { name, fields, rows, directory });
3680 }
3681 Ok(entries)
3682}
3683
3684struct Cursor<'a> {
3685 bytes: &'a [u8],
3686 at: usize,
3687}
3688impl<'a> Cursor<'a> {
3689 fn take(&mut self, len: usize) -> Result<&'a [u8]> {
3690 let end = self.at.checked_add(len).ok_or_else(|| invalid("directory offset overflow"))?;
3691 let bytes =
3692 self.bytes.get(self.at..end).ok_or_else(|| invalid("directory is truncated"))?;
3693 self.at = end;
3694 Ok(bytes)
3695 }
3696 fn u8(&mut self) -> Result<u8> {
3697 Ok(self.take(1)?[0])
3698 }
3699 fn u16(&mut self) -> Result<u16> {
3700 Ok(u16::from_le_bytes(self.take(2)?.try_into().expect("two bytes")))
3701 }
3702 fn u32(&mut self) -> Result<u32> {
3703 Ok(u32::from_le_bytes(self.take(4)?.try_into().expect("four bytes")))
3704 }
3705 fn u64(&mut self) -> Result<u64> {
3706 Ok(u64::from_le_bytes(self.take(8)?.try_into().expect("eight bytes")))
3707 }
3708 fn var_u64(&mut self) -> Result<u64> {
3709 let mut value = 0_u64;
3710 for shift in (0..=63).step_by(7) {
3711 let byte = self.u8()?;
3712 let part = u64::from(byte & 0x7f);
3713 if shift == 63 && part > 1 {
3714 return Err(invalid("frequency ordinal varint overflows"));
3715 }
3716 value |= part << shift;
3717 if byte & 0x80 == 0 {
3718 return Ok(value);
3719 }
3720 }
3721 Err(invalid("frequency ordinal varint is too long"))
3722 }
3723 fn bound(&mut self) -> Result<Option<Bound>> {
3724 Ok(match self.u8()? {
3725 0 => None,
3726 1 => Some(Bound::Int(i128::from_le_bytes(
3727 self.take(16)?.try_into().expect("sixteen bytes"),
3728 ))),
3729 2 => Some(Bound::Real(f64::from_le_bytes(
3730 self.take(8)?.try_into().expect("eight bytes"),
3731 ))),
3732 3 => {
3733 let length = self.u32()? as usize;
3734 Some(Bound::Bytes(self.take(length)?.to_vec()))
3735 }
3736 4 => {
3737 let unscaled =
3738 i128::from_le_bytes(self.take(16)?.try_into().expect("sixteen bytes"));
3739 Some(Bound::Scaled { unscaled, scale: self.u8()? })
3740 }
3741 _ => return Err(invalid("bound tag differs")),
3742 })
3743 }
3744 fn text(&mut self) -> Result<String> {
3745 let len = self.u16()? as usize;
3746 String::from_utf8(self.take(len)?.to_vec()).map_err(|_| invalid("name is not UTF-8"))
3747 }
3748}
3749
3750fn decode_directory(bytes: &[u8], size: u64) -> Result<Table> {
3751 let mut cur = Cursor { bytes, at: 0 };
3752 if cur.take(8)? != DIRECTORY {
3753 return Err(invalid("directory magic differs"));
3754 }
3755 let name = cur.text()?;
3756 let width = cur.u16()? as usize;
3757 let mut fields = Vec::with_capacity(width);
3758 for _ in 0..width {
3759 let name = cur.text()?;
3760 let ty = tag_type(cur.u8()?)?;
3761 let not_null = match cur.u8()? {
3762 0 => false,
3763 1 => true,
3764 _ => return Err(invalid("nullability flag differs")),
3765 };
3766 fields.push(Field { name, ty, not_null });
3767 }
3768 let mut dictionaries = Vec::with_capacity(width);
3769 for _ in 0..width {
3770 dictionaries.push(match cur.u8()? {
3771 0 => None,
3772 1 => {
3773 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3774 let end = page
3775 .offset
3776 .checked_add(u64::from(page.length))
3777 .ok_or_else(|| invalid("dictionary page offset overflow"))?;
3778 if page.offset < HEADER || end > size {
3783 return Err(invalid("dictionary page range is outside the file"));
3784 }
3785 Some(page)
3786 }
3787 _ => return Err(invalid("dictionary page tag differs")),
3788 });
3789 }
3790 let mut distincts = Vec::with_capacity(width);
3791 for _ in 0..width {
3792 distincts.push(match cur.u8()? {
3793 0 => None,
3794 1 => Some(cur.u64()?),
3795 _ => return Err(invalid("distinct count tag differs")),
3796 });
3797 }
3798 let rows = usize::try_from(cur.u64()?).map_err(|_| invalid("row count does not fit"))?;
3799 let count = cur.u32()? as usize;
3800 let mut stripes = Vec::with_capacity(count);
3801 let mut total = 0_usize;
3802 for _ in 0..count {
3803 let count = cur.u32()? as usize;
3804 if count == 0 || count > STRIPE_PARTS {
3805 return Err(invalid("stripe part count is outside its bound"));
3806 }
3807 let mut parts = Vec::with_capacity(count);
3808 let mut stripe_rows = 0_usize;
3809 for _ in 0..count {
3810 let rows = cur.u32()?;
3811 if rows == 0 {
3812 return Err(invalid("empty part"));
3813 }
3814 parts.push(rows);
3815 stripe_rows = stripe_rows
3816 .checked_add(rows as usize)
3817 .ok_or_else(|| invalid("stripe row count overflow"))?;
3818 }
3819 total =
3820 total.checked_add(stripe_rows).ok_or_else(|| invalid("stripe row count overflow"))?;
3821 let index = Span { offset: cur.u64()?, length: cur.u32()? };
3822 let section = index_section(count)?;
3823 let wanted = section
3824 .checked_mul(width)
3825 .and_then(|bytes| u32::try_from(bytes).ok())
3826 .ok_or_else(|| invalid("index page length overflow"))?;
3827 let end = index
3828 .offset
3829 .checked_add(u64::from(index.length))
3830 .ok_or_else(|| invalid("index page offset overflow"))?;
3831 if index.offset < HEADER || end > size || index.length != wanted {
3832 return Err(invalid("index page range is outside the file"));
3833 }
3834 let mut pages = Vec::with_capacity(width);
3835 for _ in 0..width {
3836 let offset = cur.u64()?;
3837 let length = cur.u32()?;
3838 let end = offset
3839 .checked_add(u64::from(length))
3840 .ok_or_else(|| invalid("page offset overflow"))?;
3841 if offset < HEADER || end > size || length as usize > MAX_PAGE {
3842 return Err(invalid("page range is outside the file"));
3843 }
3844 pages.push(Span { offset, length });
3845 }
3846 let mut memberships = vec![None; width];
3847 for (column, field) in fields.iter().enumerate() {
3848 if field.ty != LogicalType::Varchar {
3849 continue;
3850 }
3851 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3852 let end = page
3853 .offset
3854 .checked_add(u64::from(page.length))
3855 .ok_or_else(|| invalid("membership page offset overflow"))?;
3856 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
3857 return Err(invalid("membership page range is outside the file"));
3858 }
3859 memberships[column] = Some(page);
3860 }
3861 let mut sieves = vec![None; width];
3862 for sieve in sieves.iter_mut().take(width) {
3863 match cur.u8()? {
3864 0 => continue,
3865 1 => {}
3866 _ => return Err(invalid("a sieve page has an unknown tag")),
3867 }
3868 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3869 let end = page
3870 .offset
3871 .checked_add(u64::from(page.length))
3872 .ok_or_else(|| invalid("sieve page offset overflow"))?;
3873 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
3874 return Err(invalid("sieve page range is outside the file"));
3875 }
3876 *sieve = Some(page);
3877 }
3878 let mut part_ranges = vec![None; width];
3879 for held in part_ranges.iter_mut().take(width) {
3880 match cur.u8()? {
3881 0 => continue,
3882 1 => {}
3883 _ => return Err(invalid("a part range page has an unknown tag")),
3884 }
3885 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3886 let end = page
3887 .offset
3888 .checked_add(u64::from(page.length))
3889 .ok_or_else(|| invalid("part range page offset overflow"))?;
3890 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
3891 return Err(invalid("part range page range is outside the file"));
3892 }
3893 *held = Some(page);
3894 }
3895 let mut ranges = Vec::with_capacity(width);
3896 for column in 0..width {
3897 let low = cur.bound()?;
3898 let high = cur.bound()?;
3899 let nulls = cur.u32()? as usize;
3900 if nulls > stripe_rows {
3901 return Err(invalid("null count exceeds stripe rows"));
3902 }
3903 let exact = cur.u8()? != 0;
3904 let sum = match cur.u8()? {
3905 0 => None,
3906 1 => Some(i128::from_le_bytes(
3907 cur.take(16)?.try_into().map_err(|_| invalid("a stripe sum is truncated"))?,
3908 )),
3909 _ => return Err(invalid("a stripe sum has an unknown tag")),
3910 };
3911 let ty = &fields.get(column).ok_or_else(|| invalid("a stripe range has no column"))?.ty;
3917 let low = low.map(|bound| scaled_as(bound, ty));
3918 let high = high.map(|bound| scaled_as(bound, ty));
3919 ranges.push(Range { low, high, nulls, exact, sum });
3920 }
3921 stripes.push(Stripe {
3922 rows: stripe_rows,
3923 parts,
3924 index,
3925 pages,
3926 memberships,
3927 sieves,
3928 part_ranges,
3929 zone: Zone::from_ranges(ranges),
3930 });
3931 }
3932 if total != rows {
3933 return Err(invalid("table row count differs from stripes"));
3934 }
3935 let frequencies = if cur.at == bytes.len() {
3936 vec![None; width]
3937 } else {
3938 if cur.take(8)? != FREQUENCIES {
3939 return Err(invalid("directory extension magic differs"));
3940 }
3941 if cur.u16()? as usize != width {
3942 return Err(invalid("frequency column count differs"));
3943 }
3944 let mut frequencies = Vec::with_capacity(width);
3945 for field in &fields {
3946 let summary = match cur.u8()? {
3947 0 => None,
3948 1 => {
3949 let omitted_max = cur.u64()?;
3950 let count = cur.u32()? as usize;
3951 if count > FREQUENCY_ENTRIES {
3952 return Err(invalid("frequency entry count exceeds its bound"));
3953 }
3954 let mut entries = Vec::with_capacity(count);
3955 for _ in 0..count {
3957 let value = match cur.u8()? {
3958 0 => FrequencyValue::Null,
3959 1 => FrequencyValue::Integer(i128::from_le_bytes(
3960 cur.take(16)?.try_into().expect("sixteen bytes"),
3961 )),
3962 2 => FrequencyValue::Code(cur.u32()?),
3963 _ => return Err(invalid("frequency value tag differs")),
3964 };
3965 let valid = matches!(
3966 (&field.ty, value),
3967 (_, FrequencyValue::Null)
3968 | (LogicalType::Varchar, FrequencyValue::Code(_))
3969 | (
3970 LogicalType::TinyInt
3971 | LogicalType::SmallInt
3972 | LogicalType::Integer
3973 | LogicalType::BigInt
3974 | LogicalType::UTinyInt
3975 | LogicalType::USmallInt
3976 | LogicalType::UInteger
3977 | LogicalType::UBigInt
3978 | LogicalType::Date
3979 | LogicalType::Timestamp,
3980 FrequencyValue::Integer(_),
3981 )
3982 );
3983 if !valid {
3984 return Err(invalid("frequency value does not match its column"));
3985 }
3986 let count = cur.u64()?;
3987 if count == 0 || count > rows as u64 {
3988 return Err(invalid("frequency count is outside the table"));
3989 }
3990 entries.push(FrequencyEntry { value, count });
3991 }
3992 if entries.windows(2).any(|pair| pair[0].count < pair[1].count) {
3993 return Err(invalid("frequency entries are not descending"));
3994 }
3995 let ordinals = {
3996 let ordinal_count = cur.u32()? as usize;
3997 if ordinal_count > FREQUENCY_ORDINALS || ordinal_count > rows {
3998 return Err(invalid("frequency ordinal count exceeds its bound"));
3999 }
4000 let mut ordinals = Vec::with_capacity(ordinal_count);
4001 let mut previous = 0_u64;
4002 for at in 0..ordinal_count {
4003 let delta = cur.var_u64()?;
4004 if at != 0 && delta == 0 {
4005 return Err(invalid("frequency ordinals are not increasing"));
4006 }
4007 let ordinal = if at == 0 {
4008 delta
4009 } else {
4010 previous
4011 .checked_add(delta)
4012 .ok_or_else(|| invalid("frequency ordinal overflows"))?
4013 };
4014 if ordinal >= rows as u64 {
4015 return Err(invalid("frequency ordinal is outside the table"));
4016 }
4017 ordinals.push(ordinal);
4018 previous = ordinal;
4019 }
4020 ordinals
4021 };
4022 Some(FrequencySummary { entries, omitted_max, ordinals })
4023 }
4024 _ => return Err(invalid("frequency summary tag differs")),
4025 };
4026 frequencies.push(summary);
4027 }
4028 frequencies
4029 };
4030 if cur.at != bytes.len() {
4031 return Err(invalid("directory has trailing bytes"));
4032 }
4033 Ok(Table { name, fields, stripes, rows, dictionaries, distincts, frequencies })
4034}
4035
4036fn put_bound(out: &mut Vec<u8>, bound: Option<&Bound>) -> Result<()> {
4037 match bound {
4038 None => out.push(0),
4039 Some(Bound::Int(value)) => {
4040 out.push(1);
4041 out.extend_from_slice(&value.to_le_bytes());
4042 }
4043 Some(Bound::Real(value)) => {
4044 out.push(2);
4045 out.extend_from_slice(&value.to_le_bytes());
4046 }
4047 Some(Bound::Bytes(value)) => {
4048 out.push(3);
4049 put_u32(out, u32::try_from(value.len()).map_err(|_| invalid("bound length overflow"))?);
4050 out.extend_from_slice(value);
4051 }
4052 Some(Bound::Scaled { unscaled, scale }) => {
4053 out.push(4);
4054 out.extend_from_slice(&unscaled.to_le_bytes());
4055 out.push(*scale);
4056 }
4057 }
4058 Ok(())
4059}
4060
4061#[derive(Debug)]
4078struct Codes;
4079
4080impl chooser::Chooser for Codes {
4081 fn name(&self) -> &'static str {
4082 "codes"
4083 }
4084
4085 fn narrow_strings(
4086 &self,
4087 _values: &[&[u8]],
4088 offered: &[string::Kind],
4089 _depth: u8,
4090 ) -> Vec<string::Kind> {
4091 offered.to_vec()
4094 }
4095
4096 fn narrow_integers(
4097 &self,
4098 _values: &[i64],
4099 offered: &[integer::Kind],
4100 depth: u8,
4101 ) -> Vec<integer::Kind> {
4102 let keep: &[integer::Kind] = if depth == 0 {
4103 &[integer::Kind::Constant, integer::Kind::Packed, integer::Kind::Rle]
4104 } else {
4105 &[integer::Kind::Constant, integer::Kind::Packed]
4106 };
4107 let narrowed: Vec<integer::Kind> =
4108 offered.iter().copied().filter(|kind| keep.contains(kind)).collect();
4109 if narrowed.is_empty() { offered.to_vec() } else { narrowed }
4112 }
4113}
4114
4115#[derive(Debug)]
4127struct Fixed;
4128
4129impl chooser::Chooser for Fixed {
4130 fn name(&self) -> &'static str {
4131 "fixed"
4132 }
4133
4134 fn narrow_strings(
4135 &self,
4136 _values: &[&[u8]],
4137 offered: &[string::Kind],
4138 _depth: u8,
4139 ) -> Vec<string::Kind> {
4140 offered.to_vec()
4141 }
4142
4143 fn narrow_integers(
4144 &self,
4145 _values: &[i64],
4146 offered: &[integer::Kind],
4147 depth: u8,
4148 ) -> Vec<integer::Kind> {
4149 let keep: &[integer::Kind] = if depth == 0 {
4150 &[
4151 integer::Kind::Constant,
4152 integer::Kind::Packed,
4153 integer::Kind::Delta,
4154 integer::Kind::Rle,
4155 integer::Kind::Sparse,
4156 integer::Kind::Strided,
4157 ]
4158 } else {
4159 &[integer::Kind::Constant, integer::Kind::Packed, integer::Kind::Delta]
4160 };
4161 let narrowed: Vec<integer::Kind> =
4162 offered.iter().copied().filter(|kind| keep.contains(kind)).collect();
4163 if narrowed.is_empty() { offered.to_vec() } else { narrowed }
4164 }
4165}
4166
4167fn widened(data: &Data) -> Option<Vec<i64>> {
4174 match data {
4175 Data::Int8(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4176 Data::UInt8(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4177 Data::Int16(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4178 Data::UInt16(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4179 Data::Int32(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4180 Data::UInt32(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4181 Data::Int64(values) => Some(values.to_vec()),
4182 _ => None,
4183 }
4184}
4185
4186trait Narrow: Copy {
4193 const BIASED: (u32, u64);
4198
4199 fn narrow(value: i64) -> Self;
4201}
4202
4203#[allow(clippy::cast_sign_loss, reason = "a residue is a bit pattern and not a number")]
4220fn residue<T: Narrow>(value: i64) -> u64 {
4221 let (bits, bias) = T::BIASED;
4222 (value as u64).wrapping_add(bias) >> bits
4223}
4224
4225macro_rules! narrows {
4230 ($($ty:ty => $bias:expr),* $(,)?) => {$(
4231 impl Narrow for $ty {
4232 const BIASED: (u32, u64) = (<$ty>::BITS, $bias);
4233
4234 #[allow(
4235 clippy::cast_possible_truncation,
4236 clippy::cast_sign_loss,
4237 reason = "the caller has checked the bits this truncates away"
4238 )]
4239 fn narrow(value: i64) -> Self {
4240 value as Self
4241 }
4242 }
4243 )*};
4244}
4245
4246narrows! {
4247 i8 => 1 << 7,
4248 u8 => 0,
4249 i16 => 1 << 15,
4250 u16 => 0,
4251 i32 => 1 << 31,
4252 u32 => 0,
4253}
4254
4255fn fit<T: Narrow>(values: &[i64]) -> Result<Vec<T>> {
4268 let mut spilled = 0u64;
4269 for value in values {
4270 spilled |= residue::<T>(*value);
4271 }
4272 if spilled != 0 {
4273 return Err(invalid("page value is not of its type"));
4274 }
4275 Ok(values.iter().map(|value| T::narrow(*value)).collect())
4276}
4277
4278fn narrowed(ty: &LogicalType, values: Vec<i64>) -> Result<Data> {
4283 Ok(match ty {
4284 LogicalType::TinyInt => Data::Int8(fit::<i8>(&values)?.into()),
4285 LogicalType::UTinyInt => Data::UInt8(fit::<u8>(&values)?.into()),
4286 LogicalType::SmallInt => Data::Int16(fit::<i16>(&values)?.into()),
4287 LogicalType::USmallInt => Data::UInt16(fit::<u16>(&values)?.into()),
4288 LogicalType::Integer | LogicalType::Date => Data::Int32(fit::<i32>(&values)?.into()),
4289 LogicalType::UInteger => Data::UInt32(fit::<u32>(&values)?.into()),
4290 LogicalType::BigInt | LogicalType::Timestamp => Data::Int64(values.into()),
4291 _ => return Err(invalid("cascade codec belongs to a page that is not integers")),
4292 })
4293}
4294
4295fn plain_width(ty: &LogicalType) -> Option<usize> {
4298 Some(match ty {
4299 LogicalType::TinyInt | LogicalType::UTinyInt => 1,
4300 LogicalType::SmallInt | LogicalType::USmallInt => 2,
4301 LogicalType::Integer | LogicalType::UInteger | LogicalType::Date => 4,
4302 LogicalType::BigInt | LogicalType::Timestamp => 8,
4303 _ => return None,
4304 })
4305}
4306
4307fn cascaded(
4313 flat: &Vector,
4314 ty: &LogicalType,
4315 packed: Option<&Packed<'_>>,
4316) -> Result<Option<Vec<u8>>> {
4317 let (Some(width), Some(data)) = (plain_width(ty), flat.data()) else { return Ok(None) };
4318 let Some(values) = widened(data) else { return Ok(None) };
4319 let plain = values.len().saturating_mul(width);
4320 let best = match packed {
4321 Some(packed) => plain.min(21 + size_of_val(packed.words())),
4323 None => plain,
4324 };
4325 let out = integer::encode_with(&values, &Fixed)?;
4326 Ok((out.len() < best).then_some(out))
4327}
4328
4329fn encoded_codes(codes: &[u32]) -> Result<Option<Vec<u8>>> {
4341 let wide: Vec<i64> = codes.iter().map(|code| i64::from(*code)).collect();
4342 let coded = integer::encode_with(&wide, &Codes)?;
4343 let plain = codes.len().saturating_mul(size_of::<u32>());
4344 Ok((coded.len() < plain).then_some(coded))
4345}
4346
4347fn encode(
4348 vector: &Vector,
4349 global: Option<&mut GlobalDictionary>,
4350) -> Result<(Vec<u8>, Option<Vec<u32>>)> {
4351 let ty = vector.logical_type();
4352 let flat = vector.flatten()?;
4354 let mut out = Vec::new();
4355 let mut global_codes = None;
4356 if let Some(global) = global {
4357 let mut codes = Vec::with_capacity(flat.len());
4358 for row in 0..flat.len() {
4359 let text = flat.text_at(row).unwrap_or("");
4360 let code = global.code(text)?;
4361 global.observe(code, flat.is_null_at(row))?;
4362 codes.push(code);
4363 }
4364 global_codes = Some(codes);
4365 }
4366 let membership = global_codes.as_deref().map(unique_codes);
4367 let dictionary = if global_codes.is_none() && ty == &LogicalType::Varchar {
4368 string_dictionary(&flat)?
4369 } else {
4370 None
4371 };
4372 let packed_vector = if dictionary.is_none() && global_codes.is_none() {
4373 Some(flat.bit_packed()?)
4374 } else {
4375 None
4376 };
4377 let packed = packed_vector.as_ref().and_then(Vector::packed_parts);
4378 let coded = match global_codes.as_deref() {
4379 Some(codes) => encoded_codes(codes)?,
4380 None => None,
4381 };
4382 let cascade = if dictionary.is_none() && global_codes.is_none() {
4386 cascaded(&flat, ty, packed.as_ref())?
4387 } else {
4388 None
4389 };
4390 out.push(if coded.is_some() {
4391 4
4392 } else if cascade.is_some() {
4393 5
4394 } else if global_codes.is_some() {
4395 3
4396 } else if dictionary.is_some() {
4397 1
4398 } else if packed.is_some() {
4399 2
4400 } else {
4401 0
4402 });
4403 let nulls = flat.validity();
4404 let flag = match nulls {
4405 Validity::AllValid => 0,
4406 Validity::AllInvalid => 1,
4407 Validity::Mask(_) => 2,
4408 };
4409 out.push(flag);
4410 if flag == 2 {
4411 for group in (0..vector.len()).step_by(8) {
4412 let mut bits = 0_u8;
4413 for bit in 0..8 {
4414 if group + bit < vector.len() && !flat.is_null_at(group + bit) {
4415 bits |= 1 << bit;
4416 }
4417 }
4418 out.push(bits);
4419 }
4420 }
4421 if let Some(coded) = coded {
4422 out.extend_from_slice(&coded);
4423 return Ok((out, membership));
4424 }
4425 if let Some(cascade) = cascade {
4426 out.extend_from_slice(&cascade);
4427 return Ok((out, membership));
4428 }
4429 if let Some(codes) = global_codes {
4430 for code in codes {
4431 put_u32(&mut out, code);
4432 }
4433 return Ok((out, membership));
4434 }
4435 if let Some(dictionary) = dictionary {
4436 out.extend_from_slice(&dictionary);
4437 return Ok((out, membership));
4438 }
4439 if let Some(packed) = packed {
4440 if packed.offset() != 0 {
4441 return Err(invalid("writer received a sliced packed vector"));
4442 }
4443 out.push(u8::try_from(packed.width()).map_err(|_| invalid("packed width overflow"))?);
4444 out.extend_from_slice(&packed.base().to_le_bytes());
4445 put_u32(
4446 &mut out,
4447 u32::try_from(packed.words().len()).map_err(|_| invalid("too many packed words"))?,
4448 );
4449 for word in packed.words() {
4450 put_u64(&mut out, *word);
4451 }
4452 return Ok((out, membership));
4453 }
4454 let data = flat.data().ok_or_else(|| invalid("scalar column did not flatten"))?;
4455 match (ty, data) {
4456 (LogicalType::TinyInt, Data::Int8(values)) => {
4457 for value in &**values {
4458 out.extend_from_slice(&value.to_le_bytes());
4459 }
4460 }
4461 (LogicalType::UTinyInt, Data::UInt8(values)) => {
4462 for value in &**values {
4463 out.extend_from_slice(&value.to_le_bytes());
4464 }
4465 }
4466 (LogicalType::SmallInt, Data::Int16(values)) => {
4467 for value in &**values {
4468 out.extend_from_slice(&value.to_le_bytes());
4469 }
4470 }
4471 (LogicalType::USmallInt, Data::UInt16(values)) => {
4472 for value in &**values {
4473 out.extend_from_slice(&value.to_le_bytes());
4474 }
4475 }
4476 (LogicalType::UInteger, Data::UInt32(values)) => {
4477 for value in &**values {
4478 out.extend_from_slice(&value.to_le_bytes());
4479 }
4480 }
4481 (LogicalType::UBigInt, Data::UInt64(values)) => {
4482 for value in &**values {
4483 out.extend_from_slice(&value.to_le_bytes());
4484 }
4485 }
4486 (LogicalType::Integer | LogicalType::Date, Data::Int32(values)) => {
4487 for value in &**values {
4488 out.extend_from_slice(&value.to_le_bytes());
4489 }
4490 }
4491 (LogicalType::BigInt | LogicalType::Timestamp, Data::Int64(values)) => {
4492 for value in &**values {
4493 out.extend_from_slice(&value.to_le_bytes());
4494 }
4495 }
4496 (LogicalType::Boolean, Data::Bool(values)) => {
4497 for value in &**values {
4498 out.push(u8::from(*value));
4499 }
4500 }
4501 (LogicalType::Varchar, Data::Varlen(values)) => {
4502 let mut bytes = Vec::new();
4503 put_u32(&mut out, 0);
4504 for row in 0..vector.len() {
4505 let value = values.bytes(row).ok_or_else(|| invalid("string view is invalid"))?;
4506 bytes.extend_from_slice(value);
4507 put_u32(
4508 &mut out,
4509 u32::try_from(bytes.len())
4510 .map_err(|_| invalid("string payload exceeds 4GiB"))?,
4511 );
4512 }
4513 out.extend_from_slice(&bytes);
4514 }
4515 _ => return Err(Error::not_implemented(format!("native page for {ty}"))),
4516 }
4517 Ok((out, membership))
4518}
4519
4520fn put_varint(out: &mut Vec<u8>, mut value: u32) {
4521 while value >= 0x80 {
4522 out.push((value as u8 & 0x7f) | 0x80);
4523 value >>= 7;
4524 }
4525 out.push(value as u8);
4526}
4527
4528fn unique_codes(codes: &[u32]) -> Vec<u32> {
4530 let mut unique = codes.to_vec();
4531 unique.sort_unstable();
4532 unique.dedup();
4533 unique
4534}
4535
4536fn merged_codes(lists: Vec<Vec<u32>>) -> Vec<u32> {
4542 let mut lists = lists;
4543 while lists.len() > 1 {
4544 let mut next = Vec::with_capacity(lists.len().div_ceil(2));
4545 for pair in lists.chunks(2) {
4546 match pair {
4547 [left, right] => next.push(merged_pair(left, right)),
4548 [only] => next.push(only.clone()),
4549 _ => {}
4550 }
4551 }
4552 lists = next;
4553 }
4554 lists.pop().unwrap_or_default()
4555}
4556
4557fn merged_pair(left: &[u32], right: &[u32]) -> Vec<u32> {
4558 let mut out = Vec::with_capacity(left.len().saturating_add(right.len()));
4559 let mut at = 0;
4560 let mut to = 0;
4561 while at < left.len() && to < right.len() {
4562 match left[at].cmp(&right[to]) {
4563 Ordering::Less => {
4564 out.push(left[at]);
4565 at += 1;
4566 }
4567 Ordering::Greater => {
4568 out.push(right[to]);
4569 to += 1;
4570 }
4571 Ordering::Equal => {
4572 out.push(left[at]);
4573 at += 1;
4574 to += 1;
4575 }
4576 }
4577 }
4578 out.extend_from_slice(&left[at..]);
4579 out.extend_from_slice(&right[to..]);
4580 out
4581}
4582
4583fn merged_range(ranges: impl Iterator<Item = Range>) -> Range {
4588 let mut merged = Range::default();
4589 let mut first = true;
4590 for range in ranges {
4591 merged.nulls = merged.nulls.saturating_add(range.nulls);
4592 merged.sum = match (merged.sum.take(), range.sum) {
4596 (Some(held), Some(next)) if !first => held.checked_add(next),
4597 (_, next) if first => next,
4598 _ => None,
4599 };
4600 merged.exact = if first { range.exact } else { merged.exact && range.exact };
4601 if first {
4602 merged.low = range.low;
4603 merged.high = range.high;
4604 first = false;
4605 continue;
4606 }
4607 merged.low = match (merged.low.take(), range.low) {
4608 (Some(held), Some(next)) => Some(held.smaller(next)),
4609 _ => None,
4610 };
4611 merged.high = match (merged.high.take(), range.high) {
4612 (Some(held), Some(next)) => Some(held.larger(next)),
4613 _ => None,
4614 };
4615 }
4616 merged
4617}
4618
4619fn shortened(bound: Option<Bound>, high: bool) -> Option<Bound> {
4632 match bound {
4633 Some(Bound::Bytes(mut value)) if value.len() > PART_BOUND_BYTES => {
4634 value.truncate(PART_BOUND_BYTES);
4635 if !high {
4636 return Some(Bound::Bytes(value));
4637 }
4638 while let Some(last) = value.pop() {
4639 if last < u8::MAX {
4640 value.push(last + 1);
4641 return Some(Bound::Bytes(value));
4642 }
4643 }
4644 None
4645 }
4646 other => other,
4647 }
4648}
4649
4650fn encode_part_ranges(ranges: &[Range]) -> Result<Vec<u8>> {
4658 let mut out = Vec::new();
4659 put_u32(
4660 &mut out,
4661 u32::try_from(ranges.len()).map_err(|_| invalid("too many parts in a stripe"))?,
4662 );
4663 for range in ranges {
4664 put_bound(&mut out, shortened(range.low.clone(), false).as_ref())?;
4665 put_bound(&mut out, shortened(range.high.clone(), true).as_ref())?;
4666 put_u32(&mut out, u32::try_from(range.nulls).map_err(|_| invalid("null count overflow"))?);
4667 }
4668 Ok(out)
4669}
4670
4671fn decode_part_ranges(bytes: &[u8]) -> Result<Vec<Range>> {
4673 let mut cur = Cursor { bytes, at: 0 };
4674 let parts = cur.u32()? as usize;
4675 let mut out = Vec::new();
4676 for _ in 0..parts {
4677 let low = cur.bound()?;
4678 let high = cur.bound()?;
4679 let nulls = cur.u32()? as usize;
4680 out.push(Range { low, high, nulls, exact: false, sum: None });
4681 }
4682 Ok(out)
4683}
4684
4685fn encode_sieves<'a>(sieves: impl Iterator<Item = &'a Option<Sieve>>) -> Result<Vec<u8>> {
4686 let held: Vec<&Option<Sieve>> = sieves.collect();
4687 let mut out = Vec::new();
4688 put_u32(
4689 &mut out,
4690 u32::try_from(held.len()).map_err(|_| invalid("too many parts in a stripe"))?,
4691 );
4692 for sieve in &held {
4693 let length = sieve.as_ref().map_or(0, Sieve::len);
4694 put_u32(&mut out, u32::try_from(length).map_err(|_| invalid("sieve length overflow"))?);
4695 }
4696 for sieve in held.into_iter().flatten() {
4698 out.extend_from_slice(&sieve.to_bytes());
4699 }
4700 Ok(out)
4701}
4702
4703fn decode_sieves(bytes: &[u8]) -> Result<Vec<Option<Sieve>>> {
4709 let parts = u32::from_le_bytes(
4710 bytes
4711 .get(..4)
4712 .ok_or_else(|| invalid("sieve page is truncated"))?
4713 .try_into()
4714 .map_err(|_| invalid("sieve page is truncated"))?,
4715 ) as usize;
4716 let mut lengths = Vec::with_capacity(parts);
4717 for part in 0..parts {
4718 let at = 4 + part * 4;
4719 let field = bytes.get(at..at + 4).ok_or_else(|| invalid("sieve page is truncated"))?;
4720 lengths.push(u32::from_le_bytes(
4721 field.try_into().map_err(|_| invalid("sieve page is truncated"))?,
4722 ) as usize);
4723 }
4724 let mut at = 4 + parts * 4;
4725 let mut out = Vec::with_capacity(parts);
4726 for length in lengths {
4727 if length == 0 {
4728 out.push(None);
4729 continue;
4730 }
4731 let end = at.checked_add(length).ok_or_else(|| invalid("sieve page is truncated"))?;
4732 let field = bytes.get(at..end).ok_or_else(|| invalid("sieve page is truncated"))?;
4733 out.push(Sieve::from_bytes(field));
4734 at = end;
4735 }
4736 if at != bytes.len() {
4737 return Err(invalid("sieve page has trailing bytes"));
4738 }
4739 Ok(out)
4740}
4741
4742fn encode_membership(unique: &[u32]) -> Vec<u8> {
4748 let mut out = Vec::with_capacity(unique.len().saturating_mul(2).saturating_add(5));
4749 put_varint(&mut out, u32::try_from(unique.len()).unwrap_or(u32::MAX));
4750 let mut previous = 0;
4751 for (at, &code) in unique.iter().enumerate() {
4752 put_varint(&mut out, if at == 0 { code } else { code - previous });
4753 previous = code;
4754 }
4755 out
4756}
4757
4758fn take_varint(bytes: &[u8], at: &mut usize) -> Result<u32> {
4759 let mut value = 0_u32;
4760 for shift in (0..35).step_by(7) {
4761 let byte = *bytes.get(*at).ok_or_else(|| invalid("membership varint is truncated"))?;
4762 *at += 1;
4763 let part = u32::from(byte & 0x7f);
4764 if shift == 28 && part > 0x0f {
4765 return Err(invalid("membership varint overflow"));
4766 }
4767 value = value
4768 .checked_add(
4769 part.checked_shl(shift).ok_or_else(|| invalid("membership varint overflow"))?,
4770 )
4771 .ok_or_else(|| invalid("membership varint overflow"))?;
4772 if byte & 0x80 == 0 {
4773 return Ok(value);
4774 }
4775 }
4776 Err(invalid("membership varint is too long"))
4777}
4778
4779fn decode_membership(bytes: &[u8]) -> Result<Vec<u32>> {
4780 let mut at = 0;
4781 let count = take_varint(bytes, &mut at)? as usize;
4782 let mut codes = Vec::with_capacity(count);
4783 let mut previous = 0_u32;
4784 for index in 0..count {
4785 let delta = take_varint(bytes, &mut at)?;
4786 let code = if index == 0 {
4787 delta
4788 } else {
4789 previous.checked_add(delta).ok_or_else(|| invalid("membership code overflow"))?
4790 };
4791 if index > 0 && code <= previous {
4792 return Err(invalid("membership codes are not increasing"));
4793 }
4794 codes.push(code);
4795 previous = code;
4796 }
4797 if at != bytes.len() {
4798 return Err(invalid("membership page has trailing bytes"));
4799 }
4800 Ok(codes)
4801}
4802
4803fn string_dictionary(vector: &Vector) -> Result<Option<Vec<u8>>> {
4804 let mut by_text = HashMap::new();
4805 let mut values = Vec::new();
4806 let mut codes = Vec::with_capacity(vector.len());
4807 let mut plain_bytes = 0_usize;
4808 for row in 0..vector.len() {
4809 let text = vector.text_at(row).unwrap_or("");
4810 plain_bytes = plain_bytes.saturating_add(text.len());
4811 let code = match by_text.get(text) {
4812 Some(&code) => code,
4813 None => {
4814 let code = u32::try_from(values.len())
4815 .map_err(|_| invalid("too many dictionary values"))?;
4816 by_text.insert(text, code);
4817 values.push(text);
4818 code
4819 }
4820 };
4821 codes.push(code);
4822 }
4823 let dictionary_bytes = values.iter().map(|value| value.len()).sum::<usize>();
4824 let encoded = 8_usize
4825 .saturating_add((values.len() + 1).saturating_mul(4))
4826 .saturating_add(dictionary_bytes)
4827 .saturating_add(codes.len().saturating_mul(4));
4828 let plain = (vector.len() + 1).saturating_mul(4).saturating_add(plain_bytes);
4829 if encoded >= plain {
4830 return Ok(None);
4831 }
4832 let mut out = Vec::with_capacity(encoded);
4833 put_u32(
4834 &mut out,
4835 u32::try_from(values.len()).map_err(|_| invalid("too many dictionary values"))?,
4836 );
4837 put_u32(
4838 &mut out,
4839 u32::try_from(dictionary_bytes).map_err(|_| invalid("dictionary payload exceeds 4GiB"))?,
4840 );
4841 let mut offset = 0_u32;
4842 put_u32(&mut out, offset);
4843 for value in &values {
4844 offset = offset
4845 .checked_add(
4846 u32::try_from(value.len()).map_err(|_| invalid("dictionary value is too long"))?,
4847 )
4848 .ok_or_else(|| invalid("dictionary payload exceeds 4GiB"))?;
4849 put_u32(&mut out, offset);
4850 }
4851 for value in values {
4852 out.extend_from_slice(value.as_bytes());
4853 }
4854 for code in codes {
4855 put_u32(&mut out, code);
4856 }
4857 Ok(Some(out))
4858}
4859
4860struct EncodedDictionary {
4861 index: Vec<u8>,
4862 ranks: Vec<u8>,
4863 payload: Vec<Vec<u8>>,
4866}
4867
4868fn head(bytes: &[u8]) -> u64 {
4870 let mut word = [0; 8];
4871 let take = bytes.len().min(8);
4872 word[..take].copy_from_slice(&bytes[..take]);
4873 u64::from_be_bytes(word)
4874}
4875
4876fn rankings(dictionaries: &[Option<GlobalDictionary>]) -> Result<Vec<Vec<(u64, u32)>>> {
4884 let present =
4885 dictionaries.iter().enumerate().filter(|(_, held)| held.is_some()).map(|(at, _)| at);
4886 let present = present.collect::<Vec<_>>();
4887 let mut orders = vec![Vec::new(); dictionaries.len()];
4888 let workers = std::thread::available_parallelism()
4889 .map_or(1, usize::from)
4890 .min(MAX_FREQUENCY_WORKERS)
4891 .min(present.len());
4892 if workers <= 1 {
4893 for at in present {
4894 if let Some(dictionary) = &dictionaries[at] {
4895 orders[at] = dictionary.ranked();
4896 }
4897 }
4898 return Ok(orders);
4899 }
4900 let width = present.len().div_ceil(workers);
4901 let pieces = std::thread::scope(|scope| {
4902 present
4903 .chunks(width)
4904 .map(|columns| {
4905 scope.spawn(|| {
4906 columns
4907 .iter()
4908 .filter_map(|&at| dictionaries[at].as_ref().map(|held| (at, held.ranked())))
4909 .collect::<Vec<_>>()
4910 })
4911 })
4912 .collect::<Vec<_>>()
4913 .into_iter()
4914 .map(|handle| {
4915 handle.join().map_err(|_| Error::internal("a dictionary sort worker panicked"))
4916 })
4917 .collect::<Result<Vec<_>>>()
4918 })?;
4919 for piece in pieces {
4920 for (at, order) in piece {
4921 orders[at] = order;
4922 }
4923 }
4924 Ok(orders)
4925}
4926
4927fn encode_global_dictionary(
4928 dictionary: GlobalDictionary,
4929 order: &[(u64, u32)],
4930) -> Result<EncodedDictionary> {
4931 let values = dictionary.offsets.len() - 1;
4932 if order.len() != values {
4933 return Err(invalid("global dictionary order does not cover its values"));
4934 }
4935 let blocks = values.div_ceil(TEXT_PAYLOAD_VALUES);
4936 let payload = encode_payload(&dictionary)?;
4937 if payload.len() != blocks {
4938 return Err(invalid("global dictionary payload is not the blocks it says it is"));
4939 }
4940 let (ranks, rank_ends) = encode_ranks(order, code_width(values))?;
4941 let rank_blocks = values.div_ceil(TEXT_RANK_BLOCK);
4942 let offset_bits = offset_width(&dictionary.offsets);
4943 let mut index = Vec::with_capacity(
4944 DICTIONARY_HEADER + offset_bytes(values, offset_bits) + (blocks + rank_blocks) * 16,
4945 );
4946 put_u32(
4947 &mut index,
4948 u32::try_from(values).map_err(|_| invalid("global dictionary has too many values"))?,
4949 );
4950 put_u32(&mut index, TEXT_PAYLOAD_VALUES as u32);
4951 put_u32(
4952 &mut index,
4953 u32::try_from(blocks).map_err(|_| invalid("global dictionary has too many blocks"))?,
4954 );
4955 put_u32(&mut index, offset_bits as u32);
4956 encode_offsets(&dictionary.offsets, offset_bits, &mut index)?;
4957 let mut at = 0_u64;
4961 for block in &payload {
4962 at = at
4963 .checked_add(block.len() as u64)
4964 .ok_or_else(|| invalid("global dictionary payload overflow"))?;
4965 put_u64(&mut index, at);
4966 }
4967 for block in &payload {
4968 put_u64(&mut index, checksum(block));
4969 }
4970 if rank_ends.len() != rank_blocks {
4973 return Err(invalid("global dictionary order is not the blocks it says it is"));
4974 }
4975 for end in &rank_ends {
4976 put_u64(&mut index, *end);
4977 }
4978 let mut at = 0_usize;
4979 for end in &rank_ends {
4980 let end = usize::try_from(*end).map_err(|_| invalid("global dictionary order overflow"))?;
4981 put_u64(&mut index, checksum(&ranks[at..end]));
4982 at = end;
4983 }
4984 Ok(EncodedDictionary { index, ranks, payload })
4985}
4986
4987const PAYLOAD_SAMPLE_BLOCKS: usize = 8;
4994
4995fn payload_shapes() -> Vec<chooser::Settled> {
5021 let integers = vec![integer::Kind::Packed];
5022 [
5023 vec![string::Kind::Front, string::Kind::Lz],
5024 vec![string::Kind::Lz, string::Kind::Fsst],
5025 vec![string::Kind::Lz, string::Kind::Plain],
5026 vec![string::Kind::Fsst],
5027 vec![string::Kind::Plain],
5028 ]
5029 .into_iter()
5030 .map(|strings| chooser::Settled::new(strings, integers.clone()))
5031 .collect()
5032}
5033
5034fn encode_payload(dictionary: &GlobalDictionary) -> Result<Vec<Vec<u8>>> {
5040 let values = dictionary.offsets.len() - 1;
5041 let blocks = values.div_ceil(TEXT_PAYLOAD_VALUES);
5042 let run = |block: usize| {
5043 let first = block * TEXT_PAYLOAD_VALUES;
5044 let last = (first + TEXT_PAYLOAD_VALUES).min(values);
5045 (first..last)
5046 .map(|value| {
5047 let from = dictionary.offsets[value] as usize;
5048 let to = dictionary.offsets[value + 1] as usize;
5049 &dictionary.payload[from..to]
5050 })
5051 .collect::<Vec<_>>()
5052 };
5053 let shape = (blocks > PAYLOAD_SAMPLE_BLOCKS).then(|| settle_shape(&run, blocks)).transpose()?;
5056 let one = |block: usize| match &shape {
5057 Some(shape) => string::encode_with(&run(block), shape),
5058 None => string::encode(&run(block)),
5059 };
5060 let workers = std::thread::available_parallelism()
5061 .map_or(1, usize::from)
5062 .min(MAX_FREQUENCY_WORKERS)
5063 .min(blocks);
5064 if workers <= 1 {
5065 return (0..blocks).map(one).collect();
5066 }
5067 let next = AtomicUsize::new(0);
5068 let pieces = std::thread::scope(|scope| {
5069 (0..workers)
5070 .map(|_| {
5071 scope.spawn(|| {
5072 let mut mine = Vec::new();
5073 loop {
5074 let block = next.fetch_add(1, Atomic::Relaxed);
5075 if block >= blocks {
5076 break;
5077 }
5078 mine.push((block, one(block)?));
5079 }
5080 Ok(mine)
5081 })
5082 })
5083 .collect::<Vec<_>>()
5084 .into_iter()
5085 .map(|handle| {
5086 handle.join().map_err(|_| Error::internal("a dictionary encode worker panicked"))?
5087 })
5088 .collect::<Result<Vec<_>>>()
5089 })?;
5090 let mut payload = vec![Vec::new(); blocks];
5091 for piece in pieces {
5092 for (block, bytes) in piece {
5093 payload[block] = bytes;
5094 }
5095 }
5096 Ok(payload)
5097}
5098
5099fn settle_shape<'a>(
5107 run: &dyn Fn(usize) -> Vec<&'a [u8]>,
5108 blocks: usize,
5109) -> Result<chooser::Settled> {
5110 let last = blocks - 1;
5111 let sample = (0..PAYLOAD_SAMPLE_BLOCKS)
5112 .map(|region| run(region * last / (PAYLOAD_SAMPLE_BLOCKS - 1)))
5113 .collect::<Vec<_>>();
5114 let mut best: Option<(chooser::Settled, usize)> = None;
5115 for shape in payload_shapes() {
5116 let mut size = 0;
5117 for block in &sample {
5118 size += string::encode_with(block, &shape)?.len();
5119 }
5120 if best.as_ref().is_none_or(|(_, smallest)| size < *smallest) {
5121 best = Some((shape, size));
5122 }
5123 }
5124 best.map(|(shape, _)| shape)
5125 .ok_or_else(|| invalid("no shape applies to a global dictionary payload"))
5126}
5127
5128fn encode_ranks(order: &[(u64, u32)], code_bits: usize) -> Result<(Vec<u8>, Vec<u64>)> {
5135 let mut out = Vec::with_capacity(order.len() * 4);
5136 let mut ends = Vec::with_capacity(order.len().div_ceil(TEXT_RANK_BLOCK));
5137 let mut heads = Vec::with_capacity(TEXT_RANK_BLOCK);
5138 let mut codes = Vec::with_capacity(TEXT_RANK_BLOCK);
5139 for block in order.chunks(TEXT_RANK_BLOCK) {
5140 let base = block.first().map_or(0, |&(head, _)| head);
5143 let span = block.last().map_or(0, |&(head, _)| head.wrapping_sub(base));
5144 let width = (u64::BITS - span.leading_zeros()) as usize;
5145 heads.clear();
5146 codes.clear();
5147 for &(head, code) in block {
5148 heads.push(head.wrapping_sub(base));
5149 codes.push(u64::from(code));
5150 }
5151 put_u64(&mut out, base);
5152 out.push(width as u8);
5153 bitpack::pack_tail(&heads, width, &mut out)
5154 .map_err(|_| invalid("global dictionary heads do not pack"))?;
5155 bitpack::pack_tail(&codes, code_bits, &mut out)
5156 .map_err(|_| invalid("global dictionary codes do not pack"))?;
5157 ends.push(out.len() as u64);
5158 }
5159 Ok((out, ends))
5160}
5161
5162fn open_global_dictionary(
5169 file: Arc<File>,
5170 page: Page,
5171 ty: &LogicalType,
5172 keep_budget: usize,
5173) -> Result<Vector> {
5174 if ty != &LogicalType::Varchar {
5175 return Err(invalid("global dictionary belongs to a non-string column"));
5176 }
5177 let mut header = [0; DICTIONARY_HEADER];
5178 read_at(&file, page.offset, &mut header)?;
5179 let count = u32::from_le_bytes(header[0..4].try_into().expect("four bytes")) as usize;
5180 let per_block = u32::from_le_bytes(header[4..8].try_into().expect("four bytes")) as usize;
5181 let blocks = u32::from_le_bytes(header[8..12].try_into().expect("four bytes")) as usize;
5182 let offset_bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
5183 if per_block != TEXT_PAYLOAD_VALUES {
5184 return Err(invalid("global dictionary block width differs"));
5185 }
5186 if blocks != count.div_ceil(TEXT_PAYLOAD_VALUES) {
5187 return Err(invalid("global dictionary block count differs from its value count"));
5188 }
5189 if offset_bits > u32::BITS as usize {
5190 return Err(invalid("global dictionary packs offsets past a payload"));
5191 }
5192 let offset_len = offset_bytes(count, offset_bits);
5193 let ranks = count;
5198 let rank_blocks = ranks.div_ceil(TEXT_RANK_BLOCK);
5199 let hash_len = blocks
5202 .checked_add(rank_blocks)
5203 .and_then(|words| words.checked_mul(16))
5204 .ok_or_else(|| invalid("global dictionary block count overflow"))?;
5205 let index_len = DICTIONARY_HEADER
5206 .checked_add(offset_len)
5207 .and_then(|len| len.checked_add(hash_len))
5208 .ok_or_else(|| invalid("global dictionary header overflow"))?;
5209 if index_len > page.length as usize {
5210 return Err(invalid("global dictionary offset index exceeds its page"));
5211 }
5212 let mut index = vec![0; index_len];
5213 index[..DICTIONARY_HEADER].copy_from_slice(&header);
5214 read_at(&file, page.offset + DICTIONARY_HEADER as u64, &mut index[DICTIONARY_HEADER..])?;
5215 if checksum(&index) != page.hash {
5216 return Err(invalid("global dictionary index checksum differs"));
5217 }
5218 let offsets = index[DICTIONARY_HEADER..DICTIONARY_HEADER + offset_len].to_vec();
5219 let mut words = index[DICTIONARY_HEADER + offset_len..]
5220 .chunks_exact(8)
5221 .map(|part| u64::from_le_bytes(part.try_into().expect("eight bytes")))
5222 .collect::<Vec<_>>();
5223 let mut hashes = words.split_off(blocks);
5224 let mut rank_ends = hashes.split_off(blocks);
5225 let rank_hashes = rank_ends.split_off(rank_blocks);
5226 let ends = words;
5227 if rank_ends.windows(2).any(|pair| pair[0] >= pair[1]) {
5230 return Err(invalid("global dictionary order blocks do not rise"));
5231 }
5232 let rank_len = usize::try_from(rank_ends.last().copied().unwrap_or_default())
5233 .map_err(|_| invalid("global dictionary rank overflow"))?;
5234 let body_len = index_len
5235 .checked_add(rank_len)
5236 .ok_or_else(|| invalid("global dictionary header overflow"))?;
5237 if body_len > page.length as usize {
5238 return Err(invalid("global dictionary order exceeds its page"));
5239 }
5240 let stored_len = page.length as usize - body_len;
5243 if ends.last().copied().unwrap_or_default() as usize != stored_len
5244 || ends.windows(2).any(|pair| pair[0] > pair[1])
5245 {
5246 return Err(invalid("global dictionary blocks do not bound the payload"));
5247 }
5248 Vector::external_text(
5249 LogicalType::Varchar,
5250 Arc::new(NativeText {
5251 file,
5252 values: count,
5253 offsets,
5254 offset_bits,
5255 ranks,
5256 rank_at: page.offset + index_len as u64,
5257 rank_ends,
5258 rank_hashes,
5259 rank_blocks: (0..rank_blocks).map(|_| OnceLock::new()).collect(),
5260 code_bits: code_width(count),
5261 code_ranks: OnceLock::new(),
5262 payload: page.offset + body_len as u64,
5263 ends,
5264 hashes,
5265 blocks: (0..blocks).map(|_| OnceLock::new()).collect(),
5266 keep_budget,
5267 payload_kept: AtomicUsize::new(0),
5268 }),
5269 )
5270}
5271
5272fn decode(
5273 ty: &LogicalType,
5274 rows: usize,
5275 bytes: &[u8],
5276 global: Option<Arc<Vector>>,
5277) -> Result<Vector> {
5278 let mut cur = Cursor { bytes, at: 0 };
5279 let codec = cur.u8()?;
5280 let flag = cur.u8()?;
5281 let validity = match flag {
5282 0 => Validity::AllValid,
5283 1 => Validity::AllInvalid,
5284 2 => {
5285 let mask = cur.take(rows.div_ceil(8))?;
5286 Validity::from_iter(rows, |row| mask[row / 8] >> (row % 8) & 1 == 1)
5287 }
5288 _ => return Err(invalid("page validity tag differs")),
5289 };
5290 if codec == 1 {
5291 if ty != &LogicalType::Varchar {
5292 return Err(invalid("dictionary codec belongs to a non-string page"));
5293 }
5294 let count = cur.u32()? as usize;
5295 let payload_len = cur.u32()? as usize;
5296 let offset_bytes = cur.take(
5297 (count + 1)
5298 .checked_mul(4)
5299 .ok_or_else(|| invalid("dictionary offset count overflow"))?,
5300 )?;
5301 let offsets = offset_bytes
5302 .chunks_exact(4)
5303 .map(|part| u32::from_le_bytes(part.try_into().expect("four bytes")))
5304 .collect::<Vec<_>>();
5305 let payload = cur.take(payload_len)?.to_vec();
5306 if offsets.first() != Some(&0)
5307 || offsets.last().copied().map(|last| last as usize) != Some(payload.len())
5308 || offsets.windows(2).any(|pair| pair[0] > pair[1])
5309 {
5310 return Err(invalid("dictionary offsets do not bound the payload"));
5311 }
5312 let mut strings = StringColumn::over(Buffer::from_vec(payload));
5313 for pair in offsets.windows(2) {
5314 strings.push_in_place(pair[0] as usize, (pair[1] - pair[0]) as usize)?;
5315 }
5316 let mut codes = Vec::with_capacity(rows);
5317 for _ in 0..rows {
5318 codes.push(cur.u32()?);
5319 }
5320 if codes.iter().any(|code| *code as usize >= count) {
5321 return Err(invalid("dictionary code is out of range"));
5322 }
5323 if cur.at != bytes.len() {
5324 return Err(invalid("dictionary page has trailing bytes"));
5325 }
5326 let dictionary = Vector::flat(LogicalType::Varchar, Data::Varlen(strings))?;
5327 return Ok(Vector::dictionary(codes, dictionary)?.with_validity(validity));
5328 }
5329 if codec == 3 || codec == 4 {
5330 let dictionary = global.ok_or_else(|| invalid("global code page has no dictionary"))?;
5331 let codes = if codec == 4 {
5332 let wide = integer::decode(&bytes[cur.at..])?;
5335 if wide.len() != rows {
5336 return Err(invalid("encoded code page holds the wrong number of rows"));
5337 }
5338 let mut codes = Vec::with_capacity(wide.len());
5345 let mut seen = 0_i64;
5346 for &code in &wide {
5347 seen |= code;
5348 codes.push(code as u32);
5349 }
5350 if seen < 0 || seen > i64::from(u32::MAX) {
5351 return Err(invalid("code is not a code"));
5352 }
5353 codes
5354 } else {
5355 let mut codes = Vec::with_capacity(rows);
5356 for _ in 0..rows {
5357 codes.push(cur.u32()?);
5358 }
5359 if cur.at != bytes.len() {
5360 return Err(invalid("global code page has trailing bytes"));
5361 }
5362 codes
5363 };
5364 let highest = codes.iter().copied().max();
5365 return Ok(Vector::stable_dictionary_validated(codes, dictionary, highest)?
5366 .with_validity(validity));
5367 }
5368 if codec == 5 {
5369 let values = integer::decode(&bytes[cur.at..])?;
5371 if values.len() != rows {
5372 return Err(invalid("cascade page holds the wrong number of rows"));
5373 }
5374 let data = narrowed(ty, values)?;
5375 return Ok(Vector::flat(ty.clone(), data)?.with_validity(validity));
5376 }
5377 if codec == 2 {
5378 let width = u32::from(cur.u8()?);
5379 let base = i128::from_le_bytes(cur.take(16)?.try_into().expect("sixteen bytes"));
5380 let count = cur.u32()? as usize;
5381 let mut words = Vec::with_capacity(count);
5382 for _ in 0..count {
5383 words.push(cur.u64()?);
5384 }
5385 if cur.at != bytes.len() {
5386 return Err(invalid("packed page has trailing bytes"));
5387 }
5388 return Ok(Vector::packed(ty.clone(), words, width, base, rows)?.with_validity(validity));
5389 }
5390 if codec != 0 {
5391 return Err(invalid("page codec is unknown"));
5392 }
5393 let data = match ty {
5394 LogicalType::TinyInt => {
5395 let values = cur.take(rows)?;
5396 Data::Int8(values.iter().map(|item| *item as i8).collect::<Vec<_>>().into())
5397 }
5398 LogicalType::UTinyInt => Data::UInt8(cur.take(rows)?.to_vec().into()),
5399 LogicalType::SmallInt => {
5400 let values =
5401 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5402 Data::Int16(
5403 values
5404 .chunks_exact(2)
5405 .map(|item| i16::from_le_bytes(item.try_into().expect("two bytes")))
5406 .collect::<Vec<_>>()
5407 .into(),
5408 )
5409 }
5410 LogicalType::USmallInt => {
5411 let values =
5412 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5413 Data::UInt16(
5414 values
5415 .chunks_exact(2)
5416 .map(|item| u16::from_le_bytes(item.try_into().expect("two bytes")))
5417 .collect::<Vec<_>>()
5418 .into(),
5419 )
5420 }
5421 LogicalType::UInteger => {
5422 let values =
5423 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5424 Data::UInt32(
5425 values
5426 .chunks_exact(4)
5427 .map(|item| u32::from_le_bytes(item.try_into().expect("four bytes")))
5428 .collect::<Vec<_>>()
5429 .into(),
5430 )
5431 }
5432 LogicalType::UBigInt => {
5433 let values =
5434 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5435 Data::UInt64(
5436 values
5437 .chunks_exact(8)
5438 .map(|item| u64::from_le_bytes(item.try_into().expect("eight bytes")))
5439 .collect::<Vec<_>>()
5440 .into(),
5441 )
5442 }
5443 LogicalType::Integer | LogicalType::Date => {
5444 let values =
5445 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5446 Data::Int32(
5447 values
5448 .chunks_exact(4)
5449 .map(|item| i32::from_le_bytes(item.try_into().expect("four bytes")))
5450 .collect::<Vec<_>>()
5451 .into(),
5452 )
5453 }
5454 LogicalType::BigInt | LogicalType::Timestamp => {
5455 let values =
5456 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5457 Data::Int64(
5458 values
5459 .chunks_exact(8)
5460 .map(|item| i64::from_le_bytes(item.try_into().expect("eight bytes")))
5461 .collect::<Vec<_>>()
5462 .into(),
5463 )
5464 }
5465 LogicalType::Boolean => {
5466 let values = cur.take(rows)?;
5467 if values.iter().any(|value| *value > 1) {
5468 return Err(invalid("boolean page has another value"));
5469 }
5470 Data::Bool(values.iter().map(|value| *value == 1).collect::<Vec<_>>().into())
5471 }
5472 LogicalType::Varchar => {
5473 let offset_bytes = cur
5474 .take((rows + 1).checked_mul(4).ok_or_else(|| invalid("offset count overflow"))?)?;
5475 let offsets = offset_bytes
5476 .chunks_exact(4)
5477 .map(|part| u32::from_le_bytes(part.try_into().expect("four bytes")))
5478 .collect::<Vec<_>>();
5479 let payload = cur.take(bytes.len() - cur.at)?.to_vec();
5480 if offsets.first() != Some(&0)
5481 || offsets.last().copied().map(|last| last as usize) != Some(payload.len())
5482 || offsets.windows(2).any(|pair| pair[0] > pair[1])
5483 {
5484 return Err(invalid("string offsets do not bound the payload"));
5485 }
5486 let mut values = StringColumn::over(Buffer::from_vec(payload));
5487 for pair in offsets.windows(2) {
5488 values.push_in_place(pair[0] as usize, (pair[1] - pair[0]) as usize)?;
5489 }
5490 Data::Varlen(values)
5491 }
5492 _ => return Err(Error::not_implemented(format!("native page for {ty}"))),
5493 };
5494 if cur.at != bytes.len() {
5495 return Err(invalid("page has trailing bytes"));
5496 }
5497 Ok(Vector::flat(ty.clone(), data)?.with_validity(validity))
5498}
5499
5500#[cfg(test)]
5501mod tests {
5502 use std::fs;
5503 use std::io::{Seek, SeekFrom, Write};
5504 use std::path::PathBuf;
5505 use std::time::{SystemTime, UNIX_EPOCH};
5506
5507 use rudb_common::Value;
5508 use rudb_common::bounds::Op;
5509
5510 use super::*;
5511
5512 #[test]
5513 fn checksum_matches_fixed_vectors() {
5514 assert_eq!(checksum(b""), 0xef46_db37_51d8_e999);
5515 assert_eq!(checksum(b"a"), 0xd24e_c4f1_a98c_6e5b);
5516 assert_eq!(checksum(b"abc"), 0x44bc_2cf5_ad77_0999);
5517 }
5518
5519 fn path(label: &str) -> PathBuf {
5520 let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
5521 std::env::temp_dir().join(format!("rudb-native-{label}-{}-{stamp}.rdb", std::process::id()))
5522 }
5523
5524 #[test]
5526 fn a_read_at_an_offset_ignores_where_another_thread_left_the_cursor() {
5527 const SPANS: usize = 64;
5528 const SPAN: usize = 512;
5529 let path = path("positional");
5530 let content: Vec<u8> =
5531 (0..SPANS).flat_map(|span| std::iter::repeat_n(span as u8, SPAN)).collect();
5532 fs::write(&path, &content).expect("the file is written");
5533 let file = Arc::new(File::open(&path).expect("the file opens"));
5534 std::thread::scope(|scope| {
5535 for _ in 0..8 {
5536 let file = Arc::clone(&file);
5537 scope.spawn(move || {
5538 for _ in 0..64 {
5539 for span in 0..SPANS {
5540 let mut bytes = [0_u8; SPAN];
5541 read_at(&file, (span * SPAN) as u64, &mut bytes)
5542 .expect("the span reads");
5543 assert!(
5544 bytes.iter().all(|byte| *byte == span as u8),
5545 "span {span} came back as {}",
5546 bytes[0],
5547 );
5548 }
5549 }
5550 });
5551 }
5552 });
5553 let mut past = [0_u8; SPAN];
5554 let end = (SPANS * SPAN) as u64;
5555 let error = read_at(&file, end, &mut past).expect_err("a read past the end is refused");
5556 assert!(error.message().contains("ends before its declared length"), "{error}");
5557 drop(file);
5558 let _ = fs::remove_file(&path);
5559 }
5560
5561 #[test]
5567 fn a_writer_puts_a_page_where_it_said_it_did_wherever_the_cursor_has_got_to() {
5568 let path = path("cursor");
5569 let mut writer = Writer::create(
5570 &path,
5571 "items",
5572 vec![
5573 Field::required("id", LogicalType::Integer),
5574 Field::new("text", LogicalType::Varchar),
5575 ],
5576 )
5577 .expect("new file");
5578 writer.append(&sample()).expect("first part");
5579 writer.file.seek(SeekFrom::Start(0)).expect("the cursor goes back to the header");
5580 writer.append(&sample()).expect("second part");
5581 writer.file.seek(SeekFrom::Start(1)).expect("and somewhere useless again");
5582 writer.finish().expect("commit");
5583 let reader = Reader::open(&path).expect("reopen from disk");
5584 assert_eq!(reader.table().rows(), 6);
5585 let ids = reader.read(0, &[0]).expect("the integer page reads back");
5586 assert_eq!(ids.value_at(0, 0), Value::Integer(4));
5587 assert_eq!(ids.value_at(2, 0), Value::Integer(-2));
5588 let text = reader.read(1, &[1]).expect("the text page reads back");
5589 assert_eq!(text.value_at(1, 0), Value::Null);
5590 assert_eq!(text.value_at(2, 0), Value::Varchar("long text after a slash".into()));
5591 let end = reader.table().stripes().iter().flat_map(|stripe| {
5594 stripe
5595 .pages
5596 .iter()
5597 .map(|page| page.offset + u64::from(page.length))
5598 .chain(std::iter::once(stripe.index.offset + u64::from(stripe.index.length)))
5599 });
5600 let last = end.fold(HEADER, u64::max);
5601 let directory = fs::metadata(&path).expect("the file is there").len();
5602 assert!(last <= directory, "a page runs to {last} in a file of {directory} bytes");
5603 fs::remove_file(path).expect("remove scratch file");
5604 }
5605
5606 fn dictionary_index_len(header: &[u8; DICTIONARY_HEADER]) -> u64 {
5612 let count = u64::from(u32::from_le_bytes(header[0..4].try_into().expect("four bytes")));
5613 let blocks = u64::from(u32::from_le_bytes(header[8..12].try_into().expect("four bytes")));
5614 let bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
5615 let rank_blocks = count.div_ceil(TEXT_RANK_BLOCK as u64);
5616 DICTIONARY_HEADER as u64
5617 + offset_bytes(count as usize, bits) as u64
5618 + (blocks + rank_blocks) * 16
5619 }
5620
5621 fn last_rank_end(file: &File, offset: u64, header: &[u8; DICTIONARY_HEADER]) -> u64 {
5623 let count = u64::from(u32::from_le_bytes(header[0..4].try_into().expect("four bytes")));
5624 let blocks = u64::from(u32::from_le_bytes(header[8..12].try_into().expect("four bytes")));
5625 let bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
5626 let rank_blocks = count.div_ceil(TEXT_RANK_BLOCK as u64);
5627 let at = offset
5628 + DICTIONARY_HEADER as u64
5629 + offset_bytes(count as usize, bits) as u64
5630 + blocks * 16
5631 + (rank_blocks - 1) * 8;
5632 let mut end = [0; 8];
5633 read_at(file, at, &mut end).expect("the last rank block end");
5634 u64::from_le_bytes(end)
5635 }
5636
5637 fn sample() -> Chunk {
5638 Chunk::new(vec![
5639 Vector::from_values(
5640 LogicalType::Integer,
5641 &[Value::Integer(4), Value::Integer(9), Value::Integer(-2)],
5642 )
5643 .expect("integers"),
5644 Vector::from_values(
5645 LogicalType::Varchar,
5646 &[
5647 Value::Varchar("alpha".into()),
5648 Value::Null,
5649 Value::Varchar("long text after a slash".into()),
5650 ],
5651 )
5652 .expect("strings"),
5653 ])
5654 .expect("matching rows")
5655 }
5656
5657 fn sample_ids() -> Chunk {
5658 Chunk::new(vec![
5659 Vector::flat(LogicalType::Integer, Data::Int32(vec![7, 8, 9].into()))
5660 .expect("integers"),
5661 ])
5662 .expect("one column")
5663 }
5664
5665 #[test]
5666 fn committed_file_reopens_and_reads_only_requested_columns() {
5667 let path = path("reopen");
5668 let mut writer = Writer::create(
5669 &path,
5670 "items",
5671 vec![
5672 Field::required("id", LogicalType::Integer),
5673 Field::new("text", LogicalType::Varchar),
5674 ],
5675 )
5676 .expect("new file");
5677 writer.append(&sample()).expect("first part");
5678 writer.append(&sample()).expect("second part");
5679 writer.finish().expect("commit");
5680 let reader = Reader::open(&path).expect("reopen from disk");
5681 assert_eq!(reader.table().rows(), 6);
5682 assert_eq!(reader.table().stripes().len(), 1);
5685 assert_eq!(reader.parts(), 2);
5686 assert_eq!(reader.part_rows(0), 3);
5687 assert_eq!(reader.part_rows(1), 3);
5688 let text = reader.read(1, &[1]).expect("only text page");
5689 assert_eq!(text.width(), 1);
5690 assert_eq!(text.value_at(1, 0), Value::Null);
5691 assert_eq!(text.value_at(2, 0), Value::Varchar("long text after a slash".into()));
5692 let sparse = reader.read_sparse(1, &[1]).expect("one part without its whole page");
5693 assert_eq!(sparse.width(), 1);
5694 assert_eq!(sparse.value_at(1, 0), Value::Null);
5695 assert_eq!(sparse.value_at(2, 0), Value::Varchar("long text after a slash".into()));
5696 assert!(!reader.skips_codes(0, 1, &[0]).expect("alpha is in the stripe"));
5697 assert!(!reader.skips_codes(0, 1, &[2]).expect("long text is in the stripe"));
5698 assert!(reader.skips_codes(0, 1, &[3]).expect("unknown code is absent"));
5699 let count = reader.read(0, &[]).expect("no page is needed for count");
5700 assert_eq!(count.len(), 3);
5701 assert!(reader.skips(0, &[Probe { column: 0, op: Op::Greater, value: Bound::Int(100) }]));
5702 assert!(!reader.skips(0, &[Probe { column: 0, op: Op::Greater, value: Bound::Int(0) }]));
5703 let integers = reader.top_frequencies(0, 1).expect("valid integer synopsis").expect("kept");
5704 assert_eq!(
5705 integers,
5706 vec![(Value::Integer(-2), 2), (Value::Integer(4), 2), (Value::Integer(9), 2),]
5707 );
5708 let strings = reader.top_frequencies(1, 1).expect("valid string synopsis").expect("kept");
5709 assert_eq!(strings.len(), 3);
5710 assert!(strings.contains(&(Value::Null, 2)));
5711 assert!(strings.contains(&(Value::Varchar("alpha".into()), 2)));
5712 assert!(strings.contains(&(Value::Varchar("long text after a slash".into()), 2)));
5713 fs::remove_file(path).expect("remove scratch file");
5714 }
5715
5716 #[test]
5724 fn runs_handed_over_out_of_order_still_read_back_in_source_order() {
5725 let path = path("interleaved-runs");
5726 let mut writer =
5727 Writer::create(&path, "interleaved", vec![Field::new("v", LogicalType::BigInt)])
5728 .expect("new file");
5729 for morsel in [2_u64, 0, 3, 1] {
5730 let parts = (0..4_u64)
5731 .map(|chunk| {
5732 let first = i64::try_from(morsel * 32 + chunk * 8).expect("small");
5733 let values =
5734 (0..8_i64).map(|row| Value::BigInt(first + row)).collect::<Vec<_>>();
5735 let column =
5736 Vector::from_values(LogicalType::BigInt, &values).expect("a column");
5737 ((morsel, chunk), Chunk::new(vec![column]).expect("one column"))
5738 })
5739 .collect::<Vec<_>>();
5740 writer.append_stripe(parts).expect("a stripe");
5741 }
5742 writer.finish().expect("commit");
5743
5744 let reader = Reader::open(&path).expect("valid directory");
5745 assert_eq!(reader.table().stripes().len(), 4, "a run is a stripe of its own");
5746 assert_eq!(reader.table().rows(), 128);
5747 for part in 0..16_usize {
5748 let read = reader.read(part, &[0]).expect("a part back");
5749 for row in 0..8_usize {
5750 let want = i64::try_from(part * 8 + row).expect("small");
5751 assert_eq!(read.value_at(row, 0), Value::BigInt(want), "part {part} row {row}");
5752 }
5753 }
5754 fs::remove_file(path).expect("remove scratch file");
5755 }
5756
5757 #[test]
5760 fn runs_that_overlap_each_other_are_refused_at_commit() {
5761 let path = path("overlapping-runs");
5762 let mut writer =
5763 Writer::create(&path, "overlapping", vec![Field::new("v", LogicalType::BigInt)])
5764 .expect("new file");
5765 let one = |order: (u64, u64)| {
5766 let column =
5767 Vector::from_values(LogicalType::BigInt, &[Value::BigInt(1)]).expect("a column");
5768 (order, Chunk::new(vec![column]).expect("one column"))
5769 };
5770 writer.append_stripe(vec![one((0, 0)), one((0, 2))]).expect("a stripe");
5773 writer.append_stripe(vec![one((0, 1))]).expect("a stripe");
5774 let error = writer.finish().expect_err("the runs overlap");
5775 assert!(error.message().contains("source order"), "{error}");
5776 fs::remove_file(path).expect("remove scratch file");
5777 }
5778
5779 #[test]
5782 fn a_run_longer_than_a_stripe_is_refused() {
5783 let path = path("overlong-run");
5784 let mut writer =
5785 Writer::create(&path, "overlong", vec![Field::new("v", LogicalType::BigInt)])
5786 .expect("new file");
5787 let parts = (0..=STRIPE_PARTS)
5788 .map(|at| {
5789 let column = Vector::from_values(LogicalType::BigInt, &[Value::BigInt(1)])
5790 .expect("a column");
5791 let chunk = Chunk::new(vec![column]).expect("one column");
5792 ((0, u64::try_from(at).expect("small")), chunk)
5793 })
5794 .collect::<Vec<_>>();
5795 let error = writer.append_stripe(parts).expect_err("one part too many");
5796 assert!(error.message().contains("more parts than it holds"), "{error}");
5797 fs::remove_file(path).expect("remove scratch file");
5798 }
5799
5800 #[test]
5806 fn parts_past_the_stripe_bound_start_a_new_stripe() {
5807 let path = path("stripe-bound");
5808 let mut writer = Writer::create(
5809 &path,
5810 "items",
5811 vec![
5812 Field::required("id", LogicalType::Integer),
5813 Field::new("text", LogicalType::Varchar),
5814 ],
5815 )
5816 .expect("new file");
5817 let parts = STRIPE_PARTS * 2 + 3;
5818 for part in 0..parts {
5819 let id = part as i32;
5820 let chunk = Chunk::new(vec![
5821 Vector::from_values(
5822 LogicalType::Integer,
5823 &[Value::Integer(id), Value::Integer(-id)],
5824 )
5825 .expect("integers"),
5826 Vector::from_values(
5827 LogicalType::Varchar,
5828 &[Value::Varchar(format!("value {part}")), Value::Null],
5829 )
5830 .expect("strings"),
5831 ])
5832 .expect("matching rows");
5833 writer.append(&chunk).expect("one part");
5834 }
5835 writer.finish().expect("commit");
5836
5837 let reader = Reader::open(&path).expect("reopen from disk");
5838 assert_eq!(reader.parts(), parts);
5839 assert_eq!(reader.table().rows(), parts * 2);
5840 assert_eq!(reader.table().stripes().len(), parts.div_ceil(STRIPE_PARTS));
5841 assert_eq!(reader.table().stripes()[0].parts(), STRIPE_PARTS);
5842 assert_eq!(reader.table().stripes()[0].rows(), STRIPE_PARTS * 2);
5843 assert_eq!(reader.table().stripes()[2].parts(), 3);
5844 for part in (0..parts).rev() {
5847 let dense = reader.read(part, &[0, 1]).expect("a whole page read");
5848 let sparse = reader.read_sparse(part, &[0, 1]).expect("one part read");
5849 for chunk in [&dense, &sparse] {
5850 assert_eq!(chunk.len(), 2, "part {part} has its own row count");
5851 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
5852 assert_eq!(chunk.value_at(1, 0), Value::Integer(-(part as i32)));
5853 assert_eq!(chunk.value_at(0, 1), Value::Varchar(format!("value {part}")));
5854 assert_eq!(chunk.value_at(1, 1), Value::Null);
5855 }
5856 }
5857 let above = [Probe { column: 0, op: Op::Greater, value: Bound::Int(100) }];
5860 assert!(reader.skips(0, &above), "the first stripe stops at 63");
5861 assert!(!reader.skips(STRIPE_PARTS * 2, &above), "the third stripe reaches 130");
5862 fs::remove_file(path).expect("remove scratch file");
5863 }
5864
5865 fn scattered(n: i64) -> i64 {
5867 n.wrapping_mul(-7_046_029_254_386_353_131)
5868 }
5869
5870 #[test]
5876 fn a_part_is_skipped_when_its_sieve_does_not_hold_the_constant() {
5877 let path = path("sieve-skip");
5878 let mut writer =
5879 Writer::create(&path, "hits", vec![Field::required("id", LogicalType::BigInt)])
5880 .expect("new file");
5881 let parts = STRIPE_PARTS + 3;
5882 let per_part = 128;
5886 for part in 0..parts {
5887 let held: Vec<Value> = (0..per_part)
5888 .map(|row| Value::BigInt(scattered((part * per_part + row) as i64)))
5889 .collect();
5890 let chunk =
5891 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
5892 .expect("one column");
5893 writer.append(&chunk).expect("one part");
5894 }
5895 writer.finish().expect("commit");
5896
5897 let reader = Reader::open(&path).expect("reopen from disk");
5898 let probe = |value: i64| Probe {
5899 column: 0,
5900 op: Op::Equal,
5901 value: Bound::Int(i128::from(scattered(value))),
5902 };
5903 for wanted in [0_i64, (per_part + 1) as i64, (parts * per_part - 1) as i64] {
5904 let tests = [probe(wanted)];
5905 let kept: Vec<usize> = (0..parts).filter(|&part| !reader.skips(part, &tests)).collect();
5906 let home = wanted as usize / per_part;
5907 assert!(kept.contains(&home), "the part holding {wanted} is read");
5908 assert!(kept.len() <= 2, "{wanted} keeps {kept:?}, which is more than one stray part");
5912 }
5913 let absent = [probe((parts * per_part) as i64 + 1)];
5914 let kept = (0..parts).filter(|&part| !reader.skips(part, &absent)).count();
5915 assert!(kept <= 1, "{kept} parts of {parts} kept a value no part holds");
5916 let tests = [probe(0)];
5919 assert!(
5920 reader.table().stripes().iter().all(|stripe| !stripe.zone.skips(&tests)),
5921 "the bounds rule out no stripe at all"
5922 );
5923 fs::remove_file(path).expect("remove scratch file");
5924 }
5925
5926 #[test]
5932 fn a_part_is_skipped_when_its_own_bounds_rule_out_a_comparison_the_stripe_keeps() {
5933 let path = path("part-range-skip");
5934 let mut writer =
5935 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
5936 .expect("new file");
5937 let parts = STRIPE_PARTS + 3;
5938 let per_part = 128;
5939 for part in 0..parts {
5940 let held: Vec<Value> = (0..per_part)
5944 .map(|row| {
5945 Value::BigInt((part * 1_000) as i64 + (scattered(row as i64).rem_euclid(900)))
5946 })
5947 .collect();
5948 let chunk =
5949 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
5950 .expect("one column");
5951 writer.append(&chunk).expect("one part");
5952 }
5953 writer.finish().expect("commit");
5954
5955 let reader = Reader::open(&path).expect("reopen from disk");
5956 let under = [Probe { column: 0, op: Op::Less, value: Bound::Int(3_000) }];
5957 let kept: Vec<usize> = (0..parts).filter(|&part| !reader.skips(part, &under)).collect();
5958 assert_eq!(kept, vec![0, 1, 2], "only the three parts that start under three thousand");
5959 assert!(!reader.stripe_skips(0, &under), "the stripe reaches from zero and keeps itself");
5961 fs::remove_file(path).expect("remove scratch file");
5962 }
5963
5964 #[test]
5968 fn a_part_is_waved_through_when_its_own_bounds_pass_a_comparison_the_stripe_cannot() {
5969 let path = path("part-range-certain");
5970 let mut writer =
5971 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
5972 .expect("new file");
5973 let parts = STRIPE_PARTS + 3;
5974 let per_part = 128;
5975 for part in 0..parts {
5976 let held: Vec<Value> = (0..per_part)
5977 .map(|row| {
5978 Value::BigInt((part * 1_000) as i64 + (scattered(row as i64).rem_euclid(900)))
5979 })
5980 .collect();
5981 let chunk =
5982 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
5983 .expect("one column");
5984 writer.append(&chunk).expect("one part");
5985 }
5986 writer.finish().expect("commit");
5987
5988 let reader = Reader::open(&path).expect("reopen from disk");
5989 let under = [Probe { column: 0, op: Op::Less, value: Bound::Int(3_000) }];
5990 let waved: Vec<usize> = (0..parts).filter(|&part| reader.certain(part, &under)).collect();
5991 assert_eq!(waved, vec![0, 1, 2], "the three parts that end under three thousand");
5992 assert!(!reader.stripe_skips(0, &under), "the stripe straddles the comparison");
5995 fs::remove_file(path).expect("remove scratch file");
5996 }
5997
5998 #[test]
6001 fn a_stripe_of_one_part_writes_no_range_page_and_a_stripe_of_many_does() {
6002 for (parts, wanted) in [(1_usize, false), (STRIPE_PARTS, true)] {
6003 let path = path("part-range-page");
6004 let mut writer =
6005 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
6006 .expect("new file");
6007 for part in 0..parts {
6008 let held: Vec<Value> = (0..128)
6009 .map(|row| {
6010 Value::BigInt((part * 1_000) as i64 + scattered(row as i64).rem_euclid(900))
6011 })
6012 .collect();
6013 let chunk = Chunk::new(vec![
6014 Vector::from_values(LogicalType::BigInt, &held).expect("numbers"),
6015 ])
6016 .expect("one column");
6017 writer.append(&chunk).expect("one part");
6018 }
6019 writer.finish().expect("commit");
6020 let reader = Reader::open(&path).expect("reopen from disk");
6021 let bytes = reader.layout().columns[0].part_ranges;
6022 assert_eq!(bytes > 0, wanted, "{parts} parts wrote {bytes} bytes of ranges");
6023 fs::remove_file(path).expect("remove scratch file");
6024 }
6025 }
6026
6027 #[test]
6030 fn a_string_end_that_is_cut_down_still_covers_the_value_it_came_from() {
6031 let long = vec![b'a'; PART_BOUND_BYTES * 2];
6032 let low = shortened(Some(Bound::Bytes(long.clone())), false).expect("a low end");
6033 let high = shortened(Some(Bound::Bytes(long.clone())), true).expect("a high end");
6034 let Bound::Bytes(low) = low else { panic!("a string stays a string") };
6035 let Bound::Bytes(high) = high else { panic!("a string stays a string") };
6036 assert!(low.len() <= PART_BOUND_BYTES && high.len() <= PART_BOUND_BYTES);
6037 assert!(low.as_slice() <= long.as_slice(), "the low end is at or under the value");
6038 assert!(high.as_slice() >= long.as_slice(), "the high end is at or over the value");
6039 }
6040
6041 #[test]
6044 fn a_string_end_with_no_room_to_step_up_gives_up_the_bound() {
6045 let long = vec![u8::MAX; PART_BOUND_BYTES * 2];
6046 assert_eq!(shortened(Some(Bound::Bytes(long.clone())), true), None);
6047 let low = shortened(Some(Bound::Bytes(long)), false).expect("a low end is still a prefix");
6048 assert_eq!(low, Bound::Bytes(vec![u8::MAX; PART_BOUND_BYTES]));
6049 }
6050
6051 #[test]
6061 fn a_sieve_larger_than_the_part_it_indexes_is_not_written() {
6062 let path = path("sieve-pays");
6063 let fields = vec![
6064 Field::required("spread", LogicalType::BigInt),
6065 Field::required("repeated", LogicalType::BigInt),
6066 ];
6067 let mut writer = Writer::create(&path, "hits", fields).expect("new file");
6068 let parts = 3;
6069 let per_part = 1024;
6070 for part in 0..parts {
6071 let base = (part * per_part) as i64;
6072 let spread: Vec<Value> =
6073 (0..per_part).map(|row| Value::BigInt(scattered(base + row as i64))).collect();
6074 let repeated: Vec<Value> =
6075 (0..per_part).map(|row| Value::BigInt(scattered((row / 256) as i64))).collect();
6076 let chunk = Chunk::new(vec![
6077 Vector::from_values(LogicalType::BigInt, &spread).expect("numbers"),
6078 Vector::from_values(LogicalType::BigInt, &repeated).expect("numbers"),
6079 ])
6080 .expect("two columns");
6081 writer.append(&chunk).expect("one part");
6082 }
6083 writer.finish().expect("commit");
6084
6085 let reader = Reader::open(&path).expect("reopen from disk");
6086 let layout = reader.layout();
6087 let spread = &layout.columns[0];
6088 let repeated = &layout.columns[1];
6089 assert!(spread.sieves > 0, "a column whose parts are worth a filter keeps one");
6090 assert_eq!(
6091 repeated.sieves, 0,
6092 "a column whose filter costs more than its parts keeps none"
6093 );
6094 for column in &layout.columns {
6097 assert!(
6098 column.sieves < column.pages,
6099 "{} spends {} on sieves over {} of data",
6100 column.name,
6101 column.sieves,
6102 column.pages
6103 );
6104 }
6105 let absent = [Probe {
6107 column: 0,
6108 op: Op::Equal,
6109 value: Bound::Int(i128::from(scattered((parts * per_part) as i64 + 1))),
6110 }];
6111 assert!((0..parts).all(|part| reader.skips(part, &absent)), "no part holds it");
6112 fs::remove_file(path).expect("remove scratch file");
6113 }
6114
6115 #[test]
6121 fn a_damaged_sieve_page_is_read_through_rather_than_refused() {
6122 let path = path("sieve-damaged");
6123 let mut writer =
6124 Writer::create(&path, "hits", vec![Field::required("id", LogicalType::BigInt)])
6125 .expect("new file");
6126 let rows = 128;
6127 let held: Vec<Value> = (0..rows).map(|row| Value::BigInt(scattered(row))).collect();
6128 let chunk =
6129 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6130 .expect("one column");
6131 writer.append(&chunk).expect("one part");
6132 writer.finish().expect("commit");
6133
6134 let page =
6135 Reader::open(&path).expect("reopen").table.stripes[0].sieves[0].expect("a sieve page");
6136 let mut file = OpenOptions::new().write(true).open(&path).expect("open the sieve page");
6137 file.seek(SeekFrom::Start(page.offset + u64::from(page.length) - 1)).expect("seek");
6138 file.write_all(&[0xff]).expect("damage one byte");
6139 drop(file);
6140
6141 let reader = Reader::open(&path).expect("reopen the damaged file");
6142 let absent =
6143 [Probe { column: 0, op: Op::Equal, value: Bound::Int(i128::from(scattered(99))) }];
6144 assert!(!reader.skips(0, &absent), "a sieve that cannot be read skips nothing");
6145 assert_eq!(
6146 reader.read(0, &[0]).expect("the rows are untouched").len(),
6147 usize::try_from(rows).expect("a small count")
6148 );
6149 fs::remove_file(path).expect("remove scratch file");
6150 }
6151
6152 #[test]
6163 fn workers_that_want_the_same_stripe_read_it_once() {
6164 let path = path("single-flight");
6165 let mut writer =
6166 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6167 .expect("new file");
6168 for part in 0..STRIPE_PARTS {
6169 let id = part as i32;
6170 let chunk = Chunk::new(vec![
6171 Vector::from_values(
6172 LogicalType::Integer,
6173 &[Value::Integer(id), Value::Integer(-id)],
6174 )
6175 .expect("integers"),
6176 ])
6177 .expect("matching rows");
6178 writer.append(&chunk).expect("one part");
6179 }
6180 writer.finish().expect("commit");
6181
6182 let reader = Reader::open(&path).expect("reopen from disk");
6183 assert_eq!(reader.table().stripes().len(), 1, "one stripe is the point of the test");
6184 let barrier = std::sync::Barrier::new(8);
6185 std::thread::scope(|scope| {
6186 for worker in 0..8 {
6187 let reader = &reader;
6188 let barrier = &barrier;
6189 scope.spawn(move || {
6190 barrier.wait();
6191 for part in (worker..STRIPE_PARTS).step_by(8) {
6192 let chunk = reader.read(part, &[0]).expect("a whole page read");
6193 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6194 assert_eq!(chunk.value_at(1, 0), Value::Integer(-(part as i32)));
6195 }
6196 });
6197 }
6198 });
6199 assert_eq!(reader.pages.load(Atomic::Relaxed), 1, "one stripe, one page read, whoever won");
6200 fs::remove_file(path).expect("remove scratch file");
6201 }
6202
6203 #[test]
6216 fn opening_costs_the_same_over_a_thousand_times_the_rows() {
6217 let opened = |label: &str, rows_per_part: i32| {
6218 let path = path(label);
6219 let mut writer =
6220 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6221 .expect("new file");
6222 for part in 0..STRIPE_PARTS * 3 {
6223 let values = (0..rows_per_part)
6227 .map(|row| {
6228 Value::Integer((part as i32 * rows_per_part + row).wrapping_mul(2_654_435))
6229 })
6230 .collect::<Vec<_>>();
6231 let chunk = Chunk::new(vec![
6232 Vector::from_values(LogicalType::Integer, &values).expect("integers"),
6233 ])
6234 .expect("matching rows");
6235 writer.append(&chunk).expect("one part");
6236 }
6237 writer.finish().expect("commit");
6238 let reader = Reader::open(&path).expect("reopen from disk");
6239 let size = fs::metadata(&path).expect("the file is there").len();
6240 let out = (reader.reads(), reader.table().stripes().len(), size);
6241 fs::remove_file(path).expect("remove scratch file");
6242 out
6243 };
6244
6245 let (thin, thin_stripes, thin_size) = opened("open-thin", 1);
6246 let (fat, fat_stripes, fat_size) = opened("open-fat", 1000);
6247 assert_eq!(
6248 thin_stripes, fat_stripes,
6249 "the same stripe count is what makes this a fair ask"
6250 );
6251 assert!(
6252 fat_size > thin_size * 50,
6253 "the fat file has to actually be larger, and it is {fat_size} against {thin_size}"
6254 );
6255
6256 assert_eq!(thin.opening.reads, fat.opening.reads, "the same reads either way");
6257 assert_eq!(thin.pages, 0, "opening read a page");
6258 assert_eq!(fat.pages, 0, "opening read a page");
6259 assert_eq!(thin.indexes, 0, "opening read an index");
6260 assert_eq!(fat.indexes, 0, "opening read an index");
6261 assert!(
6264 fat.opening.bytes < thin.opening.bytes * 2,
6265 "opening the thin file read {} bytes and the fat one read {}",
6266 thin.opening.bytes,
6267 fat.opening.bytes
6268 );
6269 }
6270
6271 #[test]
6279 fn two_opens_of_one_file_cost_the_same_and_the_second_is_not_cheaper() {
6280 let path = path("open-twice");
6281 let mut writer =
6282 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6283 .expect("new file");
6284 for part in 0..STRIPE_PARTS * 3 {
6285 let chunk = Chunk::new(vec![
6286 Vector::from_values(LogicalType::Integer, &[Value::Integer(part as i32)])
6287 .expect("integers"),
6288 ])
6289 .expect("matching rows");
6290 writer.append(&chunk).expect("one part");
6291 }
6292 writer.finish().expect("commit");
6293
6294 let first = Reader::open(&path).expect("open");
6295 for part in 0..first.parts() {
6298 first.read(part, &[0]).expect("a part");
6299 }
6300 assert!(first.reads().pages > 0, "the scan has to have read something");
6301 let second = Reader::open(&path).expect("open again");
6302
6303 assert_eq!(first.reads().opening, second.reads().opening);
6304 assert_eq!(
6305 second.reads().pages,
6306 0,
6307 "the second open read a page off the back of the first"
6308 );
6309 assert_eq!(second.reads().indexes, 0, "the second open read an index it inherited");
6310 fs::remove_file(path).expect("remove scratch file");
6311 }
6312
6313 #[test]
6321 fn an_index_is_read_once_per_stripe_however_often_the_page_is_evicted() {
6322 let path = path("index-cache");
6323 let mut writer =
6324 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6325 .expect("new file");
6326 let parts = STRIPE_PARTS * (CACHED_STRIPES_PER_COLUMN + 2);
6327 for part in 0..parts {
6328 let id = part as i32;
6329 let chunk = Chunk::new(vec![
6330 Vector::from_values(LogicalType::Integer, &[Value::Integer(id)]).expect("integers"),
6331 ])
6332 .expect("matching rows");
6333 writer.append(&chunk).expect("one part");
6334 }
6335 writer.finish().expect("commit");
6336
6337 let reader = Reader::open(&path).expect("reopen from disk");
6338 let stripes = reader.table().stripes().len();
6339 assert!(stripes > CACHED_STRIPES_PER_COLUMN, "the page cache has to be too small for this");
6340 for _ in 0..2 {
6342 for part in 0..parts {
6343 let chunk = reader.read(part, &[0]).expect("a part");
6344 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6345 }
6346 }
6347 assert_eq!(reader.indexes.load(Atomic::Relaxed), stripes, "one index read per stripe");
6348 assert!(
6349 reader.pages.load(Atomic::Relaxed) > stripes,
6350 "the pages are the ones that get read again, which is what makes the index count mean \
6351 something"
6352 );
6353 fs::remove_file(path).expect("remove scratch file");
6354 }
6355
6356 #[test]
6365 fn a_worker_per_stripe_reads_its_page_once_when_the_cache_was_told_to_expect_it() {
6366 let workers = CACHED_STRIPES_PER_COLUMN + 4;
6367 let path = path("stripe-per-worker");
6368 let mut writer =
6369 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6370 .expect("new file");
6371 for part in 0..STRIPE_PARTS * workers {
6372 let chunk = Chunk::new(vec![
6373 Vector::from_values(LogicalType::Integer, &[Value::Integer(part as i32)])
6374 .expect("integers"),
6375 ])
6376 .expect("matching rows");
6377 writer.append(&chunk).expect("one part");
6378 }
6379 writer.finish().expect("commit");
6380
6381 let read = |told: bool| {
6382 let reader = Reader::open(&path).expect("reopen from disk");
6383 assert_eq!(reader.table().stripes().len(), workers, "a stripe per worker");
6384 if told {
6385 reader.keep_stripes(workers);
6386 }
6387 let barrier = std::sync::Barrier::new(workers);
6388 std::thread::scope(|scope| {
6389 for (worker, run) in reader.stripe_parts().into_iter().enumerate() {
6390 let reader = &reader;
6391 let barrier = &barrier;
6392 scope.spawn(move || {
6393 for part in run {
6394 barrier.wait();
6395 let chunk = reader.read(part, &[0]).expect("a part of my own stripe");
6396 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6397 }
6398 assert!(worker < workers);
6399 });
6400 }
6401 });
6402 reader.pages.load(Atomic::Relaxed)
6403 };
6404
6405 assert_eq!(read(true), workers, "one page read per stripe and no more");
6406 assert!(read(false) > workers, "a cache that small is read again on every part");
6407 fs::remove_file(path).expect("remove scratch file");
6408 }
6409
6410 #[test]
6415 fn a_damaged_index_page_is_an_error() {
6416 let path = path("damaged-index");
6417 let mut writer =
6418 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6419 .expect("new file");
6420 writer.append(&sample_ids()).expect("first part");
6421 writer.append(&sample_ids()).expect("second part");
6422 writer.finish().expect("commit");
6423
6424 let reader = Reader::open(&path).expect("valid directory");
6425 let index = reader.table.stripes[0].index;
6426 let mut byte = [0; 1];
6427 read_at(&reader.file, index.offset, &mut byte).expect("the first part length");
6428 let mut file = OpenOptions::new().write(true).open(&path).expect("open index page");
6429 file.seek(SeekFrom::Start(index.offset)).expect("index start");
6430 file.write_all(&[!byte[0]]).expect("damage the first part length");
6431 let error = reader.read(1, &[0]).expect_err("a damaged index must not be used");
6432 assert!(error.message().contains("index page section checksum differs"), "{error}");
6433 fs::remove_file(path).expect("remove scratch file");
6434 }
6435
6436 #[test]
6443 fn every_integer_width_round_trips_through_a_page() {
6444 let path = path("integer-widths");
6445 let columns = [
6446 (LogicalType::TinyInt, vec![Value::TinyInt(i8::MIN), Value::TinyInt(i8::MAX)]),
6447 (LogicalType::UTinyInt, vec![Value::UTinyInt(0), Value::UTinyInt(u8::MAX)]),
6448 (LogicalType::SmallInt, vec![Value::SmallInt(i16::MIN), Value::SmallInt(i16::MAX)]),
6449 (LogicalType::USmallInt, vec![Value::USmallInt(0), Value::USmallInt(u16::MAX)]),
6450 (LogicalType::Integer, vec![Value::Integer(i32::MIN), Value::Integer(i32::MAX)]),
6451 (LogicalType::UInteger, vec![Value::UInteger(0), Value::UInteger(u32::MAX)]),
6452 (LogicalType::BigInt, vec![Value::BigInt(i64::MIN), Value::BigInt(i64::MAX)]),
6453 (LogicalType::UBigInt, vec![Value::UBigInt(0), Value::UBigInt(u64::MAX)]),
6454 ];
6455 let fields = columns
6456 .iter()
6457 .enumerate()
6458 .map(|(at, (ty, _))| Field::required(format!("c{at}"), ty.clone()))
6459 .collect::<Vec<_>>();
6460 let vectors = columns
6461 .iter()
6462 .map(|(ty, values)| Vector::from_values(ty.clone(), values).expect("a vector"))
6463 .collect::<Vec<_>>();
6464 let mut writer = Writer::create(&path, "widths", fields).expect("new file");
6465 writer.append(&Chunk::new(vectors).expect("matching rows")).expect("one stripe");
6466 writer.finish().expect("commit");
6467
6468 let reader = Reader::open(&path).expect("reopen from disk");
6469 let wanted = (0..columns.len()).collect::<Vec<_>>();
6470 let read = reader.read(0, &wanted).expect("every column");
6471 assert_eq!(read.len(), 2);
6472 for (at, (ty, values)) in columns.iter().enumerate() {
6474 assert_eq!(read.value_at(0, at), values[0], "the low end of {ty}");
6475 assert_eq!(read.value_at(1, at), values[1], "the high end of {ty}");
6476 }
6477 fs::remove_file(path).expect("remove scratch file");
6478 }
6479
6480 #[test]
6481 fn numeric_frequency_candidates_keep_bounded_row_ordinals() {
6482 let path = path("frequency-ordinals");
6483 let mut writer =
6484 Writer::create(&path, "items", vec![Field::required("id", LogicalType::BigInt)])
6485 .expect("new file");
6486 let mut values = Vec::new();
6487 for leader in 0..10_i64 {
6488 values.extend(std::iter::repeat_n(leader, 100));
6489 }
6490 values.extend(1_000_i64..41_000);
6491 for part in values.chunks(1_024) {
6492 let vector = Vector::flat(LogicalType::BigInt, Data::Int64(part.to_vec().into()))
6493 .expect("big integers");
6494 writer.append(&Chunk::new(vec![vector]).expect("one column")).expect("one stripe");
6495 }
6496 writer.finish().expect("commit");
6497
6498 let reader = Reader::open(&path).expect("reopen from disk");
6499 let occurrences =
6500 reader.frequency_occurrences(0).expect("valid metadata").expect("bounded ordinals");
6501 assert!(occurrences.omitted_max < 100);
6502 assert!(occurrences.ordinals.len() <= FREQUENCY_ORDINALS);
6503 assert!(occurrences.ordinals.windows(2).all(|pair| pair[0] < pair[1]));
6504 assert_eq!(&occurrences.ordinals[..1_000], &(0_u64..1_000).collect::<Vec<_>>());
6505 fs::remove_file(path).expect("remove scratch file");
6506 }
6507
6508 #[test]
6514 fn a_file_from_another_format_says_which_format_it_is() {
6515 let older = path("older-format");
6516 let mut writer =
6517 Writer::create(&older, "items", vec![Field::new("id", LogicalType::Integer)])
6518 .expect("new file");
6519 let chunk = Chunk::new(vec![
6520 Vector::flat(LogicalType::Integer, Data::Int32(vec![1, 2, 3].into()))
6521 .expect("integers"),
6522 ])
6523 .expect("chunk");
6524 writer.append(&chunk).expect("page written");
6525 writer.finish().expect("commit");
6526
6527 let mut file = OpenOptions::new().write(true).open(&older).expect("open for the header");
6528 file.seek(SeekFrom::Start(8)).expect("the version follows the magic");
6529 file.write_all(&(FORMAT - 1).to_le_bytes()).expect("write an older version");
6530 drop(file);
6531 let complaint = Reader::open(&older).expect_err("an older format is refused").to_string();
6532 assert!(complaint.contains(&format!("format {}", FORMAT - 1)), "{complaint}");
6533 assert!(complaint.contains(&format!("format {FORMAT}")), "{complaint}");
6534
6535 let mut file = OpenOptions::new().write(true).open(&older).expect("open for the header");
6536 file.seek(SeekFrom::Start(0)).expect("the magic is first");
6537 file.write_all(b"NOTRUDB!").expect("write another engine's magic");
6538 drop(file);
6539 let complaint = Reader::open(&older).expect_err("a foreign file is refused").to_string();
6540 assert!(complaint.contains("magic"), "{complaint}");
6541 assert!(!complaint.contains("format"), "a version has nothing to do with it: {complaint}");
6542 fs::remove_file(older).expect("remove scratch file");
6543 }
6544
6545 #[test]
6546 fn an_unfinished_or_damaged_file_does_not_answer_with_partial_rows() {
6547 let unfinished = path("unfinished");
6548 let mut writer =
6549 Writer::create(&unfinished, "items", vec![Field::new("id", LogicalType::Integer)])
6550 .expect("new file");
6551 let chunk = Chunk::new(vec![
6552 Vector::flat(LogicalType::Integer, Data::Int32(vec![1, 2, 3].into()))
6553 .expect("integers"),
6554 ])
6555 .expect("chunk");
6556 writer.append(&chunk).expect("page written");
6557 drop(writer);
6558 assert!(Reader::open(&unfinished).is_err(), "no directory was committed");
6559 fs::remove_file(unfinished).expect("remove scratch file");
6560
6561 let damaged = path("damaged");
6562 let mut writer =
6563 Writer::create(&damaged, "items", vec![Field::new("id", LogicalType::Integer)])
6564 .expect("new file");
6565 writer.append(&chunk).expect("page written");
6566 writer.finish().expect("commit");
6567 let reader = Reader::open(&damaged).expect("valid directory");
6568 let mut file =
6569 OpenOptions::new().write(true).open(&damaged).expect("open for a damaged page");
6570 file.seek(SeekFrom::Start(HEADER + 1)).expect("inside first page");
6571 file.write_all(&[255]).expect("damage one byte");
6572 assert!(reader.read(0, &[0]).is_err(), "page checksum rejects corruption");
6573 fs::remove_file(damaged).expect("remove scratch file");
6574 }
6575
6576 #[test]
6577 fn damaged_lazy_dictionary_payload_is_an_error() {
6578 let path = path("damaged-dictionary");
6579 let mut writer = Writer::create(
6580 &path,
6581 "items",
6582 vec![
6583 Field::required("id", LogicalType::Integer),
6584 Field::new("text", LogicalType::Varchar),
6585 ],
6586 )
6587 .expect("new file");
6588 writer.append(&sample()).expect("stripe written");
6589 writer.finish().expect("commit");
6590
6591 let reader = Reader::open(&path).expect("valid directory");
6592 let dictionary = reader.table.dictionaries[1].expect("string dictionary page");
6593 let mut header = [0; DICTIONARY_HEADER];
6596 read_at(&reader.file, dictionary.offset, &mut header).expect("dictionary header");
6597 let index_len = dictionary_index_len(&header);
6598 let rank_len = last_rank_end(&reader.file, dictionary.offset, &header);
6599 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
6600 file.seek(SeekFrom::Start(dictionary.offset + index_len + rank_len))
6601 .expect("inside dictionary payload");
6602 file.write_all(&[255]).expect("damage dictionary payload");
6603
6604 let chunk = reader.read(0, &[1]).expect("code page and dictionary index remain valid");
6605 let error =
6606 chunk.validate_external().expect_err("payload corruption must reach the caller");
6607 assert!(error.message().contains("payload checksum differs"), "{error}");
6608 fs::remove_file(path).expect("remove scratch file");
6609 }
6610
6611 #[test]
6618 fn a_dictionary_over_many_blocks_checks_every_block_of_it() {
6619 let path = path("dictionary-blocks");
6620 let value =
6621 |row: usize| format!("{row:07} a value long enough to be worth a payload block");
6622 let parts = 30;
6623 let per_part = 1000;
6624 let mut writer =
6625 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
6626 .expect("new file");
6627 for part in 0..parts {
6628 let values = (0..per_part)
6629 .map(|row| Value::Varchar(value(part * per_part + row)))
6630 .collect::<Vec<_>>();
6631 let chunk = Chunk::new(vec![
6632 Vector::from_values(LogicalType::Varchar, &values).expect("strings"),
6633 ])
6634 .expect("matching rows");
6635 writer.append(&chunk).expect("a part");
6636 }
6637 writer.finish().expect("commit");
6638
6639 let reader = Reader::open(&path).expect("reopen from disk");
6640 let dictionary = reader.table.dictionaries[0].expect("string dictionary page");
6641 assert!(
6642 parts * per_part > TEXT_PAYLOAD_VALUES * 4,
6643 "the dictionary has to be several blocks for this to be testing anything"
6644 );
6645 for part in [0, parts - 1] {
6646 let chunk = reader.read(part, &[0]).expect("a part");
6647 chunk.validate_external().expect("every payload block checks out");
6648 assert_eq!(chunk.value_at(0, 0), Value::Varchar(value(part * per_part)));
6649 }
6650
6651 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
6652 file.seek(SeekFrom::Start(dictionary.offset + u64::from(dictionary.length) - 4))
6653 .expect("the last bytes of the page are payload");
6654 file.write_all(&[255]).expect("damage the last payload block");
6655 let reader = Reader::open(&path).expect("the directory and the index are untouched");
6656 let chunk = reader.read(parts - 1, &[0]).expect("the code page remains valid");
6657 let error = chunk.validate_external().expect_err("the damage must reach the caller");
6658 assert!(error.message().contains("payload checksum differs"), "{error}");
6659 fs::remove_file(path).expect("remove scratch file");
6660 }
6661
6662 #[test]
6672 fn values_of_different_lengths_read_back_out_of_packed_offsets() {
6673 let path = path("dictionary-offsets");
6674 let value = |row: usize| {
6675 if row % 511 == 3 { String::new() } else { "x".repeat(row % 97) + &format!("{row:05}") }
6676 };
6677 let rows = 5_000;
6678 let mut writer =
6679 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
6680 .expect("new file");
6681 let values = (0..rows).map(|row| Value::Varchar(value(row))).collect::<Vec<_>>();
6682 for part in values.chunks(1_000) {
6683 let chunk =
6684 Chunk::new(vec![Vector::from_values(LogicalType::Varchar, part).expect("strings")])
6685 .expect("matching rows");
6686 writer.append(&chunk).expect("a part");
6687 }
6688 writer.finish().expect("commit");
6689
6690 let reader = Reader::open(&path).expect("reopen from disk");
6691 assert!(
6692 rows > TEXT_PAYLOAD_VALUES * 4,
6693 "the dictionary has to be several blocks for this to be testing anything"
6694 );
6695 for part in 0..rows / 1_000 {
6696 let chunk = reader.read(part, &[0]).expect("a part");
6697 for row in 0..1_000 {
6698 let row = part * 1_000 + row;
6699 assert_eq!(
6700 chunk.value_at(row % 1_000, 0),
6701 Value::Varchar(value(row)),
6702 "value {row}"
6703 );
6704 }
6705 }
6706 fs::remove_file(path).expect("remove scratch file");
6707 }
6708
6709 #[test]
6721 fn a_global_dictionary_is_opened_once_however_many_workers_ask_at_once() {
6722 let path = path("dictionary-once");
6723 let parts = 8;
6724 let per_part = 500;
6725 let value =
6726 |row: usize| format!("{row:07} a value long enough to be worth a payload block");
6727 let mut writer =
6728 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
6729 .expect("new file");
6730 for part in 0..parts {
6731 let values = (0..per_part)
6732 .map(|row| Value::Varchar(value(part * per_part + row)))
6733 .collect::<Vec<_>>();
6734 let chunk = Chunk::new(vec![
6735 Vector::from_values(LogicalType::Varchar, &values).expect("strings"),
6736 ])
6737 .expect("matching rows");
6738 writer.append(&chunk).expect("a part");
6739 }
6740 writer.finish().expect("commit");
6741
6742 let reader = Reader::open(&path).expect("reopen from disk");
6743 assert!(reader.table.dictionaries[0].is_some(), "the column has to have one to share");
6744 assert_eq!(reader.reads().dictionaries, 0, "opening the file does not open a dictionary");
6745
6746 let workers = 16;
6747 let gate = std::sync::Barrier::new(workers);
6748 std::thread::scope(|scope| {
6749 for worker in 0..workers {
6750 let reader = reader.clone();
6751 let gate = &gate;
6752 scope.spawn(move || {
6753 gate.wait();
6754 let chunk = reader.read(worker % parts, &[0]).expect("a part");
6755 assert_eq!(
6756 chunk.value_at(0, 0),
6757 Value::Varchar(value((worker % parts) * per_part))
6758 );
6759 });
6760 }
6761 });
6762
6763 assert_eq!(reader.reads().dictionaries, 1, "sixteen workers, one dictionary, one open");
6764 fs::remove_file(path).expect("remove scratch file");
6765 }
6766
6767 #[test]
6772 fn a_damaged_sorted_order_is_an_error() {
6773 let path = path("damaged-order");
6774 let mut writer = Writer::create(
6775 &path,
6776 "items",
6777 vec![
6778 Field::required("id", LogicalType::Integer),
6779 Field::new("text", LogicalType::Varchar),
6780 ],
6781 )
6782 .expect("new file");
6783 writer.append(&sample()).expect("stripe written");
6784 writer.finish().expect("commit");
6785
6786 let reader = Reader::open(&path).expect("valid directory");
6787 let page = reader.table.dictionaries[1].expect("string dictionary page");
6788 let mut header = [0; DICTIONARY_HEADER];
6789 read_at(&reader.file, page.offset, &mut header).expect("dictionary header");
6790 let index_len = dictionary_index_len(&header);
6791 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
6792 file.seek(SeekFrom::Start(page.offset + index_len)).expect("the first head");
6793 file.write_all(&[255]).expect("damage the order");
6794
6795 let dictionary = reader.dictionary(1).expect("read").expect("a string column has one");
6796 let error = dictionary.compare_rank(0, b"anything").expect_err("a damaged order is caught");
6797 assert!(error.message().contains("rank checksum differs"), "{error}");
6798 fs::remove_file(path).expect("remove scratch file");
6799 }
6800
6801 #[test]
6805 fn a_global_dictionary_carries_the_sorted_order_of_its_values() {
6806 let spellings = ["overlong1z", "b", "", "overlong1a", "overlong", "ab", "a", "overlong1"];
6809 let path = path("dictionary-order");
6810 let mut writer =
6811 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
6812 .expect("new file");
6813 writer
6814 .append(
6815 &Chunk::new(vec![
6816 Vector::from_values(
6817 LogicalType::Varchar,
6818 &spellings.map(|text| Value::Varchar(text.into())),
6819 )
6820 .expect("strings"),
6821 ])
6822 .expect("one column"),
6823 )
6824 .expect("stripe written");
6825 writer.finish().expect("commit");
6826
6827 let reader = Reader::open(&path).expect("valid directory");
6828 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
6829 let count = dictionary.ranks().expect("a v10 file stores one");
6830 assert_eq!(count, spellings.len(), "every distinct value has a rank");
6831 let order = (0..count)
6832 .map(|rank| dictionary.code_at_rank(rank).expect("a code"))
6833 .collect::<Vec<_>>();
6834 let mut seen = order.clone();
6835 seen.sort_unstable();
6836 assert_eq!(seen, (0..spellings.len() as u32).collect::<Vec<_>>(), "a permutation of codes");
6837
6838 let ranked = order
6839 .iter()
6840 .map(|&code| {
6841 dictionary.try_bytes_at(code as usize).expect("read").expect("a value").to_vec()
6842 })
6843 .collect::<Vec<_>>();
6844 let mut expected = spellings.map(|text| text.as_bytes().to_vec()).to_vec();
6845 expected.sort();
6846 assert_eq!(ranked, expected, "rank order is value order");
6847
6848 for (rank, value) in expected.iter().enumerate() {
6851 assert_eq!(
6852 dictionary.compare_rank(rank, value).expect("compare"),
6853 Ordering::Equal,
6854 "rank {rank} is its own value"
6855 );
6856 if rank > 0 {
6857 assert_eq!(
6858 dictionary.compare_rank(rank - 1, value).expect("compare"),
6859 Ordering::Less,
6860 "rank {rank} follows the one before it"
6861 );
6862 }
6863 }
6864 fs::remove_file(path).expect("remove scratch file");
6865 }
6866
6867 #[test]
6875 fn a_dictionary_sweep_reads_every_value_and_keeps_it_under_the_budget() {
6876 let path = path("dictionary-sweep");
6877 let spellings = (0..2_500)
6880 .map(|index| Value::Varchar(format!("value {index:08} {}", "x".repeat(index % 40))))
6881 .collect::<Vec<_>>();
6882 let mut writer =
6883 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
6884 .expect("new file");
6885 for part in spellings.chunks(1_024) {
6888 writer
6889 .append(
6890 &Chunk::new(vec![
6891 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
6892 ])
6893 .expect("one column"),
6894 )
6895 .expect("stripe written");
6896 }
6897 writer.finish().expect("commit");
6898
6899 let reader = Reader::open(&path).expect("valid directory");
6900 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
6901 assert_eq!(dictionary.len(), spellings.len(), "every value is distinct");
6902
6903 let resting = dictionary.footprint();
6904 let mut swept: Vec<Vec<u8>> = Vec::new();
6905 let mut at = 0;
6906 let mut calls = 0;
6907 while at < dictionary.len() {
6908 let stopped = dictionary
6909 .sweep_text(at, dictionary.len(), &mut |index: usize, text: &[u8]| {
6910 assert_eq!(index, swept.len(), "a sweep hands its values over in order");
6911 swept.push(text.to_vec());
6912 Ok(())
6913 })
6914 .expect("a sweep reads");
6915 assert!(stopped > at, "a sweep moves");
6916 at = stopped;
6917 calls += 1;
6918 }
6919 assert_eq!(calls, 3, "a sweep hands over one block at a time");
6920 let after = dictionary.footprint();
6921 assert!(after > resting, "a sweep under the budget keeps what it decoded");
6922
6923 let read = (0..dictionary.len())
6924 .map(|code| dictionary.try_bytes_at(code).expect("read").expect("a value").to_vec())
6925 .collect::<Vec<_>>();
6926 assert_eq!(swept, read, "a sweep answers what a point read answers");
6927 assert_eq!(dictionary.footprint(), after, "a point read of a kept block decodes nothing");
6928 fs::remove_file(path).expect("remove scratch file");
6929 }
6930
6931 #[test]
6942 fn a_sweep_over_a_block_with_a_short_second_run_reads_what_a_point_read_reads() {
6943 let path = path("dictionary-sweep-short-run");
6944 let spellings = (0..2_800)
6945 .map(|index| Value::Varchar(format!("value {index:08} {}", "x".repeat(index % 40))))
6946 .collect::<Vec<_>>();
6947 let mut writer =
6948 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
6949 .expect("new file");
6950 for part in spellings.chunks(1_024) {
6951 writer
6952 .append(
6953 &Chunk::new(vec![
6954 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
6955 ])
6956 .expect("one column"),
6957 )
6958 .expect("stripe written");
6959 }
6960 writer.finish().expect("commit");
6961
6962 let reader = Reader::open(&path).expect("valid directory");
6963 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
6964 assert_eq!(dictionary.len(), spellings.len(), "every value is distinct");
6965 let last = dictionary.len() % TEXT_PAYLOAD_VALUES;
6966 assert!(last > TEXT_OFFSET_RUN, "the last block has to reach into a second run of offsets");
6967 assert!(last < TEXT_PAYLOAD_VALUES, "and that second run has to be short of a whole one");
6968
6969 let mut swept: Vec<Vec<u8>> = Vec::new();
6970 let mut at = 0;
6971 while at < dictionary.len() {
6972 let stopped = dictionary
6973 .sweep_text(at, dictionary.len(), &mut |index: usize, text: &[u8]| {
6974 assert_eq!(index, swept.len(), "a sweep hands its values over in order");
6975 swept.push(text.to_vec());
6976 Ok(())
6977 })
6978 .expect("a sweep reads");
6979 assert!(stopped > at, "a sweep moves");
6980 at = stopped;
6981 }
6982 let read = (0..dictionary.len())
6983 .map(|code| dictionary.try_bytes_at(code).expect("read").expect("a value").to_vec())
6984 .collect::<Vec<_>>();
6985 assert_eq!(swept, read, "a sweep answers what a point read answers");
6986 fs::remove_file(path).expect("remove scratch file");
6987 }
6988
6989 #[test]
6999 fn narrowing_a_page_takes_what_fits_and_refuses_what_does_not() {
7000 assert_eq!(fit::<i8>(&[]).expect("an empty page fits anything"), Vec::<i8>::new());
7001 assert_eq!(fit::<i8>(&[-128, 0, 127]).expect("the edges fit"), vec![-128_i8, 0, 127]);
7002 fit::<i8>(&[128]).expect_err("one past the top does not fit");
7003 fit::<i8>(&[-129]).expect_err("one past the bottom does not fit");
7004 assert_eq!(fit::<u8>(&[0, 255]).expect("the edges fit"), vec![0_u8, 255]);
7005 fit::<u8>(&[256]).expect_err("one past the top does not fit");
7006 fit::<u8>(&[-1]).expect_err("a negative does not fit an unsigned page");
7007 assert_eq!(
7008 fit::<i16>(&[-32_768, 0, 32_767]).expect("the edges fit"),
7009 vec![-32_768_i16, 0, 32_767]
7010 );
7011 fit::<i16>(&[32_768]).expect_err("one past the top does not fit");
7012 fit::<i16>(&[-32_769]).expect_err("one past the bottom does not fit");
7013 assert_eq!(fit::<u16>(&[0, 65_535]).expect("the edges fit"), vec![0_u16, 65_535]);
7014 fit::<u16>(&[65_536]).expect_err("one past the top does not fit");
7015 fit::<u16>(&[-1]).expect_err("a negative does not fit an unsigned page");
7016 assert_eq!(
7017 fit::<i32>(&[i64::from(i32::MIN), 0, i64::from(i32::MAX)]).expect("the edges fit"),
7018 vec![i32::MIN, 0, i32::MAX]
7019 );
7020 fit::<i32>(&[i64::from(i32::MAX) + 1]).expect_err("one past the top does not fit");
7021 fit::<i32>(&[i64::from(i32::MIN) - 1]).expect_err("one past the bottom does not fit");
7022 assert_eq!(
7023 fit::<u32>(&[0, 4_294_967_295]).expect("the edges fit"),
7024 vec![0_u32, 4_294_967_295]
7025 );
7026 fit::<u32>(&[4_294_967_296]).expect_err("one past the top does not fit");
7027 fit::<u32>(&[-1]).expect_err("a negative does not fit an unsigned page");
7028
7029 fit::<i8>(&[0, 1, 2, 128, 3]).expect_err("one bad value spoils the page");
7032 }
7033
7034 #[test]
7041 fn the_residue_agrees_with_a_checked_conversion_everywhere() {
7042 for value in -70_000_i64..70_000 {
7043 assert_eq!(fit::<i8>(&[value]).is_ok(), i8::try_from(value).is_ok(), "{value} as i8");
7044 assert_eq!(fit::<u8>(&[value]).is_ok(), u8::try_from(value).is_ok(), "{value} as u8");
7045 assert_eq!(fit::<i16>(&[value]).is_ok(), i16::try_from(value).is_ok(), "{value} i16");
7046 assert_eq!(fit::<u16>(&[value]).is_ok(), u16::try_from(value).is_ok(), "{value} u16");
7047 }
7048 let wide = [i64::MIN, i64::MIN + 1, i64::from(i32::MIN), 0, i64::from(u32::MAX), i64::MAX];
7049 for edge in wide {
7050 for step in -2_i64..=2 {
7051 let value = edge.saturating_add(step);
7052 assert_eq!(
7053 fit::<i32>(&[value]).is_ok(),
7054 i32::try_from(value).is_ok(),
7055 "{value} as i32"
7056 );
7057 assert_eq!(
7058 fit::<u32>(&[value]).is_ok(),
7059 u32::try_from(value).is_ok(),
7060 "{value} as u32"
7061 );
7062 }
7063 }
7064 }
7065
7066 #[test]
7074 fn a_dictionary_at_its_budget_sweeps_without_keeping() {
7075 let path = path("dictionary-budget");
7076 let spellings = (0..2_500)
7077 .map(|index| Value::Varchar(format!("value {index:08} {}", "y".repeat(index % 40))))
7078 .collect::<Vec<_>>();
7079 let mut writer =
7080 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7081 .expect("new file");
7082 for part in spellings.chunks(1_024) {
7083 writer
7084 .append(
7085 &Chunk::new(vec![
7086 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
7087 ])
7088 .expect("one column"),
7089 )
7090 .expect("stripe written");
7091 }
7092 writer.finish().expect("commit");
7093
7094 let reader = Reader::open(&path).expect("valid directory");
7095 let page = reader.table.dictionaries[0].expect("a string column has one");
7096 let file = Arc::clone(&reader.file);
7097 let starved = open_global_dictionary(file, page, &LogicalType::Varchar, 0)
7098 .expect("a dictionary opens whatever it may keep");
7099
7100 let resting = starved.footprint();
7101 let mut swept: Vec<Vec<u8>> = Vec::new();
7102 let mut at = 0;
7103 while at < starved.len() {
7104 at = starved
7105 .sweep_text(at, starved.len(), &mut |_index: usize, text: &[u8]| {
7106 swept.push(text.to_vec());
7107 Ok(())
7108 })
7109 .expect("a sweep reads");
7110 }
7111 assert_eq!(swept.len(), spellings.len(), "a starved sweep still reads every value");
7112 assert_eq!(starved.footprint(), resting, "and keeps no block it decoded");
7113
7114 let generous = reader.dictionary(0).expect("read").expect("a string column has one");
7115 let read = (0..generous.len())
7116 .map(|code| generous.try_bytes_at(code).expect("read").expect("a value").to_vec())
7117 .collect::<Vec<_>>();
7118 assert_eq!(swept, read, "a starved sweep answers what a point read answers");
7119 fs::remove_file(path).expect("remove scratch file");
7120 }
7121
7122 #[test]
7123 fn damaged_membership_cannot_skip_a_string_page() {
7124 let path = path("damaged-membership");
7125 let mut writer = Writer::create(
7126 &path,
7127 "items",
7128 vec![
7129 Field::required("id", LogicalType::Integer),
7130 Field::new("text", LogicalType::Varchar),
7131 ],
7132 )
7133 .expect("new file");
7134 writer.append(&sample()).expect("stripe written");
7135 writer.finish().expect("commit");
7136
7137 let reader = Reader::open(&path).expect("valid directory");
7138 let membership = reader.table.stripes[0].memberships[1].expect("string membership");
7139 let mut file = OpenOptions::new().write(true).open(&path).expect("open membership page");
7140 file.seek(SeekFrom::Start(membership.offset)).expect("membership start");
7141 file.write_all(&[255]).expect("damage membership");
7142 let error = reader.skips_codes(0, 1, &[3]).expect_err("corruption must not skip rows");
7143 assert!(error.message().contains("membership page checksum differs"), "{error}");
7144 fs::remove_file(path).expect("remove scratch file");
7145 }
7146
7147 #[test]
7148 fn membership_delta_stream_is_sorted_exact_and_bounded() {
7149 let unique = unique_codes(&[900, 4, 4, 72, 9, u32::MAX]);
7150 assert_eq!(unique, [4, 9, 72, 900, u32::MAX]);
7151 let encoded = encode_membership(&unique);
7152 assert_eq!(
7153 decode_membership(&encoded).expect("valid membership"),
7154 [4, 9, 72, 900, u32::MAX]
7155 );
7156 let merged = merged_codes(vec![vec![4, 900], vec![9, 900, u32::MAX], vec![72]]);
7159 assert_eq!(merged, [4, 9, 72, 900, u32::MAX]);
7160 assert_eq!(
7161 decode_membership(&encode_membership(&merged)).expect("valid membership"),
7162 unique
7163 );
7164 assert!(decode_membership(&[1, 0x80]).is_err(), "a truncated varint is invalid");
7165 assert!(
7166 decode_membership(&[1, 0xff, 0xff, 0xff, 0xff, 0x10]).is_err(),
7167 "a value past u32 is invalid"
7168 );
7169 }
7170
7171 #[test]
7172 fn a_global_dictionary_may_be_larger_than_one_column_page() {
7173 let dictionary = Page {
7174 offset: HEADER,
7175 length: u32::try_from(MAX_PAGE + 1).expect("the page bound fits on disk"),
7176 hash: 0,
7177 };
7178 let table = Table {
7179 name: "items".to_owned(),
7180 fields: vec![Field::new("text", LogicalType::Varchar)],
7181 stripes: Vec::new(),
7182 rows: 0,
7183 dictionaries: vec![Some(dictionary)],
7184 distincts: vec![None],
7185 frequencies: vec![None],
7186 };
7187 let directory = encode_directory(&table).expect("directory");
7188 let file_size = dictionary.offset + u64::from(dictionary.length) + 1;
7189
7190 let decoded = decode_directory(&directory, file_size).expect("large lazy dictionary");
7191 assert_eq!(decoded.dictionaries[0].expect("dictionary").length, dictionary.length);
7192 }
7193
7194 #[test]
7195 fn a_column_with_one_value_everywhere_costs_almost_nothing_a_row() {
7196 let path = path("constant-codes");
7197 let mut writer =
7198 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7199 .expect("new file");
7200 let empty = vec![Value::Varchar(String::new()); 1024];
7201 for _ in 0..4 {
7202 let column = Vector::from_values(LogicalType::Varchar, &empty).expect("strings");
7203 writer.append(&Chunk::new(vec![column]).expect("one column")).expect("a part");
7204 }
7205 writer.finish().expect("commit");
7206
7207 let reader = Reader::open(&path).expect("valid directory");
7208 let pages = reader.layout().columns.first().expect("one column").pages;
7209 assert!(pages < 256, "{pages} bytes of pages for 4,096 rows of one value");
7213 let read = reader.read(3, &[0]).expect("the last part back");
7214 assert_eq!(read.value_at(0, 0), Value::Varchar(String::new()));
7215 assert_eq!(read.value_at(1023, 0), Value::Varchar(String::new()));
7216 fs::remove_file(path).expect("remove scratch file");
7217 }
7218
7219 #[test]
7220 fn a_cascade_value_too_wide_for_its_column_is_refused_rather_than_cut() {
7221 let over = vec![i64::from(i32::MAX) + 1];
7224 let error = narrowed(&LogicalType::Integer, over).expect_err("a page that disagrees");
7225 assert!(format!("{error}").contains("not of its type"), "{error}");
7226 assert!(narrowed(&LogicalType::BigInt, vec![i64::MIN]).is_ok(), "bigint holds all of i64");
7227 assert!(narrowed(&LogicalType::Varchar, vec![0]).is_err(), "strings are not integers");
7228 }
7229
7230 #[test]
7231 fn a_code_stream_the_cascade_cannot_shrink_is_left_alone() {
7232 let mut state: u32 = 0x9e37_79b9;
7236 let spread: Vec<u32> = (0..1024)
7237 .map(|_| {
7238 state ^= state << 13;
7239 state ^= state >> 17;
7240 state ^= state << 5;
7241 state
7242 })
7243 .collect();
7244 assert_eq!(encoded_codes(&spread).expect("no failure"), None);
7245 let near: Vec<u32> = (0..1024).collect();
7246 let coded = encoded_codes(&near).expect("no failure").expect("counting up is packable");
7247 assert!(coded.len() < near.len() * 4, "{} bytes for a run of 1,024", coded.len());
7248 }
7249
7250 #[test]
7256 fn two_writes_of_the_same_rows_give_the_same_bytes() {
7257 fn written(path: &PathBuf) {
7258 let fields = (0..40)
7259 .map(|column| {
7260 let ty =
7261 if column % 4 == 0 { LogicalType::Varchar } else { LogicalType::BigInt };
7262 Field::new(format!("c{column}"), ty)
7263 })
7264 .collect::<Vec<_>>();
7265 let mut writer = Writer::create(path, "wide", fields).expect("new file");
7266 for part in 0..70_u64 {
7267 let columns = (0..40)
7268 .map(|column| {
7269 let values = (0..64_u64)
7270 .map(|row| {
7271 let seed = part.wrapping_mul(31).wrapping_add(row);
7272 if column % 4 == 0 {
7273 Value::Varchar(format!("v{}", seed % 17))
7274 } else {
7275 Value::BigInt(i64::try_from(seed % 97).expect("small"))
7276 }
7277 })
7278 .collect::<Vec<_>>();
7279 let ty = if column % 4 == 0 {
7280 LogicalType::Varchar
7281 } else {
7282 LogicalType::BigInt
7283 };
7284 Vector::from_values(ty, &values).expect("a column")
7285 })
7286 .collect::<Vec<_>>();
7287 writer.append(&Chunk::new(columns).expect("forty columns")).expect("a part");
7288 }
7289 writer.finish().expect("commit");
7290 }
7291
7292 let first = path("repeatable-one");
7293 let second = path("repeatable-two");
7294 written(&first);
7295 written(&second);
7296 let left = fs::read(&first).expect("the first file");
7297 let right = fs::read(&second).expect("the second file");
7298 assert_eq!(left.len(), right.len(), "two writes of the same rows differ in length");
7299 assert!(left == right, "two writes of the same rows differ in their bytes");
7300
7301 let reader = Reader::open(&first).expect("valid directory");
7304 assert_eq!(reader.table().rows(), 70 * 64);
7305 let read = reader.read(0, &[0, 1]).expect("the first part back");
7306 assert_eq!(read.value_at(0, 0), Value::Varchar("v0".to_owned()));
7307 assert_eq!(read.value_at(0, 1), Value::BigInt(0));
7308 fs::remove_file(first).expect("remove scratch file");
7309 fs::remove_file(second).expect("remove scratch file");
7310 }
7311
7312 fn three_tables(path: &PathBuf) {
7314 let writer = Writer::create(
7315 path,
7316 "region",
7317 vec![
7318 Field::new("r_key", LogicalType::Integer),
7319 Field::new("r_name", LogicalType::Varchar),
7320 ],
7321 )
7322 .expect("new file");
7323 let mut writer = writer;
7324 writer
7325 .append(
7326 &Chunk::new(vec![
7327 Vector::from_values(
7328 LogicalType::Integer,
7329 &[Value::Integer(0), Value::Integer(1)],
7330 )
7331 .expect("keys"),
7332 Vector::from_values(
7333 LogicalType::Varchar,
7334 &[Value::Varchar("AFRICA".to_owned()), Value::Varchar("ASIA".to_owned())],
7335 )
7336 .expect("names"),
7337 ])
7338 .expect("two columns"),
7339 )
7340 .expect("a part");
7341 let mut writer = writer
7342 .next("empty", vec![Field::new("nothing", LogicalType::BigInt)])
7343 .expect("a second table");
7344 writer
7345 .append(
7346 &Chunk::new(vec![
7347 Vector::from_values(LogicalType::BigInt, &[Value::BigInt(7)]).expect("a row"),
7348 ])
7349 .expect("one column"),
7350 )
7351 .expect("a part");
7352 let mut writer =
7353 writer.next("wide", vec![Field::new("n", LogicalType::BigInt)]).expect("a third table");
7354 for part in 0..70_i64 {
7355 let values = (0..64).map(|row| Value::BigInt(part * 64 + row)).collect::<Vec<_>>();
7356 writer
7357 .append(
7358 &Chunk::new(vec![
7359 Vector::from_values(LogicalType::BigInt, &values).expect("a column"),
7360 ])
7361 .expect("one column"),
7362 )
7363 .expect("a part");
7364 }
7365 writer.finish().expect("commit");
7366 }
7367
7368 #[test]
7369 fn three_tables_in_one_file_read_back_by_name() {
7370 let file = path("three-tables");
7371 three_tables(&file);
7372 let catalog = Catalog::open(&file).expect("a committed catalog");
7373 assert_eq!(catalog.names().collect::<Vec<_>>(), ["region", "empty", "wide"]);
7374
7375 let region = catalog.table("region").expect("the first table");
7376 assert_eq!(region.table().rows(), 2);
7377 assert_eq!(
7378 region.read(0, &[1]).expect("names").value_at(1, 0),
7379 Value::Varchar("ASIA".to_owned())
7380 );
7381
7382 let wide = catalog.table("wide").expect("the third table");
7383 assert_eq!(wide.table().rows(), 70 * 64);
7384 assert_eq!(wide.read(0, &[0]).expect("the first part").value_at(0, 0), Value::BigInt(0));
7385
7386 let empty = catalog.table("empty").expect("the second table");
7389 assert_eq!(empty.table().rows(), 1);
7390 assert_eq!(empty.read(0, &[0]).expect("the row").value_at(0, 0), Value::BigInt(7));
7391
7392 fs::remove_file(file).expect("remove scratch file");
7393 }
7394
7395 #[test]
7396 fn a_name_the_file_does_not_hold_is_an_error_rather_than_the_first_table() {
7397 let file = path("three-tables-missing");
7398 three_tables(&file);
7399 let catalog = Catalog::open(&file).expect("a committed catalog");
7400 let error = catalog.table("nation").expect_err("no such table");
7401 assert!(error.message().contains("nation"), "{}", error.message());
7402 fs::remove_file(file).expect("remove scratch file");
7403 }
7404
7405 #[test]
7406 fn a_file_of_three_tables_will_not_open_as_one() {
7407 let file = path("three-tables-unnamed");
7408 three_tables(&file);
7409 let error = Reader::open(&file).expect_err("more than one table");
7410 assert!(error.message().contains("more than one table"), "{}", error.message());
7411 fs::remove_file(file).expect("remove scratch file");
7412 }
7413
7414 #[test]
7415 fn two_tables_of_one_name_are_refused_before_anything_is_committed() {
7416 let file = path("two-of-a-name");
7417 let writer = Writer::create(&file, "t", vec![Field::new("a", LogicalType::BigInt)])
7418 .expect("new file");
7419 let error = writer
7420 .next("t", vec![Field::new("a", LogicalType::BigInt)])
7421 .expect_err("the same name twice");
7422 assert!(error.message().contains("same name"), "{}", error.message());
7423 fs::remove_file(file).expect("remove scratch file");
7424 }
7425
7426 #[test]
7427 fn opening_the_catalog_reads_no_table_directory() {
7428 let file = path("catalog-only");
7429 three_tables(&file);
7430 let catalog = Catalog::open(&file).expect("a committed catalog");
7431 assert_eq!(catalog.opening.reads, 2, "opening the catalog read more than the slot");
7434 assert_eq!(catalog.names().len(), 3);
7435 fs::remove_file(file).expect("remove scratch file");
7436 }
7437}