1#![forbid(unsafe_code)]
34
35use std::cmp::Ordering;
36use std::collections::{HashMap, VecDeque};
37use std::fs::{File, OpenOptions};
38use std::io::{Read, Seek, SeekFrom};
39use std::mem::{size_of, size_of_val};
40use std::path::Path;
41use std::slice;
42use std::sync::atomic::{AtomicUsize, Ordering as Atomic};
43use std::sync::{Arc, Mutex, OnceLock};
44
45use rudb_common::bounds::{Bound, Op, scaled_as};
46use rudb_common::{Error, Field, LogicalType, PhysicalType, Result, Value};
47use rudb_encoding::{bitpack, chooser, integer, string};
48use rudb_storage::sieve::Sieve;
49use rudb_storage::{Probe, Range, Zone};
50use rudb_vector::string::StringColumn;
51use rudb_vector::validity::Validity;
52use rudb_vector::{Buffer, Chunk, Data, Packed, TextSource, Vector, search_below};
53
54mod zones;
55
56pub use zones::{Common, Stripes, distincts};
57
58const MAGIC: &[u8; 8] = b"RUDBNV10";
59const DIRECTORY: &[u8; 8] = b"RUDBDI10";
60const CATALOG: &[u8; 8] = b"RUDBCA10";
61const FORMAT: u32 = 22;
62const HEADER: u64 = 80;
63const SLOT_BYTES: usize = 28;
64const MAX_PAGE: usize = 256 * 1024 * 1024;
65const MAX_DIRECTORY: usize = 128 * 1024 * 1024;
66const FREQUENCIES: &[u8; 8] = b"RUDBFQ2\0";
67const FREQUENCY_CANDIDATES: usize = 32_768;
68const FREQUENCY_ENTRIES: usize = 512;
69const FREQUENCY_BUILD_RANK: usize = 10;
70const FREQUENCY_ORDINALS: usize = 65_536;
71const MAX_FREQUENCY_WORKERS: usize = 32;
78
79const MAX_ENCODE_WORKERS: usize = 32;
86
87const SIEVE_BUDGET: usize = 8 * 1024;
95
96const PART_BOUND_BYTES: usize = 24;
105
106fn io(error: std::io::Error) -> Error {
107 Error::io(error.to_string())
108}
109
110fn invalid(message: &str) -> Error {
111 Error::invalid_input(format!("invalid rudb native file: {message}"))
112}
113
114fn sum(counts: impl Iterator<Item = u64>) -> u64 {
116 counts.fold(0, u64::saturating_add)
117}
118
119fn span_bytes(spans: &[Span], at: usize) -> u64 {
121 spans.get(at).map_or(0, |span| u64::from(span.length))
122}
123
124fn page_bytes(pages: &[Option<Page>], at: usize) -> u64 {
126 pages.get(at).and_then(Option::as_ref).map_or(0, Page::bytes)
127}
128
129fn checksum(bytes: &[u8]) -> u64 {
130 const P1: u64 = 11_400_714_785_074_694_791;
131 const P2: u64 = 14_029_467_366_897_019_727;
132 const P3: u64 = 1_609_587_929_392_839_161;
133 const P4: u64 = 9_650_029_242_287_828_579;
134 const P5: u64 = 2_870_177_450_012_600_261;
135 let round = |state: u64, word: u64| {
136 state.wrapping_add(word.wrapping_mul(P2)).rotate_left(31).wrapping_mul(P1)
137 };
138 let merge = |state: u64, lane: u64| (state ^ round(0, lane)).wrapping_mul(P1).wrapping_add(P4);
139 let word =
140 |at: usize| u64::from_le_bytes(bytes[at..at + 8].try_into().expect("eight checksum bytes"));
141
142 let mut at = 0;
143 let mut hash = if bytes.len() >= 32 {
144 let mut one = P1.wrapping_add(P2);
145 let mut two = P2;
146 let mut three = 0;
147 let mut four = 0_u64.wrapping_sub(P1);
148 while at + 32 <= bytes.len() {
149 one = round(one, word(at));
150 two = round(two, word(at + 8));
151 three = round(three, word(at + 16));
152 four = round(four, word(at + 24));
153 at += 32;
154 }
155 let combined = one
156 .rotate_left(1)
157 .wrapping_add(two.rotate_left(7))
158 .wrapping_add(three.rotate_left(12))
159 .wrapping_add(four.rotate_left(18));
160 merge(merge(merge(merge(combined, one), two), three), four)
161 } else {
162 P5
163 };
164 hash = hash.wrapping_add(bytes.len() as u64);
165 while at + 8 <= bytes.len() {
166 hash ^= round(0, word(at));
167 hash = hash.rotate_left(27).wrapping_mul(P1).wrapping_add(P4);
168 at += 8;
169 }
170 if at + 4 <= bytes.len() {
171 let tail = u32::from_le_bytes(bytes[at..at + 4].try_into().expect("four checksum bytes"));
172 hash ^= u64::from(tail).wrapping_mul(P1);
173 hash = hash.rotate_left(23).wrapping_mul(P2).wrapping_add(P3);
174 at += 4;
175 }
176 while at < bytes.len() {
177 hash ^= u64::from(bytes[at]).wrapping_mul(P5);
178 hash = hash.rotate_left(11).wrapping_mul(P1);
179 at += 1;
180 }
181 hash ^= hash >> 33;
182 hash = hash.wrapping_mul(P2);
183 hash ^= hash >> 29;
184 hash = hash.wrapping_mul(P3);
185 hash ^ (hash >> 32)
186}
187
188#[derive(Debug, Clone, Copy)]
189struct Slot {
190 offset: u64,
191 length: u32,
192 generation: u64,
193 hash: u64,
194}
195
196impl Slot {
197 fn bytes(self) -> [u8; SLOT_BYTES] {
198 let mut result = [0; SLOT_BYTES];
199 result[..8].copy_from_slice(&self.offset.to_le_bytes());
200 result[8..12].copy_from_slice(&self.length.to_le_bytes());
201 result[12..20].copy_from_slice(&self.generation.to_le_bytes());
202 result[20..28].copy_from_slice(&self.hash.to_le_bytes());
203 result
204 }
205
206 fn read(bytes: &[u8]) -> Self {
207 Self {
208 offset: u64::from_le_bytes(bytes[..8].try_into().expect("eight bytes")),
209 length: u32::from_le_bytes(bytes[8..12].try_into().expect("four bytes")),
210 generation: u64::from_le_bytes(bytes[12..20].try_into().expect("eight bytes")),
211 hash: u64::from_le_bytes(bytes[20..28].try_into().expect("eight bytes")),
212 }
213 }
214}
215
216#[derive(Debug, Clone, Copy)]
217struct Page {
218 offset: u64,
219 length: u32,
220 hash: u64,
221}
222
223impl Page {
224 fn bytes(&self) -> u64 {
226 u64::from(self.length)
227 }
228}
229
230#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
231enum FrequencyValue {
232 Null,
233 Integer(i128),
234 Code(u32),
235}
236
237#[derive(Debug, Clone)]
238struct FrequencyEntry {
239 value: FrequencyValue,
240 count: u64,
241}
242
243#[derive(Debug, Clone)]
248struct FrequencySummary {
249 entries: Vec<FrequencyEntry>,
250 omitted_max: u64,
251 ordinals: Vec<u64>,
252}
253
254#[derive(Debug, Clone, PartialEq, Eq)]
256pub struct FrequencyOccurrences {
257 pub omitted_max: u64,
259 pub ordinals: Vec<u64>,
261}
262
263#[derive(Debug, Clone, Copy, Default)]
270struct Span {
271 offset: u64,
272 length: u32,
273}
274
275#[derive(Debug, Clone)]
277pub struct Stripe {
278 rows: usize,
279 parts: Vec<u32>,
282 index: Span,
286 pages: Vec<Span>,
287 memberships: Vec<Option<Page>>,
288 sieves: Vec<Option<Page>>,
291 part_ranges: Vec<Option<Page>>,
302 zone: Zone,
303}
304
305impl Stripe {
306 #[must_use]
308 pub fn rows(&self) -> usize {
309 self.rows
310 }
311
312 #[must_use]
314 pub fn parts(&self) -> usize {
315 self.parts.len()
316 }
317
318 #[must_use]
324 pub fn zone(&self) -> &Zone {
325 &self.zone
326 }
327}
328
329#[derive(Debug, Clone)]
331pub struct Table {
332 name: String,
333 fields: Vec<Field>,
334 stripes: Vec<Stripe>,
335 rows: usize,
336 dictionaries: Vec<Option<Page>>,
337 frequencies: Vec<Option<FrequencySummary>>,
338 distincts: Vec<Option<u64>>,
348}
349
350impl Table {
351 #[must_use]
353 pub fn name(&self) -> &str {
354 &self.name
355 }
356
357 #[must_use]
359 pub fn fields(&self) -> &[Field] {
360 &self.fields
361 }
362
363 #[must_use]
365 pub fn rows(&self) -> usize {
366 self.rows
367 }
368
369 #[must_use]
371 pub fn stripes(&self) -> &[Stripe] {
372 &self.stripes
373 }
374}
375
376#[derive(Debug, Clone)]
388struct Entry {
389 name: String,
390 fields: Vec<Field>,
391 rows: usize,
392 directory: Page,
394}
395
396#[derive(Debug, Clone)]
398pub struct ColumnLayout {
399 pub name: String,
401 pub kind: String,
403 pub pages: u64,
405 pub memberships: u64,
407 pub sieves: u64,
409 pub part_ranges: u64,
411 pub dictionary: u64,
413}
414
415impl ColumnLayout {
416 #[must_use]
418 pub fn total(&self) -> u64 {
419 self.pages
420 .saturating_add(self.memberships)
421 .saturating_add(self.sieves)
422 .saturating_add(self.part_ranges)
423 .saturating_add(self.dictionary)
424 }
425}
426
427#[derive(Debug, Clone)]
438pub struct Layout {
439 pub file: u64,
441 pub rows: usize,
443 pub stripes: usize,
445 pub parts: usize,
447 pub columns: Vec<ColumnLayout>,
449 pub indexes: u64,
452 pub directory: u64,
454 pub header: u64,
456}
457
458impl Layout {
459 #[must_use]
461 pub fn columns_total(&self) -> u64 {
462 self.columns.iter().map(ColumnLayout::total).fold(0, u64::saturating_add)
463 }
464
465 #[must_use]
471 pub fn unaccounted(&self) -> u64 {
472 self.file
473 .saturating_sub(self.columns_total())
474 .saturating_sub(self.indexes)
475 .saturating_sub(self.directory)
476 .saturating_sub(self.header)
477 }
478}
479
480#[derive(Debug)]
482struct GlobalDictionary {
483 primary: HashMap<u64, u32>,
484 collisions: HashMap<u64, Vec<u32>>,
485 offsets: Vec<u32>,
486 payload: Vec<u8>,
487 counts: Vec<u64>,
488 nulls: u64,
489}
490
491impl GlobalDictionary {
492 fn new() -> Self {
493 Self {
494 primary: HashMap::new(),
495 collisions: HashMap::new(),
496 offsets: vec![0],
497 payload: Vec::new(),
498 counts: Vec::new(),
499 nulls: 0,
500 }
501 }
502
503 fn bytes(&self, code: u32) -> Option<&[u8]> {
504 let start = *self.offsets.get(code as usize)? as usize;
505 let end = *self.offsets.get(code as usize + 1)? as usize;
506 self.payload.get(start..end)
507 }
508
509 fn code(&mut self, text: &str) -> Result<u32> {
510 let hash = checksum(text.as_bytes());
511 if let Some(&code) = self.primary.get(&hash) {
512 if self.bytes(code) == Some(text.as_bytes()) {
513 return Ok(code);
514 }
515 if let Some(codes) = self.collisions.get(&hash) {
516 if let Some(code) =
517 codes.iter().copied().find(|&code| self.bytes(code) == Some(text.as_bytes()))
518 {
519 return Ok(code);
520 }
521 }
522 let code = self.insert(text)?;
523 self.collisions.entry(hash).or_default().push(code);
524 return Ok(code);
525 }
526 let code = self.insert(text)?;
527 self.primary.insert(hash, code);
528 Ok(code)
529 }
530
531 fn insert(&mut self, text: &str) -> Result<u32> {
532 let code = u32::try_from(self.offsets.len() - 1)
533 .map_err(|_| invalid("global dictionary has too many values"))?;
534 self.payload.extend_from_slice(text.as_bytes());
535 self.offsets.push(
536 u32::try_from(self.payload.len())
537 .map_err(|_| invalid("global dictionary payload exceeds 4 GiB"))?,
538 );
539 self.counts.push(0);
540 Ok(code)
541 }
542
543 fn ranked(&self) -> Vec<(u64, u32)> {
563 let count = self.offsets.len() - 1;
564 let mut ranked = (0..count)
565 .map(|code| {
566 let code = code as u32;
567 (head(self.bytes(code).unwrap_or_default()), code)
568 })
569 .collect::<Vec<_>>();
570 ranked.sort_unstable_by(|left, right| {
571 left.0.cmp(&right.0).then_with(|| self.bytes(left.1).cmp(&self.bytes(right.1)))
572 });
573 ranked
574 }
575
576 fn observe(&mut self, code: u32, null: bool) -> Result<()> {
577 if null {
578 self.nulls = self.nulls.saturating_add(1);
579 return Ok(());
580 }
581 let count = self
582 .counts
583 .get_mut(code as usize)
584 .ok_or_else(|| invalid("global dictionary count code is out of range"))?;
585 *count = count.saturating_add(1);
586 Ok(())
587 }
588}
589
590#[derive(Debug)]
598pub struct Writer {
599 file: File,
600 at: u64,
608 table: Table,
609 generation: u64,
610 order: Vec<((u64, u64), (u64, u64))>,
613 next_order: u64,
614 dictionaries: Vec<Option<GlobalDictionary>>,
615 pending: Vec<PendingChunk>,
616 closed: Vec<Entry>,
618}
619
620#[derive(Debug)]
628struct PendingChunk {
629 order: (u64, u64),
630 chunk: Chunk,
631}
632
633#[derive(Debug)]
639struct ColumnStripe {
640 pages: Vec<Vec<u8>>,
641 codes: Vec<Option<Vec<u32>>>,
642 sieves: Vec<Option<Sieve>>,
643 ranges: Vec<Range>,
644}
645
646fn weight(ty: &LogicalType) -> usize {
654 match ty {
655 LogicalType::Varchar | LogicalType::Blob => 64,
656 LogicalType::BigInt
657 | LogicalType::UBigInt
658 | LogicalType::Timestamp
659 | LogicalType::Double
660 | LogicalType::Decimal { .. } => 8,
661 LogicalType::Integer | LogicalType::UInteger | LogicalType::Date | LogicalType::Float => 4,
662 LogicalType::SmallInt | LogicalType::USmallInt => 2,
663 _ => 1,
664 }
665}
666
667pub const STRIPE_PARTS: usize = 64;
674
675const INDEX_ENTRY: usize = size_of::<u32>() + size_of::<u64>();
677
678fn index_section(parts: usize) -> Result<usize> {
680 parts
681 .checked_mul(INDEX_ENTRY)
682 .and_then(|bytes| bytes.checked_add(size_of::<u64>()))
683 .ok_or_else(|| invalid("index page length overflow"))
684}
685
686impl Writer {
687 pub fn create(
693 path: impl AsRef<Path>,
694 name: impl Into<String>,
695 fields: Vec<Field>,
696 ) -> Result<Self> {
697 for field in &fields {
698 type_tag(&field.ty)?;
699 }
700 let file =
701 OpenOptions::new().write(true).read(true).create_new(true).open(path).map_err(io)?;
702 let mut header = [0; HEADER as usize];
703 header[..8].copy_from_slice(MAGIC);
704 header[8..12].copy_from_slice(&FORMAT.to_le_bytes());
705 write_at(&file, 0, &header)?;
706 Ok(Self {
707 file,
708 at: HEADER,
709 dictionaries: fields
710 .iter()
711 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
712 .collect(),
713 table: Table {
714 name: name.into(),
715 dictionaries: vec![None; fields.len()],
716 distincts: vec![None; fields.len()],
717 fields,
718 stripes: Vec::new(),
719 rows: 0,
720 frequencies: Vec::new(),
721 },
722 generation: 1,
723 order: Vec::new(),
724 next_order: 0,
725 pending: Vec::with_capacity(STRIPE_PARTS),
726 closed: Vec::new(),
727 })
728 }
729
730 pub fn next(mut self, name: impl Into<String>, fields: Vec<Field>) -> Result<Self> {
741 for field in &fields {
742 type_tag(&field.ty)?;
743 }
744 let name = name.into();
745 let entry = self.close()?;
746 if self.closed.iter().chain(std::iter::once(&entry)).any(|held| held.name == name) {
747 return Err(invalid("two tables in one native file have the same name"));
748 }
749 let Self { file, at, generation, mut closed, .. } = self;
750 closed.push(entry);
751 Ok(Self {
752 file,
753 at,
754 generation,
755 closed,
756 dictionaries: fields
757 .iter()
758 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
759 .collect(),
760 table: Table {
761 name,
762 dictionaries: vec![None; fields.len()],
763 distincts: vec![None; fields.len()],
764 fields,
765 stripes: Vec::new(),
766 rows: 0,
767 frequencies: Vec::new(),
768 },
769 order: Vec::new(),
770 next_order: 0,
771 pending: Vec::with_capacity(STRIPE_PARTS),
772 })
773 }
774
775 fn put(&mut self, bytes: &[u8]) -> Result<()> {
780 write_at(&self.file, self.at, bytes)?;
781 self.at = self
782 .at
783 .checked_add(bytes.len() as u64)
784 .ok_or_else(|| invalid("native file length overflow"))?;
785 Ok(())
786 }
787
788 pub fn append(&mut self, chunk: &Chunk) -> Result<()> {
794 let order = (self.next_order, 0);
795 self.next_order = self.next_order.saturating_add(1);
796 self.append_at(order, chunk)
797 }
798
799 pub fn append_at(&mut self, order: (u64, u64), chunk: &Chunk) -> Result<()> {
810 if chunk.is_empty() {
811 return Ok(());
812 }
813 self.admit(chunk)?;
814 if self.pending.last().is_some_and(|last| last.order > order) {
815 self.flush_pending()?;
816 }
817 self.pending.push(PendingChunk { order, chunk: chunk.clone() });
822 if self.pending.len() == STRIPE_PARTS {
823 self.flush_pending()?;
824 }
825 Ok(())
826 }
827
828 pub fn append_stripe(&mut self, parts: Vec<((u64, u64), Chunk)>) -> Result<()> {
844 if parts.len() > STRIPE_PARTS {
845 return Err(invalid("a stripe was handed more parts than it holds"));
846 }
847 self.flush_pending()?;
850 for (order, chunk) in parts {
851 if chunk.is_empty() {
852 continue;
853 }
854 self.admit(&chunk)?;
855 self.pending.push(PendingChunk { order, chunk });
856 }
857 self.flush_pending()
858 }
859
860 fn admit(&mut self, chunk: &Chunk) -> Result<()> {
862 if chunk.width() != self.table.fields.len() {
863 return Err(invalid("chunk width differs from table schema"));
864 }
865 for (index, field) in self.table.fields.iter().enumerate() {
866 if chunk.column(index)?.logical_type() != &field.ty {
867 return Err(invalid("chunk type differs from table schema"));
868 }
869 }
870 self.table.rows = self
871 .table
872 .rows
873 .checked_add(chunk.len())
874 .ok_or_else(|| invalid("row count overflow"))?;
875 Ok(())
876 }
877
878 fn encode_column(
886 index: usize,
887 held: &[PendingChunk],
888 mut dictionary: Option<&mut GlobalDictionary>,
889 ) -> Result<ColumnStripe> {
890 let mut stripe = ColumnStripe {
891 pages: Vec::with_capacity(held.len()),
892 codes: Vec::with_capacity(held.len()),
893 sieves: Vec::with_capacity(held.len()),
894 ranges: Vec::with_capacity(held.len()),
895 };
896 for pending in held {
897 let column = pending.chunk.column(index)?;
898 let (bytes, unique) = encode(column, dictionary.as_deref_mut())?;
899 if bytes.len() > MAX_PAGE {
900 return Err(invalid("column page exceeds the configured bound"));
901 }
902 let range = Range::of(column);
905 let sieve = match dictionary {
918 Some(_) => None,
919 None => Sieve::of(column, &range, SIEVE_BUDGET)
920 .filter(|sieve| sieve.len() < bytes.len()),
921 };
922 stripe.pages.push(bytes);
923 stripe.codes.push(unique);
924 stripe.sieves.push(sieve);
925 stripe.ranges.push(range);
926 }
927 Ok(stripe)
928 }
929
930 fn encode_columns(&mut self, held: &[PendingChunk]) -> Result<Vec<ColumnStripe>> {
939 let width = self.table.fields.len();
940 let workers = std::thread::available_parallelism()
941 .map_or(1, usize::from)
942 .min(MAX_ENCODE_WORKERS)
943 .min(width);
944 if workers <= 1 || held.len() <= 1 {
945 return self
946 .dictionaries
947 .iter_mut()
948 .enumerate()
949 .map(|(index, dictionary)| Self::encode_column(index, held, dictionary.as_mut()))
950 .collect();
951 }
952 let mut jobs: Vec<(usize, Option<GlobalDictionary>)> =
955 std::mem::take(&mut self.dictionaries).into_iter().enumerate().collect();
956 jobs.sort_by_key(|(index, _)| weight(&self.table.fields[*index].ty));
958 let queue = Mutex::new(jobs);
959 let pieces = std::thread::scope(|scope| {
960 (0..workers)
961 .map(|_| {
962 scope.spawn(|| {
963 let mut mine = Vec::new();
964 loop {
965 let taken = queue
966 .lock()
967 .map_err(|_| Error::internal("a native encode worker panicked"))?
968 .pop();
969 let Some((index, mut dictionary)) = taken else { break };
970 let encoded = Self::encode_column(index, held, dictionary.as_mut())?;
971 mine.push((index, dictionary, encoded));
972 }
973 Ok(mine)
974 })
975 })
976 .collect::<Vec<_>>()
977 .into_iter()
978 .map(|handle| {
979 handle.join().map_err(|_| Error::internal("a native encode worker panicked"))?
980 })
981 .collect::<Result<Vec<_>>>()
982 })?;
983 let mut dictionaries: Vec<Option<GlobalDictionary>> = (0..width).map(|_| None).collect();
984 let mut encoded: Vec<Option<ColumnStripe>> = (0..width).map(|_| None).collect();
985 for piece in pieces {
986 for (index, dictionary, stripe) in piece {
987 dictionaries[index] = dictionary;
988 encoded[index] = Some(stripe);
989 }
990 }
991 self.dictionaries = dictionaries;
992 encoded
993 .into_iter()
994 .map(|stripe| stripe.ok_or_else(|| Error::internal("a column was never encoded")))
995 .collect()
996 }
997
998 fn flush_pending(&mut self) -> Result<()> {
1000 if self.pending.is_empty() {
1001 return Ok(());
1002 }
1003 let width = self.table.fields.len();
1004 let mut held = std::mem::take(&mut self.pending);
1007 let parts = held.len();
1008 let encoded = self.encode_columns(&held)?;
1009 let mut pages = Vec::with_capacity(width);
1010 let mut memberships = vec![None; width];
1011 let mut ranges = Vec::with_capacity(width);
1012 let mut index = Vec::with_capacity(width.saturating_mul(index_section(parts)?));
1013 for stripe in &encoded {
1014 let offset = self.at;
1015 let section = index.len();
1016 let mut length = 0_usize;
1017 for bytes in &stripe.pages {
1018 write_at(&self.file, self.at + length as u64, bytes)?;
1019 put_u32(
1020 &mut index,
1021 u32::try_from(bytes.len()).map_err(|_| invalid("part length overflow"))?,
1022 );
1023 put_u64(&mut index, checksum(bytes));
1024 length = length
1025 .checked_add(bytes.len())
1026 .ok_or_else(|| invalid("column page length overflow"))?;
1027 }
1028 let hash = checksum(&index[section..]);
1029 put_u64(&mut index, hash);
1030 if length > MAX_PAGE {
1031 return Err(invalid("column page exceeds the configured bound"));
1032 }
1033 self.at = self
1034 .at
1035 .checked_add(length as u64)
1036 .ok_or_else(|| invalid("native file length overflow"))?;
1037 pages.push(Span {
1038 offset,
1039 length: u32::try_from(length).map_err(|_| invalid("page length overflow"))?,
1040 });
1041 ranges.push(merged_range(stripe.ranges.iter().cloned()));
1042 }
1043 for (membership, stripe) in memberships.iter_mut().zip(&encoded) {
1044 if stripe.codes.iter().all(Option::is_none) {
1045 continue;
1046 }
1047 let lists = stripe
1048 .codes
1049 .iter()
1050 .map(|codes| codes.clone().unwrap_or_default())
1051 .collect::<Vec<_>>();
1052 let bytes = encode_membership(&merged_codes(lists));
1053 let offset = self.at;
1054 self.put(&bytes)?;
1055 *membership = Some(Page {
1056 offset,
1057 length: u32::try_from(bytes.len())
1058 .map_err(|_| invalid("membership page length overflow"))?,
1059 hash: checksum(&bytes),
1060 });
1061 }
1062 let mut sieves = vec![None; width];
1063 for (page, stripe) in sieves.iter_mut().zip(&encoded) {
1064 if stripe.sieves.iter().all(Option::is_none) {
1065 continue;
1066 }
1067 let bytes = encode_sieves(stripe.sieves.iter())?;
1068 let offset = self.at;
1069 self.put(&bytes)?;
1070 *page = Some(Page {
1071 offset,
1072 length: u32::try_from(bytes.len())
1073 .map_err(|_| invalid("sieve page length overflow"))?,
1074 hash: checksum(&bytes),
1075 });
1076 }
1077 let mut part_ranges = vec![None; width];
1083 if parts > 1 {
1084 for ((page, stripe), span) in part_ranges.iter_mut().zip(&encoded).zip(&pages) {
1085 let bytes = encode_part_ranges(&stripe.ranges)?;
1086 if bytes.len() >= span.length as usize {
1087 continue;
1088 }
1089 let offset = self.at;
1090 self.put(&bytes)?;
1091 *page = Some(Page {
1092 offset,
1093 length: u32::try_from(bytes.len())
1094 .map_err(|_| invalid("part range page length overflow"))?,
1095 hash: checksum(&bytes),
1096 });
1097 }
1098 }
1099 let offset = self.at;
1100 self.put(&index)?;
1101 let index = Span {
1102 offset,
1103 length: u32::try_from(index.len())
1104 .map_err(|_| invalid("index page length overflow"))?,
1105 };
1106 let mut rows = 0_usize;
1107 let mut lengths = Vec::with_capacity(parts);
1108 let mut span = None;
1109 for pending in held.drain(..) {
1110 let part = pending.chunk.len();
1111 rows = rows.checked_add(part).ok_or_else(|| invalid("row count overflow"))?;
1112 lengths.push(u32::try_from(part).map_err(|_| invalid("part row count overflow"))?);
1113 span = Some(
1114 span.map_or((pending.order, pending.order), |(first, _)| (first, pending.order)),
1115 );
1116 }
1117 self.order.push(span.ok_or_else(|| invalid("a stripe was flushed with no parts"))?);
1118 self.table.stripes.push(Stripe {
1119 rows,
1120 parts: lengths,
1121 index,
1122 pages,
1123 memberships,
1124 sieves,
1125 part_ranges,
1126 zone: Zone::from_ranges(ranges),
1127 });
1128 self.pending = held;
1130 Ok(())
1131 }
1132
1133 fn numeric_frequency(&self, column: usize) -> Result<Option<FrequencySummary>> {
1137 let ty = &self.table.fields[column].ty;
1138 if !matches!(
1139 ty,
1140 LogicalType::TinyInt
1141 | LogicalType::SmallInt
1142 | LogicalType::Integer
1143 | LogicalType::BigInt
1144 | LogicalType::UTinyInt
1145 | LogicalType::USmallInt
1146 | LogicalType::UInteger
1147 | LogicalType::UBigInt
1148 | LogicalType::Date
1149 | LogicalType::Timestamp
1150 ) {
1151 return Ok(None);
1152 }
1153 let mut candidates: HashMap<FrequencyValue, u32> = HashMap::new();
1154 let mut decrements = 0_u64;
1155 self.visit_numeric(column, |_, value| {
1156 if let Some(count) = candidates.get_mut(&value) {
1157 *count = count.saturating_add(1);
1158 } else if candidates.len() < FREQUENCY_CANDIDATES {
1159 candidates.insert(value, 1);
1160 } else {
1161 candidates.retain(|_, count| {
1162 *count -= 1;
1163 *count != 0
1164 });
1165 decrements = decrements.saturating_add(1);
1166 }
1167 })?;
1168 let (exact, ordinals) = if decrements == 0 {
1169 (
1170 candidates
1171 .into_iter()
1172 .map(|(value, count)| (value, u64::from(count)))
1173 .collect::<HashMap<_, _>>(),
1174 Vec::new(),
1175 )
1176 } else {
1177 let mut lower = candidates.values().copied().collect::<Vec<_>>();
1178 lower.sort_unstable_by(|left, right| right.cmp(left));
1179 if lower.len() < FREQUENCY_BUILD_RANK
1180 || u64::from(lower[FREQUENCY_BUILD_RANK - 1]) <= decrements
1181 {
1182 return Ok(None);
1183 }
1184 let mut exact =
1185 candidates.into_keys().map(|value| (value, 0_u64)).collect::<HashMap<_, _>>();
1186 let mut ordinals = Vec::new();
1187 let mut exceeded = false;
1188 self.visit_numeric(column, |ordinal, value| {
1189 if let Some(count) = exact.get_mut(&value) {
1190 *count = count.saturating_add(1);
1191 if !exceeded {
1192 if ordinals.len() < FREQUENCY_ORDINALS {
1193 ordinals.push(ordinal);
1194 } else {
1195 ordinals.clear();
1196 exceeded = true;
1197 }
1198 }
1199 }
1200 })?;
1201 (exact, ordinals)
1202 };
1203 let mut entries = exact
1204 .into_iter()
1205 .map(|(value, count)| FrequencyEntry { value, count })
1206 .collect::<Vec<_>>();
1207 entries.sort_unstable_by(|left, right| {
1208 right.count.cmp(&left.count).then_with(|| frequency_order(left.value, right.value))
1209 });
1210 let omitted_max =
1211 entries.get(FREQUENCY_ENTRIES).map_or(decrements, |entry| decrements.max(entry.count));
1212 entries.truncate(FREQUENCY_ENTRIES);
1213 Ok(Some(FrequencySummary { entries, omitted_max, ordinals }))
1214 }
1215
1216 fn visit_numeric(
1217 &self,
1218 column: usize,
1219 mut visit: impl FnMut(u64, FrequencyValue),
1220 ) -> Result<()> {
1221 let ty = &self.table.fields[column].ty;
1222 let mut start = 0_u64;
1223 for stripe in &self.table.stripes {
1224 let spans = read_index(&self.file, stripe, column)?;
1225 let page = stripe.pages[column];
1226 let mut bytes = vec![0; page.length as usize];
1227 read_at(&self.file, page.offset, &mut bytes)?;
1228 for (span, &rows) in spans.iter().zip(&stripe.parts) {
1229 let part = part_bytes(&bytes, *span)?;
1230 if checksum(part) != span.hash {
1231 return Err(invalid("column page checksum differs while building frequencies"));
1232 }
1233 let rows = rows as usize;
1234 let vector = decode(ty, rows, part, None)?;
1235 for row in 0..rows {
1237 let value = if vector.is_null_at(row) {
1238 FrequencyValue::Null
1239 } else {
1240 let widened = match vector.signed_at(row) {
1244 Some(value) => Some(value),
1245 None => match vector.value_at(row) {
1246 Value::UTinyInt(value) => Some(i128::from(value)),
1247 Value::USmallInt(value) => Some(i128::from(value)),
1248 Value::UInteger(value) => Some(i128::from(value)),
1249 Value::UBigInt(value) => Some(i128::from(value)),
1250 _ => None,
1251 },
1252 };
1253 FrequencyValue::Integer(widened.ok_or_else(|| {
1254 invalid("numeric frequency page did not contain an integer value")
1255 })?)
1256 };
1257 visit(start.saturating_add(row as u64), value);
1258 }
1259 start = start.saturating_add(rows as u64);
1260 }
1261 }
1262 Ok(())
1263 }
1264
1265 fn numeric_frequencies(&self) -> Result<Vec<Option<FrequencySummary>>> {
1273 let mut columns = self
1274 .table
1275 .fields
1276 .iter()
1277 .enumerate()
1278 .filter_map(|(column, field)| {
1279 matches!(
1280 field.ty,
1281 LogicalType::TinyInt
1282 | LogicalType::SmallInt
1283 | LogicalType::Integer
1284 | LogicalType::BigInt
1285 | LogicalType::UTinyInt
1286 | LogicalType::USmallInt
1287 | LogicalType::UInteger
1288 | LogicalType::UBigInt
1289 | LogicalType::Date
1290 | LogicalType::Timestamp
1291 )
1292 .then_some(column)
1293 })
1294 .collect::<Vec<_>>();
1295 let workers = std::thread::available_parallelism()
1296 .map_or(1, usize::from)
1297 .min(MAX_FREQUENCY_WORKERS)
1298 .min(columns.len());
1299 if workers <= 1 {
1300 let mut frequencies = vec![None; self.table.fields.len()];
1301 for column in columns {
1302 frequencies[column] = self.numeric_frequency(column)?;
1303 }
1304 return Ok(frequencies);
1305 }
1306 columns.sort_by_key(|&column| weight(&self.table.fields[column].ty));
1309 let queue = Mutex::new(columns);
1310 let pieces = std::thread::scope(|scope| {
1311 (0..workers)
1312 .map(|_| {
1313 scope.spawn(|| {
1314 let mut mine = Vec::new();
1315 loop {
1316 let taken = queue
1317 .lock()
1318 .map_err(|_| Error::internal("a native frequency worker panicked"))?
1319 .pop();
1320 let Some(column) = taken else { break };
1321 mine.push((column, self.numeric_frequency(column)?));
1322 }
1323 Ok(mine)
1324 })
1325 })
1326 .collect::<Vec<_>>()
1327 .into_iter()
1328 .map(|handle| {
1329 handle
1330 .join()
1331 .map_err(|_| Error::internal("a native frequency worker panicked"))?
1332 })
1333 .collect::<Result<Vec<_>>>()
1334 })?;
1335 let mut frequencies = vec![None; self.table.fields.len()];
1336 for piece in pieces {
1337 for (column, summary) in piece {
1338 frequencies[column] = summary;
1339 }
1340 }
1341 Ok(frequencies)
1342 }
1343
1344 fn close(&mut self) -> Result<Entry> {
1355 self.flush_pending()?;
1356 let mut stripes = std::mem::take(&mut self.order)
1357 .into_iter()
1358 .zip(std::mem::take(&mut self.table.stripes))
1359 .collect::<Vec<_>>();
1360 stripes.sort_by_key(|(order, _)| order.0);
1361 let mut previous: Option<(u64, u64)> = None;
1362 for ((first, last), _) in &stripes {
1363 if previous.is_some_and(|previous| previous >= *first) {
1364 return Err(invalid("chunks did not arrive in source order"));
1365 }
1366 previous = Some(*last);
1367 }
1368 self.table.stripes = stripes.into_iter().map(|(_, stripe)| stripe).collect();
1369 self.table.frequencies = self.numeric_frequencies()?;
1370 let dictionaries = std::mem::take(&mut self.dictionaries);
1371 let orders = rankings(&dictionaries)?;
1372 for (index, (dictionary, order)) in dictionaries.into_iter().zip(orders).enumerate() {
1373 let Some(dictionary) = dictionary else { continue };
1374 self.table.distincts[index] =
1378 Some(dictionary.counts.iter().filter(|count| **count != 0).count() as u64);
1379 self.table.frequencies[index] = Some(code_frequency(&dictionary));
1380 let encoded = encode_global_dictionary(dictionary, &order)?;
1381 let offset = self.at;
1382 self.put(&encoded.index)?;
1383 self.put(&encoded.ranks)?;
1384 for block in &encoded.payload {
1385 self.put(block)?;
1386 }
1387 let payload_len =
1388 encoded.payload.iter().try_fold(0_usize, |len, block| len.checked_add(block.len()));
1389 let length = payload_len
1390 .and_then(|len| len.checked_add(encoded.index.len()))
1391 .and_then(|len| len.checked_add(encoded.ranks.len()))
1392 .ok_or_else(|| invalid("dictionary page length overflow"))?;
1393 self.table.dictionaries[index] = Some(Page {
1394 offset,
1395 length: u32::try_from(length)
1396 .map_err(|_| invalid("dictionary page length overflow"))?,
1397 hash: checksum(&encoded.index),
1398 });
1399 }
1400 let directory = encode_directory(&self.table)?;
1401 if directory.len() > MAX_DIRECTORY {
1402 return Err(invalid("directory exceeds the configured bound"));
1403 }
1404 let offset = self.at;
1405 self.put(&directory)?;
1406 Ok(Entry {
1407 name: self.table.name.clone(),
1408 fields: self.table.fields.clone(),
1409 rows: self.table.rows,
1410 directory: Page {
1411 offset,
1412 length: u32::try_from(directory.len())
1413 .map_err(|_| invalid("directory length overflow"))?,
1414 hash: checksum(&directory),
1415 },
1416 })
1417 }
1418
1419 pub fn finish(mut self) -> Result<Table> {
1429 let entry = self.close()?;
1430 let mut tables = std::mem::take(&mut self.closed);
1431 tables.push(entry);
1432 let catalog = encode_catalog(&tables)?;
1433 if catalog.len() > MAX_DIRECTORY {
1434 return Err(invalid("catalog exceeds the configured bound"));
1435 }
1436 let offset = self.at;
1437 self.put(&catalog)?;
1438 self.file.sync_all().map_err(io)?;
1442 let slot = Slot {
1443 offset,
1444 length: u32::try_from(catalog.len()).map_err(|_| invalid("catalog length overflow"))?,
1445 generation: self.generation,
1446 hash: checksum(&catalog),
1447 };
1448 write_at(&self.file, 16, &slot.bytes())?;
1451 self.file.sync_all().map_err(io)?;
1452 Ok(self.table)
1453 }
1454}
1455
1456#[derive(Debug, Clone)]
1458pub struct Reader {
1459 file: Arc<File>,
1460 table: Arc<Table>,
1461 dictionaries: Arc<Vec<OnceLock<Arc<Vector>>>>,
1462 loading: Arc<Vec<Mutex<()>>>,
1471 opened: Arc<AtomicUsize>,
1475 sieves: Arc<Vec<Vec<SieveSlot>>>,
1479 part_ranges: Arc<Vec<Vec<RangeSlot>>>,
1482 places: Arc<Vec<Place>>,
1484 cache: Arc<Vec<Mutex<Cached>>>,
1485 pages: Arc<AtomicUsize>,
1488 indexes: Arc<AtomicUsize>,
1491 kept: Arc<AtomicUsize>,
1494 size: u64,
1496 directory: u64,
1498 opening: Opening,
1500}
1501
1502#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1514pub struct Opening {
1515 pub reads: u32,
1518 pub bytes: u64,
1520}
1521
1522#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1524pub struct Reads {
1525 pub opening: Opening,
1527 pub pages: usize,
1529 pub indexes: usize,
1531 pub dictionaries: usize,
1534}
1535
1536#[derive(Debug, Clone, Copy)]
1538struct Place {
1539 stripe: u32,
1540 part: u32,
1541 rows: u32,
1542}
1543
1544#[derive(Debug, Clone, Copy)]
1546struct PartSpan {
1547 start: usize,
1548 length: usize,
1549 hash: u64,
1550}
1551
1552#[derive(Debug, Clone)]
1558struct CachedColumn {
1559 stripe: usize,
1560 index: Arc<Vec<PartSpan>>,
1561 page: Option<Arc<Vec<u8>>>,
1562}
1563
1564#[derive(Debug, Default)]
1584struct Cached {
1585 pages: Vec<Option<Arc<Vec<u8>>>>,
1586 order: VecDeque<usize>,
1587 loading: Vec<usize>,
1588 index: Vec<Option<Arc<Vec<PartSpan>>>>,
1589}
1590
1591const CACHED_STRIPES_PER_COLUMN: usize = 4;
1603
1604type SieveSlot = OnceLock<Arc<Vec<Option<Sieve>>>>;
1606
1607type RangeSlot = OnceLock<Arc<Vec<Range>>>;
1608
1609#[derive(Debug)]
1610struct NativeText {
1611 file: Arc<File>,
1612 values: usize,
1614 offsets: Vec<u8>,
1623 offset_bits: usize,
1626 ranks: usize,
1628 rank_at: u64,
1632 rank_ends: Vec<u64>,
1636 rank_hashes: Vec<u64>,
1637 rank_blocks: Vec<OnceLock<Result<Vec<u8>>>>,
1638 code_bits: usize,
1641 code_ranks: OnceLock<Option<Vec<u32>>>,
1648 payload: u64,
1649 ends: Vec<u64>,
1652 hashes: Vec<u64>,
1653 blocks: Vec<OnceLock<Result<Vec<u8>>>>,
1655 keep_budget: usize,
1658 payload_kept: AtomicUsize,
1666 searched: Mutex<HashMap<Vec<u8>, (usize, bool)>>,
1683}
1684
1685const TEXT_SEARCH_MEMO: usize = 64;
1690
1691const TEXT_PAYLOAD_VALUES: usize = 1024;
1707
1708const TEXT_KEEP_BUDGET: usize = 256 * 1024 * 1024;
1729
1730const TEXT_OFFSET_RUN: usize = 512;
1737
1738const DICTIONARY_HEADER: usize = 16;
1741
1742const TEXT_RANK_BLOCK: usize = 512;
1753
1754const RANK_BLOCK_HEADER: usize = size_of::<u64>() + 1;
1768
1769impl NativeText {
1770 fn payload_block(&self, block: usize) -> Result<Option<&[u8]>> {
1777 let Some(slot) = self.blocks.get(block) else { return Ok(None) };
1778 let bytes = slot.get_or_init(|| self.decode_block(block)).as_ref().map_err(Clone::clone)?;
1779 Ok(Some(bytes.as_slice()))
1780 }
1781
1782 fn decode_block(&self, block: usize) -> Result<Vec<u8>> {
1787 let start = if block == 0 { 0 } else { self.ends[block - 1] };
1788 let end = self.ends[block];
1789 let len = end
1790 .checked_sub(start)
1791 .ok_or_else(|| invalid("global dictionary block ends before it starts"))?;
1792 let mut stored = vec![
1793 0;
1794 usize::try_from(len).map_err(|_| invalid(
1795 "global dictionary block does not fit in memory"
1796 ))?
1797 ];
1798 read_at(&self.file, self.payload + start, &mut stored)?;
1799 if checksum(&stored) != self.hashes[block] {
1800 return Err(invalid("global dictionary payload checksum differs"));
1801 }
1802 let first = block * TEXT_PAYLOAD_VALUES;
1803 let last = (first + TEXT_PAYLOAD_VALUES).min(self.values);
1804 let want = self.end_within(last - 1)? as usize;
1805 let values = string::decode_flat(&stored)?;
1806 if values.len() != last - first {
1807 return Err(invalid("global dictionary block holds the wrong value count"));
1808 }
1809 let bytes = values.into_bytes();
1810 if bytes.len() != want {
1811 return Err(invalid("global dictionary block decodes to the wrong length"));
1812 }
1813 Ok(bytes)
1814 }
1815
1816 fn end_within(&self, index: usize) -> Result<u32> {
1818 let run = index / TEXT_OFFSET_RUN;
1819 let bytes = self
1820 .offsets
1821 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
1822 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1823 let end = bitpack::tail_at(bytes, self.offset_bits, index % TEXT_OFFSET_RUN)
1824 .map_err(|_| invalid("global dictionary offsets are short"))?;
1825 u32::try_from(end).map_err(|_| invalid("global dictionary offset is past the payload"))
1826 }
1827
1828 fn ends_within(&self, first: usize, last: usize) -> Result<Vec<u64>> {
1841 let mut ends = Vec::with_capacity(last.saturating_sub(first));
1842 let mut at = first;
1843 while at < last {
1844 let run = at / TEXT_OFFSET_RUN;
1845 let stop = ((run + 1) * TEXT_OFFSET_RUN).min(last);
1846 let held = self.values.saturating_sub(run * TEXT_OFFSET_RUN).min(TEXT_OFFSET_RUN);
1847 let bytes = self
1848 .offsets
1849 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
1850 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1851 let run_ends = bitpack::unpack_tail(bytes, self.offset_bits, held)
1852 .map_err(|_| invalid("global dictionary offsets are short"))?;
1853 let within = run_ends
1854 .get(at % TEXT_OFFSET_RUN..stop - run * TEXT_OFFSET_RUN)
1855 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1856 ends.extend_from_slice(within);
1857 at = stop;
1858 }
1859 Ok(ends)
1860 }
1861
1862 fn start_within(&self, index: usize) -> Result<u32> {
1865 if index % TEXT_PAYLOAD_VALUES == 0 { Ok(0) } else { self.end_within(index - 1) }
1866 }
1867
1868 fn span_within(&self, index: usize) -> Result<(u32, u32)> {
1876 let within = index % TEXT_OFFSET_RUN;
1877 let (start, end) = if within == 0 {
1878 (self.start_within(index)?, self.end_within(index)?)
1879 } else {
1880 let run = index / TEXT_OFFSET_RUN;
1881 let bytes = self
1882 .offsets
1883 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
1884 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1885 let (start, end) = bitpack::tail_pair(bytes, self.offset_bits, within)
1886 .map_err(|_| invalid("global dictionary offsets are short"))?;
1887 let ends = u32::try_from(end)
1888 .map_err(|_| invalid("global dictionary offset is past the payload"))?;
1889 let starts = u32::try_from(start)
1890 .map_err(|_| invalid("global dictionary offset is past the payload"))?;
1891 (starts, ends)
1892 };
1893 if start > end {
1894 return Err(invalid("global dictionary value ends before it starts"));
1895 }
1896 Ok((start, end))
1897 }
1898
1899 fn rank_parts(&self, rank: usize) -> Result<(&[u8], usize)> {
1906 let slot = self
1907 .rank_blocks
1908 .get(rank / TEXT_RANK_BLOCK)
1909 .ok_or_else(|| invalid("global dictionary rank is past the order"))?;
1910 let block = slot
1911 .get_or_init(|| {
1912 let which = rank / TEXT_RANK_BLOCK;
1913 let start = if which == 0 { 0 } else { self.rank_ends[which - 1] };
1914 let end = self.rank_ends[which];
1915 let mut bytes = vec![0; (end - start) as usize];
1916 read_at(&self.file, self.rank_at + start, &mut bytes)?;
1917 if checksum(&bytes)
1918 != *self
1919 .rank_hashes
1920 .get(rank / TEXT_RANK_BLOCK)
1921 .ok_or_else(|| invalid("global dictionary rank block has no checksum"))?
1922 {
1923 return Err(invalid("global dictionary rank checksum differs"));
1924 }
1925 Ok(bytes)
1926 })
1927 .as_ref()
1928 .map_err(Clone::clone)?;
1929 Ok((block.as_slice(), rank % TEXT_RANK_BLOCK))
1930 }
1931
1932 fn head_at(&self, rank: usize) -> Result<u64> {
1934 let (block, within) = self.rank_parts(rank)?;
1935 let (base, width, packed) = rank_heads(block)?;
1936 let above = bitpack::tail_at(packed, width, within)
1937 .map_err(|_| invalid("global dictionary rank block is short of heads"))?;
1938 Ok(base.wrapping_add(above))
1939 }
1940
1941 fn rank_codes<'block>(&self, block: &'block [u8], count: usize) -> Result<&'block [u8]> {
1943 let (_, width, packed) = rank_heads(block)?;
1944 packed
1945 .get(bitpack::tail_len(count, width)..)
1946 .ok_or_else(|| invalid("global dictionary rank block is short of codes"))
1947 }
1948
1949 fn rank_block_len(&self, rank: usize) -> usize {
1951 let first = rank / TEXT_RANK_BLOCK * TEXT_RANK_BLOCK;
1952 TEXT_RANK_BLOCK.min(self.ranks - first)
1953 }
1954}
1955
1956fn rank_heads(block: &[u8]) -> Result<(u64, usize, &[u8])> {
1958 let header = block
1959 .get(..RANK_BLOCK_HEADER)
1960 .ok_or_else(|| invalid("global dictionary rank block is short"))?;
1961 let base = u64::from_le_bytes(header[..8].try_into().expect("eight bytes"));
1962 let width = header[8] as usize;
1963 if width > 64 {
1964 return Err(invalid("global dictionary rank block packs heads past a word"));
1965 }
1966 Ok((base, width, &block[RANK_BLOCK_HEADER..]))
1967}
1968
1969fn offset_width(offsets: &[u32]) -> usize {
1976 let values = offsets.len() - 1;
1977 let mut span = 0;
1978 for first in (0..values).step_by(TEXT_PAYLOAD_VALUES) {
1979 let last = (first + TEXT_PAYLOAD_VALUES).min(values);
1980 span = span.max(offsets[last] - offsets[first]);
1981 }
1982 (u32::BITS - span.leading_zeros()) as usize
1983}
1984
1985fn offset_bytes(values: usize, bits: usize) -> usize {
1988 let full = values / TEXT_OFFSET_RUN;
1989 let rest = values % TEXT_OFFSET_RUN;
1990 full * TEXT_OFFSET_RUN / 8 * bits + bitpack::tail_len(rest, bits)
1991}
1992
1993fn encode_offsets(offsets: &[u32], bits: usize, out: &mut Vec<u8>) -> Result<()> {
1995 let values = offsets.len() - 1;
1996 let mut run = Vec::with_capacity(TEXT_OFFSET_RUN);
1997 for first in (0..values).step_by(TEXT_OFFSET_RUN) {
1998 let last = (first + TEXT_OFFSET_RUN).min(values);
1999 let base = offsets[first / TEXT_PAYLOAD_VALUES * TEXT_PAYLOAD_VALUES];
2000 run.clear();
2001 run.extend((first..last).map(|value| u64::from(offsets[value + 1] - base)));
2002 bitpack::pack_tail(&run, bits, out)
2003 .map_err(|_| invalid("global dictionary offsets do not pack"))?;
2004 }
2005 Ok(())
2006}
2007
2008fn code_width(values: usize) -> usize {
2010 match u64::try_from(values).unwrap_or(u64::MAX) {
2011 0 | 1 => 0,
2012 last => (u64::BITS - (last - 1).leading_zeros()) as usize,
2013 }
2014}
2015
2016impl TextSource for NativeText {
2017 fn len(&self) -> usize {
2018 self.values
2019 }
2020
2021 fn bytes_at(&self, index: usize) -> Result<Option<&[u8]>> {
2022 if index >= self.values {
2023 return Ok(None);
2024 }
2025 let (start, end) = self.span_within(index)?;
2026 if start == end {
2027 return Ok(Some(&[]));
2028 }
2029 let block = index / TEXT_PAYLOAD_VALUES;
2032 let Some(bytes) = self.payload_block(block)? else { return Ok(None) };
2033 Ok(bytes.get(start as usize..end as usize))
2034 }
2035
2036 fn bytes_len_at(&self, index: usize) -> Result<Option<usize>> {
2037 if index >= self.values {
2038 return Ok(None);
2039 }
2040 let (start, end) = self.span_within(index)?;
2041 Ok(Some((end - start) as usize))
2042 }
2043
2044 fn sweep(
2057 &self,
2058 first: usize,
2059 limit: usize,
2060 body: &mut dyn FnMut(usize, &[u8]) -> Result<()>,
2061 ) -> Result<usize> {
2062 let limit = limit.min(self.values);
2063 if first >= limit {
2064 return Ok(first);
2065 }
2066 let block = first / TEXT_PAYLOAD_VALUES;
2067 let last = ((block + 1) * TEXT_PAYLOAD_VALUES).min(limit);
2068 let decoded;
2069 let bytes: &[u8] = match self.blocks.get(block).and_then(OnceLock::get) {
2070 Some(Ok(kept)) => kept,
2071 _ if self.payload_kept.load(Atomic::Relaxed) < self.keep_budget => {
2072 let kept = self
2073 .payload_block(block)?
2074 .ok_or_else(|| invalid("global dictionary block is past the payload"))?;
2075 self.payload_kept.fetch_add(kept.len(), Atomic::Relaxed);
2076 kept
2077 }
2078 _ => {
2079 decoded = self.decode_block(block)?;
2080 &decoded
2081 }
2082 };
2083 let ends = self.ends_within(first, last)?;
2084 if ends.len() != last - first {
2085 return Err(invalid("global dictionary offsets are short"));
2086 }
2087 let mut start = u64::from(self.start_within(first)?);
2088 for (index, &end) in (first..last).zip(&ends) {
2091 let value = usize::try_from(start)
2092 .ok()
2093 .zip(usize::try_from(end).ok())
2094 .and_then(|(from, to)| bytes.get(from..to))
2095 .ok_or_else(|| invalid("global dictionary value is past its block"))?;
2096 body(index, value)?;
2097 start = end;
2098 }
2099 Ok(last)
2100 }
2101
2102 fn ranks(&self) -> Option<usize> {
2103 (self.ranks > 0).then_some(self.ranks)
2104 }
2105
2106 fn below(&self, ranks: usize, wanted: &[u8]) -> Result<(usize, bool)> {
2114 let mut memo = self.searched.lock().map_err(|_| invalid("a poisoned dictionary search"))?;
2115 if let Some(&answer) = memo.get(wanted) {
2116 return Ok(answer);
2117 }
2118 let answer = search_below(self, ranks, wanted)?;
2119 if memo.len() >= TEXT_SEARCH_MEMO {
2120 memo.clear();
2121 }
2122 memo.insert(wanted.to_vec(), answer);
2123 Ok(answer)
2124 }
2125
2126 fn compare_rank(&self, rank: usize, wanted: &[u8]) -> Result<Ordering> {
2127 let settled = self.head_at(rank)?.cmp(&head(wanted));
2131 if settled != Ordering::Equal {
2132 return Ok(settled);
2133 }
2134 let code = self.code_at_rank(rank)?;
2135 let bytes = self
2136 .bytes_at(code as usize)?
2137 .ok_or_else(|| invalid("global dictionary order names a code it does not have"))?;
2138 Ok(bytes.cmp(wanted))
2139 }
2140
2141 fn code_at_rank(&self, rank: usize) -> Result<u32> {
2142 let (block, within) = self.rank_parts(rank)?;
2143 let codes = self.rank_codes(block, self.rank_block_len(rank))?;
2144 let code = bitpack::tail_at(codes, self.code_bits, within)
2145 .map_err(|_| invalid("global dictionary rank block is short of codes"))?;
2146 let code = u32::try_from(code)
2147 .map_err(|_| invalid("global dictionary order names a code it does not have"))?;
2148 if code as usize >= self.len() {
2149 return Err(invalid("global dictionary order names a code it does not have"));
2150 }
2151 Ok(code)
2152 }
2153
2154 fn code_ranks(&self) -> Option<&[u32]> {
2155 if self.ranks == 0 || self.ranks != self.len() {
2159 return None;
2160 }
2161 self.code_ranks
2162 .get_or_init(|| {
2163 let mut ranks = vec![u32::MAX; self.ranks];
2164 for first in (0..self.ranks).step_by(TEXT_RANK_BLOCK) {
2167 let (block, _) = self.rank_parts(first).ok()?;
2168 let count = self.rank_block_len(first);
2169 let codes = self.rank_codes(block, count).ok()?;
2170 for (within, code) in bitpack::unpack_tail(codes, self.code_bits, count)
2171 .ok()?
2172 .into_iter()
2173 .enumerate()
2174 {
2175 let code = usize::try_from(code).ok()?;
2176 *ranks.get_mut(code)? = u32::try_from(first + within).ok()?;
2177 }
2178 }
2179 if ranks.contains(&u32::MAX) {
2180 return None;
2181 }
2182 Some(ranks)
2183 })
2184 .as_deref()
2185 }
2186
2187 fn footprint(&self) -> usize {
2188 self.offsets.capacity()
2189 + self
2190 .code_ranks
2191 .get()
2192 .and_then(Option::as_ref)
2193 .map_or(0, |ranks| ranks.capacity() * size_of::<u32>())
2194 + self.rank_hashes.capacity() * size_of::<u64>()
2195 + self.rank_ends.capacity() * size_of::<u64>()
2196 + self.rank_blocks.capacity() * size_of::<OnceLock<Result<Vec<u8>>>>()
2197 + self
2198 .rank_blocks
2199 .iter()
2200 .filter_map(OnceLock::get)
2201 .filter_map(|result| result.as_ref().ok())
2202 .map(Vec::capacity)
2203 .sum::<usize>()
2204 + self.blocks.capacity() * size_of::<OnceLock<Result<Vec<u8>>>>()
2205 + self.hashes.capacity() * size_of::<u64>()
2206 + self.ends.capacity() * size_of::<u64>()
2207 + self
2208 .blocks
2209 .iter()
2210 .filter_map(OnceLock::get)
2211 .filter_map(|result| result.as_ref().ok())
2212 .map(Vec::capacity)
2213 .sum::<usize>()
2214 }
2215}
2216
2217fn places(table: &Table) -> Result<Vec<Place>> {
2219 let mut places = Vec::with_capacity(table.stripes.len().saturating_mul(STRIPE_PARTS));
2220 for (at, stripe) in table.stripes.iter().enumerate() {
2221 let index = u32::try_from(at).map_err(|_| invalid("too many stripes"))?;
2222 for (part, &rows) in stripe.parts.iter().enumerate() {
2223 places.push(Place {
2224 stripe: index,
2225 part: u32::try_from(part).map_err(|_| invalid("too many parts in a stripe"))?,
2226 rows,
2227 });
2228 }
2229 }
2230 Ok(places)
2231}
2232
2233fn read_index(file: &File, stripe: &Stripe, column: usize) -> Result<Vec<PartSpan>> {
2238 let parts = stripe.parts.len();
2239 let section = index_section(parts)?;
2240 let at = column.checked_mul(section).ok_or_else(|| invalid("index page offset overflow"))?;
2241 let end = at.checked_add(section).ok_or_else(|| invalid("index page offset overflow"))?;
2242 if end > stripe.index.length as usize {
2243 return Err(invalid("index page is shorter than its columns"));
2244 }
2245 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
2246 let mut bytes = vec![0; section];
2247 let offset = stripe
2248 .index
2249 .offset
2250 .checked_add(at as u64)
2251 .ok_or_else(|| invalid("index page offset overflow"))?;
2252 read_at(file, offset, &mut bytes)?;
2253 let entries = section - size_of::<u64>();
2254 let stored = u64::from_le_bytes(bytes[entries..].try_into().expect("eight bytes"));
2255 if checksum(&bytes[..entries]) != stored {
2256 return Err(invalid(&format!(
2259 "index page section checksum differs, column {column} of {parts} parts at {offset}, \
2260 wanted {stored:016x} and got {:016x}",
2261 checksum(&bytes[..entries]),
2262 )));
2263 }
2264 let mut spans = Vec::with_capacity(parts);
2265 let mut start = 0_usize;
2266 for part in 0..parts {
2267 let at = part * INDEX_ENTRY;
2268 let length = u32::from_le_bytes(bytes[at..at + 4].try_into().expect("four bytes")) as usize;
2269 let hash = u64::from_le_bytes(bytes[at + 4..at + 12].try_into().expect("eight bytes"));
2270 spans.push(PartSpan { start, length, hash });
2271 start = start.checked_add(length).ok_or_else(|| invalid("column page length overflow"))?;
2272 }
2273 if start != page.length as usize {
2274 return Err(invalid("column page length differs from its index"));
2275 }
2276 Ok(spans)
2277}
2278
2279fn part_bytes(page: &[u8], span: PartSpan) -> Result<&[u8]> {
2281 let end = span.start.checked_add(span.length).ok_or_else(|| invalid("part range overflow"))?;
2282 page.get(span.start..end).ok_or_else(|| invalid("part exceeds its column page"))
2283}
2284
2285fn remember(cached: &mut Cached, held: &CachedColumn, kept: usize) {
2290 if let Some(slot) = cached.index.get_mut(held.stripe) {
2291 if slot.is_none() {
2292 *slot = Some(Arc::clone(&held.index));
2293 }
2294 }
2295 let Some(page) = held.page.clone() else { return };
2296 let Some(slot) = cached.pages.get_mut(held.stripe) else { return };
2297 if slot.is_none() {
2298 cached.order.push_back(held.stripe);
2299 }
2300 *slot = Some(page);
2301 while cached.order.len() > kept.max(1) {
2302 let Some(oldest) = cached.order.pop_front() else { break };
2303 if let Some(slot) = cached.pages.get_mut(oldest) {
2304 *slot = None;
2305 }
2306 }
2307}
2308
2309#[derive(Debug, Clone)]
2318pub struct Catalog {
2319 file: Arc<File>,
2320 size: u64,
2321 entries: Arc<Vec<Entry>>,
2322 opening: Opening,
2323}
2324
2325impl Catalog {
2326 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
2332 let (file, size, bytes, opening) = slot_bytes(path)?;
2333 let entries = decode_catalog(&bytes, size)?;
2334 Ok(Self { file: Arc::new(file), size, entries: Arc::new(entries), opening })
2335 }
2336
2337 pub fn names(&self) -> impl ExactSizeIterator<Item = &str> {
2339 self.entries.iter().map(|entry| entry.name.as_str())
2340 }
2341
2342 #[must_use]
2344 pub fn len(&self) -> usize {
2345 self.entries.len()
2346 }
2347
2348 #[must_use]
2350 pub fn is_empty(&self) -> bool {
2351 self.entries.is_empty()
2352 }
2353
2354 pub fn table(&self, name: &str) -> Result<Reader> {
2360 let entry = self
2361 .entries
2362 .iter()
2363 .find(|entry| entry.name == name)
2364 .ok_or_else(|| invalid(&format!("the file holds no table called {name}")))?;
2365 let mut bytes = vec![0; entry.directory.length as usize];
2366 read_at(&self.file, entry.directory.offset, &mut bytes)?;
2367 if checksum(&bytes) != entry.directory.hash {
2368 return Err(invalid(&format!("the directory of table {name} does not checksum")));
2369 }
2370 let mut opening = self.opening;
2371 opening.reads += 1;
2372 opening.bytes += u64::from(entry.directory.length);
2373 Reader::build(
2374 Arc::clone(&self.file),
2375 self.size,
2376 decode_directory(&bytes, self.size)?,
2377 u64::from(entry.directory.length),
2378 opening,
2379 )
2380 }
2381}
2382
2383fn slot_bytes(path: impl AsRef<Path>) -> Result<(File, u64, Vec<u8>, Opening)> {
2388 let mut file = File::open(path).map_err(io)?;
2389 let size = file.metadata().map_err(io)?.len();
2390 if size < HEADER {
2391 return Err(invalid("file is shorter than its header"));
2392 }
2393 let mut header = [0; HEADER as usize];
2394 file.read_exact(&mut header).map_err(io)?;
2395 let mut opening = Opening { reads: 1, bytes: HEADER };
2396 let version = u32::from_le_bytes([header[8], header[9], header[10], header[11]]);
2397 if &header[..8] != MAGIC {
2402 return Err(invalid("the header does not begin with a rudb native magic"));
2403 }
2404 if version != FORMAT {
2405 return Err(invalid(&format!(
2406 "the file is format {version} and this build reads format {FORMAT}, so it has to \
2407 be written again"
2408 )));
2409 }
2410 let mut selected = None;
2411 for start in [16, 16 + SLOT_BYTES] {
2412 let slot = Slot::read(&header[start..start + SLOT_BYTES]);
2413 if slot.generation == 0 || slot.length == 0 || slot.length as usize > MAX_DIRECTORY {
2414 continue;
2415 }
2416 let Some(end) = slot.offset.checked_add(u64::from(slot.length)) else { continue };
2417 if slot.offset < HEADER || end > size {
2418 continue;
2419 }
2420 let mut bytes = vec![0; slot.length as usize];
2421 file.seek(SeekFrom::Start(slot.offset)).map_err(io)?;
2422 file.read_exact(&mut bytes).map_err(io)?;
2423 opening.reads += 1;
2424 opening.bytes += u64::from(slot.length);
2425 if checksum(&bytes) == slot.hash
2426 && selected
2427 .as_ref()
2428 .is_none_or(|(old, _): &(Slot, Vec<u8>)| old.generation < slot.generation)
2429 {
2430 selected = Some((slot, bytes));
2431 }
2432 }
2433 let (_, bytes) = selected.ok_or_else(|| invalid("no committed directory slot is valid"))?;
2434 Ok((file, size, bytes, opening))
2435}
2436
2437impl Reader {
2438 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
2445 let catalog = Catalog::open(path)?;
2446 let mut names = catalog.names();
2447 let name = names.next().ok_or_else(|| invalid("the file holds no table"))?.to_string();
2448 if names.next().is_some() {
2449 return Err(invalid(
2450 "the file holds more than one table, so it has to be opened by name",
2451 ));
2452 }
2453 catalog.table(&name)
2454 }
2455
2456 fn build(
2458 file: Arc<File>,
2459 size: u64,
2460 table: Table,
2461 directory: u64,
2462 opening: Opening,
2463 ) -> Result<Self> {
2464 let places = places(&table)?;
2465 let dictionaries = (0..table.fields.len()).map(|_| OnceLock::new()).collect();
2466 let table_fields = table.fields.len();
2467 let stripes = table.stripes.len();
2468 let cache = (0..table.fields.len())
2469 .map(|_| {
2470 Mutex::new(Cached {
2471 pages: (0..stripes).map(|_| None).collect(),
2472 index: (0..stripes).map(|_| None).collect(),
2473 ..Cached::default()
2474 })
2475 })
2476 .collect::<Vec<_>>();
2477 let sieves: Vec<Vec<SieveSlot>> = (0..table.fields.len())
2478 .map(|_| table.stripes.iter().map(|_| OnceLock::new()).collect())
2479 .collect();
2480 let part_ranges: Vec<Vec<RangeSlot>> = (0..table.fields.len())
2481 .map(|_| table.stripes.iter().map(|_| OnceLock::new()).collect())
2482 .collect();
2483 Ok(Self {
2484 file,
2485 table: Arc::new(table),
2486 dictionaries: Arc::new(dictionaries),
2487 loading: Arc::new((0..table_fields).map(|_| Mutex::new(())).collect()),
2488 opened: Arc::new(AtomicUsize::new(0)),
2489 sieves: Arc::new(sieves),
2490 part_ranges: Arc::new(part_ranges),
2491 places: Arc::new(places),
2492 cache: Arc::new(cache),
2493 pages: Arc::new(AtomicUsize::new(0)),
2494 indexes: Arc::new(AtomicUsize::new(0)),
2495 kept: Arc::new(AtomicUsize::new(CACHED_STRIPES_PER_COLUMN)),
2496 size,
2497 directory,
2498 opening,
2499 })
2500 }
2501
2502 #[must_use]
2509 pub fn reads(&self) -> Reads {
2510 Reads {
2511 opening: self.opening,
2512 pages: self.pages.load(Atomic::Relaxed),
2513 indexes: self.indexes.load(Atomic::Relaxed),
2514 dictionaries: self.opened.load(Atomic::Relaxed),
2515 }
2516 }
2517
2518 #[must_use]
2523 pub fn layout(&self) -> Layout {
2524 let table = &self.table;
2525 let stripes = table.stripes.as_slice();
2526 let columns = table
2527 .fields
2528 .iter()
2529 .enumerate()
2530 .map(|(at, field)| ColumnLayout {
2531 name: field.name.clone(),
2532 kind: field.ty.to_string(),
2533 pages: sum(stripes.iter().map(|stripe| span_bytes(&stripe.pages, at))),
2534 memberships: sum(stripes.iter().map(|stripe| page_bytes(&stripe.memberships, at))),
2535 sieves: sum(stripes.iter().map(|stripe| page_bytes(&stripe.sieves, at))),
2536 part_ranges: sum(stripes.iter().map(|stripe| page_bytes(&stripe.part_ranges, at))),
2537 dictionary: page_bytes(&table.dictionaries, at),
2538 })
2539 .collect();
2540 Layout {
2541 file: self.size,
2542 rows: table.rows,
2543 stripes: stripes.len(),
2544 parts: self.places.len(),
2545 columns,
2546 indexes: sum(stripes.iter().map(|stripe| u64::from(stripe.index.length))),
2547 directory: self.directory,
2548 header: HEADER,
2549 }
2550 }
2551
2552 #[must_use]
2554 pub fn parts(&self) -> usize {
2555 self.places.len()
2556 }
2557
2558 #[must_use]
2565 pub fn stripe_parts(&self) -> Vec<std::ops::Range<usize>> {
2566 let mut runs = Vec::with_capacity(self.table.stripes.len());
2567 let mut start = 0;
2568 for stripe in &self.table.stripes {
2569 let end = start + stripe.parts.len();
2570 runs.push(start..end);
2571 start = end;
2572 }
2573 runs
2574 }
2575
2576 #[must_use]
2581 pub fn stripe_rows(&self, stripe: usize) -> usize {
2582 self.table.stripes.get(stripe).map_or(0, |held| held.rows)
2583 }
2584
2585 pub fn keep_stripes(&self, stripes: usize) {
2592 self.kept.fetch_max(stripes, Atomic::Relaxed);
2593 }
2594
2595 #[must_use]
2597 pub fn part_rows(&self, at: usize) -> usize {
2598 self.places.get(at).map_or(0, |place| place.rows as usize)
2599 }
2600
2601 #[must_use]
2603 pub fn table(&self) -> &Table {
2604 &self.table
2605 }
2606
2607 pub fn top_frequencies(&self, column: usize, top: usize) -> Result<Option<Vec<(Value, u64)>>> {
2616 let field = self
2617 .table
2618 .fields
2619 .get(column)
2620 .ok_or_else(|| invalid("frequency column index out of range"))?;
2621 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2622 return Ok(None);
2623 };
2624 if top == 0 || summary.entries.len() < top {
2625 return Ok(None);
2626 }
2627 let boundary = summary.entries[top - 1].count;
2628 if boundary <= summary.omitted_max {
2629 return Ok(None);
2630 }
2631 self.decode_frequencies(column, &field.ty, &summary.entries).map(Some)
2632 }
2633
2634 pub fn exact_frequencies(&self, column: usize) -> Result<Option<Vec<(Value, u64)>>> {
2654 let field = self
2655 .table
2656 .fields
2657 .get(column)
2658 .ok_or_else(|| invalid("frequency column index out of range"))?;
2659 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2660 return Ok(None);
2661 };
2662 if summary.omitted_max > 0 {
2663 return Ok(None);
2664 }
2665 self.decode_frequencies(column, &field.ty, &summary.entries).map(Some)
2666 }
2667
2668 fn decode_frequencies(
2670 &self,
2671 column: usize,
2672 ty: &LogicalType,
2673 entries: &[FrequencyEntry],
2674 ) -> Result<Vec<(Value, u64)>> {
2675 let dictionary = if *ty == LogicalType::Varchar { self.dictionary(column)? } else { None };
2676 let mut out = Vec::with_capacity(entries.len());
2677 for entry in entries {
2678 let value = match entry.value {
2679 FrequencyValue::Null => Value::Null,
2680 FrequencyValue::Integer(value) => match *ty {
2681 LogicalType::TinyInt => Value::TinyInt(
2682 i8::try_from(value)
2683 .map_err(|_| invalid("frequency TINYINT is out of range"))?,
2684 ),
2685 LogicalType::UTinyInt => Value::UTinyInt(
2686 u8::try_from(value)
2687 .map_err(|_| invalid("frequency UTINYINT is out of range"))?,
2688 ),
2689 LogicalType::USmallInt => Value::USmallInt(
2690 u16::try_from(value)
2691 .map_err(|_| invalid("frequency USMALLINT is out of range"))?,
2692 ),
2693 LogicalType::UInteger => Value::UInteger(
2694 u32::try_from(value)
2695 .map_err(|_| invalid("frequency UINTEGER is out of range"))?,
2696 ),
2697 LogicalType::UBigInt => Value::UBigInt(
2698 u64::try_from(value)
2699 .map_err(|_| invalid("frequency UBIGINT is out of range"))?,
2700 ),
2701 LogicalType::SmallInt => Value::SmallInt(
2702 i16::try_from(value)
2703 .map_err(|_| invalid("frequency SMALLINT is out of range"))?,
2704 ),
2705 LogicalType::Integer => Value::Integer(
2706 i32::try_from(value)
2707 .map_err(|_| invalid("frequency INTEGER is out of range"))?,
2708 ),
2709 LogicalType::BigInt => Value::BigInt(
2710 i64::try_from(value)
2711 .map_err(|_| invalid("frequency BIGINT is out of range"))?,
2712 ),
2713 LogicalType::Date => Value::Date(
2714 i32::try_from(value)
2715 .map_err(|_| invalid("frequency DATE is out of range"))?,
2716 ),
2717 LogicalType::Timestamp => Value::Timestamp(
2718 i64::try_from(value)
2719 .map_err(|_| invalid("frequency TIMESTAMP is out of range"))?,
2720 ),
2721 _ => return Err(invalid("integer frequency belongs to another type")),
2722 },
2723 FrequencyValue::Code(code) => dictionary
2724 .as_ref()
2725 .ok_or_else(|| invalid("frequency code has no dictionary"))?
2726 .try_value_at(code as usize)?,
2727 };
2728 out.push((value, entry.count));
2729 }
2730 Ok(out)
2731 }
2732
2733 pub fn frequency_occurrences(&self, column: usize) -> Result<Option<FrequencyOccurrences>> {
2743 self.table
2744 .fields
2745 .get(column)
2746 .ok_or_else(|| invalid("frequency column index out of range"))?;
2747 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2748 return Ok(None);
2749 };
2750 if summary.ordinals.is_empty() {
2751 return Ok(None);
2752 }
2753 Ok(Some(FrequencyOccurrences {
2754 omitted_max: summary.omitted_max,
2755 ordinals: summary.ordinals.clone(),
2756 }))
2757 }
2758
2759 pub fn distinct_values(&self, column: usize) -> Result<Option<u64>> {
2783 self.table
2784 .distincts
2785 .get(column)
2786 .copied()
2787 .ok_or_else(|| invalid("distinct column index out of range"))
2788 }
2789
2790 pub fn null_count(&self, column: usize) -> Result<u64> {
2801 if column >= self.table.fields.len() {
2802 return Err(invalid("null count column index out of range"));
2803 }
2804 let mut nulls = 0_u64;
2805 for stripe in &self.table.stripes {
2806 let range = stripe
2807 .zone
2808 .column(column)
2809 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
2810 nulls = nulls
2811 .checked_add(range.nulls as u64)
2812 .ok_or_else(|| invalid("null count overflow"))?;
2813 }
2814 Ok(nulls)
2815 }
2816
2817 pub fn text_extremes(&self, column: usize) -> Result<Option<(Value, Value)>> {
2832 if self.null_count(column)? > 0 {
2833 return Ok(None);
2834 }
2835 let Some(dictionary) = self.dictionary(column)? else { return Ok(None) };
2836 let Some(ranks) = dictionary.ranks() else { return Ok(None) };
2837 if ranks == 0 {
2838 return Ok(None);
2839 }
2840 let low = text_at_rank(&dictionary, 0)?;
2841 let high = text_at_rank(&dictionary, ranks - 1)?;
2842 Ok(Some((low, high)))
2843 }
2844
2845 pub fn exact_extremes(&self, column: usize) -> Result<Option<(Bound, Bound)>> {
2868 if column >= self.table.fields.len() {
2869 return Err(invalid("extremes column index out of range"));
2870 }
2871 let mut low: Option<Bound> = None;
2872 let mut high: Option<Bound> = None;
2873 for stripe in &self.table.stripes {
2874 let range = stripe
2875 .zone
2876 .column(column)
2877 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
2878 if !range.exact {
2879 return Ok(None);
2880 }
2881 let (Some(small), Some(large)) = (range.low.as_ref(), range.high.as_ref()) else {
2886 if stripe.rows > range.nulls {
2887 return Ok(None);
2888 }
2889 continue;
2890 };
2891 low = Some(low.map_or_else(|| small.clone(), |held| held.smaller(small.clone())));
2892 high = Some(high.map_or_else(|| large.clone(), |held| held.larger(large.clone())));
2893 }
2894 Ok(low.zip(high))
2895 }
2896
2897 pub fn exact_sum(&self, column: usize) -> Result<Option<(i128, u64)>> {
2910 if column >= self.table.fields.len() {
2911 return Err(invalid("sum column index out of range"));
2912 }
2913 let mut total = 0_i128;
2914 let mut rows = 0_u64;
2915 for stripe in &self.table.stripes {
2916 let range = stripe
2917 .zone
2918 .column(column)
2919 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
2920 let Some(part) = range.sum else { return Ok(None) };
2921 let Some(sum) = total.checked_add(part) else { return Ok(None) };
2922 total = sum;
2923 rows = rows.saturating_add(stripe.rows as u64 - range.nulls as u64);
2924 }
2925 Ok(Some((total, rows)))
2926 }
2927
2928 fn dictionary(&self, column: usize) -> Result<Option<Arc<Vector>>> {
2937 let Some(page) = self.table.dictionaries[column] else { return Ok(None) };
2938 if let Some(dictionary) = self.dictionaries[column].get() {
2939 return Ok(Some(Arc::clone(dictionary)));
2940 }
2941 let _queued = self.loading[column].lock().map_err(|_| invalid("a poisoned dictionary"))?;
2942 if let Some(dictionary) = self.dictionaries[column].get() {
2943 return Ok(Some(Arc::clone(dictionary)));
2944 }
2945 self.opened.fetch_add(1, Atomic::Relaxed);
2946 let dictionary = Arc::new(open_global_dictionary(
2947 Arc::clone(&self.file),
2948 page,
2949 &self.table.fields[column].ty,
2950 TEXT_KEEP_BUDGET,
2951 )?);
2952 let _ = self.dictionaries[column].set(Arc::clone(&dictionary));
2953 Ok(Some(dictionary))
2954 }
2955
2956 pub fn read(&self, part: usize, columns: &[usize]) -> Result<Chunk> {
2965 self.read_impl(part, columns, true)
2966 }
2967
2968 pub fn read_sparse(&self, part: usize, columns: &[usize]) -> Result<Chunk> {
2978 self.read_impl(part, columns, false)
2979 }
2980
2981 pub fn skips_codes(&self, part: usize, column: usize, candidates: &[u32]) -> Result<bool> {
2988 if candidates.is_empty() {
2989 return Ok(true);
2990 }
2991 if candidates.windows(2).any(|pair| pair[0] >= pair[1]) {
2992 return Err(Error::internal("native code candidates are not sorted and unique"));
2993 }
2994 let stripe = self.stripe_of(part)?;
2995 let Some(page) = stripe.memberships.get(column).copied().flatten() else {
2996 return Ok(false);
2997 };
2998 let mut bytes = vec![0; page.length as usize];
2999 read_at(&self.file, page.offset, &mut bytes)?;
3000 if checksum(&bytes) != page.hash {
3001 return Err(invalid("membership page checksum differs"));
3002 }
3003 let codes = decode_membership(&bytes)?;
3004 let mut left = 0;
3005 let mut right = 0;
3006 while left < codes.len() && right < candidates.len() {
3007 match codes[left].cmp(&candidates[right]) {
3008 Ordering::Less => left += 1,
3009 Ordering::Greater => right += 1,
3010 Ordering::Equal => return Ok(false),
3011 }
3012 }
3013 Ok(true)
3014 }
3015
3016 fn stripe_of(&self, part: usize) -> Result<&Stripe> {
3017 let place = self.places.get(part).ok_or_else(|| invalid("part index out of range"))?;
3018 self.table
3019 .stripes
3020 .get(place.stripe as usize)
3021 .ok_or_else(|| invalid("stripe index out of range"))
3022 }
3023
3024 fn held(&self, at: usize, stripe: &Stripe, column: usize, whole: bool) -> Result<CachedColumn> {
3041 let cache = self.cache.get(column).ok_or_else(|| invalid("column index out of range"))?;
3042 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3043 let known = cached.index.get(at).and_then(Clone::clone);
3044 let page = cached.pages.get(at).and_then(Clone::clone);
3045 if let Some(index) = known.clone() {
3046 if !whole || page.is_some() {
3047 return Ok(CachedColumn { stripe: at, index, page });
3048 }
3049 }
3050 if cached.loading.contains(&at) {
3051 drop(cached);
3052 if let Some(index) = known {
3056 return Ok(CachedColumn { stripe: at, index, page: None });
3057 }
3058 let held = self.page_of(stripe, column, at, false, None)?;
3059 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3060 remember(&mut cached, &held, self.kept.load(Atomic::Relaxed));
3061 return Ok(held);
3062 }
3063 cached.loading.push(at);
3064 drop(cached);
3065
3066 let read = self.page_of(stripe, column, at, whole, known);
3067
3068 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
3072 if let Some(position) = cached.loading.iter().position(|loading| *loading == at) {
3073 cached.loading.remove(position);
3074 }
3075 let held = read?;
3076 remember(&mut cached, &held, self.kept.load(Atomic::Relaxed));
3077 Ok(held)
3078 }
3079
3080 fn page_of(
3086 &self,
3087 stripe: &Stripe,
3088 column: usize,
3089 at: usize,
3090 whole: bool,
3091 known: Option<Arc<Vec<PartSpan>>>,
3092 ) -> Result<CachedColumn> {
3093 let index = match known {
3094 Some(index) => index,
3095 None => {
3096 self.indexes.fetch_add(1, Atomic::Relaxed);
3097 Arc::new(read_index(&self.file, stripe, column)?)
3098 }
3099 };
3100 let page = if whole {
3101 self.pages.fetch_add(1, Atomic::Relaxed);
3102 let span = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
3103 let mut bytes = vec![0; span.length as usize];
3104 read_at(&self.file, span.offset, &mut bytes)?;
3105 Some(Arc::new(bytes))
3106 } else {
3107 None
3108 };
3109 Ok(CachedColumn { stripe: at, index, page })
3110 }
3111
3112 fn read_impl(&self, at: usize, columns: &[usize], whole: bool) -> Result<Chunk> {
3113 let place = *self.places.get(at).ok_or_else(|| invalid("part index out of range"))?;
3114 let index = place.stripe as usize;
3115 let stripe =
3116 self.table.stripes.get(index).ok_or_else(|| invalid("stripe index out of range"))?;
3117 let rows = place.rows as usize;
3118 let mut picked = Vec::with_capacity(columns.len());
3119 for &column in columns {
3120 let field = self
3121 .table
3122 .fields
3123 .get(column)
3124 .ok_or_else(|| invalid("column index out of range"))?;
3125 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
3126 let held = self.held(index, stripe, column, whole)?;
3127 let span = *held
3128 .index
3129 .get(place.part as usize)
3130 .ok_or_else(|| invalid("part index out of range"))?;
3131 let owned;
3132 let bytes = match &held.page {
3133 Some(held) => part_bytes(held, span)?,
3134 None => {
3135 let offset = page
3136 .offset
3137 .checked_add(span.start as u64)
3138 .ok_or_else(|| invalid("part range overflow"))?;
3139 let mut bytes = vec![0; span.length];
3140 read_at(&self.file, offset, &mut bytes)?;
3141 owned = bytes;
3142 &owned
3143 }
3144 };
3145 if checksum(bytes) != span.hash {
3146 return Err(invalid(&format!(
3147 "column page checksum differs, column {column} part {} at {}+{} of {} bytes, \
3148 wanted {:016x} and got {:016x}",
3149 place.part,
3150 page.offset,
3151 span.start,
3152 span.length,
3153 span.hash,
3154 checksum(bytes),
3155 )));
3156 }
3157 let dictionary = self.dictionary(column)?;
3158 picked.push(decode(&field.ty, rows, bytes, dictionary)?.into_pages());
3164 }
3165 Chunk::with_rows(picked, rows)
3166 }
3167
3168 #[must_use]
3184 pub fn skips(&self, part: usize, probes: &[Probe]) -> bool {
3185 let Some(place) = self.places.get(part).copied() else { return false };
3186 let Some(stripe) = self.table.stripes.get(place.stripe as usize) else { return false };
3187 if stripe.zone.skips(probes) {
3188 return true;
3189 }
3190 probes.iter().any(|probe| self.outside(place, probe) || self.sifted(place, probe))
3191 }
3192
3193 fn outside(&self, place: Place, probe: &Probe) -> bool {
3199 match self.stripe_part_ranges(place.stripe as usize, probe.column) {
3200 Some(ranges) => ranges
3201 .get(place.part as usize)
3202 .is_some_and(|range| range.excludes(probe.op, &probe.value)),
3203 None => false,
3204 }
3205 }
3206
3207 fn stripe_part_ranges(&self, stripe: usize, column: usize) -> Option<&[Range]> {
3213 let slot = self.part_ranges.get(column)?.get(stripe)?;
3214 if let Some(held) = slot.get() {
3215 return Some(held);
3216 }
3217 let page = self.table.stripes.get(stripe)?.part_ranges.get(column).copied().flatten()?;
3218 let mut bytes = vec![0; page.length as usize];
3219 read_at(&self.file, page.offset, &mut bytes).ok()?;
3220 if checksum(&bytes) != page.hash {
3221 return None;
3222 }
3223 let ranges = Arc::new(decode_part_ranges(&bytes).ok()?);
3224 let _ = slot.set(ranges);
3225 slot.get().map(|held| held.as_slice())
3226 }
3227
3228 #[must_use]
3245 pub fn certain(&self, part: usize, probes: &[Probe]) -> bool {
3246 let Some(place) = self.places.get(part).copied() else { return false };
3247 let Some(stripe) = self.table.stripes.get(place.stripe as usize) else { return false };
3248 if stripe.zone.certain(probes) {
3249 return true;
3250 }
3251 probes
3252 .iter()
3253 .all(|probe| stripe.zone.certain(slice::from_ref(probe)) || self.inside(place, probe))
3254 }
3255
3256 fn inside(&self, place: Place, probe: &Probe) -> bool {
3262 match self.stripe_part_ranges(place.stripe as usize, probe.column) {
3263 Some(ranges) => ranges
3264 .get(place.part as usize)
3265 .is_some_and(|range| range.certain(probe.op, &probe.value)),
3266 None => false,
3267 }
3268 }
3269
3270 #[must_use]
3281 pub fn stripe_skips(&self, stripe: usize, probes: &[Probe]) -> bool {
3282 self.table.stripes.get(stripe).is_some_and(|held| held.zone.skips(probes))
3283 }
3284
3285 fn sifted(&self, place: Place, probe: &Probe) -> bool {
3291 if probe.op != Op::Equal {
3292 return false;
3293 }
3294 match self.stripe_sieves(place.stripe as usize, probe.column) {
3295 Some(sieves) => sieves
3296 .get(place.part as usize)
3297 .and_then(Option::as_ref)
3298 .is_some_and(|sieve| sieve.excludes(&probe.value)),
3299 None => false,
3300 }
3301 }
3302
3303 fn stripe_sieves(&self, stripe: usize, column: usize) -> Option<&[Option<Sieve>]> {
3310 let slot = self.sieves.get(column)?.get(stripe)?;
3311 if let Some(held) = slot.get() {
3312 return Some(held);
3313 }
3314 let page = self.table.stripes.get(stripe)?.sieves.get(column).copied().flatten()?;
3315 let mut bytes = vec![0; page.length as usize];
3316 read_at(&self.file, page.offset, &mut bytes).ok()?;
3317 if checksum(&bytes) != page.hash {
3318 return None;
3319 }
3320 let sieves = Arc::new(decode_sieves(&bytes).ok()?);
3321 let _ = slot.set(sieves);
3322 slot.get().map(|held| held.as_slice())
3323 }
3324}
3325
3326fn text_at_rank(dictionary: &Vector, rank: usize) -> Result<Value> {
3328 let code = dictionary.code_at_rank(rank)? as usize;
3329 let text = dictionary
3330 .try_text_at(code)?
3331 .ok_or_else(|| invalid("global dictionary order names a code it does not have"))?;
3332 Ok(Value::Varchar(text.into()))
3333}
3334
3335#[cfg(unix)]
3340fn write_at(file: &File, mut offset: u64, mut bytes: &[u8]) -> Result<()> {
3341 use std::os::unix::fs::FileExt;
3342 while !bytes.is_empty() {
3343 let written = file.write_at(bytes, offset).map_err(io)?;
3344 if written == 0 {
3345 return Err(invalid("a write to the native file wrote nothing"));
3346 }
3347 offset += written as u64;
3348 bytes = &bytes[written..];
3349 }
3350 Ok(())
3351}
3352
3353#[cfg(windows)]
3355fn write_at(file: &File, mut offset: u64, mut bytes: &[u8]) -> Result<()> {
3356 use std::os::windows::fs::FileExt;
3357 while !bytes.is_empty() {
3358 let written = file.seek_write(bytes, offset).map_err(io)?;
3359 if written == 0 {
3360 return Err(invalid("a write to the native file wrote nothing"));
3361 }
3362 offset += written as u64;
3363 bytes = &bytes[written..];
3364 }
3365 Ok(())
3366}
3367
3368#[cfg(not(any(unix, windows)))]
3370fn write_at(file: &File, offset: u64, bytes: &[u8]) -> Result<()> {
3371 use std::io::Write;
3372 let mut file = file.try_clone().map_err(io)?;
3373 file.seek(SeekFrom::Start(offset)).map_err(io)?;
3374 file.write_all(bytes).map_err(io)
3375}
3376
3377#[cfg(unix)]
3387fn read_at(file: &File, mut offset: u64, mut bytes: &mut [u8]) -> Result<()> {
3388 use std::os::unix::fs::FileExt;
3389 while !bytes.is_empty() {
3390 let read = file.read_at(bytes, offset).map_err(io)?;
3391 if read == 0 {
3392 return Err(invalid("column page ends before its declared length"));
3393 }
3394 offset += read as u64;
3395 bytes = &mut bytes[read..];
3396 }
3397 Ok(())
3398}
3399
3400#[cfg(windows)]
3406fn read_at(file: &File, mut offset: u64, mut bytes: &mut [u8]) -> Result<()> {
3407 use std::os::windows::fs::FileExt;
3408 while !bytes.is_empty() {
3409 let read = file.seek_read(bytes, offset).map_err(io)?;
3410 if read == 0 {
3411 return Err(invalid("column page ends before its declared length"));
3412 }
3413 offset += read as u64;
3414 bytes = &mut bytes[read..];
3415 }
3416 Ok(())
3417}
3418
3419#[cfg(not(any(unix, windows)))]
3424fn read_at(file: &File, offset: u64, bytes: &mut [u8]) -> Result<()> {
3425 let mut file = file.try_clone().map_err(io)?;
3426 file.seek(SeekFrom::Start(offset)).map_err(io)?;
3427 file.read_exact(bytes).map_err(io)
3428}
3429
3430fn type_tag(ty: &LogicalType) -> Result<u8> {
3431 match ty {
3432 LogicalType::SmallInt => Ok(1),
3433 LogicalType::Integer => Ok(2),
3434 LogicalType::BigInt => Ok(3),
3435 LogicalType::Varchar => Ok(4),
3436 LogicalType::Date => Ok(5),
3437 LogicalType::Timestamp => Ok(6),
3438 LogicalType::Boolean => Ok(7),
3439 LogicalType::TinyInt => Ok(8),
3440 LogicalType::UTinyInt => Ok(9),
3441 LogicalType::USmallInt => Ok(10),
3442 LogicalType::UInteger => Ok(11),
3443 LogicalType::UBigInt => Ok(12),
3444 LogicalType::Decimal { .. } => Ok(13),
3445 _ => Err(Error::not_implemented(format!("native storage for {ty}"))),
3446 }
3447}
3448
3449fn put_type(out: &mut Vec<u8>, ty: &LogicalType) -> Result<()> {
3455 out.push(type_tag(ty)?);
3456 if let LogicalType::Decimal { width, scale } = ty {
3457 out.push(*width);
3458 out.push(*scale);
3459 }
3460 Ok(())
3461}
3462
3463fn read_type(cur: &mut Cursor<'_>) -> Result<LogicalType> {
3465 let tag = cur.u8()?;
3466 if tag == 13 {
3467 let width = cur.u8()?;
3468 let scale = cur.u8()?;
3469 return LogicalType::decimal(width, scale)
3470 .map_err(|_| invalid("decimal column width and scale are not a decimal"));
3471 }
3472 tag_type(tag)
3473}
3474
3475fn tag_type(tag: u8) -> Result<LogicalType> {
3476 match tag {
3477 1 => Ok(LogicalType::SmallInt),
3478 2 => Ok(LogicalType::Integer),
3479 3 => Ok(LogicalType::BigInt),
3480 4 => Ok(LogicalType::Varchar),
3481 5 => Ok(LogicalType::Date),
3482 6 => Ok(LogicalType::Timestamp),
3483 7 => Ok(LogicalType::Boolean),
3484 8 => Ok(LogicalType::TinyInt),
3485 9 => Ok(LogicalType::UTinyInt),
3486 10 => Ok(LogicalType::USmallInt),
3487 11 => Ok(LogicalType::UInteger),
3488 12 => Ok(LogicalType::UBigInt),
3489 _ => Err(invalid("column type tag is unknown")),
3490 }
3491}
3492
3493fn put_u16(out: &mut Vec<u8>, value: u16) {
3494 out.extend_from_slice(&value.to_le_bytes());
3495}
3496fn put_u32(out: &mut Vec<u8>, value: u32) {
3497 out.extend_from_slice(&value.to_le_bytes());
3498}
3499fn put_u64(out: &mut Vec<u8>, value: u64) {
3500 out.extend_from_slice(&value.to_le_bytes());
3501}
3502fn put_var_u64(out: &mut Vec<u8>, mut value: u64) {
3503 while value >= 0x80 {
3504 out.push((value as u8 & 0x7f) | 0x80);
3505 value >>= 7;
3506 }
3507 out.push(value as u8);
3508}
3509
3510fn frequency_order(left: FrequencyValue, right: FrequencyValue) -> Ordering {
3511 match (left, right) {
3512 (FrequencyValue::Null, FrequencyValue::Null) => Ordering::Equal,
3513 (FrequencyValue::Null, _) => Ordering::Less,
3514 (_, FrequencyValue::Null) => Ordering::Greater,
3515 (FrequencyValue::Integer(left), FrequencyValue::Integer(right)) => left.cmp(&right),
3516 (FrequencyValue::Code(left), FrequencyValue::Code(right)) => left.cmp(&right),
3517 (FrequencyValue::Integer(_), FrequencyValue::Code(_)) => Ordering::Less,
3518 (FrequencyValue::Code(_), FrequencyValue::Integer(_)) => Ordering::Greater,
3519 }
3520}
3521
3522fn code_frequency(dictionary: &GlobalDictionary) -> FrequencySummary {
3523 let mut entries = dictionary
3524 .counts
3525 .iter()
3526 .enumerate()
3527 .filter(|(_, count)| **count != 0)
3528 .map(|(code, &count)| FrequencyEntry { value: FrequencyValue::Code(code as u32), count })
3529 .collect::<Vec<_>>();
3530 if dictionary.nulls != 0 {
3531 entries.push(FrequencyEntry { value: FrequencyValue::Null, count: dictionary.nulls });
3532 }
3533 entries.sort_unstable_by(|left, right| {
3534 right.count.cmp(&left.count).then_with(|| frequency_order(left.value, right.value))
3535 });
3536 let omitted_max = entries.get(FREQUENCY_ENTRIES).map_or(0, |entry| entry.count);
3537 entries.truncate(FREQUENCY_ENTRIES);
3538 FrequencySummary { entries, omitted_max, ordinals: Vec::new() }
3539}
3540
3541fn encode_directory(table: &Table) -> Result<Vec<u8>> {
3542 let mut out = DIRECTORY.to_vec();
3543 let name = table.name.as_bytes();
3544 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("table name too long"))?);
3545 out.extend_from_slice(name);
3546 put_u16(&mut out, u16::try_from(table.fields.len()).map_err(|_| invalid("too many columns"))?);
3547 for field in &table.fields {
3548 let name = field.name.as_bytes();
3549 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("column name too long"))?);
3550 out.extend_from_slice(name);
3551 put_type(&mut out, &field.ty)?;
3552 out.push(u8::from(field.not_null));
3553 }
3554 for dictionary in &table.dictionaries {
3555 match dictionary {
3556 None => out.push(0),
3557 Some(page) => {
3558 out.push(1);
3559 put_u64(&mut out, page.offset);
3560 put_u32(&mut out, page.length);
3561 put_u64(&mut out, page.hash);
3562 }
3563 }
3564 }
3565 for distinct in &table.distincts {
3566 match distinct {
3567 None => out.push(0),
3568 Some(count) => {
3569 out.push(1);
3570 put_u64(&mut out, *count);
3571 }
3572 }
3573 }
3574 put_u64(&mut out, u64::try_from(table.rows).map_err(|_| invalid("row count overflow"))?);
3575 put_u32(&mut out, u32::try_from(table.stripes.len()).map_err(|_| invalid("too many stripes"))?);
3576 for stripe in &table.stripes {
3577 put_u32(
3578 &mut out,
3579 u32::try_from(stripe.parts.len()).map_err(|_| invalid("too many parts in a stripe"))?,
3580 );
3581 for &rows in &stripe.parts {
3582 put_u32(&mut out, rows);
3583 }
3584 put_u64(&mut out, stripe.index.offset);
3585 put_u32(&mut out, stripe.index.length);
3586 for page in &stripe.pages {
3587 put_u64(&mut out, page.offset);
3588 put_u32(&mut out, page.length);
3589 }
3590 for (field, membership) in table.fields.iter().zip(&stripe.memberships) {
3591 if field.ty != LogicalType::Varchar {
3592 continue;
3593 }
3594 let page =
3595 membership.ok_or_else(|| invalid("string page has no code membership index"))?;
3596 put_u64(&mut out, page.offset);
3597 put_u32(&mut out, page.length);
3598 put_u64(&mut out, page.hash);
3599 }
3600 for sieve in &stripe.sieves {
3601 match sieve {
3602 None => out.push(0),
3603 Some(page) => {
3604 out.push(1);
3605 put_u64(&mut out, page.offset);
3606 put_u32(&mut out, page.length);
3607 put_u64(&mut out, page.hash);
3608 }
3609 }
3610 }
3611 for held in &stripe.part_ranges {
3612 match held {
3613 None => out.push(0),
3614 Some(page) => {
3615 out.push(1);
3616 put_u64(&mut out, page.offset);
3617 put_u32(&mut out, page.length);
3618 put_u64(&mut out, page.hash);
3619 }
3620 }
3621 }
3622 for range in stripe.zone.columns() {
3623 put_bound(&mut out, range.low.as_ref())?;
3624 put_bound(&mut out, range.high.as_ref())?;
3625 put_u32(
3626 &mut out,
3627 u32::try_from(range.nulls).map_err(|_| invalid("null count overflow"))?,
3628 );
3629 out.push(u8::from(range.exact));
3630 match range.sum {
3631 None => out.push(0),
3632 Some(total) => {
3633 out.push(1);
3634 out.extend_from_slice(&total.to_le_bytes());
3635 }
3636 }
3637 }
3638 }
3639 out.extend_from_slice(FREQUENCIES);
3640 put_u16(
3641 &mut out,
3642 u16::try_from(table.frequencies.len())
3643 .map_err(|_| invalid("too many frequency columns"))?,
3644 );
3645 for summary in &table.frequencies {
3646 let Some(summary) = summary else {
3647 out.push(0);
3648 continue;
3649 };
3650 out.push(1);
3651 put_u64(&mut out, summary.omitted_max);
3652 put_u32(
3653 &mut out,
3654 u32::try_from(summary.entries.len())
3655 .map_err(|_| invalid("too many frequency entries"))?,
3656 );
3657 for entry in &summary.entries {
3658 match entry.value {
3659 FrequencyValue::Null => out.push(0),
3660 FrequencyValue::Integer(value) => {
3661 out.push(1);
3662 out.extend_from_slice(&value.to_le_bytes());
3663 }
3664 FrequencyValue::Code(value) => {
3665 out.push(2);
3666 put_u32(&mut out, value);
3667 }
3668 }
3669 put_u64(&mut out, entry.count);
3670 }
3671 put_u32(
3672 &mut out,
3673 u32::try_from(summary.ordinals.len())
3674 .map_err(|_| invalid("too many frequency ordinals"))?,
3675 );
3676 let mut previous = 0_u64;
3677 for (at, &ordinal) in summary.ordinals.iter().enumerate() {
3678 let delta = if at == 0 {
3679 ordinal
3680 } else {
3681 ordinal
3682 .checked_sub(previous)
3683 .ok_or_else(|| invalid("frequency ordinals are not ordered"))?
3684 };
3685 if at != 0 && delta == 0 {
3686 return Err(invalid("frequency ordinals are not unique"));
3687 }
3688 put_var_u64(&mut out, delta);
3689 previous = ordinal;
3690 }
3691 }
3692 Ok(out)
3693}
3694
3695fn encode_catalog(entries: &[Entry]) -> Result<Vec<u8>> {
3701 let mut out = CATALOG.to_vec();
3702 put_u32(&mut out, u32::try_from(entries.len()).map_err(|_| invalid("too many tables"))?);
3703 for entry in entries {
3704 let name = entry.name.as_bytes();
3705 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("table name too long"))?);
3706 out.extend_from_slice(name);
3707 put_u64(&mut out, u64::try_from(entry.rows).map_err(|_| invalid("row count overflow"))?);
3708 put_u16(
3709 &mut out,
3710 u16::try_from(entry.fields.len()).map_err(|_| invalid("too many columns"))?,
3711 );
3712 for field in &entry.fields {
3713 let name = field.name.as_bytes();
3714 put_u16(
3715 &mut out,
3716 u16::try_from(name.len()).map_err(|_| invalid("column name too long"))?,
3717 );
3718 out.extend_from_slice(name);
3719 put_type(&mut out, &field.ty)?;
3720 out.push(u8::from(field.not_null));
3721 }
3722 put_u64(&mut out, entry.directory.offset);
3723 put_u32(&mut out, entry.directory.length);
3724 put_u64(&mut out, entry.directory.hash);
3725 }
3726 Ok(out)
3727}
3728
3729fn decode_catalog(bytes: &[u8], size: u64) -> Result<Vec<Entry>> {
3732 let mut cur = Cursor { bytes, at: 0 };
3733 if cur.take(8)? != CATALOG {
3734 return Err(invalid("catalog magic differs"));
3735 }
3736 let count = cur.u32()? as usize;
3737 let mut entries: Vec<Entry> = Vec::with_capacity(count.min(1024));
3738 for _ in 0..count {
3739 let name = cur.text()?;
3740 let rows = usize::try_from(cur.u64()?).map_err(|_| invalid("row count does not fit"))?;
3741 let width = cur.u16()? as usize;
3742 let mut fields = Vec::with_capacity(width);
3743 for _ in 0..width {
3744 let name = cur.text()?;
3745 let ty = read_type(&mut cur)?;
3746 let not_null = match cur.u8()? {
3747 0 => false,
3748 1 => true,
3749 _ => return Err(invalid("nullability flag differs")),
3750 };
3751 fields.push(Field { name, ty, not_null });
3752 }
3753 let directory = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3754 let end = directory
3755 .offset
3756 .checked_add(u64::from(directory.length))
3757 .ok_or_else(|| invalid("table directory offset overflow"))?;
3758 if directory.offset < HEADER
3759 || end > size
3760 || directory.length as usize > MAX_DIRECTORY
3761 || directory.length == 0
3762 {
3763 return Err(invalid("table directory range is outside the file"));
3764 }
3765 if entries.iter().any(|held| held.name == name) {
3766 return Err(invalid("two tables in the catalog have the same name"));
3767 }
3768 entries.push(Entry { name, fields, rows, directory });
3769 }
3770 Ok(entries)
3771}
3772
3773struct Cursor<'a> {
3774 bytes: &'a [u8],
3775 at: usize,
3776}
3777impl<'a> Cursor<'a> {
3778 fn take(&mut self, len: usize) -> Result<&'a [u8]> {
3779 let end = self.at.checked_add(len).ok_or_else(|| invalid("directory offset overflow"))?;
3780 let bytes =
3781 self.bytes.get(self.at..end).ok_or_else(|| invalid("directory is truncated"))?;
3782 self.at = end;
3783 Ok(bytes)
3784 }
3785 fn u8(&mut self) -> Result<u8> {
3786 Ok(self.take(1)?[0])
3787 }
3788 fn u16(&mut self) -> Result<u16> {
3789 Ok(u16::from_le_bytes(self.take(2)?.try_into().expect("two bytes")))
3790 }
3791 fn u32(&mut self) -> Result<u32> {
3792 Ok(u32::from_le_bytes(self.take(4)?.try_into().expect("four bytes")))
3793 }
3794 fn u64(&mut self) -> Result<u64> {
3795 Ok(u64::from_le_bytes(self.take(8)?.try_into().expect("eight bytes")))
3796 }
3797 fn var_u64(&mut self) -> Result<u64> {
3798 let mut value = 0_u64;
3799 for shift in (0..=63).step_by(7) {
3800 let byte = self.u8()?;
3801 let part = u64::from(byte & 0x7f);
3802 if shift == 63 && part > 1 {
3803 return Err(invalid("frequency ordinal varint overflows"));
3804 }
3805 value |= part << shift;
3806 if byte & 0x80 == 0 {
3807 return Ok(value);
3808 }
3809 }
3810 Err(invalid("frequency ordinal varint is too long"))
3811 }
3812 fn bound(&mut self) -> Result<Option<Bound>> {
3813 Ok(match self.u8()? {
3814 0 => None,
3815 1 => Some(Bound::Int(i128::from_le_bytes(
3816 self.take(16)?.try_into().expect("sixteen bytes"),
3817 ))),
3818 2 => Some(Bound::Real(f64::from_le_bytes(
3819 self.take(8)?.try_into().expect("eight bytes"),
3820 ))),
3821 3 => {
3822 let length = self.u32()? as usize;
3823 Some(Bound::Bytes(self.take(length)?.to_vec()))
3824 }
3825 4 => {
3826 let unscaled =
3827 i128::from_le_bytes(self.take(16)?.try_into().expect("sixteen bytes"));
3828 Some(Bound::Scaled { unscaled, scale: self.u8()? })
3829 }
3830 _ => return Err(invalid("bound tag differs")),
3831 })
3832 }
3833 fn text(&mut self) -> Result<String> {
3834 let len = self.u16()? as usize;
3835 String::from_utf8(self.take(len)?.to_vec()).map_err(|_| invalid("name is not UTF-8"))
3836 }
3837}
3838
3839fn decode_directory(bytes: &[u8], size: u64) -> Result<Table> {
3840 let mut cur = Cursor { bytes, at: 0 };
3841 if cur.take(8)? != DIRECTORY {
3842 return Err(invalid("directory magic differs"));
3843 }
3844 let name = cur.text()?;
3845 let width = cur.u16()? as usize;
3846 let mut fields = Vec::with_capacity(width);
3847 for _ in 0..width {
3848 let name = cur.text()?;
3849 let ty = read_type(&mut cur)?;
3850 let not_null = match cur.u8()? {
3851 0 => false,
3852 1 => true,
3853 _ => return Err(invalid("nullability flag differs")),
3854 };
3855 fields.push(Field { name, ty, not_null });
3856 }
3857 let mut dictionaries = Vec::with_capacity(width);
3858 for _ in 0..width {
3859 dictionaries.push(match cur.u8()? {
3860 0 => None,
3861 1 => {
3862 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3863 let end = page
3864 .offset
3865 .checked_add(u64::from(page.length))
3866 .ok_or_else(|| invalid("dictionary page offset overflow"))?;
3867 if page.offset < HEADER || end > size {
3872 return Err(invalid("dictionary page range is outside the file"));
3873 }
3874 Some(page)
3875 }
3876 _ => return Err(invalid("dictionary page tag differs")),
3877 });
3878 }
3879 let mut distincts = Vec::with_capacity(width);
3880 for _ in 0..width {
3881 distincts.push(match cur.u8()? {
3882 0 => None,
3883 1 => Some(cur.u64()?),
3884 _ => return Err(invalid("distinct count tag differs")),
3885 });
3886 }
3887 let rows = usize::try_from(cur.u64()?).map_err(|_| invalid("row count does not fit"))?;
3888 let count = cur.u32()? as usize;
3889 let mut stripes = Vec::with_capacity(count);
3890 let mut total = 0_usize;
3891 for _ in 0..count {
3892 let count = cur.u32()? as usize;
3893 if count == 0 || count > STRIPE_PARTS {
3894 return Err(invalid("stripe part count is outside its bound"));
3895 }
3896 let mut parts = Vec::with_capacity(count);
3897 let mut stripe_rows = 0_usize;
3898 for _ in 0..count {
3899 let rows = cur.u32()?;
3900 if rows == 0 {
3901 return Err(invalid("empty part"));
3902 }
3903 parts.push(rows);
3904 stripe_rows = stripe_rows
3905 .checked_add(rows as usize)
3906 .ok_or_else(|| invalid("stripe row count overflow"))?;
3907 }
3908 total =
3909 total.checked_add(stripe_rows).ok_or_else(|| invalid("stripe row count overflow"))?;
3910 let index = Span { offset: cur.u64()?, length: cur.u32()? };
3911 let section = index_section(count)?;
3912 let wanted = section
3913 .checked_mul(width)
3914 .and_then(|bytes| u32::try_from(bytes).ok())
3915 .ok_or_else(|| invalid("index page length overflow"))?;
3916 let end = index
3917 .offset
3918 .checked_add(u64::from(index.length))
3919 .ok_or_else(|| invalid("index page offset overflow"))?;
3920 if index.offset < HEADER || end > size || index.length != wanted {
3921 return Err(invalid("index page range is outside the file"));
3922 }
3923 let mut pages = Vec::with_capacity(width);
3924 for _ in 0..width {
3925 let offset = cur.u64()?;
3926 let length = cur.u32()?;
3927 let end = offset
3928 .checked_add(u64::from(length))
3929 .ok_or_else(|| invalid("page offset overflow"))?;
3930 if offset < HEADER || end > size || length as usize > MAX_PAGE {
3931 return Err(invalid("page range is outside the file"));
3932 }
3933 pages.push(Span { offset, length });
3934 }
3935 let mut memberships = vec![None; width];
3936 for (column, field) in fields.iter().enumerate() {
3937 if field.ty != LogicalType::Varchar {
3938 continue;
3939 }
3940 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3941 let end = page
3942 .offset
3943 .checked_add(u64::from(page.length))
3944 .ok_or_else(|| invalid("membership page offset overflow"))?;
3945 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
3946 return Err(invalid("membership page range is outside the file"));
3947 }
3948 memberships[column] = Some(page);
3949 }
3950 let mut sieves = vec![None; width];
3951 for sieve in sieves.iter_mut().take(width) {
3952 match cur.u8()? {
3953 0 => continue,
3954 1 => {}
3955 _ => return Err(invalid("a sieve page has an unknown tag")),
3956 }
3957 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3958 let end = page
3959 .offset
3960 .checked_add(u64::from(page.length))
3961 .ok_or_else(|| invalid("sieve page offset overflow"))?;
3962 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
3963 return Err(invalid("sieve page range is outside the file"));
3964 }
3965 *sieve = Some(page);
3966 }
3967 let mut part_ranges = vec![None; width];
3968 for held in part_ranges.iter_mut().take(width) {
3969 match cur.u8()? {
3970 0 => continue,
3971 1 => {}
3972 _ => return Err(invalid("a part range page has an unknown tag")),
3973 }
3974 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3975 let end = page
3976 .offset
3977 .checked_add(u64::from(page.length))
3978 .ok_or_else(|| invalid("part range page offset overflow"))?;
3979 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
3980 return Err(invalid("part range page range is outside the file"));
3981 }
3982 *held = Some(page);
3983 }
3984 let mut ranges = Vec::with_capacity(width);
3985 for column in 0..width {
3986 let low = cur.bound()?;
3987 let high = cur.bound()?;
3988 let nulls = cur.u32()? as usize;
3989 if nulls > stripe_rows {
3990 return Err(invalid("null count exceeds stripe rows"));
3991 }
3992 let exact = cur.u8()? != 0;
3993 let sum = match cur.u8()? {
3994 0 => None,
3995 1 => Some(i128::from_le_bytes(
3996 cur.take(16)?.try_into().map_err(|_| invalid("a stripe sum is truncated"))?,
3997 )),
3998 _ => return Err(invalid("a stripe sum has an unknown tag")),
3999 };
4000 let ty = &fields.get(column).ok_or_else(|| invalid("a stripe range has no column"))?.ty;
4006 let low = low.map(|bound| scaled_as(bound, ty));
4007 let high = high.map(|bound| scaled_as(bound, ty));
4008 ranges.push(Range { low, high, nulls, exact, sum });
4009 }
4010 stripes.push(Stripe {
4011 rows: stripe_rows,
4012 parts,
4013 index,
4014 pages,
4015 memberships,
4016 sieves,
4017 part_ranges,
4018 zone: Zone::from_ranges(ranges),
4019 });
4020 }
4021 if total != rows {
4022 return Err(invalid("table row count differs from stripes"));
4023 }
4024 let frequencies = if cur.at == bytes.len() {
4025 vec![None; width]
4026 } else {
4027 if cur.take(8)? != FREQUENCIES {
4028 return Err(invalid("directory extension magic differs"));
4029 }
4030 if cur.u16()? as usize != width {
4031 return Err(invalid("frequency column count differs"));
4032 }
4033 let mut frequencies = Vec::with_capacity(width);
4034 for field in &fields {
4035 let summary = match cur.u8()? {
4036 0 => None,
4037 1 => {
4038 let omitted_max = cur.u64()?;
4039 let count = cur.u32()? as usize;
4040 if count > FREQUENCY_ENTRIES {
4041 return Err(invalid("frequency entry count exceeds its bound"));
4042 }
4043 let mut entries = Vec::with_capacity(count);
4044 for _ in 0..count {
4046 let value = match cur.u8()? {
4047 0 => FrequencyValue::Null,
4048 1 => FrequencyValue::Integer(i128::from_le_bytes(
4049 cur.take(16)?.try_into().expect("sixteen bytes"),
4050 )),
4051 2 => FrequencyValue::Code(cur.u32()?),
4052 _ => return Err(invalid("frequency value tag differs")),
4053 };
4054 let valid = matches!(
4055 (&field.ty, value),
4056 (_, FrequencyValue::Null)
4057 | (LogicalType::Varchar, FrequencyValue::Code(_))
4058 | (
4059 LogicalType::TinyInt
4060 | LogicalType::SmallInt
4061 | LogicalType::Integer
4062 | LogicalType::BigInt
4063 | LogicalType::UTinyInt
4064 | LogicalType::USmallInt
4065 | LogicalType::UInteger
4066 | LogicalType::UBigInt
4067 | LogicalType::Date
4068 | LogicalType::Timestamp,
4069 FrequencyValue::Integer(_),
4070 )
4071 );
4072 if !valid {
4073 return Err(invalid("frequency value does not match its column"));
4074 }
4075 let count = cur.u64()?;
4076 if count == 0 || count > rows as u64 {
4077 return Err(invalid("frequency count is outside the table"));
4078 }
4079 entries.push(FrequencyEntry { value, count });
4080 }
4081 if entries.windows(2).any(|pair| pair[0].count < pair[1].count) {
4082 return Err(invalid("frequency entries are not descending"));
4083 }
4084 let ordinals = {
4085 let ordinal_count = cur.u32()? as usize;
4086 if ordinal_count > FREQUENCY_ORDINALS || ordinal_count > rows {
4087 return Err(invalid("frequency ordinal count exceeds its bound"));
4088 }
4089 let mut ordinals = Vec::with_capacity(ordinal_count);
4090 let mut previous = 0_u64;
4091 for at in 0..ordinal_count {
4092 let delta = cur.var_u64()?;
4093 if at != 0 && delta == 0 {
4094 return Err(invalid("frequency ordinals are not increasing"));
4095 }
4096 let ordinal = if at == 0 {
4097 delta
4098 } else {
4099 previous
4100 .checked_add(delta)
4101 .ok_or_else(|| invalid("frequency ordinal overflows"))?
4102 };
4103 if ordinal >= rows as u64 {
4104 return Err(invalid("frequency ordinal is outside the table"));
4105 }
4106 ordinals.push(ordinal);
4107 previous = ordinal;
4108 }
4109 ordinals
4110 };
4111 Some(FrequencySummary { entries, omitted_max, ordinals })
4112 }
4113 _ => return Err(invalid("frequency summary tag differs")),
4114 };
4115 frequencies.push(summary);
4116 }
4117 frequencies
4118 };
4119 if cur.at != bytes.len() {
4120 return Err(invalid("directory has trailing bytes"));
4121 }
4122 Ok(Table { name, fields, stripes, rows, dictionaries, distincts, frequencies })
4123}
4124
4125fn put_bound(out: &mut Vec<u8>, bound: Option<&Bound>) -> Result<()> {
4126 match bound {
4127 None => out.push(0),
4128 Some(Bound::Int(value)) => {
4129 out.push(1);
4130 out.extend_from_slice(&value.to_le_bytes());
4131 }
4132 Some(Bound::Real(value)) => {
4133 out.push(2);
4134 out.extend_from_slice(&value.to_le_bytes());
4135 }
4136 Some(Bound::Bytes(value)) => {
4137 out.push(3);
4138 put_u32(out, u32::try_from(value.len()).map_err(|_| invalid("bound length overflow"))?);
4139 out.extend_from_slice(value);
4140 }
4141 Some(Bound::Scaled { unscaled, scale }) => {
4142 out.push(4);
4143 out.extend_from_slice(&unscaled.to_le_bytes());
4144 out.push(*scale);
4145 }
4146 }
4147 Ok(())
4148}
4149
4150#[derive(Debug)]
4167struct Codes;
4168
4169impl chooser::Chooser for Codes {
4170 fn name(&self) -> &'static str {
4171 "codes"
4172 }
4173
4174 fn narrow_strings(
4175 &self,
4176 _values: &[&[u8]],
4177 offered: &[string::Kind],
4178 _depth: u8,
4179 ) -> Vec<string::Kind> {
4180 offered.to_vec()
4183 }
4184
4185 fn narrow_integers(
4186 &self,
4187 _values: &[i64],
4188 offered: &[integer::Kind],
4189 depth: u8,
4190 ) -> Vec<integer::Kind> {
4191 let keep: &[integer::Kind] = if depth == 0 {
4192 &[integer::Kind::Constant, integer::Kind::Packed, integer::Kind::Rle]
4193 } else {
4194 &[integer::Kind::Constant, integer::Kind::Packed]
4195 };
4196 let narrowed: Vec<integer::Kind> =
4197 offered.iter().copied().filter(|kind| keep.contains(kind)).collect();
4198 if narrowed.is_empty() { offered.to_vec() } else { narrowed }
4201 }
4202}
4203
4204#[derive(Debug)]
4216struct Fixed;
4217
4218impl chooser::Chooser for Fixed {
4219 fn name(&self) -> &'static str {
4220 "fixed"
4221 }
4222
4223 fn narrow_strings(
4224 &self,
4225 _values: &[&[u8]],
4226 offered: &[string::Kind],
4227 _depth: u8,
4228 ) -> Vec<string::Kind> {
4229 offered.to_vec()
4230 }
4231
4232 fn narrow_integers(
4233 &self,
4234 _values: &[i64],
4235 offered: &[integer::Kind],
4236 depth: u8,
4237 ) -> Vec<integer::Kind> {
4238 let keep: &[integer::Kind] = if depth == 0 {
4239 &[
4240 integer::Kind::Constant,
4241 integer::Kind::Packed,
4242 integer::Kind::Delta,
4243 integer::Kind::Rle,
4244 integer::Kind::Sparse,
4245 integer::Kind::Strided,
4246 ]
4247 } else {
4248 &[integer::Kind::Constant, integer::Kind::Packed, integer::Kind::Delta]
4249 };
4250 let narrowed: Vec<integer::Kind> =
4251 offered.iter().copied().filter(|kind| keep.contains(kind)).collect();
4252 if narrowed.is_empty() { offered.to_vec() } else { narrowed }
4253 }
4254}
4255
4256fn widened(data: &Data) -> Option<Vec<i64>> {
4263 match data {
4264 Data::Int8(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4265 Data::UInt8(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4266 Data::Int16(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4267 Data::UInt16(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4268 Data::Int32(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4269 Data::UInt32(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
4270 Data::Int64(values) => Some(values.to_vec()),
4271 _ => None,
4272 }
4273}
4274
4275trait Narrow: Copy {
4282 const BIASED: (u32, u64);
4287
4288 fn narrow(value: i64) -> Self;
4290}
4291
4292#[allow(clippy::cast_sign_loss, reason = "a residue is a bit pattern and not a number")]
4309fn residue<T: Narrow>(value: i64) -> u64 {
4310 let (bits, bias) = T::BIASED;
4311 (value as u64).wrapping_add(bias) >> bits
4312}
4313
4314macro_rules! narrows {
4319 ($($ty:ty => $bias:expr),* $(,)?) => {$(
4320 impl Narrow for $ty {
4321 const BIASED: (u32, u64) = (<$ty>::BITS, $bias);
4322
4323 #[allow(
4324 clippy::cast_possible_truncation,
4325 clippy::cast_sign_loss,
4326 reason = "the caller has checked the bits this truncates away"
4327 )]
4328 fn narrow(value: i64) -> Self {
4329 value as Self
4330 }
4331 }
4332 )*};
4333}
4334
4335narrows! {
4336 i8 => 1 << 7,
4337 u8 => 0,
4338 i16 => 1 << 15,
4339 u16 => 0,
4340 i32 => 1 << 31,
4341 u32 => 0,
4342}
4343
4344fn fit<T: Narrow>(values: &[i64]) -> Result<Vec<T>> {
4357 let mut spilled = 0u64;
4358 for value in values {
4359 spilled |= residue::<T>(*value);
4360 }
4361 if spilled != 0 {
4362 return Err(invalid("page value is not of its type"));
4363 }
4364 Ok(values.iter().map(|value| T::narrow(*value)).collect())
4365}
4366
4367fn narrowed(ty: &LogicalType, values: Vec<i64>) -> Result<Data> {
4372 Ok(match ty {
4373 LogicalType::TinyInt => Data::Int8(fit::<i8>(&values)?.into()),
4374 LogicalType::UTinyInt => Data::UInt8(fit::<u8>(&values)?.into()),
4375 LogicalType::SmallInt => Data::Int16(fit::<i16>(&values)?.into()),
4376 LogicalType::USmallInt => Data::UInt16(fit::<u16>(&values)?.into()),
4377 LogicalType::Integer | LogicalType::Date => Data::Int32(fit::<i32>(&values)?.into()),
4378 LogicalType::UInteger => Data::UInt32(fit::<u32>(&values)?.into()),
4379 LogicalType::BigInt | LogicalType::Timestamp => Data::Int64(values.into()),
4380 LogicalType::Decimal { .. } => match ty.physical() {
4383 PhysicalType::Int16 => Data::Int16(fit::<i16>(&values)?.into()),
4384 PhysicalType::Int32 => Data::Int32(fit::<i32>(&values)?.into()),
4385 PhysicalType::Int64 => Data::Int64(values.into()),
4386 _ => return Err(invalid("cascade codec belongs to a decimal that is not an integer")),
4387 },
4388 _ => return Err(invalid("cascade codec belongs to a page that is not integers")),
4389 })
4390}
4391
4392fn plain_width(ty: &LogicalType) -> Option<usize> {
4395 Some(match ty {
4396 LogicalType::TinyInt | LogicalType::UTinyInt => 1,
4397 LogicalType::SmallInt | LogicalType::USmallInt => 2,
4398 LogicalType::Integer | LogicalType::UInteger | LogicalType::Date => 4,
4399 LogicalType::BigInt | LogicalType::Timestamp => 8,
4400 LogicalType::Decimal { .. } => match ty.physical() {
4401 PhysicalType::Int16 => 2,
4402 PhysicalType::Int32 => 4,
4403 PhysicalType::Int64 => 8,
4404 _ => return None,
4407 },
4408 _ => return None,
4409 })
4410}
4411
4412fn cascaded(
4418 flat: &Vector,
4419 ty: &LogicalType,
4420 packed: Option<&Packed<'_>>,
4421) -> Result<Option<Vec<u8>>> {
4422 let (Some(width), Some(data)) = (plain_width(ty), flat.data()) else { return Ok(None) };
4423 let Some(values) = widened(data) else { return Ok(None) };
4424 let plain = values.len().saturating_mul(width);
4425 let best = match packed {
4426 Some(packed) => plain.min(21 + size_of_val(packed.words())),
4428 None => plain,
4429 };
4430 let out = integer::encode_with(&values, &Fixed)?;
4431 Ok((out.len() < best).then_some(out))
4432}
4433
4434fn encoded_codes(codes: &[u32]) -> Result<Option<Vec<u8>>> {
4446 let wide: Vec<i64> = codes.iter().map(|code| i64::from(*code)).collect();
4447 let coded = integer::encode_with(&wide, &Codes)?;
4448 let plain = codes.len().saturating_mul(size_of::<u32>());
4449 Ok((coded.len() < plain).then_some(coded))
4450}
4451
4452fn encode(
4453 vector: &Vector,
4454 global: Option<&mut GlobalDictionary>,
4455) -> Result<(Vec<u8>, Option<Vec<u32>>)> {
4456 let ty = vector.logical_type();
4457 let flat = vector.flatten()?;
4459 let mut out = Vec::new();
4460 let mut global_codes = None;
4461 if let Some(global) = global {
4462 let mut codes = Vec::with_capacity(flat.len());
4463 for row in 0..flat.len() {
4464 let text = flat.text_at(row).unwrap_or("");
4465 let code = global.code(text)?;
4466 global.observe(code, flat.is_null_at(row))?;
4467 codes.push(code);
4468 }
4469 global_codes = Some(codes);
4470 }
4471 let membership = global_codes.as_deref().map(unique_codes);
4472 let dictionary = if global_codes.is_none() && ty == &LogicalType::Varchar {
4473 string_dictionary(&flat)?
4474 } else {
4475 None
4476 };
4477 let packed_vector = if dictionary.is_none() && global_codes.is_none() {
4478 Some(flat.bit_packed()?)
4479 } else {
4480 None
4481 };
4482 let packed = packed_vector.as_ref().and_then(Vector::packed_parts);
4483 let coded = match global_codes.as_deref() {
4484 Some(codes) => encoded_codes(codes)?,
4485 None => None,
4486 };
4487 let cascade = if dictionary.is_none() && global_codes.is_none() {
4491 cascaded(&flat, ty, packed.as_ref())?
4492 } else {
4493 None
4494 };
4495 out.push(if coded.is_some() {
4496 4
4497 } else if cascade.is_some() {
4498 5
4499 } else if global_codes.is_some() {
4500 3
4501 } else if dictionary.is_some() {
4502 1
4503 } else if packed.is_some() {
4504 2
4505 } else {
4506 0
4507 });
4508 let nulls = flat.validity();
4509 let flag = match nulls {
4510 Validity::AllValid => 0,
4511 Validity::AllInvalid => 1,
4512 Validity::Mask(_) => 2,
4513 };
4514 out.push(flag);
4515 if flag == 2 {
4516 for group in (0..vector.len()).step_by(8) {
4517 let mut bits = 0_u8;
4518 for bit in 0..8 {
4519 if group + bit < vector.len() && !flat.is_null_at(group + bit) {
4520 bits |= 1 << bit;
4521 }
4522 }
4523 out.push(bits);
4524 }
4525 }
4526 if let Some(coded) = coded {
4527 out.extend_from_slice(&coded);
4528 return Ok((out, membership));
4529 }
4530 if let Some(cascade) = cascade {
4531 out.extend_from_slice(&cascade);
4532 return Ok((out, membership));
4533 }
4534 if let Some(codes) = global_codes {
4535 for code in codes {
4536 put_u32(&mut out, code);
4537 }
4538 return Ok((out, membership));
4539 }
4540 if let Some(dictionary) = dictionary {
4541 out.extend_from_slice(&dictionary);
4542 return Ok((out, membership));
4543 }
4544 if let Some(packed) = packed {
4545 if packed.offset() != 0 {
4546 return Err(invalid("writer received a sliced packed vector"));
4547 }
4548 out.push(u8::try_from(packed.width()).map_err(|_| invalid("packed width overflow"))?);
4549 out.extend_from_slice(&packed.base().to_le_bytes());
4550 put_u32(
4551 &mut out,
4552 u32::try_from(packed.words().len()).map_err(|_| invalid("too many packed words"))?,
4553 );
4554 for word in packed.words() {
4555 put_u64(&mut out, *word);
4556 }
4557 return Ok((out, membership));
4558 }
4559 let data = flat.data().ok_or_else(|| invalid("scalar column did not flatten"))?;
4560 match (ty, data) {
4561 (LogicalType::TinyInt, Data::Int8(values)) => {
4562 for value in &**values {
4563 out.extend_from_slice(&value.to_le_bytes());
4564 }
4565 }
4566 (LogicalType::UTinyInt, Data::UInt8(values)) => {
4567 for value in &**values {
4568 out.extend_from_slice(&value.to_le_bytes());
4569 }
4570 }
4571 (LogicalType::SmallInt, Data::Int16(values)) => {
4572 for value in &**values {
4573 out.extend_from_slice(&value.to_le_bytes());
4574 }
4575 }
4576 (LogicalType::USmallInt, Data::UInt16(values)) => {
4577 for value in &**values {
4578 out.extend_from_slice(&value.to_le_bytes());
4579 }
4580 }
4581 (LogicalType::UInteger, Data::UInt32(values)) => {
4582 for value in &**values {
4583 out.extend_from_slice(&value.to_le_bytes());
4584 }
4585 }
4586 (LogicalType::UBigInt, Data::UInt64(values)) => {
4587 for value in &**values {
4588 out.extend_from_slice(&value.to_le_bytes());
4589 }
4590 }
4591 (LogicalType::Integer | LogicalType::Date, Data::Int32(values)) => {
4592 for value in &**values {
4593 out.extend_from_slice(&value.to_le_bytes());
4594 }
4595 }
4596 (LogicalType::BigInt | LogicalType::Timestamp, Data::Int64(values)) => {
4597 for value in &**values {
4598 out.extend_from_slice(&value.to_le_bytes());
4599 }
4600 }
4601 (LogicalType::Boolean, Data::Bool(values)) => {
4602 for value in &**values {
4603 out.push(u8::from(*value));
4604 }
4605 }
4606 (LogicalType::Decimal { .. }, Data::Int16(values)) => {
4609 for value in &**values {
4610 out.extend_from_slice(&value.to_le_bytes());
4611 }
4612 }
4613 (LogicalType::Decimal { .. }, Data::Int32(values)) => {
4614 for value in &**values {
4615 out.extend_from_slice(&value.to_le_bytes());
4616 }
4617 }
4618 (LogicalType::Decimal { .. }, Data::Int64(values)) => {
4619 for value in &**values {
4620 out.extend_from_slice(&value.to_le_bytes());
4621 }
4622 }
4623 (LogicalType::Decimal { .. }, Data::Int128(values)) => {
4624 for value in &**values {
4625 out.extend_from_slice(&value.to_le_bytes());
4626 }
4627 }
4628 (LogicalType::Varchar, Data::Varlen(values)) => {
4629 let mut bytes = Vec::new();
4630 put_u32(&mut out, 0);
4631 for row in 0..vector.len() {
4632 let value = values.bytes(row).ok_or_else(|| invalid("string view is invalid"))?;
4633 bytes.extend_from_slice(value);
4634 put_u32(
4635 &mut out,
4636 u32::try_from(bytes.len())
4637 .map_err(|_| invalid("string payload exceeds 4GiB"))?,
4638 );
4639 }
4640 out.extend_from_slice(&bytes);
4641 }
4642 _ => return Err(Error::not_implemented(format!("native page for {ty}"))),
4643 }
4644 Ok((out, membership))
4645}
4646
4647fn put_varint(out: &mut Vec<u8>, mut value: u32) {
4648 while value >= 0x80 {
4649 out.push((value as u8 & 0x7f) | 0x80);
4650 value >>= 7;
4651 }
4652 out.push(value as u8);
4653}
4654
4655fn unique_codes(codes: &[u32]) -> Vec<u32> {
4657 let mut unique = codes.to_vec();
4658 unique.sort_unstable();
4659 unique.dedup();
4660 unique
4661}
4662
4663fn merged_codes(lists: Vec<Vec<u32>>) -> Vec<u32> {
4669 let mut lists = lists;
4670 while lists.len() > 1 {
4671 let mut next = Vec::with_capacity(lists.len().div_ceil(2));
4672 for pair in lists.chunks(2) {
4673 match pair {
4674 [left, right] => next.push(merged_pair(left, right)),
4675 [only] => next.push(only.clone()),
4676 _ => {}
4677 }
4678 }
4679 lists = next;
4680 }
4681 lists.pop().unwrap_or_default()
4682}
4683
4684fn merged_pair(left: &[u32], right: &[u32]) -> Vec<u32> {
4685 let mut out = Vec::with_capacity(left.len().saturating_add(right.len()));
4686 let mut at = 0;
4687 let mut to = 0;
4688 while at < left.len() && to < right.len() {
4689 match left[at].cmp(&right[to]) {
4690 Ordering::Less => {
4691 out.push(left[at]);
4692 at += 1;
4693 }
4694 Ordering::Greater => {
4695 out.push(right[to]);
4696 to += 1;
4697 }
4698 Ordering::Equal => {
4699 out.push(left[at]);
4700 at += 1;
4701 to += 1;
4702 }
4703 }
4704 }
4705 out.extend_from_slice(&left[at..]);
4706 out.extend_from_slice(&right[to..]);
4707 out
4708}
4709
4710fn merged_range(ranges: impl Iterator<Item = Range>) -> Range {
4715 let mut merged = Range::default();
4716 let mut first = true;
4717 for range in ranges {
4718 merged.nulls = merged.nulls.saturating_add(range.nulls);
4719 merged.sum = match (merged.sum.take(), range.sum) {
4723 (Some(held), Some(next)) if !first => held.checked_add(next),
4724 (_, next) if first => next,
4725 _ => None,
4726 };
4727 merged.exact = if first { range.exact } else { merged.exact && range.exact };
4728 if first {
4729 merged.low = range.low;
4730 merged.high = range.high;
4731 first = false;
4732 continue;
4733 }
4734 merged.low = match (merged.low.take(), range.low) {
4735 (Some(held), Some(next)) => Some(held.smaller(next)),
4736 _ => None,
4737 };
4738 merged.high = match (merged.high.take(), range.high) {
4739 (Some(held), Some(next)) => Some(held.larger(next)),
4740 _ => None,
4741 };
4742 }
4743 merged
4744}
4745
4746fn shortened(bound: Option<Bound>, high: bool) -> Option<Bound> {
4759 match bound {
4760 Some(Bound::Bytes(mut value)) if value.len() > PART_BOUND_BYTES => {
4761 value.truncate(PART_BOUND_BYTES);
4762 if !high {
4763 return Some(Bound::Bytes(value));
4764 }
4765 while let Some(last) = value.pop() {
4766 if last < u8::MAX {
4767 value.push(last + 1);
4768 return Some(Bound::Bytes(value));
4769 }
4770 }
4771 None
4772 }
4773 other => other,
4774 }
4775}
4776
4777fn encode_part_ranges(ranges: &[Range]) -> Result<Vec<u8>> {
4785 let mut out = Vec::new();
4786 put_u32(
4787 &mut out,
4788 u32::try_from(ranges.len()).map_err(|_| invalid("too many parts in a stripe"))?,
4789 );
4790 for range in ranges {
4791 put_bound(&mut out, shortened(range.low.clone(), false).as_ref())?;
4792 put_bound(&mut out, shortened(range.high.clone(), true).as_ref())?;
4793 put_u32(&mut out, u32::try_from(range.nulls).map_err(|_| invalid("null count overflow"))?);
4794 }
4795 Ok(out)
4796}
4797
4798fn decode_part_ranges(bytes: &[u8]) -> Result<Vec<Range>> {
4800 let mut cur = Cursor { bytes, at: 0 };
4801 let parts = cur.u32()? as usize;
4802 let mut out = Vec::new();
4803 for _ in 0..parts {
4804 let low = cur.bound()?;
4805 let high = cur.bound()?;
4806 let nulls = cur.u32()? as usize;
4807 out.push(Range { low, high, nulls, exact: false, sum: None });
4808 }
4809 Ok(out)
4810}
4811
4812fn encode_sieves<'a>(sieves: impl Iterator<Item = &'a Option<Sieve>>) -> Result<Vec<u8>> {
4813 let held: Vec<&Option<Sieve>> = sieves.collect();
4814 let mut out = Vec::new();
4815 put_u32(
4816 &mut out,
4817 u32::try_from(held.len()).map_err(|_| invalid("too many parts in a stripe"))?,
4818 );
4819 for sieve in &held {
4820 let length = sieve.as_ref().map_or(0, Sieve::len);
4821 put_u32(&mut out, u32::try_from(length).map_err(|_| invalid("sieve length overflow"))?);
4822 }
4823 for sieve in held.into_iter().flatten() {
4825 out.extend_from_slice(&sieve.to_bytes());
4826 }
4827 Ok(out)
4828}
4829
4830fn decode_sieves(bytes: &[u8]) -> Result<Vec<Option<Sieve>>> {
4836 let parts = u32::from_le_bytes(
4837 bytes
4838 .get(..4)
4839 .ok_or_else(|| invalid("sieve page is truncated"))?
4840 .try_into()
4841 .map_err(|_| invalid("sieve page is truncated"))?,
4842 ) as usize;
4843 let mut lengths = Vec::with_capacity(parts);
4844 for part in 0..parts {
4845 let at = 4 + part * 4;
4846 let field = bytes.get(at..at + 4).ok_or_else(|| invalid("sieve page is truncated"))?;
4847 lengths.push(u32::from_le_bytes(
4848 field.try_into().map_err(|_| invalid("sieve page is truncated"))?,
4849 ) as usize);
4850 }
4851 let mut at = 4 + parts * 4;
4852 let mut out = Vec::with_capacity(parts);
4853 for length in lengths {
4854 if length == 0 {
4855 out.push(None);
4856 continue;
4857 }
4858 let end = at.checked_add(length).ok_or_else(|| invalid("sieve page is truncated"))?;
4859 let field = bytes.get(at..end).ok_or_else(|| invalid("sieve page is truncated"))?;
4860 out.push(Sieve::from_bytes(field));
4861 at = end;
4862 }
4863 if at != bytes.len() {
4864 return Err(invalid("sieve page has trailing bytes"));
4865 }
4866 Ok(out)
4867}
4868
4869fn encode_membership(unique: &[u32]) -> Vec<u8> {
4875 let mut out = Vec::with_capacity(unique.len().saturating_mul(2).saturating_add(5));
4876 put_varint(&mut out, u32::try_from(unique.len()).unwrap_or(u32::MAX));
4877 let mut previous = 0;
4878 for (at, &code) in unique.iter().enumerate() {
4879 put_varint(&mut out, if at == 0 { code } else { code - previous });
4880 previous = code;
4881 }
4882 out
4883}
4884
4885fn take_varint(bytes: &[u8], at: &mut usize) -> Result<u32> {
4886 let mut value = 0_u32;
4887 for shift in (0..35).step_by(7) {
4888 let byte = *bytes.get(*at).ok_or_else(|| invalid("membership varint is truncated"))?;
4889 *at += 1;
4890 let part = u32::from(byte & 0x7f);
4891 if shift == 28 && part > 0x0f {
4892 return Err(invalid("membership varint overflow"));
4893 }
4894 value = value
4895 .checked_add(
4896 part.checked_shl(shift).ok_or_else(|| invalid("membership varint overflow"))?,
4897 )
4898 .ok_or_else(|| invalid("membership varint overflow"))?;
4899 if byte & 0x80 == 0 {
4900 return Ok(value);
4901 }
4902 }
4903 Err(invalid("membership varint is too long"))
4904}
4905
4906fn decode_membership(bytes: &[u8]) -> Result<Vec<u32>> {
4907 let mut at = 0;
4908 let count = take_varint(bytes, &mut at)? as usize;
4909 let mut codes = Vec::with_capacity(count);
4910 let mut previous = 0_u32;
4911 for index in 0..count {
4912 let delta = take_varint(bytes, &mut at)?;
4913 let code = if index == 0 {
4914 delta
4915 } else {
4916 previous.checked_add(delta).ok_or_else(|| invalid("membership code overflow"))?
4917 };
4918 if index > 0 && code <= previous {
4919 return Err(invalid("membership codes are not increasing"));
4920 }
4921 codes.push(code);
4922 previous = code;
4923 }
4924 if at != bytes.len() {
4925 return Err(invalid("membership page has trailing bytes"));
4926 }
4927 Ok(codes)
4928}
4929
4930fn string_dictionary(vector: &Vector) -> Result<Option<Vec<u8>>> {
4931 let mut by_text = HashMap::new();
4932 let mut values = Vec::new();
4933 let mut codes = Vec::with_capacity(vector.len());
4934 let mut plain_bytes = 0_usize;
4935 for row in 0..vector.len() {
4936 let text = vector.text_at(row).unwrap_or("");
4937 plain_bytes = plain_bytes.saturating_add(text.len());
4938 let code = match by_text.get(text) {
4939 Some(&code) => code,
4940 None => {
4941 let code = u32::try_from(values.len())
4942 .map_err(|_| invalid("too many dictionary values"))?;
4943 by_text.insert(text, code);
4944 values.push(text);
4945 code
4946 }
4947 };
4948 codes.push(code);
4949 }
4950 let dictionary_bytes = values.iter().map(|value| value.len()).sum::<usize>();
4951 let encoded = 8_usize
4952 .saturating_add((values.len() + 1).saturating_mul(4))
4953 .saturating_add(dictionary_bytes)
4954 .saturating_add(codes.len().saturating_mul(4));
4955 let plain = (vector.len() + 1).saturating_mul(4).saturating_add(plain_bytes);
4956 if encoded >= plain {
4957 return Ok(None);
4958 }
4959 let mut out = Vec::with_capacity(encoded);
4960 put_u32(
4961 &mut out,
4962 u32::try_from(values.len()).map_err(|_| invalid("too many dictionary values"))?,
4963 );
4964 put_u32(
4965 &mut out,
4966 u32::try_from(dictionary_bytes).map_err(|_| invalid("dictionary payload exceeds 4GiB"))?,
4967 );
4968 let mut offset = 0_u32;
4969 put_u32(&mut out, offset);
4970 for value in &values {
4971 offset = offset
4972 .checked_add(
4973 u32::try_from(value.len()).map_err(|_| invalid("dictionary value is too long"))?,
4974 )
4975 .ok_or_else(|| invalid("dictionary payload exceeds 4GiB"))?;
4976 put_u32(&mut out, offset);
4977 }
4978 for value in values {
4979 out.extend_from_slice(value.as_bytes());
4980 }
4981 for code in codes {
4982 put_u32(&mut out, code);
4983 }
4984 Ok(Some(out))
4985}
4986
4987struct EncodedDictionary {
4988 index: Vec<u8>,
4989 ranks: Vec<u8>,
4990 payload: Vec<Vec<u8>>,
4993}
4994
4995fn head(bytes: &[u8]) -> u64 {
4997 let mut word = [0; 8];
4998 let take = bytes.len().min(8);
4999 word[..take].copy_from_slice(&bytes[..take]);
5000 u64::from_be_bytes(word)
5001}
5002
5003fn rankings(dictionaries: &[Option<GlobalDictionary>]) -> Result<Vec<Vec<(u64, u32)>>> {
5011 let present =
5012 dictionaries.iter().enumerate().filter(|(_, held)| held.is_some()).map(|(at, _)| at);
5013 let present = present.collect::<Vec<_>>();
5014 let mut orders = vec![Vec::new(); dictionaries.len()];
5015 let workers = std::thread::available_parallelism()
5016 .map_or(1, usize::from)
5017 .min(MAX_FREQUENCY_WORKERS)
5018 .min(present.len());
5019 if workers <= 1 {
5020 for at in present {
5021 if let Some(dictionary) = &dictionaries[at] {
5022 orders[at] = dictionary.ranked();
5023 }
5024 }
5025 return Ok(orders);
5026 }
5027 let width = present.len().div_ceil(workers);
5028 let pieces = std::thread::scope(|scope| {
5029 present
5030 .chunks(width)
5031 .map(|columns| {
5032 scope.spawn(|| {
5033 columns
5034 .iter()
5035 .filter_map(|&at| dictionaries[at].as_ref().map(|held| (at, held.ranked())))
5036 .collect::<Vec<_>>()
5037 })
5038 })
5039 .collect::<Vec<_>>()
5040 .into_iter()
5041 .map(|handle| {
5042 handle.join().map_err(|_| Error::internal("a dictionary sort worker panicked"))
5043 })
5044 .collect::<Result<Vec<_>>>()
5045 })?;
5046 for piece in pieces {
5047 for (at, order) in piece {
5048 orders[at] = order;
5049 }
5050 }
5051 Ok(orders)
5052}
5053
5054fn encode_global_dictionary(
5055 dictionary: GlobalDictionary,
5056 order: &[(u64, u32)],
5057) -> Result<EncodedDictionary> {
5058 let values = dictionary.offsets.len() - 1;
5059 if order.len() != values {
5060 return Err(invalid("global dictionary order does not cover its values"));
5061 }
5062 let blocks = values.div_ceil(TEXT_PAYLOAD_VALUES);
5063 let payload = encode_payload(&dictionary)?;
5064 if payload.len() != blocks {
5065 return Err(invalid("global dictionary payload is not the blocks it says it is"));
5066 }
5067 let (ranks, rank_ends) = encode_ranks(order, code_width(values))?;
5068 let rank_blocks = values.div_ceil(TEXT_RANK_BLOCK);
5069 let offset_bits = offset_width(&dictionary.offsets);
5070 let mut index = Vec::with_capacity(
5071 DICTIONARY_HEADER + offset_bytes(values, offset_bits) + (blocks + rank_blocks) * 16,
5072 );
5073 put_u32(
5074 &mut index,
5075 u32::try_from(values).map_err(|_| invalid("global dictionary has too many values"))?,
5076 );
5077 put_u32(&mut index, TEXT_PAYLOAD_VALUES as u32);
5078 put_u32(
5079 &mut index,
5080 u32::try_from(blocks).map_err(|_| invalid("global dictionary has too many blocks"))?,
5081 );
5082 put_u32(&mut index, offset_bits as u32);
5083 encode_offsets(&dictionary.offsets, offset_bits, &mut index)?;
5084 let mut at = 0_u64;
5088 for block in &payload {
5089 at = at
5090 .checked_add(block.len() as u64)
5091 .ok_or_else(|| invalid("global dictionary payload overflow"))?;
5092 put_u64(&mut index, at);
5093 }
5094 for block in &payload {
5095 put_u64(&mut index, checksum(block));
5096 }
5097 if rank_ends.len() != rank_blocks {
5100 return Err(invalid("global dictionary order is not the blocks it says it is"));
5101 }
5102 for end in &rank_ends {
5103 put_u64(&mut index, *end);
5104 }
5105 let mut at = 0_usize;
5106 for end in &rank_ends {
5107 let end = usize::try_from(*end).map_err(|_| invalid("global dictionary order overflow"))?;
5108 put_u64(&mut index, checksum(&ranks[at..end]));
5109 at = end;
5110 }
5111 Ok(EncodedDictionary { index, ranks, payload })
5112}
5113
5114const PAYLOAD_SAMPLE_BLOCKS: usize = 8;
5121
5122fn payload_shapes() -> Vec<chooser::Settled> {
5148 let integers = vec![integer::Kind::Packed];
5149 [
5150 vec![string::Kind::Front, string::Kind::Lz],
5151 vec![string::Kind::Lz, string::Kind::Fsst],
5152 vec![string::Kind::Lz, string::Kind::Plain],
5153 vec![string::Kind::Fsst],
5154 vec![string::Kind::Plain],
5155 ]
5156 .into_iter()
5157 .map(|strings| chooser::Settled::new(strings, integers.clone()))
5158 .collect()
5159}
5160
5161fn encode_payload(dictionary: &GlobalDictionary) -> Result<Vec<Vec<u8>>> {
5167 let values = dictionary.offsets.len() - 1;
5168 let blocks = values.div_ceil(TEXT_PAYLOAD_VALUES);
5169 let run = |block: usize| {
5170 let first = block * TEXT_PAYLOAD_VALUES;
5171 let last = (first + TEXT_PAYLOAD_VALUES).min(values);
5172 (first..last)
5173 .map(|value| {
5174 let from = dictionary.offsets[value] as usize;
5175 let to = dictionary.offsets[value + 1] as usize;
5176 &dictionary.payload[from..to]
5177 })
5178 .collect::<Vec<_>>()
5179 };
5180 let shape = (blocks > PAYLOAD_SAMPLE_BLOCKS).then(|| settle_shape(&run, blocks)).transpose()?;
5183 let one = |block: usize| match &shape {
5184 Some(shape) => string::encode_with(&run(block), shape),
5185 None => string::encode(&run(block)),
5186 };
5187 let workers = std::thread::available_parallelism()
5188 .map_or(1, usize::from)
5189 .min(MAX_FREQUENCY_WORKERS)
5190 .min(blocks);
5191 if workers <= 1 {
5192 return (0..blocks).map(one).collect();
5193 }
5194 let next = AtomicUsize::new(0);
5195 let pieces = std::thread::scope(|scope| {
5196 (0..workers)
5197 .map(|_| {
5198 scope.spawn(|| {
5199 let mut mine = Vec::new();
5200 loop {
5201 let block = next.fetch_add(1, Atomic::Relaxed);
5202 if block >= blocks {
5203 break;
5204 }
5205 mine.push((block, one(block)?));
5206 }
5207 Ok(mine)
5208 })
5209 })
5210 .collect::<Vec<_>>()
5211 .into_iter()
5212 .map(|handle| {
5213 handle.join().map_err(|_| Error::internal("a dictionary encode worker panicked"))?
5214 })
5215 .collect::<Result<Vec<_>>>()
5216 })?;
5217 let mut payload = vec![Vec::new(); blocks];
5218 for piece in pieces {
5219 for (block, bytes) in piece {
5220 payload[block] = bytes;
5221 }
5222 }
5223 Ok(payload)
5224}
5225
5226fn settle_shape<'a>(
5234 run: &dyn Fn(usize) -> Vec<&'a [u8]>,
5235 blocks: usize,
5236) -> Result<chooser::Settled> {
5237 let last = blocks - 1;
5238 let sample = (0..PAYLOAD_SAMPLE_BLOCKS)
5239 .map(|region| run(region * last / (PAYLOAD_SAMPLE_BLOCKS - 1)))
5240 .collect::<Vec<_>>();
5241 let mut best: Option<(chooser::Settled, usize)> = None;
5242 for shape in payload_shapes() {
5243 let mut size = 0;
5244 for block in &sample {
5245 size += string::encode_with(block, &shape)?.len();
5246 }
5247 if best.as_ref().is_none_or(|(_, smallest)| size < *smallest) {
5248 best = Some((shape, size));
5249 }
5250 }
5251 best.map(|(shape, _)| shape)
5252 .ok_or_else(|| invalid("no shape applies to a global dictionary payload"))
5253}
5254
5255fn encode_ranks(order: &[(u64, u32)], code_bits: usize) -> Result<(Vec<u8>, Vec<u64>)> {
5262 let mut out = Vec::with_capacity(order.len() * 4);
5263 let mut ends = Vec::with_capacity(order.len().div_ceil(TEXT_RANK_BLOCK));
5264 let mut heads = Vec::with_capacity(TEXT_RANK_BLOCK);
5265 let mut codes = Vec::with_capacity(TEXT_RANK_BLOCK);
5266 for block in order.chunks(TEXT_RANK_BLOCK) {
5267 let base = block.first().map_or(0, |&(head, _)| head);
5270 let span = block.last().map_or(0, |&(head, _)| head.wrapping_sub(base));
5271 let width = (u64::BITS - span.leading_zeros()) as usize;
5272 heads.clear();
5273 codes.clear();
5274 for &(head, code) in block {
5275 heads.push(head.wrapping_sub(base));
5276 codes.push(u64::from(code));
5277 }
5278 put_u64(&mut out, base);
5279 out.push(width as u8);
5280 bitpack::pack_tail(&heads, width, &mut out)
5281 .map_err(|_| invalid("global dictionary heads do not pack"))?;
5282 bitpack::pack_tail(&codes, code_bits, &mut out)
5283 .map_err(|_| invalid("global dictionary codes do not pack"))?;
5284 ends.push(out.len() as u64);
5285 }
5286 Ok((out, ends))
5287}
5288
5289fn open_global_dictionary(
5296 file: Arc<File>,
5297 page: Page,
5298 ty: &LogicalType,
5299 keep_budget: usize,
5300) -> Result<Vector> {
5301 if ty != &LogicalType::Varchar {
5302 return Err(invalid("global dictionary belongs to a non-string column"));
5303 }
5304 let mut header = [0; DICTIONARY_HEADER];
5305 read_at(&file, page.offset, &mut header)?;
5306 let count = u32::from_le_bytes(header[0..4].try_into().expect("four bytes")) as usize;
5307 let per_block = u32::from_le_bytes(header[4..8].try_into().expect("four bytes")) as usize;
5308 let blocks = u32::from_le_bytes(header[8..12].try_into().expect("four bytes")) as usize;
5309 let offset_bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
5310 if per_block != TEXT_PAYLOAD_VALUES {
5311 return Err(invalid("global dictionary block width differs"));
5312 }
5313 if blocks != count.div_ceil(TEXT_PAYLOAD_VALUES) {
5314 return Err(invalid("global dictionary block count differs from its value count"));
5315 }
5316 if offset_bits > u32::BITS as usize {
5317 return Err(invalid("global dictionary packs offsets past a payload"));
5318 }
5319 let offset_len = offset_bytes(count, offset_bits);
5320 let ranks = count;
5325 let rank_blocks = ranks.div_ceil(TEXT_RANK_BLOCK);
5326 let hash_len = blocks
5329 .checked_add(rank_blocks)
5330 .and_then(|words| words.checked_mul(16))
5331 .ok_or_else(|| invalid("global dictionary block count overflow"))?;
5332 let index_len = DICTIONARY_HEADER
5333 .checked_add(offset_len)
5334 .and_then(|len| len.checked_add(hash_len))
5335 .ok_or_else(|| invalid("global dictionary header overflow"))?;
5336 if index_len > page.length as usize {
5337 return Err(invalid("global dictionary offset index exceeds its page"));
5338 }
5339 let mut index = vec![0; index_len];
5340 index[..DICTIONARY_HEADER].copy_from_slice(&header);
5341 read_at(&file, page.offset + DICTIONARY_HEADER as u64, &mut index[DICTIONARY_HEADER..])?;
5342 if checksum(&index) != page.hash {
5343 return Err(invalid("global dictionary index checksum differs"));
5344 }
5345 let offsets = index[DICTIONARY_HEADER..DICTIONARY_HEADER + offset_len].to_vec();
5346 let mut words = index[DICTIONARY_HEADER + offset_len..]
5347 .chunks_exact(8)
5348 .map(|part| u64::from_le_bytes(part.try_into().expect("eight bytes")))
5349 .collect::<Vec<_>>();
5350 let mut hashes = words.split_off(blocks);
5351 let mut rank_ends = hashes.split_off(blocks);
5352 let rank_hashes = rank_ends.split_off(rank_blocks);
5353 let ends = words;
5354 if rank_ends.windows(2).any(|pair| pair[0] >= pair[1]) {
5357 return Err(invalid("global dictionary order blocks do not rise"));
5358 }
5359 let rank_len = usize::try_from(rank_ends.last().copied().unwrap_or_default())
5360 .map_err(|_| invalid("global dictionary rank overflow"))?;
5361 let body_len = index_len
5362 .checked_add(rank_len)
5363 .ok_or_else(|| invalid("global dictionary header overflow"))?;
5364 if body_len > page.length as usize {
5365 return Err(invalid("global dictionary order exceeds its page"));
5366 }
5367 let stored_len = page.length as usize - body_len;
5370 if ends.last().copied().unwrap_or_default() as usize != stored_len
5371 || ends.windows(2).any(|pair| pair[0] > pair[1])
5372 {
5373 return Err(invalid("global dictionary blocks do not bound the payload"));
5374 }
5375 Vector::external_text(
5376 LogicalType::Varchar,
5377 Arc::new(NativeText {
5378 file,
5379 values: count,
5380 offsets,
5381 offset_bits,
5382 ranks,
5383 rank_at: page.offset + index_len as u64,
5384 rank_ends,
5385 rank_hashes,
5386 rank_blocks: (0..rank_blocks).map(|_| OnceLock::new()).collect(),
5387 code_bits: code_width(count),
5388 code_ranks: OnceLock::new(),
5389 payload: page.offset + body_len as u64,
5390 ends,
5391 hashes,
5392 blocks: (0..blocks).map(|_| OnceLock::new()).collect(),
5393 keep_budget,
5394 payload_kept: AtomicUsize::new(0),
5395 searched: Mutex::new(HashMap::new()),
5396 }),
5397 )
5398}
5399
5400fn decode(
5401 ty: &LogicalType,
5402 rows: usize,
5403 bytes: &[u8],
5404 global: Option<Arc<Vector>>,
5405) -> Result<Vector> {
5406 let mut cur = Cursor { bytes, at: 0 };
5407 let codec = cur.u8()?;
5408 let flag = cur.u8()?;
5409 let validity = match flag {
5410 0 => Validity::AllValid,
5411 1 => Validity::AllInvalid,
5412 2 => {
5413 let mask = cur.take(rows.div_ceil(8))?;
5414 Validity::from_iter(rows, |row| mask[row / 8] >> (row % 8) & 1 == 1)
5415 }
5416 _ => return Err(invalid("page validity tag differs")),
5417 };
5418 if codec == 1 {
5419 if ty != &LogicalType::Varchar {
5420 return Err(invalid("dictionary codec belongs to a non-string page"));
5421 }
5422 let count = cur.u32()? as usize;
5423 let payload_len = cur.u32()? as usize;
5424 let offset_bytes = cur.take(
5425 (count + 1)
5426 .checked_mul(4)
5427 .ok_or_else(|| invalid("dictionary offset count overflow"))?,
5428 )?;
5429 let offsets = offset_bytes
5430 .chunks_exact(4)
5431 .map(|part| u32::from_le_bytes(part.try_into().expect("four bytes")))
5432 .collect::<Vec<_>>();
5433 let payload = cur.take(payload_len)?.to_vec();
5434 if offsets.first() != Some(&0)
5435 || offsets.last().copied().map(|last| last as usize) != Some(payload.len())
5436 || offsets.windows(2).any(|pair| pair[0] > pair[1])
5437 {
5438 return Err(invalid("dictionary offsets do not bound the payload"));
5439 }
5440 let mut strings = StringColumn::over(Buffer::from_vec(payload).into_page());
5443 for pair in offsets.windows(2) {
5444 strings.push_in_place(pair[0] as usize, (pair[1] - pair[0]) as usize)?;
5445 }
5446 let mut codes = Vec::with_capacity(rows);
5447 for _ in 0..rows {
5448 codes.push(cur.u32()?);
5449 }
5450 if codes.iter().any(|code| *code as usize >= count) {
5451 return Err(invalid("dictionary code is out of range"));
5452 }
5453 if cur.at != bytes.len() {
5454 return Err(invalid("dictionary page has trailing bytes"));
5455 }
5456 let dictionary = Vector::flat(LogicalType::Varchar, Data::Varlen(strings))?;
5457 return Ok(Vector::dictionary(codes, dictionary)?.with_validity(validity));
5458 }
5459 if codec == 3 || codec == 4 {
5460 let dictionary = global.ok_or_else(|| invalid("global code page has no dictionary"))?;
5461 let codes = if codec == 4 {
5462 let wide = integer::decode(&bytes[cur.at..])?;
5465 if wide.len() != rows {
5466 return Err(invalid("encoded code page holds the wrong number of rows"));
5467 }
5468 let mut codes = Vec::with_capacity(wide.len());
5475 let mut seen = 0_i64;
5476 for &code in &wide {
5477 seen |= code;
5478 codes.push(code as u32);
5479 }
5480 if seen < 0 || seen > i64::from(u32::MAX) {
5481 return Err(invalid("code is not a code"));
5482 }
5483 codes
5484 } else {
5485 let mut codes = Vec::with_capacity(rows);
5486 for _ in 0..rows {
5487 codes.push(cur.u32()?);
5488 }
5489 if cur.at != bytes.len() {
5490 return Err(invalid("global code page has trailing bytes"));
5491 }
5492 codes
5493 };
5494 let highest = codes.iter().copied().max();
5495 return Ok(Vector::stable_dictionary_validated(codes, dictionary, highest)?
5496 .with_validity(validity));
5497 }
5498 if codec == 5 {
5499 let values = integer::decode(&bytes[cur.at..])?;
5501 if values.len() != rows {
5502 return Err(invalid("cascade page holds the wrong number of rows"));
5503 }
5504 let data = narrowed(ty, values)?;
5505 return Ok(Vector::flat(ty.clone(), data)?.with_validity(validity));
5506 }
5507 if codec == 2 {
5508 let width = u32::from(cur.u8()?);
5509 let base = i128::from_le_bytes(cur.take(16)?.try_into().expect("sixteen bytes"));
5510 let count = cur.u32()? as usize;
5511 let mut words = Vec::with_capacity(count);
5512 for _ in 0..count {
5513 words.push(cur.u64()?);
5514 }
5515 if cur.at != bytes.len() {
5516 return Err(invalid("packed page has trailing bytes"));
5517 }
5518 return Ok(Vector::packed(ty.clone(), words, width, base, rows)?.with_validity(validity));
5519 }
5520 if codec != 0 {
5521 return Err(invalid("page codec is unknown"));
5522 }
5523 let data = match ty {
5524 LogicalType::TinyInt => {
5525 let values = cur.take(rows)?;
5526 Data::Int8(values.iter().map(|item| *item as i8).collect::<Vec<_>>().into())
5527 }
5528 LogicalType::UTinyInt => Data::UInt8(cur.take(rows)?.to_vec().into()),
5529 LogicalType::SmallInt => {
5530 let values =
5531 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5532 Data::Int16(
5533 values
5534 .chunks_exact(2)
5535 .map(|item| i16::from_le_bytes(item.try_into().expect("two bytes")))
5536 .collect::<Vec<_>>()
5537 .into(),
5538 )
5539 }
5540 LogicalType::USmallInt => {
5541 let values =
5542 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5543 Data::UInt16(
5544 values
5545 .chunks_exact(2)
5546 .map(|item| u16::from_le_bytes(item.try_into().expect("two bytes")))
5547 .collect::<Vec<_>>()
5548 .into(),
5549 )
5550 }
5551 LogicalType::UInteger => {
5552 let values =
5553 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5554 Data::UInt32(
5555 values
5556 .chunks_exact(4)
5557 .map(|item| u32::from_le_bytes(item.try_into().expect("four bytes")))
5558 .collect::<Vec<_>>()
5559 .into(),
5560 )
5561 }
5562 LogicalType::UBigInt => {
5563 let values =
5564 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5565 Data::UInt64(
5566 values
5567 .chunks_exact(8)
5568 .map(|item| u64::from_le_bytes(item.try_into().expect("eight bytes")))
5569 .collect::<Vec<_>>()
5570 .into(),
5571 )
5572 }
5573 LogicalType::Integer | LogicalType::Date => {
5574 let values =
5575 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5576 Data::Int32(
5577 values
5578 .chunks_exact(4)
5579 .map(|item| i32::from_le_bytes(item.try_into().expect("four bytes")))
5580 .collect::<Vec<_>>()
5581 .into(),
5582 )
5583 }
5584 LogicalType::BigInt | LogicalType::Timestamp => {
5585 let values =
5586 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5587 Data::Int64(
5588 values
5589 .chunks_exact(8)
5590 .map(|item| i64::from_le_bytes(item.try_into().expect("eight bytes")))
5591 .collect::<Vec<_>>()
5592 .into(),
5593 )
5594 }
5595 LogicalType::Boolean => {
5596 let values = cur.take(rows)?;
5597 if values.iter().any(|value| *value > 1) {
5598 return Err(invalid("boolean page has another value"));
5599 }
5600 Data::Bool(values.iter().map(|value| *value == 1).collect::<Vec<_>>().into())
5601 }
5602 LogicalType::Decimal { .. } => match ty.physical() {
5605 PhysicalType::Int16 => {
5606 let values =
5607 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5608 Data::Int16(
5609 values
5610 .chunks_exact(2)
5611 .map(|item| i16::from_le_bytes(item.try_into().expect("two bytes")))
5612 .collect::<Vec<_>>()
5613 .into(),
5614 )
5615 }
5616 PhysicalType::Int32 => {
5617 let values =
5618 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5619 Data::Int32(
5620 values
5621 .chunks_exact(4)
5622 .map(|item| i32::from_le_bytes(item.try_into().expect("four bytes")))
5623 .collect::<Vec<_>>()
5624 .into(),
5625 )
5626 }
5627 PhysicalType::Int64 => {
5628 let values =
5629 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5630 Data::Int64(
5631 values
5632 .chunks_exact(8)
5633 .map(|item| i64::from_le_bytes(item.try_into().expect("eight bytes")))
5634 .collect::<Vec<_>>()
5635 .into(),
5636 )
5637 }
5638 _ => {
5639 let values =
5640 cur.take(rows.checked_mul(16).ok_or_else(|| invalid("page size overflow"))?)?;
5641 Data::Int128(
5642 values
5643 .chunks_exact(16)
5644 .map(|item| i128::from_le_bytes(item.try_into().expect("sixteen bytes")))
5645 .collect::<Vec<_>>()
5646 .into(),
5647 )
5648 }
5649 },
5650 LogicalType::Varchar => {
5651 let offset_bytes = cur
5652 .take((rows + 1).checked_mul(4).ok_or_else(|| invalid("offset count overflow"))?)?;
5653 let offsets = offset_bytes
5654 .chunks_exact(4)
5655 .map(|part| u32::from_le_bytes(part.try_into().expect("four bytes")))
5656 .collect::<Vec<_>>();
5657 let payload = cur.take(bytes.len() - cur.at)?.to_vec();
5658 if offsets.first() != Some(&0)
5659 || offsets.last().copied().map(|last| last as usize) != Some(payload.len())
5660 || offsets.windows(2).any(|pair| pair[0] > pair[1])
5661 {
5662 return Err(invalid("string offsets do not bound the payload"));
5663 }
5664 let mut values = StringColumn::over(Buffer::from_vec(payload).into_page());
5668 for pair in offsets.windows(2) {
5669 values.push_in_place(pair[0] as usize, (pair[1] - pair[0]) as usize)?;
5670 }
5671 Data::Varlen(values)
5672 }
5673 _ => return Err(Error::not_implemented(format!("native page for {ty}"))),
5674 };
5675 if cur.at != bytes.len() {
5676 return Err(invalid("page has trailing bytes"));
5677 }
5678 Ok(Vector::flat(ty.clone(), data)?.with_validity(validity))
5679}
5680
5681#[cfg(test)]
5682mod tests {
5683 use std::fs;
5684 use std::io::{Seek, SeekFrom, Write};
5685 use std::path::PathBuf;
5686 use std::time::{SystemTime, UNIX_EPOCH};
5687
5688 use rudb_common::Stat;
5689 use rudb_common::Value;
5690 use rudb_common::bounds::{Frequencies, Op, Zones};
5691 use rudb_common::stat::Provenance;
5692
5693 use super::*;
5694
5695 #[test]
5696 fn checksum_matches_fixed_vectors() {
5697 assert_eq!(checksum(b""), 0xef46_db37_51d8_e999);
5698 assert_eq!(checksum(b"a"), 0xd24e_c4f1_a98c_6e5b);
5699 assert_eq!(checksum(b"abc"), 0x44bc_2cf5_ad77_0999);
5700 }
5701
5702 fn path(label: &str) -> PathBuf {
5703 let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
5704 std::env::temp_dir().join(format!("rudb-native-{label}-{}-{stamp}.rdb", std::process::id()))
5705 }
5706
5707 #[test]
5709 fn a_read_at_an_offset_ignores_where_another_thread_left_the_cursor() {
5710 const SPANS: usize = 64;
5711 const SPAN: usize = 512;
5712 let path = path("positional");
5713 let content: Vec<u8> =
5714 (0..SPANS).flat_map(|span| std::iter::repeat_n(span as u8, SPAN)).collect();
5715 fs::write(&path, &content).expect("the file is written");
5716 let file = Arc::new(File::open(&path).expect("the file opens"));
5717 std::thread::scope(|scope| {
5718 for _ in 0..8 {
5719 let file = Arc::clone(&file);
5720 scope.spawn(move || {
5721 for _ in 0..64 {
5722 for span in 0..SPANS {
5723 let mut bytes = [0_u8; SPAN];
5724 read_at(&file, (span * SPAN) as u64, &mut bytes)
5725 .expect("the span reads");
5726 assert!(
5727 bytes.iter().all(|byte| *byte == span as u8),
5728 "span {span} came back as {}",
5729 bytes[0],
5730 );
5731 }
5732 }
5733 });
5734 }
5735 });
5736 let mut past = [0_u8; SPAN];
5737 let end = (SPANS * SPAN) as u64;
5738 let error = read_at(&file, end, &mut past).expect_err("a read past the end is refused");
5739 assert!(error.message().contains("ends before its declared length"), "{error}");
5740 drop(file);
5741 let _ = fs::remove_file(&path);
5742 }
5743
5744 #[test]
5750 fn a_writer_puts_a_page_where_it_said_it_did_wherever_the_cursor_has_got_to() {
5751 let path = path("cursor");
5752 let mut writer = Writer::create(
5753 &path,
5754 "items",
5755 vec![
5756 Field::required("id", LogicalType::Integer),
5757 Field::new("text", LogicalType::Varchar),
5758 ],
5759 )
5760 .expect("new file");
5761 writer.append(&sample()).expect("first part");
5762 writer.file.seek(SeekFrom::Start(0)).expect("the cursor goes back to the header");
5763 writer.append(&sample()).expect("second part");
5764 writer.file.seek(SeekFrom::Start(1)).expect("and somewhere useless again");
5765 writer.finish().expect("commit");
5766 let reader = Reader::open(&path).expect("reopen from disk");
5767 assert_eq!(reader.table().rows(), 6);
5768 let ids = reader.read(0, &[0]).expect("the integer page reads back");
5769 assert_eq!(ids.value_at(0, 0), Value::Integer(4));
5770 assert_eq!(ids.value_at(2, 0), Value::Integer(-2));
5771 let text = reader.read(1, &[1]).expect("the text page reads back");
5772 assert_eq!(text.value_at(1, 0), Value::Null);
5773 assert_eq!(text.value_at(2, 0), Value::Varchar("long text after a slash".into()));
5774 let end = reader.table().stripes().iter().flat_map(|stripe| {
5777 stripe
5778 .pages
5779 .iter()
5780 .map(|page| page.offset + u64::from(page.length))
5781 .chain(std::iter::once(stripe.index.offset + u64::from(stripe.index.length)))
5782 });
5783 let last = end.fold(HEADER, u64::max);
5784 let directory = fs::metadata(&path).expect("the file is there").len();
5785 assert!(last <= directory, "a page runs to {last} in a file of {directory} bytes");
5786 fs::remove_file(path).expect("remove scratch file");
5787 }
5788
5789 fn dictionary_index_len(header: &[u8; DICTIONARY_HEADER]) -> u64 {
5795 let count = u64::from(u32::from_le_bytes(header[0..4].try_into().expect("four bytes")));
5796 let blocks = u64::from(u32::from_le_bytes(header[8..12].try_into().expect("four bytes")));
5797 let bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
5798 let rank_blocks = count.div_ceil(TEXT_RANK_BLOCK as u64);
5799 DICTIONARY_HEADER as u64
5800 + offset_bytes(count as usize, bits) as u64
5801 + (blocks + rank_blocks) * 16
5802 }
5803
5804 fn last_rank_end(file: &File, offset: u64, header: &[u8; DICTIONARY_HEADER]) -> u64 {
5806 let count = u64::from(u32::from_le_bytes(header[0..4].try_into().expect("four bytes")));
5807 let blocks = u64::from(u32::from_le_bytes(header[8..12].try_into().expect("four bytes")));
5808 let bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
5809 let rank_blocks = count.div_ceil(TEXT_RANK_BLOCK as u64);
5810 let at = offset
5811 + DICTIONARY_HEADER as u64
5812 + offset_bytes(count as usize, bits) as u64
5813 + blocks * 16
5814 + (rank_blocks - 1) * 8;
5815 let mut end = [0; 8];
5816 read_at(file, at, &mut end).expect("the last rank block end");
5817 u64::from_le_bytes(end)
5818 }
5819
5820 fn sample() -> Chunk {
5821 Chunk::new(vec![
5822 Vector::from_values(
5823 LogicalType::Integer,
5824 &[Value::Integer(4), Value::Integer(9), Value::Integer(-2)],
5825 )
5826 .expect("integers"),
5827 Vector::from_values(
5828 LogicalType::Varchar,
5829 &[
5830 Value::Varchar("alpha".into()),
5831 Value::Null,
5832 Value::Varchar("long text after a slash".into()),
5833 ],
5834 )
5835 .expect("strings"),
5836 ])
5837 .expect("matching rows")
5838 }
5839
5840 fn sample_ids() -> Chunk {
5841 Chunk::new(vec![
5842 Vector::flat(LogicalType::Integer, Data::Int32(vec![7, 8, 9].into()))
5843 .expect("integers"),
5844 ])
5845 .expect("one column")
5846 }
5847
5848 #[test]
5849 fn the_planner_gets_the_null_count_off_the_same_directory_the_bounds_are_in() {
5850 let path = path("nulls_for_the_planner");
5853 let mut writer =
5854 Writer::create(&path, "items", vec![Field::new("a", LogicalType::Integer)])
5855 .expect("new file");
5856 let rows = Chunk::new(vec![
5857 Vector::from_values(
5858 LogicalType::Integer,
5859 &[
5860 Value::Integer(4),
5861 Value::Null,
5862 Value::Integer(9),
5863 Value::Null,
5864 Value::Integer(1),
5865 Value::Integer(2),
5866 ],
5867 )
5868 .expect("integers"),
5869 ])
5870 .expect("one column");
5871 writer.append(&rows).expect("the only part");
5872 writer.finish().expect("commit");
5873 let reader = Reader::open(&path).expect("reopen from disk");
5874 let stripes = Stripes::new(reader);
5875 let column = stripes.column("a").expect("the file has that column");
5876 assert_eq!(stripes.nulls(column), Stat::exact(2, Provenance::NullCount));
5877 assert_eq!(stripes.nulls(column + 1), Stat::Unknown);
5880 fs::remove_file(&path).expect("clean up");
5881 }
5882
5883 #[test]
5884 fn the_planner_gets_a_row_count_per_value_off_a_complete_synopsis() {
5885 let path = path("frequencies_for_the_planner");
5890 let mut writer =
5891 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
5892 .expect("new file");
5893 let rows = Chunk::new(vec![
5894 Vector::from_values(
5895 LogicalType::Integer,
5896 &[
5897 Value::Integer(4),
5898 Value::Integer(4),
5899 Value::Integer(4),
5900 Value::Integer(9),
5901 Value::Integer(9),
5902 Value::Integer(1),
5903 ],
5904 )
5905 .expect("integers"),
5906 ])
5907 .expect("one column");
5908 writer.append(&rows).expect("the only part");
5909 writer.finish().expect("commit");
5910 let reader = Reader::open(&path).expect("reopen from disk");
5911 let common = Common::new(reader);
5912 assert_eq!(common.rows(), 6);
5913 let column = common.column("id").expect("the file has that column");
5914 assert_eq!(common.column("nothing"), None);
5915 assert_eq!(
5916 common.rows_with(column, &Bound::Int(4)),
5917 Stat::exact(3, Provenance::FrequencySynopsis)
5918 );
5919 assert_eq!(
5921 common.rows_with(column, &Bound::Int(7)),
5922 Stat::exact(0, Provenance::FrequencySynopsis)
5923 );
5924 assert_eq!(common.rows_with(column, &Bound::Bytes(b"four".to_vec())), Stat::Unknown);
5927 fs::remove_file(&path).expect("clean up");
5928 }
5929
5930 #[test]
5931 fn committed_file_reopens_and_reads_only_requested_columns() {
5932 let path = path("reopen");
5933 let mut writer = Writer::create(
5934 &path,
5935 "items",
5936 vec![
5937 Field::required("id", LogicalType::Integer),
5938 Field::new("text", LogicalType::Varchar),
5939 ],
5940 )
5941 .expect("new file");
5942 writer.append(&sample()).expect("first part");
5943 writer.append(&sample()).expect("second part");
5944 writer.finish().expect("commit");
5945 let reader = Reader::open(&path).expect("reopen from disk");
5946 assert_eq!(reader.table().rows(), 6);
5947 assert_eq!(reader.table().stripes().len(), 1);
5950 assert_eq!(reader.parts(), 2);
5951 assert_eq!(reader.part_rows(0), 3);
5952 assert_eq!(reader.part_rows(1), 3);
5953 let text = reader.read(1, &[1]).expect("only text page");
5954 assert_eq!(text.width(), 1);
5955 assert_eq!(text.value_at(1, 0), Value::Null);
5956 assert_eq!(text.value_at(2, 0), Value::Varchar("long text after a slash".into()));
5957 let sparse = reader.read_sparse(1, &[1]).expect("one part without its whole page");
5958 assert_eq!(sparse.width(), 1);
5959 assert_eq!(sparse.value_at(1, 0), Value::Null);
5960 assert_eq!(sparse.value_at(2, 0), Value::Varchar("long text after a slash".into()));
5961 assert!(!reader.skips_codes(0, 1, &[0]).expect("alpha is in the stripe"));
5962 assert!(!reader.skips_codes(0, 1, &[2]).expect("long text is in the stripe"));
5963 assert!(reader.skips_codes(0, 1, &[3]).expect("unknown code is absent"));
5964 let count = reader.read(0, &[]).expect("no page is needed for count");
5965 assert_eq!(count.len(), 3);
5966 assert!(reader.skips(0, &[Probe { column: 0, op: Op::Greater, value: Bound::Int(100) }]));
5967 assert!(!reader.skips(0, &[Probe { column: 0, op: Op::Greater, value: Bound::Int(0) }]));
5968 let integers = reader.top_frequencies(0, 1).expect("valid integer synopsis").expect("kept");
5969 assert_eq!(
5970 integers,
5971 vec![(Value::Integer(-2), 2), (Value::Integer(4), 2), (Value::Integer(9), 2),]
5972 );
5973 let strings = reader.top_frequencies(1, 1).expect("valid string synopsis").expect("kept");
5974 assert_eq!(strings.len(), 3);
5975 assert!(strings.contains(&(Value::Null, 2)));
5976 assert!(strings.contains(&(Value::Varchar("alpha".into()), 2)));
5977 assert!(strings.contains(&(Value::Varchar("long text after a slash".into()), 2)));
5978 fs::remove_file(path).expect("remove scratch file");
5979 }
5980
5981 #[test]
5989 fn runs_handed_over_out_of_order_still_read_back_in_source_order() {
5990 let path = path("interleaved-runs");
5991 let mut writer =
5992 Writer::create(&path, "interleaved", vec![Field::new("v", LogicalType::BigInt)])
5993 .expect("new file");
5994 for morsel in [2_u64, 0, 3, 1] {
5995 let parts = (0..4_u64)
5996 .map(|chunk| {
5997 let first = i64::try_from(morsel * 32 + chunk * 8).expect("small");
5998 let values =
5999 (0..8_i64).map(|row| Value::BigInt(first + row)).collect::<Vec<_>>();
6000 let column =
6001 Vector::from_values(LogicalType::BigInt, &values).expect("a column");
6002 ((morsel, chunk), Chunk::new(vec![column]).expect("one column"))
6003 })
6004 .collect::<Vec<_>>();
6005 writer.append_stripe(parts).expect("a stripe");
6006 }
6007 writer.finish().expect("commit");
6008
6009 let reader = Reader::open(&path).expect("valid directory");
6010 assert_eq!(reader.table().stripes().len(), 4, "a run is a stripe of its own");
6011 assert_eq!(reader.table().rows(), 128);
6012 for part in 0..16_usize {
6013 let read = reader.read(part, &[0]).expect("a part back");
6014 for row in 0..8_usize {
6015 let want = i64::try_from(part * 8 + row).expect("small");
6016 assert_eq!(read.value_at(row, 0), Value::BigInt(want), "part {part} row {row}");
6017 }
6018 }
6019 fs::remove_file(path).expect("remove scratch file");
6020 }
6021
6022 #[test]
6025 fn runs_that_overlap_each_other_are_refused_at_commit() {
6026 let path = path("overlapping-runs");
6027 let mut writer =
6028 Writer::create(&path, "overlapping", vec![Field::new("v", LogicalType::BigInt)])
6029 .expect("new file");
6030 let one = |order: (u64, u64)| {
6031 let column =
6032 Vector::from_values(LogicalType::BigInt, &[Value::BigInt(1)]).expect("a column");
6033 (order, Chunk::new(vec![column]).expect("one column"))
6034 };
6035 writer.append_stripe(vec![one((0, 0)), one((0, 2))]).expect("a stripe");
6038 writer.append_stripe(vec![one((0, 1))]).expect("a stripe");
6039 let error = writer.finish().expect_err("the runs overlap");
6040 assert!(error.message().contains("source order"), "{error}");
6041 fs::remove_file(path).expect("remove scratch file");
6042 }
6043
6044 #[test]
6047 fn a_run_longer_than_a_stripe_is_refused() {
6048 let path = path("overlong-run");
6049 let mut writer =
6050 Writer::create(&path, "overlong", vec![Field::new("v", LogicalType::BigInt)])
6051 .expect("new file");
6052 let parts = (0..=STRIPE_PARTS)
6053 .map(|at| {
6054 let column = Vector::from_values(LogicalType::BigInt, &[Value::BigInt(1)])
6055 .expect("a column");
6056 let chunk = Chunk::new(vec![column]).expect("one column");
6057 ((0, u64::try_from(at).expect("small")), chunk)
6058 })
6059 .collect::<Vec<_>>();
6060 let error = writer.append_stripe(parts).expect_err("one part too many");
6061 assert!(error.message().contains("more parts than it holds"), "{error}");
6062 fs::remove_file(path).expect("remove scratch file");
6063 }
6064
6065 #[test]
6071 fn parts_past_the_stripe_bound_start_a_new_stripe() {
6072 let path = path("stripe-bound");
6073 let mut writer = Writer::create(
6074 &path,
6075 "items",
6076 vec![
6077 Field::required("id", LogicalType::Integer),
6078 Field::new("text", LogicalType::Varchar),
6079 ],
6080 )
6081 .expect("new file");
6082 let parts = STRIPE_PARTS * 2 + 3;
6083 for part in 0..parts {
6084 let id = part as i32;
6085 let chunk = Chunk::new(vec![
6086 Vector::from_values(
6087 LogicalType::Integer,
6088 &[Value::Integer(id), Value::Integer(-id)],
6089 )
6090 .expect("integers"),
6091 Vector::from_values(
6092 LogicalType::Varchar,
6093 &[Value::Varchar(format!("value {part}")), Value::Null],
6094 )
6095 .expect("strings"),
6096 ])
6097 .expect("matching rows");
6098 writer.append(&chunk).expect("one part");
6099 }
6100 writer.finish().expect("commit");
6101
6102 let reader = Reader::open(&path).expect("reopen from disk");
6103 assert_eq!(reader.parts(), parts);
6104 assert_eq!(reader.table().rows(), parts * 2);
6105 assert_eq!(reader.table().stripes().len(), parts.div_ceil(STRIPE_PARTS));
6106 assert_eq!(reader.table().stripes()[0].parts(), STRIPE_PARTS);
6107 assert_eq!(reader.table().stripes()[0].rows(), STRIPE_PARTS * 2);
6108 assert_eq!(reader.table().stripes()[2].parts(), 3);
6109 for part in (0..parts).rev() {
6112 let dense = reader.read(part, &[0, 1]).expect("a whole page read");
6113 let sparse = reader.read_sparse(part, &[0, 1]).expect("one part read");
6114 for chunk in [&dense, &sparse] {
6115 assert_eq!(chunk.len(), 2, "part {part} has its own row count");
6116 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6117 assert_eq!(chunk.value_at(1, 0), Value::Integer(-(part as i32)));
6118 assert_eq!(chunk.value_at(0, 1), Value::Varchar(format!("value {part}")));
6119 assert_eq!(chunk.value_at(1, 1), Value::Null);
6120 }
6121 }
6122 let above = [Probe { column: 0, op: Op::Greater, value: Bound::Int(100) }];
6125 assert!(reader.skips(0, &above), "the first stripe stops at 63");
6126 assert!(!reader.skips(STRIPE_PARTS * 2, &above), "the third stripe reaches 130");
6127 fs::remove_file(path).expect("remove scratch file");
6128 }
6129
6130 fn scattered(n: i64) -> i64 {
6132 n.wrapping_mul(-7_046_029_254_386_353_131)
6133 }
6134
6135 #[test]
6141 fn a_part_is_skipped_when_its_sieve_does_not_hold_the_constant() {
6142 let path = path("sieve-skip");
6143 let mut writer =
6144 Writer::create(&path, "hits", vec![Field::required("id", LogicalType::BigInt)])
6145 .expect("new file");
6146 let parts = STRIPE_PARTS + 3;
6147 let per_part = 128;
6151 for part in 0..parts {
6152 let held: Vec<Value> = (0..per_part)
6153 .map(|row| Value::BigInt(scattered((part * per_part + row) as i64)))
6154 .collect();
6155 let chunk =
6156 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6157 .expect("one column");
6158 writer.append(&chunk).expect("one part");
6159 }
6160 writer.finish().expect("commit");
6161
6162 let reader = Reader::open(&path).expect("reopen from disk");
6163 let probe = |value: i64| Probe {
6164 column: 0,
6165 op: Op::Equal,
6166 value: Bound::Int(i128::from(scattered(value))),
6167 };
6168 for wanted in [0_i64, (per_part + 1) as i64, (parts * per_part - 1) as i64] {
6169 let tests = [probe(wanted)];
6170 let kept: Vec<usize> = (0..parts).filter(|&part| !reader.skips(part, &tests)).collect();
6171 let home = wanted as usize / per_part;
6172 assert!(kept.contains(&home), "the part holding {wanted} is read");
6173 assert!(kept.len() <= 2, "{wanted} keeps {kept:?}, which is more than one stray part");
6177 }
6178 let absent = [probe((parts * per_part) as i64 + 1)];
6179 let kept = (0..parts).filter(|&part| !reader.skips(part, &absent)).count();
6180 assert!(kept <= 1, "{kept} parts of {parts} kept a value no part holds");
6181 let tests = [probe(0)];
6184 assert!(
6185 reader.table().stripes().iter().all(|stripe| !stripe.zone.skips(&tests)),
6186 "the bounds rule out no stripe at all"
6187 );
6188 fs::remove_file(path).expect("remove scratch file");
6189 }
6190
6191 #[test]
6197 fn a_part_is_skipped_when_its_own_bounds_rule_out_a_comparison_the_stripe_keeps() {
6198 let path = path("part-range-skip");
6199 let mut writer =
6200 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
6201 .expect("new file");
6202 let parts = STRIPE_PARTS + 3;
6203 let per_part = 128;
6204 for part in 0..parts {
6205 let held: Vec<Value> = (0..per_part)
6209 .map(|row| {
6210 Value::BigInt((part * 1_000) as i64 + (scattered(row as i64).rem_euclid(900)))
6211 })
6212 .collect();
6213 let chunk =
6214 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6215 .expect("one column");
6216 writer.append(&chunk).expect("one part");
6217 }
6218 writer.finish().expect("commit");
6219
6220 let reader = Reader::open(&path).expect("reopen from disk");
6221 let under = [Probe { column: 0, op: Op::Less, value: Bound::Int(3_000) }];
6222 let kept: Vec<usize> = (0..parts).filter(|&part| !reader.skips(part, &under)).collect();
6223 assert_eq!(kept, vec![0, 1, 2], "only the three parts that start under three thousand");
6224 assert!(!reader.stripe_skips(0, &under), "the stripe reaches from zero and keeps itself");
6226 fs::remove_file(path).expect("remove scratch file");
6227 }
6228
6229 #[test]
6233 fn a_part_is_waved_through_when_its_own_bounds_pass_a_comparison_the_stripe_cannot() {
6234 let path = path("part-range-certain");
6235 let mut writer =
6236 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
6237 .expect("new file");
6238 let parts = STRIPE_PARTS + 3;
6239 let per_part = 128;
6240 for part in 0..parts {
6241 let held: Vec<Value> = (0..per_part)
6242 .map(|row| {
6243 Value::BigInt((part * 1_000) as i64 + (scattered(row as i64).rem_euclid(900)))
6244 })
6245 .collect();
6246 let chunk =
6247 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6248 .expect("one column");
6249 writer.append(&chunk).expect("one part");
6250 }
6251 writer.finish().expect("commit");
6252
6253 let reader = Reader::open(&path).expect("reopen from disk");
6254 let under = [Probe { column: 0, op: Op::Less, value: Bound::Int(3_000) }];
6255 let waved: Vec<usize> = (0..parts).filter(|&part| reader.certain(part, &under)).collect();
6256 assert_eq!(waved, vec![0, 1, 2], "the three parts that end under three thousand");
6257 assert!(!reader.stripe_skips(0, &under), "the stripe straddles the comparison");
6260 fs::remove_file(path).expect("remove scratch file");
6261 }
6262
6263 #[test]
6266 fn a_stripe_of_one_part_writes_no_range_page_and_a_stripe_of_many_does() {
6267 for (parts, wanted) in [(1_usize, false), (STRIPE_PARTS, true)] {
6268 let path = path("part-range-page");
6269 let mut writer =
6270 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
6271 .expect("new file");
6272 for part in 0..parts {
6273 let held: Vec<Value> = (0..128)
6274 .map(|row| {
6275 Value::BigInt((part * 1_000) as i64 + scattered(row as i64).rem_euclid(900))
6276 })
6277 .collect();
6278 let chunk = Chunk::new(vec![
6279 Vector::from_values(LogicalType::BigInt, &held).expect("numbers"),
6280 ])
6281 .expect("one column");
6282 writer.append(&chunk).expect("one part");
6283 }
6284 writer.finish().expect("commit");
6285 let reader = Reader::open(&path).expect("reopen from disk");
6286 let bytes = reader.layout().columns[0].part_ranges;
6287 assert_eq!(bytes > 0, wanted, "{parts} parts wrote {bytes} bytes of ranges");
6288 fs::remove_file(path).expect("remove scratch file");
6289 }
6290 }
6291
6292 #[test]
6295 fn a_string_end_that_is_cut_down_still_covers_the_value_it_came_from() {
6296 let long = vec![b'a'; PART_BOUND_BYTES * 2];
6297 let low = shortened(Some(Bound::Bytes(long.clone())), false).expect("a low end");
6298 let high = shortened(Some(Bound::Bytes(long.clone())), true).expect("a high end");
6299 let Bound::Bytes(low) = low else { panic!("a string stays a string") };
6300 let Bound::Bytes(high) = high else { panic!("a string stays a string") };
6301 assert!(low.len() <= PART_BOUND_BYTES && high.len() <= PART_BOUND_BYTES);
6302 assert!(low.as_slice() <= long.as_slice(), "the low end is at or under the value");
6303 assert!(high.as_slice() >= long.as_slice(), "the high end is at or over the value");
6304 }
6305
6306 #[test]
6309 fn a_string_end_with_no_room_to_step_up_gives_up_the_bound() {
6310 let long = vec![u8::MAX; PART_BOUND_BYTES * 2];
6311 assert_eq!(shortened(Some(Bound::Bytes(long.clone())), true), None);
6312 let low = shortened(Some(Bound::Bytes(long)), false).expect("a low end is still a prefix");
6313 assert_eq!(low, Bound::Bytes(vec![u8::MAX; PART_BOUND_BYTES]));
6314 }
6315
6316 #[test]
6326 fn a_sieve_larger_than_the_part_it_indexes_is_not_written() {
6327 let path = path("sieve-pays");
6328 let fields = vec![
6329 Field::required("spread", LogicalType::BigInt),
6330 Field::required("repeated", LogicalType::BigInt),
6331 ];
6332 let mut writer = Writer::create(&path, "hits", fields).expect("new file");
6333 let parts = 3;
6334 let per_part = 1024;
6335 for part in 0..parts {
6336 let base = (part * per_part) as i64;
6337 let spread: Vec<Value> =
6338 (0..per_part).map(|row| Value::BigInt(scattered(base + row as i64))).collect();
6339 let repeated: Vec<Value> =
6340 (0..per_part).map(|row| Value::BigInt(scattered((row / 256) as i64))).collect();
6341 let chunk = Chunk::new(vec![
6342 Vector::from_values(LogicalType::BigInt, &spread).expect("numbers"),
6343 Vector::from_values(LogicalType::BigInt, &repeated).expect("numbers"),
6344 ])
6345 .expect("two columns");
6346 writer.append(&chunk).expect("one part");
6347 }
6348 writer.finish().expect("commit");
6349
6350 let reader = Reader::open(&path).expect("reopen from disk");
6351 let layout = reader.layout();
6352 let spread = &layout.columns[0];
6353 let repeated = &layout.columns[1];
6354 assert!(spread.sieves > 0, "a column whose parts are worth a filter keeps one");
6355 assert_eq!(
6356 repeated.sieves, 0,
6357 "a column whose filter costs more than its parts keeps none"
6358 );
6359 for column in &layout.columns {
6362 assert!(
6363 column.sieves < column.pages,
6364 "{} spends {} on sieves over {} of data",
6365 column.name,
6366 column.sieves,
6367 column.pages
6368 );
6369 }
6370 let absent = [Probe {
6372 column: 0,
6373 op: Op::Equal,
6374 value: Bound::Int(i128::from(scattered((parts * per_part) as i64 + 1))),
6375 }];
6376 assert!((0..parts).all(|part| reader.skips(part, &absent)), "no part holds it");
6377 fs::remove_file(path).expect("remove scratch file");
6378 }
6379
6380 #[test]
6386 fn a_damaged_sieve_page_is_read_through_rather_than_refused() {
6387 let path = path("sieve-damaged");
6388 let mut writer =
6389 Writer::create(&path, "hits", vec![Field::required("id", LogicalType::BigInt)])
6390 .expect("new file");
6391 let rows = 128;
6392 let held: Vec<Value> = (0..rows).map(|row| Value::BigInt(scattered(row))).collect();
6393 let chunk =
6394 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
6395 .expect("one column");
6396 writer.append(&chunk).expect("one part");
6397 writer.finish().expect("commit");
6398
6399 let page =
6400 Reader::open(&path).expect("reopen").table.stripes[0].sieves[0].expect("a sieve page");
6401 let mut file = OpenOptions::new().write(true).open(&path).expect("open the sieve page");
6402 file.seek(SeekFrom::Start(page.offset + u64::from(page.length) - 1)).expect("seek");
6403 file.write_all(&[0xff]).expect("damage one byte");
6404 drop(file);
6405
6406 let reader = Reader::open(&path).expect("reopen the damaged file");
6407 let absent =
6408 [Probe { column: 0, op: Op::Equal, value: Bound::Int(i128::from(scattered(99))) }];
6409 assert!(!reader.skips(0, &absent), "a sieve that cannot be read skips nothing");
6410 assert_eq!(
6411 reader.read(0, &[0]).expect("the rows are untouched").len(),
6412 usize::try_from(rows).expect("a small count")
6413 );
6414 fs::remove_file(path).expect("remove scratch file");
6415 }
6416
6417 #[test]
6428 fn workers_that_want_the_same_stripe_read_it_once() {
6429 let path = path("single-flight");
6430 let mut writer =
6431 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6432 .expect("new file");
6433 for part in 0..STRIPE_PARTS {
6434 let id = part as i32;
6435 let chunk = Chunk::new(vec![
6436 Vector::from_values(
6437 LogicalType::Integer,
6438 &[Value::Integer(id), Value::Integer(-id)],
6439 )
6440 .expect("integers"),
6441 ])
6442 .expect("matching rows");
6443 writer.append(&chunk).expect("one part");
6444 }
6445 writer.finish().expect("commit");
6446
6447 let reader = Reader::open(&path).expect("reopen from disk");
6448 assert_eq!(reader.table().stripes().len(), 1, "one stripe is the point of the test");
6449 let barrier = std::sync::Barrier::new(8);
6450 std::thread::scope(|scope| {
6451 for worker in 0..8 {
6452 let reader = &reader;
6453 let barrier = &barrier;
6454 scope.spawn(move || {
6455 barrier.wait();
6456 for part in (worker..STRIPE_PARTS).step_by(8) {
6457 let chunk = reader.read(part, &[0]).expect("a whole page read");
6458 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6459 assert_eq!(chunk.value_at(1, 0), Value::Integer(-(part as i32)));
6460 }
6461 });
6462 }
6463 });
6464 assert_eq!(reader.pages.load(Atomic::Relaxed), 1, "one stripe, one page read, whoever won");
6465 fs::remove_file(path).expect("remove scratch file");
6466 }
6467
6468 #[test]
6481 fn opening_costs_the_same_over_a_thousand_times_the_rows() {
6482 let opened = |label: &str, rows_per_part: i32| {
6483 let path = path(label);
6484 let mut writer =
6485 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6486 .expect("new file");
6487 for part in 0..STRIPE_PARTS * 3 {
6488 let values = (0..rows_per_part)
6492 .map(|row| {
6493 Value::Integer((part as i32 * rows_per_part + row).wrapping_mul(2_654_435))
6494 })
6495 .collect::<Vec<_>>();
6496 let chunk = Chunk::new(vec![
6497 Vector::from_values(LogicalType::Integer, &values).expect("integers"),
6498 ])
6499 .expect("matching rows");
6500 writer.append(&chunk).expect("one part");
6501 }
6502 writer.finish().expect("commit");
6503 let reader = Reader::open(&path).expect("reopen from disk");
6504 let size = fs::metadata(&path).expect("the file is there").len();
6505 let out = (reader.reads(), reader.table().stripes().len(), size);
6506 fs::remove_file(path).expect("remove scratch file");
6507 out
6508 };
6509
6510 let (thin, thin_stripes, thin_size) = opened("open-thin", 1);
6511 let (fat, fat_stripes, fat_size) = opened("open-fat", 1000);
6512 assert_eq!(
6513 thin_stripes, fat_stripes,
6514 "the same stripe count is what makes this a fair ask"
6515 );
6516 assert!(
6517 fat_size > thin_size * 50,
6518 "the fat file has to actually be larger, and it is {fat_size} against {thin_size}"
6519 );
6520
6521 assert_eq!(thin.opening.reads, fat.opening.reads, "the same reads either way");
6522 assert_eq!(thin.pages, 0, "opening read a page");
6523 assert_eq!(fat.pages, 0, "opening read a page");
6524 assert_eq!(thin.indexes, 0, "opening read an index");
6525 assert_eq!(fat.indexes, 0, "opening read an index");
6526 assert!(
6529 fat.opening.bytes < thin.opening.bytes * 2,
6530 "opening the thin file read {} bytes and the fat one read {}",
6531 thin.opening.bytes,
6532 fat.opening.bytes
6533 );
6534 }
6535
6536 #[test]
6544 fn two_opens_of_one_file_cost_the_same_and_the_second_is_not_cheaper() {
6545 let path = path("open-twice");
6546 let mut writer =
6547 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6548 .expect("new file");
6549 for part in 0..STRIPE_PARTS * 3 {
6550 let chunk = Chunk::new(vec![
6551 Vector::from_values(LogicalType::Integer, &[Value::Integer(part as i32)])
6552 .expect("integers"),
6553 ])
6554 .expect("matching rows");
6555 writer.append(&chunk).expect("one part");
6556 }
6557 writer.finish().expect("commit");
6558
6559 let first = Reader::open(&path).expect("open");
6560 for part in 0..first.parts() {
6563 first.read(part, &[0]).expect("a part");
6564 }
6565 assert!(first.reads().pages > 0, "the scan has to have read something");
6566 let second = Reader::open(&path).expect("open again");
6567
6568 assert_eq!(first.reads().opening, second.reads().opening);
6569 assert_eq!(
6570 second.reads().pages,
6571 0,
6572 "the second open read a page off the back of the first"
6573 );
6574 assert_eq!(second.reads().indexes, 0, "the second open read an index it inherited");
6575 fs::remove_file(path).expect("remove scratch file");
6576 }
6577
6578 #[test]
6586 fn an_index_is_read_once_per_stripe_however_often_the_page_is_evicted() {
6587 let path = path("index-cache");
6588 let mut writer =
6589 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6590 .expect("new file");
6591 let parts = STRIPE_PARTS * (CACHED_STRIPES_PER_COLUMN + 2);
6592 for part in 0..parts {
6593 let id = part as i32;
6594 let chunk = Chunk::new(vec![
6595 Vector::from_values(LogicalType::Integer, &[Value::Integer(id)]).expect("integers"),
6596 ])
6597 .expect("matching rows");
6598 writer.append(&chunk).expect("one part");
6599 }
6600 writer.finish().expect("commit");
6601
6602 let reader = Reader::open(&path).expect("reopen from disk");
6603 let stripes = reader.table().stripes().len();
6604 assert!(stripes > CACHED_STRIPES_PER_COLUMN, "the page cache has to be too small for this");
6605 for _ in 0..2 {
6607 for part in 0..parts {
6608 let chunk = reader.read(part, &[0]).expect("a part");
6609 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6610 }
6611 }
6612 assert_eq!(reader.indexes.load(Atomic::Relaxed), stripes, "one index read per stripe");
6613 assert!(
6614 reader.pages.load(Atomic::Relaxed) > stripes,
6615 "the pages are the ones that get read again, which is what makes the index count mean \
6616 something"
6617 );
6618 fs::remove_file(path).expect("remove scratch file");
6619 }
6620
6621 #[test]
6630 fn a_worker_per_stripe_reads_its_page_once_when_the_cache_was_told_to_expect_it() {
6631 let workers = CACHED_STRIPES_PER_COLUMN + 4;
6632 let path = path("stripe-per-worker");
6633 let mut writer =
6634 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6635 .expect("new file");
6636 for part in 0..STRIPE_PARTS * workers {
6637 let chunk = Chunk::new(vec![
6638 Vector::from_values(LogicalType::Integer, &[Value::Integer(part as i32)])
6639 .expect("integers"),
6640 ])
6641 .expect("matching rows");
6642 writer.append(&chunk).expect("one part");
6643 }
6644 writer.finish().expect("commit");
6645
6646 let read = |told: bool| {
6647 let reader = Reader::open(&path).expect("reopen from disk");
6648 assert_eq!(reader.table().stripes().len(), workers, "a stripe per worker");
6649 if told {
6650 reader.keep_stripes(workers);
6651 }
6652 let barrier = std::sync::Barrier::new(workers);
6653 std::thread::scope(|scope| {
6654 for (worker, run) in reader.stripe_parts().into_iter().enumerate() {
6655 let reader = &reader;
6656 let barrier = &barrier;
6657 scope.spawn(move || {
6658 for part in run {
6659 barrier.wait();
6660 let chunk = reader.read(part, &[0]).expect("a part of my own stripe");
6661 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6662 }
6663 assert!(worker < workers);
6664 });
6665 }
6666 });
6667 reader.pages.load(Atomic::Relaxed)
6668 };
6669
6670 assert_eq!(read(true), workers, "one page read per stripe and no more");
6671 assert!(read(false) > workers, "a cache that small is read again on every part");
6672 fs::remove_file(path).expect("remove scratch file");
6673 }
6674
6675 #[test]
6680 fn a_damaged_index_page_is_an_error() {
6681 let path = path("damaged-index");
6682 let mut writer =
6683 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6684 .expect("new file");
6685 writer.append(&sample_ids()).expect("first part");
6686 writer.append(&sample_ids()).expect("second part");
6687 writer.finish().expect("commit");
6688
6689 let reader = Reader::open(&path).expect("valid directory");
6690 let index = reader.table.stripes[0].index;
6691 let mut byte = [0; 1];
6692 read_at(&reader.file, index.offset, &mut byte).expect("the first part length");
6693 let mut file = OpenOptions::new().write(true).open(&path).expect("open index page");
6694 file.seek(SeekFrom::Start(index.offset)).expect("index start");
6695 file.write_all(&[!byte[0]]).expect("damage the first part length");
6696 let error = reader.read(1, &[0]).expect_err("a damaged index must not be used");
6697 assert!(error.message().contains("index page section checksum differs"), "{error}");
6698 fs::remove_file(path).expect("remove scratch file");
6699 }
6700
6701 #[test]
6708 fn every_integer_width_round_trips_through_a_page() {
6709 let path = path("integer-widths");
6710 let columns = [
6711 (LogicalType::TinyInt, vec![Value::TinyInt(i8::MIN), Value::TinyInt(i8::MAX)]),
6712 (LogicalType::UTinyInt, vec![Value::UTinyInt(0), Value::UTinyInt(u8::MAX)]),
6713 (LogicalType::SmallInt, vec![Value::SmallInt(i16::MIN), Value::SmallInt(i16::MAX)]),
6714 (LogicalType::USmallInt, vec![Value::USmallInt(0), Value::USmallInt(u16::MAX)]),
6715 (LogicalType::Integer, vec![Value::Integer(i32::MIN), Value::Integer(i32::MAX)]),
6716 (LogicalType::UInteger, vec![Value::UInteger(0), Value::UInteger(u32::MAX)]),
6717 (LogicalType::BigInt, vec![Value::BigInt(i64::MIN), Value::BigInt(i64::MAX)]),
6718 (LogicalType::UBigInt, vec![Value::UBigInt(0), Value::UBigInt(u64::MAX)]),
6719 ];
6720 let fields = columns
6721 .iter()
6722 .enumerate()
6723 .map(|(at, (ty, _))| Field::required(format!("c{at}"), ty.clone()))
6724 .collect::<Vec<_>>();
6725 let vectors = columns
6726 .iter()
6727 .map(|(ty, values)| Vector::from_values(ty.clone(), values).expect("a vector"))
6728 .collect::<Vec<_>>();
6729 let mut writer = Writer::create(&path, "widths", fields).expect("new file");
6730 writer.append(&Chunk::new(vectors).expect("matching rows")).expect("one stripe");
6731 writer.finish().expect("commit");
6732
6733 let reader = Reader::open(&path).expect("reopen from disk");
6734 let wanted = (0..columns.len()).collect::<Vec<_>>();
6735 let read = reader.read(0, &wanted).expect("every column");
6736 assert_eq!(read.len(), 2);
6737 for (at, (ty, values)) in columns.iter().enumerate() {
6739 assert_eq!(read.value_at(0, at), values[0], "the low end of {ty}");
6740 assert_eq!(read.value_at(1, at), values[1], "the high end of {ty}");
6741 }
6742 fs::remove_file(path).expect("remove scratch file");
6743 }
6744
6745 #[test]
6746 fn numeric_frequency_candidates_keep_bounded_row_ordinals() {
6747 let path = path("frequency-ordinals");
6748 let mut writer =
6749 Writer::create(&path, "items", vec![Field::required("id", LogicalType::BigInt)])
6750 .expect("new file");
6751 let mut values = Vec::new();
6752 for leader in 0..10_i64 {
6753 values.extend(std::iter::repeat_n(leader, 100));
6754 }
6755 values.extend(1_000_i64..41_000);
6756 for part in values.chunks(1_024) {
6757 let vector = Vector::flat(LogicalType::BigInt, Data::Int64(part.to_vec().into()))
6758 .expect("big integers");
6759 writer.append(&Chunk::new(vec![vector]).expect("one column")).expect("one stripe");
6760 }
6761 writer.finish().expect("commit");
6762
6763 let reader = Reader::open(&path).expect("reopen from disk");
6764 let occurrences =
6765 reader.frequency_occurrences(0).expect("valid metadata").expect("bounded ordinals");
6766 assert!(occurrences.omitted_max < 100);
6767 assert!(occurrences.ordinals.len() <= FREQUENCY_ORDINALS);
6768 assert!(occurrences.ordinals.windows(2).all(|pair| pair[0] < pair[1]));
6769 assert_eq!(&occurrences.ordinals[..1_000], &(0_u64..1_000).collect::<Vec<_>>());
6770 fs::remove_file(path).expect("remove scratch file");
6771 }
6772
6773 #[test]
6779 fn a_file_from_another_format_says_which_format_it_is() {
6780 let older = path("older-format");
6781 let mut writer =
6782 Writer::create(&older, "items", vec![Field::new("id", LogicalType::Integer)])
6783 .expect("new file");
6784 let chunk = Chunk::new(vec![
6785 Vector::flat(LogicalType::Integer, Data::Int32(vec![1, 2, 3].into()))
6786 .expect("integers"),
6787 ])
6788 .expect("chunk");
6789 writer.append(&chunk).expect("page written");
6790 writer.finish().expect("commit");
6791
6792 let mut file = OpenOptions::new().write(true).open(&older).expect("open for the header");
6793 file.seek(SeekFrom::Start(8)).expect("the version follows the magic");
6794 file.write_all(&(FORMAT - 1).to_le_bytes()).expect("write an older version");
6795 drop(file);
6796 let complaint = Reader::open(&older).expect_err("an older format is refused").to_string();
6797 assert!(complaint.contains(&format!("format {}", FORMAT - 1)), "{complaint}");
6798 assert!(complaint.contains(&format!("format {FORMAT}")), "{complaint}");
6799
6800 let mut file = OpenOptions::new().write(true).open(&older).expect("open for the header");
6801 file.seek(SeekFrom::Start(0)).expect("the magic is first");
6802 file.write_all(b"NOTRUDB!").expect("write another engine's magic");
6803 drop(file);
6804 let complaint = Reader::open(&older).expect_err("a foreign file is refused").to_string();
6805 assert!(complaint.contains("magic"), "{complaint}");
6806 assert!(!complaint.contains("format"), "a version has nothing to do with it: {complaint}");
6807 fs::remove_file(older).expect("remove scratch file");
6808 }
6809
6810 #[test]
6811 fn an_unfinished_or_damaged_file_does_not_answer_with_partial_rows() {
6812 let unfinished = path("unfinished");
6813 let mut writer =
6814 Writer::create(&unfinished, "items", vec![Field::new("id", LogicalType::Integer)])
6815 .expect("new file");
6816 let chunk = Chunk::new(vec![
6817 Vector::flat(LogicalType::Integer, Data::Int32(vec![1, 2, 3].into()))
6818 .expect("integers"),
6819 ])
6820 .expect("chunk");
6821 writer.append(&chunk).expect("page written");
6822 drop(writer);
6823 assert!(Reader::open(&unfinished).is_err(), "no directory was committed");
6824 fs::remove_file(unfinished).expect("remove scratch file");
6825
6826 let damaged = path("damaged");
6827 let mut writer =
6828 Writer::create(&damaged, "items", vec![Field::new("id", LogicalType::Integer)])
6829 .expect("new file");
6830 writer.append(&chunk).expect("page written");
6831 writer.finish().expect("commit");
6832 let reader = Reader::open(&damaged).expect("valid directory");
6833 let mut file =
6834 OpenOptions::new().write(true).open(&damaged).expect("open for a damaged page");
6835 file.seek(SeekFrom::Start(HEADER + 1)).expect("inside first page");
6836 file.write_all(&[255]).expect("damage one byte");
6837 assert!(reader.read(0, &[0]).is_err(), "page checksum rejects corruption");
6838 fs::remove_file(damaged).expect("remove scratch file");
6839 }
6840
6841 #[test]
6842 fn damaged_lazy_dictionary_payload_is_an_error() {
6843 let path = path("damaged-dictionary");
6844 let mut writer = Writer::create(
6845 &path,
6846 "items",
6847 vec![
6848 Field::required("id", LogicalType::Integer),
6849 Field::new("text", LogicalType::Varchar),
6850 ],
6851 )
6852 .expect("new file");
6853 writer.append(&sample()).expect("stripe written");
6854 writer.finish().expect("commit");
6855
6856 let reader = Reader::open(&path).expect("valid directory");
6857 let dictionary = reader.table.dictionaries[1].expect("string dictionary page");
6858 let mut header = [0; DICTIONARY_HEADER];
6861 read_at(&reader.file, dictionary.offset, &mut header).expect("dictionary header");
6862 let index_len = dictionary_index_len(&header);
6863 let rank_len = last_rank_end(&reader.file, dictionary.offset, &header);
6864 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
6865 file.seek(SeekFrom::Start(dictionary.offset + index_len + rank_len))
6866 .expect("inside dictionary payload");
6867 file.write_all(&[255]).expect("damage dictionary payload");
6868
6869 let chunk = reader.read(0, &[1]).expect("code page and dictionary index remain valid");
6870 let error =
6871 chunk.validate_external().expect_err("payload corruption must reach the caller");
6872 assert!(error.message().contains("payload checksum differs"), "{error}");
6873 fs::remove_file(path).expect("remove scratch file");
6874 }
6875
6876 #[test]
6883 fn a_dictionary_over_many_blocks_checks_every_block_of_it() {
6884 let path = path("dictionary-blocks");
6885 let value =
6886 |row: usize| format!("{row:07} a value long enough to be worth a payload block");
6887 let parts = 30;
6888 let per_part = 1000;
6889 let mut writer =
6890 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
6891 .expect("new file");
6892 for part in 0..parts {
6893 let values = (0..per_part)
6894 .map(|row| Value::Varchar(value(part * per_part + row)))
6895 .collect::<Vec<_>>();
6896 let chunk = Chunk::new(vec![
6897 Vector::from_values(LogicalType::Varchar, &values).expect("strings"),
6898 ])
6899 .expect("matching rows");
6900 writer.append(&chunk).expect("a part");
6901 }
6902 writer.finish().expect("commit");
6903
6904 let reader = Reader::open(&path).expect("reopen from disk");
6905 let dictionary = reader.table.dictionaries[0].expect("string dictionary page");
6906 assert!(
6907 parts * per_part > TEXT_PAYLOAD_VALUES * 4,
6908 "the dictionary has to be several blocks for this to be testing anything"
6909 );
6910 for part in [0, parts - 1] {
6911 let chunk = reader.read(part, &[0]).expect("a part");
6912 chunk.validate_external().expect("every payload block checks out");
6913 assert_eq!(chunk.value_at(0, 0), Value::Varchar(value(part * per_part)));
6914 }
6915
6916 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
6917 file.seek(SeekFrom::Start(dictionary.offset + u64::from(dictionary.length) - 4))
6918 .expect("the last bytes of the page are payload");
6919 file.write_all(&[255]).expect("damage the last payload block");
6920 let reader = Reader::open(&path).expect("the directory and the index are untouched");
6921 let chunk = reader.read(parts - 1, &[0]).expect("the code page remains valid");
6922 let error = chunk.validate_external().expect_err("the damage must reach the caller");
6923 assert!(error.message().contains("payload checksum differs"), "{error}");
6924 fs::remove_file(path).expect("remove scratch file");
6925 }
6926
6927 #[test]
6937 fn values_of_different_lengths_read_back_out_of_packed_offsets() {
6938 let path = path("dictionary-offsets");
6939 let value = |row: usize| {
6940 if row % 511 == 3 { String::new() } else { "x".repeat(row % 97) + &format!("{row:05}") }
6941 };
6942 let rows = 5_000;
6943 let mut writer =
6944 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
6945 .expect("new file");
6946 let values = (0..rows).map(|row| Value::Varchar(value(row))).collect::<Vec<_>>();
6947 for part in values.chunks(1_000) {
6948 let chunk =
6949 Chunk::new(vec![Vector::from_values(LogicalType::Varchar, part).expect("strings")])
6950 .expect("matching rows");
6951 writer.append(&chunk).expect("a part");
6952 }
6953 writer.finish().expect("commit");
6954
6955 let reader = Reader::open(&path).expect("reopen from disk");
6956 assert!(
6957 rows > TEXT_PAYLOAD_VALUES * 4,
6958 "the dictionary has to be several blocks for this to be testing anything"
6959 );
6960 for part in 0..rows / 1_000 {
6961 let chunk = reader.read(part, &[0]).expect("a part");
6962 for row in 0..1_000 {
6963 let row = part * 1_000 + row;
6964 assert_eq!(
6965 chunk.value_at(row % 1_000, 0),
6966 Value::Varchar(value(row)),
6967 "value {row}"
6968 );
6969 }
6970 }
6971 fs::remove_file(path).expect("remove scratch file");
6972 }
6973
6974 #[test]
6986 fn a_global_dictionary_is_opened_once_however_many_workers_ask_at_once() {
6987 let path = path("dictionary-once");
6988 let parts = 8;
6989 let per_part = 500;
6990 let value =
6991 |row: usize| format!("{row:07} a value long enough to be worth a payload block");
6992 let mut writer =
6993 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
6994 .expect("new file");
6995 for part in 0..parts {
6996 let values = (0..per_part)
6997 .map(|row| Value::Varchar(value(part * per_part + row)))
6998 .collect::<Vec<_>>();
6999 let chunk = Chunk::new(vec![
7000 Vector::from_values(LogicalType::Varchar, &values).expect("strings"),
7001 ])
7002 .expect("matching rows");
7003 writer.append(&chunk).expect("a part");
7004 }
7005 writer.finish().expect("commit");
7006
7007 let reader = Reader::open(&path).expect("reopen from disk");
7008 assert!(reader.table.dictionaries[0].is_some(), "the column has to have one to share");
7009 assert_eq!(reader.reads().dictionaries, 0, "opening the file does not open a dictionary");
7010
7011 let workers = 16;
7012 let gate = std::sync::Barrier::new(workers);
7013 std::thread::scope(|scope| {
7014 for worker in 0..workers {
7015 let reader = reader.clone();
7016 let gate = &gate;
7017 scope.spawn(move || {
7018 gate.wait();
7019 let chunk = reader.read(worker % parts, &[0]).expect("a part");
7020 assert_eq!(
7021 chunk.value_at(0, 0),
7022 Value::Varchar(value((worker % parts) * per_part))
7023 );
7024 });
7025 }
7026 });
7027
7028 assert_eq!(reader.reads().dictionaries, 1, "sixteen workers, one dictionary, one open");
7029 fs::remove_file(path).expect("remove scratch file");
7030 }
7031
7032 #[test]
7037 fn a_damaged_sorted_order_is_an_error() {
7038 let path = path("damaged-order");
7039 let mut writer = Writer::create(
7040 &path,
7041 "items",
7042 vec![
7043 Field::required("id", LogicalType::Integer),
7044 Field::new("text", LogicalType::Varchar),
7045 ],
7046 )
7047 .expect("new file");
7048 writer.append(&sample()).expect("stripe written");
7049 writer.finish().expect("commit");
7050
7051 let reader = Reader::open(&path).expect("valid directory");
7052 let page = reader.table.dictionaries[1].expect("string dictionary page");
7053 let mut header = [0; DICTIONARY_HEADER];
7054 read_at(&reader.file, page.offset, &mut header).expect("dictionary header");
7055 let index_len = dictionary_index_len(&header);
7056 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
7057 file.seek(SeekFrom::Start(page.offset + index_len)).expect("the first head");
7058 file.write_all(&[255]).expect("damage the order");
7059
7060 let dictionary = reader.dictionary(1).expect("read").expect("a string column has one");
7061 let error = dictionary.compare_rank(0, b"anything").expect_err("a damaged order is caught");
7062 assert!(error.message().contains("rank checksum differs"), "{error}");
7063 fs::remove_file(path).expect("remove scratch file");
7064 }
7065
7066 #[test]
7070 fn a_global_dictionary_carries_the_sorted_order_of_its_values() {
7071 let spellings = ["overlong1z", "b", "", "overlong1a", "overlong", "ab", "a", "overlong1"];
7074 let path = path("dictionary-order");
7075 let mut writer =
7076 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7077 .expect("new file");
7078 writer
7079 .append(
7080 &Chunk::new(vec![
7081 Vector::from_values(
7082 LogicalType::Varchar,
7083 &spellings.map(|text| Value::Varchar(text.into())),
7084 )
7085 .expect("strings"),
7086 ])
7087 .expect("one column"),
7088 )
7089 .expect("stripe written");
7090 writer.finish().expect("commit");
7091
7092 let reader = Reader::open(&path).expect("valid directory");
7093 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
7094 let count = dictionary.ranks().expect("a v10 file stores one");
7095 assert_eq!(count, spellings.len(), "every distinct value has a rank");
7096 let order = (0..count)
7097 .map(|rank| dictionary.code_at_rank(rank).expect("a code"))
7098 .collect::<Vec<_>>();
7099 let mut seen = order.clone();
7100 seen.sort_unstable();
7101 assert_eq!(seen, (0..spellings.len() as u32).collect::<Vec<_>>(), "a permutation of codes");
7102
7103 let ranked = order
7104 .iter()
7105 .map(|&code| {
7106 dictionary.try_bytes_at(code as usize).expect("read").expect("a value").to_vec()
7107 })
7108 .collect::<Vec<_>>();
7109 let mut expected = spellings.map(|text| text.as_bytes().to_vec()).to_vec();
7110 expected.sort();
7111 assert_eq!(ranked, expected, "rank order is value order");
7112
7113 for (rank, value) in expected.iter().enumerate() {
7116 assert_eq!(
7117 dictionary.compare_rank(rank, value).expect("compare"),
7118 Ordering::Equal,
7119 "rank {rank} is its own value"
7120 );
7121 if rank > 0 {
7122 assert_eq!(
7123 dictionary.compare_rank(rank - 1, value).expect("compare"),
7124 Ordering::Less,
7125 "rank {rank} follows the one before it"
7126 );
7127 }
7128 }
7129 fs::remove_file(path).expect("remove scratch file");
7130 }
7131
7132 #[test]
7140 fn a_dictionary_sweep_reads_every_value_and_keeps_it_under_the_budget() {
7141 let path = path("dictionary-sweep");
7142 let spellings = (0..2_500)
7145 .map(|index| Value::Varchar(format!("value {index:08} {}", "x".repeat(index % 40))))
7146 .collect::<Vec<_>>();
7147 let mut writer =
7148 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7149 .expect("new file");
7150 for part in spellings.chunks(1_024) {
7153 writer
7154 .append(
7155 &Chunk::new(vec![
7156 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
7157 ])
7158 .expect("one column"),
7159 )
7160 .expect("stripe written");
7161 }
7162 writer.finish().expect("commit");
7163
7164 let reader = Reader::open(&path).expect("valid directory");
7165 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
7166 assert_eq!(dictionary.len(), spellings.len(), "every value is distinct");
7167
7168 let resting = dictionary.footprint();
7169 let mut swept: Vec<Vec<u8>> = Vec::new();
7170 let mut at = 0;
7171 let mut calls = 0;
7172 while at < dictionary.len() {
7173 let stopped = dictionary
7174 .sweep_text(at, dictionary.len(), &mut |index: usize, text: &[u8]| {
7175 assert_eq!(index, swept.len(), "a sweep hands its values over in order");
7176 swept.push(text.to_vec());
7177 Ok(())
7178 })
7179 .expect("a sweep reads");
7180 assert!(stopped > at, "a sweep moves");
7181 at = stopped;
7182 calls += 1;
7183 }
7184 assert_eq!(calls, 3, "a sweep hands over one block at a time");
7185 let after = dictionary.footprint();
7186 assert!(after > resting, "a sweep under the budget keeps what it decoded");
7187
7188 let read = (0..dictionary.len())
7189 .map(|code| dictionary.try_bytes_at(code).expect("read").expect("a value").to_vec())
7190 .collect::<Vec<_>>();
7191 assert_eq!(swept, read, "a sweep answers what a point read answers");
7192 assert_eq!(dictionary.footprint(), after, "a point read of a kept block decodes nothing");
7193 fs::remove_file(path).expect("remove scratch file");
7194 }
7195
7196 #[test]
7207 fn a_sweep_over_a_block_with_a_short_second_run_reads_what_a_point_read_reads() {
7208 let path = path("dictionary-sweep-short-run");
7209 let spellings = (0..2_800)
7210 .map(|index| Value::Varchar(format!("value {index:08} {}", "x".repeat(index % 40))))
7211 .collect::<Vec<_>>();
7212 let mut writer =
7213 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7214 .expect("new file");
7215 for part in spellings.chunks(1_024) {
7216 writer
7217 .append(
7218 &Chunk::new(vec![
7219 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
7220 ])
7221 .expect("one column"),
7222 )
7223 .expect("stripe written");
7224 }
7225 writer.finish().expect("commit");
7226
7227 let reader = Reader::open(&path).expect("valid directory");
7228 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
7229 assert_eq!(dictionary.len(), spellings.len(), "every value is distinct");
7230 let last = dictionary.len() % TEXT_PAYLOAD_VALUES;
7231 assert!(last > TEXT_OFFSET_RUN, "the last block has to reach into a second run of offsets");
7232 assert!(last < TEXT_PAYLOAD_VALUES, "and that second run has to be short of a whole one");
7233
7234 let mut swept: Vec<Vec<u8>> = Vec::new();
7235 let mut at = 0;
7236 while at < dictionary.len() {
7237 let stopped = dictionary
7238 .sweep_text(at, dictionary.len(), &mut |index: usize, text: &[u8]| {
7239 assert_eq!(index, swept.len(), "a sweep hands its values over in order");
7240 swept.push(text.to_vec());
7241 Ok(())
7242 })
7243 .expect("a sweep reads");
7244 assert!(stopped > at, "a sweep moves");
7245 at = stopped;
7246 }
7247 let read = (0..dictionary.len())
7248 .map(|code| dictionary.try_bytes_at(code).expect("read").expect("a value").to_vec())
7249 .collect::<Vec<_>>();
7250 assert_eq!(swept, read, "a sweep answers what a point read answers");
7251 fs::remove_file(path).expect("remove scratch file");
7252 }
7253
7254 #[test]
7264 fn narrowing_a_page_takes_what_fits_and_refuses_what_does_not() {
7265 assert_eq!(fit::<i8>(&[]).expect("an empty page fits anything"), Vec::<i8>::new());
7266 assert_eq!(fit::<i8>(&[-128, 0, 127]).expect("the edges fit"), vec![-128_i8, 0, 127]);
7267 fit::<i8>(&[128]).expect_err("one past the top does not fit");
7268 fit::<i8>(&[-129]).expect_err("one past the bottom does not fit");
7269 assert_eq!(fit::<u8>(&[0, 255]).expect("the edges fit"), vec![0_u8, 255]);
7270 fit::<u8>(&[256]).expect_err("one past the top does not fit");
7271 fit::<u8>(&[-1]).expect_err("a negative does not fit an unsigned page");
7272 assert_eq!(
7273 fit::<i16>(&[-32_768, 0, 32_767]).expect("the edges fit"),
7274 vec![-32_768_i16, 0, 32_767]
7275 );
7276 fit::<i16>(&[32_768]).expect_err("one past the top does not fit");
7277 fit::<i16>(&[-32_769]).expect_err("one past the bottom does not fit");
7278 assert_eq!(fit::<u16>(&[0, 65_535]).expect("the edges fit"), vec![0_u16, 65_535]);
7279 fit::<u16>(&[65_536]).expect_err("one past the top does not fit");
7280 fit::<u16>(&[-1]).expect_err("a negative does not fit an unsigned page");
7281 assert_eq!(
7282 fit::<i32>(&[i64::from(i32::MIN), 0, i64::from(i32::MAX)]).expect("the edges fit"),
7283 vec![i32::MIN, 0, i32::MAX]
7284 );
7285 fit::<i32>(&[i64::from(i32::MAX) + 1]).expect_err("one past the top does not fit");
7286 fit::<i32>(&[i64::from(i32::MIN) - 1]).expect_err("one past the bottom does not fit");
7287 assert_eq!(
7288 fit::<u32>(&[0, 4_294_967_295]).expect("the edges fit"),
7289 vec![0_u32, 4_294_967_295]
7290 );
7291 fit::<u32>(&[4_294_967_296]).expect_err("one past the top does not fit");
7292 fit::<u32>(&[-1]).expect_err("a negative does not fit an unsigned page");
7293
7294 fit::<i8>(&[0, 1, 2, 128, 3]).expect_err("one bad value spoils the page");
7297 }
7298
7299 #[test]
7306 fn the_residue_agrees_with_a_checked_conversion_everywhere() {
7307 for value in -70_000_i64..70_000 {
7308 assert_eq!(fit::<i8>(&[value]).is_ok(), i8::try_from(value).is_ok(), "{value} as i8");
7309 assert_eq!(fit::<u8>(&[value]).is_ok(), u8::try_from(value).is_ok(), "{value} as u8");
7310 assert_eq!(fit::<i16>(&[value]).is_ok(), i16::try_from(value).is_ok(), "{value} i16");
7311 assert_eq!(fit::<u16>(&[value]).is_ok(), u16::try_from(value).is_ok(), "{value} u16");
7312 }
7313 let wide = [i64::MIN, i64::MIN + 1, i64::from(i32::MIN), 0, i64::from(u32::MAX), i64::MAX];
7314 for edge in wide {
7315 for step in -2_i64..=2 {
7316 let value = edge.saturating_add(step);
7317 assert_eq!(
7318 fit::<i32>(&[value]).is_ok(),
7319 i32::try_from(value).is_ok(),
7320 "{value} as i32"
7321 );
7322 assert_eq!(
7323 fit::<u32>(&[value]).is_ok(),
7324 u32::try_from(value).is_ok(),
7325 "{value} as u32"
7326 );
7327 }
7328 }
7329 }
7330
7331 #[test]
7339 fn a_dictionary_at_its_budget_sweeps_without_keeping() {
7340 let path = path("dictionary-budget");
7341 let spellings = (0..2_500)
7342 .map(|index| Value::Varchar(format!("value {index:08} {}", "y".repeat(index % 40))))
7343 .collect::<Vec<_>>();
7344 let mut writer =
7345 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7346 .expect("new file");
7347 for part in spellings.chunks(1_024) {
7348 writer
7349 .append(
7350 &Chunk::new(vec![
7351 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
7352 ])
7353 .expect("one column"),
7354 )
7355 .expect("stripe written");
7356 }
7357 writer.finish().expect("commit");
7358
7359 let reader = Reader::open(&path).expect("valid directory");
7360 let page = reader.table.dictionaries[0].expect("a string column has one");
7361 let file = Arc::clone(&reader.file);
7362 let starved = open_global_dictionary(file, page, &LogicalType::Varchar, 0)
7363 .expect("a dictionary opens whatever it may keep");
7364
7365 let resting = starved.footprint();
7366 let mut swept: Vec<Vec<u8>> = Vec::new();
7367 let mut at = 0;
7368 while at < starved.len() {
7369 at = starved
7370 .sweep_text(at, starved.len(), &mut |_index: usize, text: &[u8]| {
7371 swept.push(text.to_vec());
7372 Ok(())
7373 })
7374 .expect("a sweep reads");
7375 }
7376 assert_eq!(swept.len(), spellings.len(), "a starved sweep still reads every value");
7377 assert_eq!(starved.footprint(), resting, "and keeps no block it decoded");
7378
7379 let generous = reader.dictionary(0).expect("read").expect("a string column has one");
7380 let read = (0..generous.len())
7381 .map(|code| generous.try_bytes_at(code).expect("read").expect("a value").to_vec())
7382 .collect::<Vec<_>>();
7383 assert_eq!(swept, read, "a starved sweep answers what a point read answers");
7384 fs::remove_file(path).expect("remove scratch file");
7385 }
7386
7387 #[test]
7388 fn damaged_membership_cannot_skip_a_string_page() {
7389 let path = path("damaged-membership");
7390 let mut writer = Writer::create(
7391 &path,
7392 "items",
7393 vec![
7394 Field::required("id", LogicalType::Integer),
7395 Field::new("text", LogicalType::Varchar),
7396 ],
7397 )
7398 .expect("new file");
7399 writer.append(&sample()).expect("stripe written");
7400 writer.finish().expect("commit");
7401
7402 let reader = Reader::open(&path).expect("valid directory");
7403 let membership = reader.table.stripes[0].memberships[1].expect("string membership");
7404 let mut file = OpenOptions::new().write(true).open(&path).expect("open membership page");
7405 file.seek(SeekFrom::Start(membership.offset)).expect("membership start");
7406 file.write_all(&[255]).expect("damage membership");
7407 let error = reader.skips_codes(0, 1, &[3]).expect_err("corruption must not skip rows");
7408 assert!(error.message().contains("membership page checksum differs"), "{error}");
7409 fs::remove_file(path).expect("remove scratch file");
7410 }
7411
7412 #[test]
7413 fn membership_delta_stream_is_sorted_exact_and_bounded() {
7414 let unique = unique_codes(&[900, 4, 4, 72, 9, u32::MAX]);
7415 assert_eq!(unique, [4, 9, 72, 900, u32::MAX]);
7416 let encoded = encode_membership(&unique);
7417 assert_eq!(
7418 decode_membership(&encoded).expect("valid membership"),
7419 [4, 9, 72, 900, u32::MAX]
7420 );
7421 let merged = merged_codes(vec![vec![4, 900], vec![9, 900, u32::MAX], vec![72]]);
7424 assert_eq!(merged, [4, 9, 72, 900, u32::MAX]);
7425 assert_eq!(
7426 decode_membership(&encode_membership(&merged)).expect("valid membership"),
7427 unique
7428 );
7429 assert!(decode_membership(&[1, 0x80]).is_err(), "a truncated varint is invalid");
7430 assert!(
7431 decode_membership(&[1, 0xff, 0xff, 0xff, 0xff, 0x10]).is_err(),
7432 "a value past u32 is invalid"
7433 );
7434 }
7435
7436 #[test]
7437 fn a_global_dictionary_may_be_larger_than_one_column_page() {
7438 let dictionary = Page {
7439 offset: HEADER,
7440 length: u32::try_from(MAX_PAGE + 1).expect("the page bound fits on disk"),
7441 hash: 0,
7442 };
7443 let table = Table {
7444 name: "items".to_owned(),
7445 fields: vec![Field::new("text", LogicalType::Varchar)],
7446 stripes: Vec::new(),
7447 rows: 0,
7448 dictionaries: vec![Some(dictionary)],
7449 distincts: vec![None],
7450 frequencies: vec![None],
7451 };
7452 let directory = encode_directory(&table).expect("directory");
7453 let file_size = dictionary.offset + u64::from(dictionary.length) + 1;
7454
7455 let decoded = decode_directory(&directory, file_size).expect("large lazy dictionary");
7456 assert_eq!(decoded.dictionaries[0].expect("dictionary").length, dictionary.length);
7457 }
7458
7459 #[test]
7460 fn a_column_with_one_value_everywhere_costs_almost_nothing_a_row() {
7461 let path = path("constant-codes");
7462 let mut writer =
7463 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
7464 .expect("new file");
7465 let empty = vec![Value::Varchar(String::new()); 1024];
7466 for _ in 0..4 {
7467 let column = Vector::from_values(LogicalType::Varchar, &empty).expect("strings");
7468 writer.append(&Chunk::new(vec![column]).expect("one column")).expect("a part");
7469 }
7470 writer.finish().expect("commit");
7471
7472 let reader = Reader::open(&path).expect("valid directory");
7473 let pages = reader.layout().columns.first().expect("one column").pages;
7474 assert!(pages < 256, "{pages} bytes of pages for 4,096 rows of one value");
7478 let read = reader.read(3, &[0]).expect("the last part back");
7479 assert_eq!(read.value_at(0, 0), Value::Varchar(String::new()));
7480 assert_eq!(read.value_at(1023, 0), Value::Varchar(String::new()));
7481 fs::remove_file(path).expect("remove scratch file");
7482 }
7483
7484 #[test]
7485 fn a_cascade_value_too_wide_for_its_column_is_refused_rather_than_cut() {
7486 let over = vec![i64::from(i32::MAX) + 1];
7489 let error = narrowed(&LogicalType::Integer, over).expect_err("a page that disagrees");
7490 assert!(format!("{error}").contains("not of its type"), "{error}");
7491 assert!(narrowed(&LogicalType::BigInt, vec![i64::MIN]).is_ok(), "bigint holds all of i64");
7492 assert!(narrowed(&LogicalType::Varchar, vec![0]).is_err(), "strings are not integers");
7493 }
7494
7495 #[test]
7496 fn a_code_stream_the_cascade_cannot_shrink_is_left_alone() {
7497 let mut state: u32 = 0x9e37_79b9;
7501 let spread: Vec<u32> = (0..1024)
7502 .map(|_| {
7503 state ^= state << 13;
7504 state ^= state >> 17;
7505 state ^= state << 5;
7506 state
7507 })
7508 .collect();
7509 assert_eq!(encoded_codes(&spread).expect("no failure"), None);
7510 let near: Vec<u32> = (0..1024).collect();
7511 let coded = encoded_codes(&near).expect("no failure").expect("counting up is packable");
7512 assert!(coded.len() < near.len() * 4, "{} bytes for a run of 1,024", coded.len());
7513 }
7514
7515 #[test]
7521 fn two_writes_of_the_same_rows_give_the_same_bytes() {
7522 fn written(path: &PathBuf) {
7523 let fields = (0..40)
7524 .map(|column| {
7525 let ty =
7526 if column % 4 == 0 { LogicalType::Varchar } else { LogicalType::BigInt };
7527 Field::new(format!("c{column}"), ty)
7528 })
7529 .collect::<Vec<_>>();
7530 let mut writer = Writer::create(path, "wide", fields).expect("new file");
7531 for part in 0..70_u64 {
7532 let columns = (0..40)
7533 .map(|column| {
7534 let values = (0..64_u64)
7535 .map(|row| {
7536 let seed = part.wrapping_mul(31).wrapping_add(row);
7537 if column % 4 == 0 {
7538 Value::Varchar(format!("v{}", seed % 17))
7539 } else {
7540 Value::BigInt(i64::try_from(seed % 97).expect("small"))
7541 }
7542 })
7543 .collect::<Vec<_>>();
7544 let ty = if column % 4 == 0 {
7545 LogicalType::Varchar
7546 } else {
7547 LogicalType::BigInt
7548 };
7549 Vector::from_values(ty, &values).expect("a column")
7550 })
7551 .collect::<Vec<_>>();
7552 writer.append(&Chunk::new(columns).expect("forty columns")).expect("a part");
7553 }
7554 writer.finish().expect("commit");
7555 }
7556
7557 let first = path("repeatable-one");
7558 let second = path("repeatable-two");
7559 written(&first);
7560 written(&second);
7561 let left = fs::read(&first).expect("the first file");
7562 let right = fs::read(&second).expect("the second file");
7563 assert_eq!(left.len(), right.len(), "two writes of the same rows differ in length");
7564 assert!(left == right, "two writes of the same rows differ in their bytes");
7565
7566 let reader = Reader::open(&first).expect("valid directory");
7569 assert_eq!(reader.table().rows(), 70 * 64);
7570 let read = reader.read(0, &[0, 1]).expect("the first part back");
7571 assert_eq!(read.value_at(0, 0), Value::Varchar("v0".to_owned()));
7572 assert_eq!(read.value_at(0, 1), Value::BigInt(0));
7573 fs::remove_file(first).expect("remove scratch file");
7574 fs::remove_file(second).expect("remove scratch file");
7575 }
7576
7577 fn three_tables(path: &PathBuf) {
7579 let writer = Writer::create(
7580 path,
7581 "region",
7582 vec![
7583 Field::new("r_key", LogicalType::Integer),
7584 Field::new("r_name", LogicalType::Varchar),
7585 ],
7586 )
7587 .expect("new file");
7588 let mut writer = writer;
7589 writer
7590 .append(
7591 &Chunk::new(vec![
7592 Vector::from_values(
7593 LogicalType::Integer,
7594 &[Value::Integer(0), Value::Integer(1)],
7595 )
7596 .expect("keys"),
7597 Vector::from_values(
7598 LogicalType::Varchar,
7599 &[Value::Varchar("AFRICA".to_owned()), Value::Varchar("ASIA".to_owned())],
7600 )
7601 .expect("names"),
7602 ])
7603 .expect("two columns"),
7604 )
7605 .expect("a part");
7606 let mut writer = writer
7607 .next("empty", vec![Field::new("nothing", LogicalType::BigInt)])
7608 .expect("a second table");
7609 writer
7610 .append(
7611 &Chunk::new(vec![
7612 Vector::from_values(LogicalType::BigInt, &[Value::BigInt(7)]).expect("a row"),
7613 ])
7614 .expect("one column"),
7615 )
7616 .expect("a part");
7617 let mut writer =
7618 writer.next("wide", vec![Field::new("n", LogicalType::BigInt)]).expect("a third table");
7619 for part in 0..70_i64 {
7620 let values = (0..64).map(|row| Value::BigInt(part * 64 + row)).collect::<Vec<_>>();
7621 writer
7622 .append(
7623 &Chunk::new(vec![
7624 Vector::from_values(LogicalType::BigInt, &values).expect("a column"),
7625 ])
7626 .expect("one column"),
7627 )
7628 .expect("a part");
7629 }
7630 writer.finish().expect("commit");
7631 }
7632
7633 #[test]
7634 fn three_tables_in_one_file_read_back_by_name() {
7635 let file = path("three-tables");
7636 three_tables(&file);
7637 let catalog = Catalog::open(&file).expect("a committed catalog");
7638 assert_eq!(catalog.names().collect::<Vec<_>>(), ["region", "empty", "wide"]);
7639
7640 let region = catalog.table("region").expect("the first table");
7641 assert_eq!(region.table().rows(), 2);
7642 assert_eq!(
7643 region.read(0, &[1]).expect("names").value_at(1, 0),
7644 Value::Varchar("ASIA".to_owned())
7645 );
7646
7647 let wide = catalog.table("wide").expect("the third table");
7648 assert_eq!(wide.table().rows(), 70 * 64);
7649 assert_eq!(wide.read(0, &[0]).expect("the first part").value_at(0, 0), Value::BigInt(0));
7650
7651 let empty = catalog.table("empty").expect("the second table");
7654 assert_eq!(empty.table().rows(), 1);
7655 assert_eq!(empty.read(0, &[0]).expect("the row").value_at(0, 0), Value::BigInt(7));
7656
7657 fs::remove_file(file).expect("remove scratch file");
7658 }
7659
7660 #[test]
7661 fn a_name_the_file_does_not_hold_is_an_error_rather_than_the_first_table() {
7662 let file = path("three-tables-missing");
7663 three_tables(&file);
7664 let catalog = Catalog::open(&file).expect("a committed catalog");
7665 let error = catalog.table("nation").expect_err("no such table");
7666 assert!(error.message().contains("nation"), "{}", error.message());
7667 fs::remove_file(file).expect("remove scratch file");
7668 }
7669
7670 #[test]
7671 fn a_file_of_three_tables_will_not_open_as_one() {
7672 let file = path("three-tables-unnamed");
7673 three_tables(&file);
7674 let error = Reader::open(&file).expect_err("more than one table");
7675 assert!(error.message().contains("more than one table"), "{}", error.message());
7676 fs::remove_file(file).expect("remove scratch file");
7677 }
7678
7679 #[test]
7681 fn decimals_of_every_storage_width_round_trip() {
7682 let file = path("decimals");
7683 let widths = [(4_u8, 2_u8), (9, 2), (18, 4), (38, 6)];
7684 let fields = widths
7685 .iter()
7686 .enumerate()
7687 .map(|(index, (width, scale))| {
7688 Field::new(
7689 format!("d{index}"),
7690 LogicalType::decimal(*width, *scale).expect("a decimal type"),
7691 )
7692 })
7693 .collect::<Vec<_>>();
7694 let mut writer = Writer::create(&file, "money", fields).expect("new file");
7695 let rows: [i128; 3] = [-1234, 0, 999];
7696 let columns = widths
7697 .iter()
7698 .map(|(width, scale)| {
7699 let values = rows
7700 .iter()
7701 .map(|unscaled| Value::Decimal {
7702 unscaled: *unscaled,
7703 width: *width,
7704 scale: *scale,
7705 })
7706 .collect::<Vec<_>>();
7707 Vector::from_values(
7708 LogicalType::decimal(*width, *scale).expect("a decimal type"),
7709 &values,
7710 )
7711 .expect("a decimal column")
7712 })
7713 .collect::<Vec<_>>();
7714 writer.append(&Chunk::new(columns).expect("four columns")).expect("a part");
7715 writer.finish().expect("commit");
7716
7717 let reader = Reader::open(&file).expect("a committed file");
7718 for (index, (width, scale)) in widths.iter().enumerate() {
7719 assert_eq!(
7720 reader.table().fields()[index].ty,
7721 LogicalType::decimal(*width, *scale).expect("a decimal type"),
7722 "column {index} came back as another type"
7723 );
7724 let column = reader.read(0, &[index]).expect("the column");
7725 for (row, unscaled) in rows.iter().enumerate() {
7726 assert_eq!(
7727 column.value_at(row, 0),
7728 Value::Decimal { unscaled: *unscaled, width: *width, scale: *scale },
7729 "column {index} row {row}"
7730 );
7731 }
7732 }
7733 fs::remove_file(file).expect("remove scratch file");
7734 }
7735
7736 #[test]
7737 fn two_tables_of_one_name_are_refused_before_anything_is_committed() {
7738 let file = path("two-of-a-name");
7739 let writer = Writer::create(&file, "t", vec![Field::new("a", LogicalType::BigInt)])
7740 .expect("new file");
7741 let error = writer
7742 .next("t", vec![Field::new("a", LogicalType::BigInt)])
7743 .expect_err("the same name twice");
7744 assert!(error.message().contains("same name"), "{}", error.message());
7745 fs::remove_file(file).expect("remove scratch file");
7746 }
7747
7748 #[test]
7749 fn opening_the_catalog_reads_no_table_directory() {
7750 let file = path("catalog-only");
7751 three_tables(&file);
7752 let catalog = Catalog::open(&file).expect("a committed catalog");
7753 assert_eq!(catalog.opening.reads, 2, "opening the catalog read more than the slot");
7756 assert_eq!(catalog.names().len(), 3);
7757 fs::remove_file(file).expect("remove scratch file");
7758 }
7759}