1#![forbid(unsafe_code)]
29
30use std::cmp::Ordering;
31use std::collections::{HashMap, VecDeque};
32use std::fs::{File, OpenOptions};
33use std::io::{Read, Seek, SeekFrom};
34use std::mem::{size_of, size_of_val};
35use std::path::Path;
36use std::sync::atomic::{AtomicUsize, Ordering as Atomic};
37use std::sync::{Arc, Mutex, OnceLock};
38
39use rudb_common::bounds::{Bound, Op, scaled_as};
40use rudb_common::{Error, Field, LogicalType, Result, Value};
41use rudb_encoding::{bitpack, chooser, integer, string};
42use rudb_storage::sieve::Sieve;
43use rudb_storage::{Probe, Range, Zone};
44use rudb_vector::string::StringColumn;
45use rudb_vector::validity::Validity;
46use rudb_vector::{Buffer, Chunk, Data, Packed, TextSource, Vector};
47
48const MAGIC: &[u8; 8] = b"RUDBNV10";
49const DIRECTORY: &[u8; 8] = b"RUDBDI10";
50const FORMAT: u32 = 21;
51const HEADER: u64 = 80;
52const SLOT_BYTES: usize = 28;
53const MAX_PAGE: usize = 256 * 1024 * 1024;
54const MAX_DIRECTORY: usize = 128 * 1024 * 1024;
55const FREQUENCIES: &[u8; 8] = b"RUDBFQ2\0";
56const FREQUENCY_CANDIDATES: usize = 32_768;
57const FREQUENCY_ENTRIES: usize = 512;
58const FREQUENCY_BUILD_RANK: usize = 10;
59const FREQUENCY_ORDINALS: usize = 65_536;
60const MAX_FREQUENCY_WORKERS: usize = 32;
67
68const MAX_ENCODE_WORKERS: usize = 32;
75
76const SIEVE_BUDGET: usize = 8 * 1024;
84
85const PART_BOUND_BYTES: usize = 24;
94
95fn io(error: std::io::Error) -> Error {
96 Error::io(error.to_string())
97}
98
99fn invalid(message: &str) -> Error {
100 Error::invalid_input(format!("invalid rudb native file: {message}"))
101}
102
103fn sum(counts: impl Iterator<Item = u64>) -> u64 {
105 counts.fold(0, u64::saturating_add)
106}
107
108fn span_bytes(spans: &[Span], at: usize) -> u64 {
110 spans.get(at).map_or(0, |span| u64::from(span.length))
111}
112
113fn page_bytes(pages: &[Option<Page>], at: usize) -> u64 {
115 pages.get(at).and_then(Option::as_ref).map_or(0, Page::bytes)
116}
117
118fn checksum(bytes: &[u8]) -> u64 {
119 const P1: u64 = 11_400_714_785_074_694_791;
120 const P2: u64 = 14_029_467_366_897_019_727;
121 const P3: u64 = 1_609_587_929_392_839_161;
122 const P4: u64 = 9_650_029_242_287_828_579;
123 const P5: u64 = 2_870_177_450_012_600_261;
124 let round = |state: u64, word: u64| {
125 state.wrapping_add(word.wrapping_mul(P2)).rotate_left(31).wrapping_mul(P1)
126 };
127 let merge = |state: u64, lane: u64| (state ^ round(0, lane)).wrapping_mul(P1).wrapping_add(P4);
128 let word =
129 |at: usize| u64::from_le_bytes(bytes[at..at + 8].try_into().expect("eight checksum bytes"));
130
131 let mut at = 0;
132 let mut hash = if bytes.len() >= 32 {
133 let mut one = P1.wrapping_add(P2);
134 let mut two = P2;
135 let mut three = 0;
136 let mut four = 0_u64.wrapping_sub(P1);
137 while at + 32 <= bytes.len() {
138 one = round(one, word(at));
139 two = round(two, word(at + 8));
140 three = round(three, word(at + 16));
141 four = round(four, word(at + 24));
142 at += 32;
143 }
144 let combined = one
145 .rotate_left(1)
146 .wrapping_add(two.rotate_left(7))
147 .wrapping_add(three.rotate_left(12))
148 .wrapping_add(four.rotate_left(18));
149 merge(merge(merge(merge(combined, one), two), three), four)
150 } else {
151 P5
152 };
153 hash = hash.wrapping_add(bytes.len() as u64);
154 while at + 8 <= bytes.len() {
155 hash ^= round(0, word(at));
156 hash = hash.rotate_left(27).wrapping_mul(P1).wrapping_add(P4);
157 at += 8;
158 }
159 if at + 4 <= bytes.len() {
160 let tail = u32::from_le_bytes(bytes[at..at + 4].try_into().expect("four checksum bytes"));
161 hash ^= u64::from(tail).wrapping_mul(P1);
162 hash = hash.rotate_left(23).wrapping_mul(P2).wrapping_add(P3);
163 at += 4;
164 }
165 while at < bytes.len() {
166 hash ^= u64::from(bytes[at]).wrapping_mul(P5);
167 hash = hash.rotate_left(11).wrapping_mul(P1);
168 at += 1;
169 }
170 hash ^= hash >> 33;
171 hash = hash.wrapping_mul(P2);
172 hash ^= hash >> 29;
173 hash = hash.wrapping_mul(P3);
174 hash ^ (hash >> 32)
175}
176
177#[derive(Debug, Clone, Copy)]
178struct Slot {
179 offset: u64,
180 length: u32,
181 generation: u64,
182 hash: u64,
183}
184
185impl Slot {
186 fn bytes(self) -> [u8; SLOT_BYTES] {
187 let mut result = [0; SLOT_BYTES];
188 result[..8].copy_from_slice(&self.offset.to_le_bytes());
189 result[8..12].copy_from_slice(&self.length.to_le_bytes());
190 result[12..20].copy_from_slice(&self.generation.to_le_bytes());
191 result[20..28].copy_from_slice(&self.hash.to_le_bytes());
192 result
193 }
194
195 fn read(bytes: &[u8]) -> Self {
196 Self {
197 offset: u64::from_le_bytes(bytes[..8].try_into().expect("eight bytes")),
198 length: u32::from_le_bytes(bytes[8..12].try_into().expect("four bytes")),
199 generation: u64::from_le_bytes(bytes[12..20].try_into().expect("eight bytes")),
200 hash: u64::from_le_bytes(bytes[20..28].try_into().expect("eight bytes")),
201 }
202 }
203}
204
205#[derive(Debug, Clone, Copy)]
206struct Page {
207 offset: u64,
208 length: u32,
209 hash: u64,
210}
211
212impl Page {
213 fn bytes(&self) -> u64 {
215 u64::from(self.length)
216 }
217}
218
219#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
220enum FrequencyValue {
221 Null,
222 Integer(i128),
223 Code(u32),
224}
225
226#[derive(Debug, Clone)]
227struct FrequencyEntry {
228 value: FrequencyValue,
229 count: u64,
230}
231
232#[derive(Debug, Clone)]
237struct FrequencySummary {
238 entries: Vec<FrequencyEntry>,
239 omitted_max: u64,
240 ordinals: Vec<u64>,
241}
242
243#[derive(Debug, Clone, PartialEq, Eq)]
245pub struct FrequencyOccurrences {
246 pub omitted_max: u64,
248 pub ordinals: Vec<u64>,
250}
251
252#[derive(Debug, Clone, Copy, Default)]
259struct Span {
260 offset: u64,
261 length: u32,
262}
263
264#[derive(Debug, Clone)]
266pub struct Stripe {
267 rows: usize,
268 parts: Vec<u32>,
271 index: Span,
275 pages: Vec<Span>,
276 memberships: Vec<Option<Page>>,
277 sieves: Vec<Option<Page>>,
280 part_ranges: Vec<Option<Page>>,
291 zone: Zone,
292}
293
294impl Stripe {
295 #[must_use]
297 pub fn rows(&self) -> usize {
298 self.rows
299 }
300
301 #[must_use]
303 pub fn parts(&self) -> usize {
304 self.parts.len()
305 }
306}
307
308#[derive(Debug, Clone)]
310pub struct Table {
311 name: String,
312 fields: Vec<Field>,
313 stripes: Vec<Stripe>,
314 rows: usize,
315 dictionaries: Vec<Option<Page>>,
316 frequencies: Vec<Option<FrequencySummary>>,
317 distincts: Vec<Option<u64>>,
327}
328
329impl Table {
330 #[must_use]
332 pub fn name(&self) -> &str {
333 &self.name
334 }
335
336 #[must_use]
338 pub fn fields(&self) -> &[Field] {
339 &self.fields
340 }
341
342 #[must_use]
344 pub fn rows(&self) -> usize {
345 self.rows
346 }
347
348 #[must_use]
350 pub fn stripes(&self) -> &[Stripe] {
351 &self.stripes
352 }
353}
354
355#[derive(Debug, Clone)]
357pub struct ColumnLayout {
358 pub name: String,
360 pub kind: String,
362 pub pages: u64,
364 pub memberships: u64,
366 pub sieves: u64,
368 pub part_ranges: u64,
370 pub dictionary: u64,
372}
373
374impl ColumnLayout {
375 #[must_use]
377 pub fn total(&self) -> u64 {
378 self.pages
379 .saturating_add(self.memberships)
380 .saturating_add(self.sieves)
381 .saturating_add(self.part_ranges)
382 .saturating_add(self.dictionary)
383 }
384}
385
386#[derive(Debug, Clone)]
397pub struct Layout {
398 pub file: u64,
400 pub rows: usize,
402 pub stripes: usize,
404 pub parts: usize,
406 pub columns: Vec<ColumnLayout>,
408 pub indexes: u64,
411 pub directory: u64,
413 pub header: u64,
415}
416
417impl Layout {
418 #[must_use]
420 pub fn columns_total(&self) -> u64 {
421 self.columns.iter().map(ColumnLayout::total).fold(0, u64::saturating_add)
422 }
423
424 #[must_use]
430 pub fn unaccounted(&self) -> u64 {
431 self.file
432 .saturating_sub(self.columns_total())
433 .saturating_sub(self.indexes)
434 .saturating_sub(self.directory)
435 .saturating_sub(self.header)
436 }
437}
438
439#[derive(Debug)]
441struct GlobalDictionary {
442 primary: HashMap<u64, u32>,
443 collisions: HashMap<u64, Vec<u32>>,
444 offsets: Vec<u32>,
445 payload: Vec<u8>,
446 counts: Vec<u64>,
447 nulls: u64,
448}
449
450impl GlobalDictionary {
451 fn new() -> Self {
452 Self {
453 primary: HashMap::new(),
454 collisions: HashMap::new(),
455 offsets: vec![0],
456 payload: Vec::new(),
457 counts: Vec::new(),
458 nulls: 0,
459 }
460 }
461
462 fn bytes(&self, code: u32) -> Option<&[u8]> {
463 let start = *self.offsets.get(code as usize)? as usize;
464 let end = *self.offsets.get(code as usize + 1)? as usize;
465 self.payload.get(start..end)
466 }
467
468 fn code(&mut self, text: &str) -> Result<u32> {
469 let hash = checksum(text.as_bytes());
470 if let Some(&code) = self.primary.get(&hash) {
471 if self.bytes(code) == Some(text.as_bytes()) {
472 return Ok(code);
473 }
474 if let Some(codes) = self.collisions.get(&hash) {
475 if let Some(code) =
476 codes.iter().copied().find(|&code| self.bytes(code) == Some(text.as_bytes()))
477 {
478 return Ok(code);
479 }
480 }
481 let code = self.insert(text)?;
482 self.collisions.entry(hash).or_default().push(code);
483 return Ok(code);
484 }
485 let code = self.insert(text)?;
486 self.primary.insert(hash, code);
487 Ok(code)
488 }
489
490 fn insert(&mut self, text: &str) -> Result<u32> {
491 let code = u32::try_from(self.offsets.len() - 1)
492 .map_err(|_| invalid("global dictionary has too many values"))?;
493 self.payload.extend_from_slice(text.as_bytes());
494 self.offsets.push(
495 u32::try_from(self.payload.len())
496 .map_err(|_| invalid("global dictionary payload exceeds 4 GiB"))?,
497 );
498 self.counts.push(0);
499 Ok(code)
500 }
501
502 fn ranked(&self) -> Vec<(u64, u32)> {
522 let count = self.offsets.len() - 1;
523 let mut ranked = (0..count)
524 .map(|code| {
525 let code = code as u32;
526 (head(self.bytes(code).unwrap_or_default()), code)
527 })
528 .collect::<Vec<_>>();
529 ranked.sort_unstable_by(|left, right| {
530 left.0.cmp(&right.0).then_with(|| self.bytes(left.1).cmp(&self.bytes(right.1)))
531 });
532 ranked
533 }
534
535 fn observe(&mut self, code: u32, null: bool) -> Result<()> {
536 if null {
537 self.nulls = self.nulls.saturating_add(1);
538 return Ok(());
539 }
540 let count = self
541 .counts
542 .get_mut(code as usize)
543 .ok_or_else(|| invalid("global dictionary count code is out of range"))?;
544 *count = count.saturating_add(1);
545 Ok(())
546 }
547}
548
549#[derive(Debug)]
551pub struct Writer {
552 file: File,
553 at: u64,
561 table: Table,
562 generation: u64,
563 order: Vec<((u64, u64), (u64, u64))>,
566 next_order: u64,
567 dictionaries: Vec<Option<GlobalDictionary>>,
568 pending: Vec<PendingChunk>,
569}
570
571#[derive(Debug)]
579struct PendingChunk {
580 order: (u64, u64),
581 chunk: Chunk,
582}
583
584#[derive(Debug)]
590struct ColumnStripe {
591 pages: Vec<Vec<u8>>,
592 codes: Vec<Option<Vec<u32>>>,
593 sieves: Vec<Option<Sieve>>,
594 ranges: Vec<Range>,
595}
596
597fn weight(ty: &LogicalType) -> usize {
605 match ty {
606 LogicalType::Varchar | LogicalType::Blob => 64,
607 LogicalType::BigInt
608 | LogicalType::UBigInt
609 | LogicalType::Timestamp
610 | LogicalType::Double
611 | LogicalType::Decimal { .. } => 8,
612 LogicalType::Integer | LogicalType::UInteger | LogicalType::Date | LogicalType::Float => 4,
613 LogicalType::SmallInt | LogicalType::USmallInt => 2,
614 _ => 1,
615 }
616}
617
618pub const STRIPE_PARTS: usize = 64;
625
626const INDEX_ENTRY: usize = size_of::<u32>() + size_of::<u64>();
628
629fn index_section(parts: usize) -> Result<usize> {
631 parts
632 .checked_mul(INDEX_ENTRY)
633 .and_then(|bytes| bytes.checked_add(size_of::<u64>()))
634 .ok_or_else(|| invalid("index page length overflow"))
635}
636
637impl Writer {
638 pub fn create(
644 path: impl AsRef<Path>,
645 name: impl Into<String>,
646 fields: Vec<Field>,
647 ) -> Result<Self> {
648 for field in &fields {
649 type_tag(&field.ty)?;
650 }
651 let file =
652 OpenOptions::new().write(true).read(true).create_new(true).open(path).map_err(io)?;
653 let mut header = [0; HEADER as usize];
654 header[..8].copy_from_slice(MAGIC);
655 header[8..12].copy_from_slice(&FORMAT.to_le_bytes());
656 write_at(&file, 0, &header)?;
657 Ok(Self {
658 file,
659 at: HEADER,
660 dictionaries: fields
661 .iter()
662 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
663 .collect(),
664 table: Table {
665 name: name.into(),
666 dictionaries: vec![None; fields.len()],
667 distincts: vec![None; fields.len()],
668 fields,
669 stripes: Vec::new(),
670 rows: 0,
671 frequencies: Vec::new(),
672 },
673 generation: 1,
674 order: Vec::new(),
675 next_order: 0,
676 pending: Vec::with_capacity(STRIPE_PARTS),
677 })
678 }
679
680 fn put(&mut self, bytes: &[u8]) -> Result<()> {
685 write_at(&self.file, self.at, bytes)?;
686 self.at = self
687 .at
688 .checked_add(bytes.len() as u64)
689 .ok_or_else(|| invalid("native file length overflow"))?;
690 Ok(())
691 }
692
693 pub fn append(&mut self, chunk: &Chunk) -> Result<()> {
699 let order = (self.next_order, 0);
700 self.next_order = self.next_order.saturating_add(1);
701 self.append_at(order, chunk)
702 }
703
704 pub fn append_at(&mut self, order: (u64, u64), chunk: &Chunk) -> Result<()> {
715 if chunk.is_empty() {
716 return Ok(());
717 }
718 self.admit(chunk)?;
719 if self.pending.last().is_some_and(|last| last.order > order) {
720 self.flush_pending()?;
721 }
722 self.pending.push(PendingChunk { order, chunk: chunk.clone() });
727 if self.pending.len() == STRIPE_PARTS {
728 self.flush_pending()?;
729 }
730 Ok(())
731 }
732
733 pub fn append_stripe(&mut self, parts: Vec<((u64, u64), Chunk)>) -> Result<()> {
749 if parts.len() > STRIPE_PARTS {
750 return Err(invalid("a stripe was handed more parts than it holds"));
751 }
752 self.flush_pending()?;
755 for (order, chunk) in parts {
756 if chunk.is_empty() {
757 continue;
758 }
759 self.admit(&chunk)?;
760 self.pending.push(PendingChunk { order, chunk });
761 }
762 self.flush_pending()
763 }
764
765 fn admit(&mut self, chunk: &Chunk) -> Result<()> {
767 if chunk.width() != self.table.fields.len() {
768 return Err(invalid("chunk width differs from table schema"));
769 }
770 for (index, field) in self.table.fields.iter().enumerate() {
771 if chunk.column(index)?.logical_type() != &field.ty {
772 return Err(invalid("chunk type differs from table schema"));
773 }
774 }
775 self.table.rows = self
776 .table
777 .rows
778 .checked_add(chunk.len())
779 .ok_or_else(|| invalid("row count overflow"))?;
780 Ok(())
781 }
782
783 fn encode_column(
791 index: usize,
792 held: &[PendingChunk],
793 mut dictionary: Option<&mut GlobalDictionary>,
794 ) -> Result<ColumnStripe> {
795 let mut stripe = ColumnStripe {
796 pages: Vec::with_capacity(held.len()),
797 codes: Vec::with_capacity(held.len()),
798 sieves: Vec::with_capacity(held.len()),
799 ranges: Vec::with_capacity(held.len()),
800 };
801 for pending in held {
802 let column = pending.chunk.column(index)?;
803 let (bytes, unique) = encode(column, dictionary.as_deref_mut())?;
804 if bytes.len() > MAX_PAGE {
805 return Err(invalid("column page exceeds the configured bound"));
806 }
807 let range = Range::of(column);
810 let sieve = match dictionary {
823 Some(_) => None,
824 None => Sieve::of(column, &range, SIEVE_BUDGET)
825 .filter(|sieve| sieve.len() < bytes.len()),
826 };
827 stripe.pages.push(bytes);
828 stripe.codes.push(unique);
829 stripe.sieves.push(sieve);
830 stripe.ranges.push(range);
831 }
832 Ok(stripe)
833 }
834
835 fn encode_columns(&mut self, held: &[PendingChunk]) -> Result<Vec<ColumnStripe>> {
844 let width = self.table.fields.len();
845 let workers = std::thread::available_parallelism()
846 .map_or(1, usize::from)
847 .min(MAX_ENCODE_WORKERS)
848 .min(width);
849 if workers <= 1 || held.len() <= 1 {
850 return self
851 .dictionaries
852 .iter_mut()
853 .enumerate()
854 .map(|(index, dictionary)| Self::encode_column(index, held, dictionary.as_mut()))
855 .collect();
856 }
857 let mut jobs: Vec<(usize, Option<GlobalDictionary>)> =
860 std::mem::take(&mut self.dictionaries).into_iter().enumerate().collect();
861 jobs.sort_by_key(|(index, _)| weight(&self.table.fields[*index].ty));
863 let queue = Mutex::new(jobs);
864 let pieces = std::thread::scope(|scope| {
865 (0..workers)
866 .map(|_| {
867 scope.spawn(|| {
868 let mut mine = Vec::new();
869 loop {
870 let taken = queue
871 .lock()
872 .map_err(|_| Error::internal("a native encode worker panicked"))?
873 .pop();
874 let Some((index, mut dictionary)) = taken else { break };
875 let encoded = Self::encode_column(index, held, dictionary.as_mut())?;
876 mine.push((index, dictionary, encoded));
877 }
878 Ok(mine)
879 })
880 })
881 .collect::<Vec<_>>()
882 .into_iter()
883 .map(|handle| {
884 handle.join().map_err(|_| Error::internal("a native encode worker panicked"))?
885 })
886 .collect::<Result<Vec<_>>>()
887 })?;
888 let mut dictionaries: Vec<Option<GlobalDictionary>> = (0..width).map(|_| None).collect();
889 let mut encoded: Vec<Option<ColumnStripe>> = (0..width).map(|_| None).collect();
890 for piece in pieces {
891 for (index, dictionary, stripe) in piece {
892 dictionaries[index] = dictionary;
893 encoded[index] = Some(stripe);
894 }
895 }
896 self.dictionaries = dictionaries;
897 encoded
898 .into_iter()
899 .map(|stripe| stripe.ok_or_else(|| Error::internal("a column was never encoded")))
900 .collect()
901 }
902
903 fn flush_pending(&mut self) -> Result<()> {
905 if self.pending.is_empty() {
906 return Ok(());
907 }
908 let width = self.table.fields.len();
909 let mut held = std::mem::take(&mut self.pending);
912 let parts = held.len();
913 let encoded = self.encode_columns(&held)?;
914 let mut pages = Vec::with_capacity(width);
915 let mut memberships = vec![None; width];
916 let mut ranges = Vec::with_capacity(width);
917 let mut index = Vec::with_capacity(width.saturating_mul(index_section(parts)?));
918 for stripe in &encoded {
919 let offset = self.at;
920 let section = index.len();
921 let mut length = 0_usize;
922 for bytes in &stripe.pages {
923 write_at(&self.file, self.at + length as u64, bytes)?;
924 put_u32(
925 &mut index,
926 u32::try_from(bytes.len()).map_err(|_| invalid("part length overflow"))?,
927 );
928 put_u64(&mut index, checksum(bytes));
929 length = length
930 .checked_add(bytes.len())
931 .ok_or_else(|| invalid("column page length overflow"))?;
932 }
933 let hash = checksum(&index[section..]);
934 put_u64(&mut index, hash);
935 if length > MAX_PAGE {
936 return Err(invalid("column page exceeds the configured bound"));
937 }
938 self.at = self
939 .at
940 .checked_add(length as u64)
941 .ok_or_else(|| invalid("native file length overflow"))?;
942 pages.push(Span {
943 offset,
944 length: u32::try_from(length).map_err(|_| invalid("page length overflow"))?,
945 });
946 ranges.push(merged_range(stripe.ranges.iter().cloned()));
947 }
948 for (membership, stripe) in memberships.iter_mut().zip(&encoded) {
949 if stripe.codes.iter().all(Option::is_none) {
950 continue;
951 }
952 let lists = stripe
953 .codes
954 .iter()
955 .map(|codes| codes.clone().unwrap_or_default())
956 .collect::<Vec<_>>();
957 let bytes = encode_membership(&merged_codes(lists));
958 let offset = self.at;
959 self.put(&bytes)?;
960 *membership = Some(Page {
961 offset,
962 length: u32::try_from(bytes.len())
963 .map_err(|_| invalid("membership page length overflow"))?,
964 hash: checksum(&bytes),
965 });
966 }
967 let mut sieves = vec![None; width];
968 for (page, stripe) in sieves.iter_mut().zip(&encoded) {
969 if stripe.sieves.iter().all(Option::is_none) {
970 continue;
971 }
972 let bytes = encode_sieves(stripe.sieves.iter())?;
973 let offset = self.at;
974 self.put(&bytes)?;
975 *page = Some(Page {
976 offset,
977 length: u32::try_from(bytes.len())
978 .map_err(|_| invalid("sieve page length overflow"))?,
979 hash: checksum(&bytes),
980 });
981 }
982 let mut part_ranges = vec![None; width];
988 if parts > 1 {
989 for ((page, stripe), span) in part_ranges.iter_mut().zip(&encoded).zip(&pages) {
990 let bytes = encode_part_ranges(&stripe.ranges)?;
991 if bytes.len() >= span.length as usize {
992 continue;
993 }
994 let offset = self.at;
995 self.put(&bytes)?;
996 *page = Some(Page {
997 offset,
998 length: u32::try_from(bytes.len())
999 .map_err(|_| invalid("part range page length overflow"))?,
1000 hash: checksum(&bytes),
1001 });
1002 }
1003 }
1004 let offset = self.at;
1005 self.put(&index)?;
1006 let index = Span {
1007 offset,
1008 length: u32::try_from(index.len())
1009 .map_err(|_| invalid("index page length overflow"))?,
1010 };
1011 let mut rows = 0_usize;
1012 let mut lengths = Vec::with_capacity(parts);
1013 let mut span = None;
1014 for pending in held.drain(..) {
1015 let part = pending.chunk.len();
1016 rows = rows.checked_add(part).ok_or_else(|| invalid("row count overflow"))?;
1017 lengths.push(u32::try_from(part).map_err(|_| invalid("part row count overflow"))?);
1018 span = Some(
1019 span.map_or((pending.order, pending.order), |(first, _)| (first, pending.order)),
1020 );
1021 }
1022 self.order.push(span.ok_or_else(|| invalid("a stripe was flushed with no parts"))?);
1023 self.table.stripes.push(Stripe {
1024 rows,
1025 parts: lengths,
1026 index,
1027 pages,
1028 memberships,
1029 sieves,
1030 part_ranges,
1031 zone: Zone::from_ranges(ranges),
1032 });
1033 self.pending = held;
1035 Ok(())
1036 }
1037
1038 fn numeric_frequency(&self, column: usize) -> Result<Option<FrequencySummary>> {
1042 let ty = &self.table.fields[column].ty;
1043 if !matches!(
1044 ty,
1045 LogicalType::TinyInt
1046 | LogicalType::SmallInt
1047 | LogicalType::Integer
1048 | LogicalType::BigInt
1049 | LogicalType::UTinyInt
1050 | LogicalType::USmallInt
1051 | LogicalType::UInteger
1052 | LogicalType::UBigInt
1053 | LogicalType::Date
1054 | LogicalType::Timestamp
1055 ) {
1056 return Ok(None);
1057 }
1058 let mut candidates: HashMap<FrequencyValue, u32> = HashMap::new();
1059 let mut decrements = 0_u64;
1060 self.visit_numeric(column, |_, value| {
1061 if let Some(count) = candidates.get_mut(&value) {
1062 *count = count.saturating_add(1);
1063 } else if candidates.len() < FREQUENCY_CANDIDATES {
1064 candidates.insert(value, 1);
1065 } else {
1066 candidates.retain(|_, count| {
1067 *count -= 1;
1068 *count != 0
1069 });
1070 decrements = decrements.saturating_add(1);
1071 }
1072 })?;
1073 let (exact, ordinals) = if decrements == 0 {
1074 (
1075 candidates
1076 .into_iter()
1077 .map(|(value, count)| (value, u64::from(count)))
1078 .collect::<HashMap<_, _>>(),
1079 Vec::new(),
1080 )
1081 } else {
1082 let mut lower = candidates.values().copied().collect::<Vec<_>>();
1083 lower.sort_unstable_by(|left, right| right.cmp(left));
1084 if lower.len() < FREQUENCY_BUILD_RANK
1085 || u64::from(lower[FREQUENCY_BUILD_RANK - 1]) <= decrements
1086 {
1087 return Ok(None);
1088 }
1089 let mut exact =
1090 candidates.into_keys().map(|value| (value, 0_u64)).collect::<HashMap<_, _>>();
1091 let mut ordinals = Vec::new();
1092 let mut exceeded = false;
1093 self.visit_numeric(column, |ordinal, value| {
1094 if let Some(count) = exact.get_mut(&value) {
1095 *count = count.saturating_add(1);
1096 if !exceeded {
1097 if ordinals.len() < FREQUENCY_ORDINALS {
1098 ordinals.push(ordinal);
1099 } else {
1100 ordinals.clear();
1101 exceeded = true;
1102 }
1103 }
1104 }
1105 })?;
1106 (exact, ordinals)
1107 };
1108 let mut entries = exact
1109 .into_iter()
1110 .map(|(value, count)| FrequencyEntry { value, count })
1111 .collect::<Vec<_>>();
1112 entries.sort_unstable_by(|left, right| {
1113 right.count.cmp(&left.count).then_with(|| frequency_order(left.value, right.value))
1114 });
1115 let omitted_max =
1116 entries.get(FREQUENCY_ENTRIES).map_or(decrements, |entry| decrements.max(entry.count));
1117 entries.truncate(FREQUENCY_ENTRIES);
1118 Ok(Some(FrequencySummary { entries, omitted_max, ordinals }))
1119 }
1120
1121 fn visit_numeric(
1122 &self,
1123 column: usize,
1124 mut visit: impl FnMut(u64, FrequencyValue),
1125 ) -> Result<()> {
1126 let ty = &self.table.fields[column].ty;
1127 let mut start = 0_u64;
1128 for stripe in &self.table.stripes {
1129 let spans = read_index(&self.file, stripe, column)?;
1130 let page = stripe.pages[column];
1131 let mut bytes = vec![0; page.length as usize];
1132 read_at(&self.file, page.offset, &mut bytes)?;
1133 for (span, &rows) in spans.iter().zip(&stripe.parts) {
1134 let part = part_bytes(&bytes, *span)?;
1135 if checksum(part) != span.hash {
1136 return Err(invalid("column page checksum differs while building frequencies"));
1137 }
1138 let rows = rows as usize;
1139 let vector = decode(ty, rows, part, None)?;
1140 for row in 0..rows {
1142 let value = if vector.is_null_at(row) {
1143 FrequencyValue::Null
1144 } else {
1145 let widened = match vector.signed_at(row) {
1149 Some(value) => Some(value),
1150 None => match vector.value_at(row) {
1151 Value::UTinyInt(value) => Some(i128::from(value)),
1152 Value::USmallInt(value) => Some(i128::from(value)),
1153 Value::UInteger(value) => Some(i128::from(value)),
1154 Value::UBigInt(value) => Some(i128::from(value)),
1155 _ => None,
1156 },
1157 };
1158 FrequencyValue::Integer(widened.ok_or_else(|| {
1159 invalid("numeric frequency page did not contain an integer value")
1160 })?)
1161 };
1162 visit(start.saturating_add(row as u64), value);
1163 }
1164 start = start.saturating_add(rows as u64);
1165 }
1166 }
1167 Ok(())
1168 }
1169
1170 fn numeric_frequencies(&self) -> Result<Vec<Option<FrequencySummary>>> {
1178 let mut columns = self
1179 .table
1180 .fields
1181 .iter()
1182 .enumerate()
1183 .filter_map(|(column, field)| {
1184 matches!(
1185 field.ty,
1186 LogicalType::TinyInt
1187 | LogicalType::SmallInt
1188 | LogicalType::Integer
1189 | LogicalType::BigInt
1190 | LogicalType::UTinyInt
1191 | LogicalType::USmallInt
1192 | LogicalType::UInteger
1193 | LogicalType::UBigInt
1194 | LogicalType::Date
1195 | LogicalType::Timestamp
1196 )
1197 .then_some(column)
1198 })
1199 .collect::<Vec<_>>();
1200 let workers = std::thread::available_parallelism()
1201 .map_or(1, usize::from)
1202 .min(MAX_FREQUENCY_WORKERS)
1203 .min(columns.len());
1204 if workers <= 1 {
1205 let mut frequencies = vec![None; self.table.fields.len()];
1206 for column in columns {
1207 frequencies[column] = self.numeric_frequency(column)?;
1208 }
1209 return Ok(frequencies);
1210 }
1211 columns.sort_by_key(|&column| weight(&self.table.fields[column].ty));
1214 let queue = Mutex::new(columns);
1215 let pieces = std::thread::scope(|scope| {
1216 (0..workers)
1217 .map(|_| {
1218 scope.spawn(|| {
1219 let mut mine = Vec::new();
1220 loop {
1221 let taken = queue
1222 .lock()
1223 .map_err(|_| Error::internal("a native frequency worker panicked"))?
1224 .pop();
1225 let Some(column) = taken else { break };
1226 mine.push((column, self.numeric_frequency(column)?));
1227 }
1228 Ok(mine)
1229 })
1230 })
1231 .collect::<Vec<_>>()
1232 .into_iter()
1233 .map(|handle| {
1234 handle
1235 .join()
1236 .map_err(|_| Error::internal("a native frequency worker panicked"))?
1237 })
1238 .collect::<Result<Vec<_>>>()
1239 })?;
1240 let mut frequencies = vec![None; self.table.fields.len()];
1241 for piece in pieces {
1242 for (column, summary) in piece {
1243 frequencies[column] = summary;
1244 }
1245 }
1246 Ok(frequencies)
1247 }
1248
1249 pub fn finish(mut self) -> Result<Table> {
1255 self.flush_pending()?;
1256 let mut stripes = std::mem::take(&mut self.order)
1257 .into_iter()
1258 .zip(std::mem::take(&mut self.table.stripes))
1259 .collect::<Vec<_>>();
1260 stripes.sort_by_key(|(order, _)| order.0);
1261 let mut previous: Option<(u64, u64)> = None;
1262 for ((first, last), _) in &stripes {
1263 if previous.is_some_and(|previous| previous >= *first) {
1264 return Err(invalid("chunks did not arrive in source order"));
1265 }
1266 previous = Some(*last);
1267 }
1268 self.table.stripes = stripes.into_iter().map(|(_, stripe)| stripe).collect();
1269 self.table.frequencies = self.numeric_frequencies()?;
1270 let dictionaries = std::mem::take(&mut self.dictionaries);
1271 let orders = rankings(&dictionaries)?;
1272 for (index, (dictionary, order)) in dictionaries.into_iter().zip(orders).enumerate() {
1273 let Some(dictionary) = dictionary else { continue };
1274 self.table.distincts[index] =
1278 Some(dictionary.counts.iter().filter(|count| **count != 0).count() as u64);
1279 self.table.frequencies[index] = Some(code_frequency(&dictionary));
1280 let encoded = encode_global_dictionary(dictionary, &order)?;
1281 let offset = self.at;
1282 self.put(&encoded.index)?;
1283 self.put(&encoded.ranks)?;
1284 for block in &encoded.payload {
1285 self.put(block)?;
1286 }
1287 let payload_len =
1288 encoded.payload.iter().try_fold(0_usize, |len, block| len.checked_add(block.len()));
1289 let length = payload_len
1290 .and_then(|len| len.checked_add(encoded.index.len()))
1291 .and_then(|len| len.checked_add(encoded.ranks.len()))
1292 .ok_or_else(|| invalid("dictionary page length overflow"))?;
1293 self.table.dictionaries[index] = Some(Page {
1294 offset,
1295 length: u32::try_from(length)
1296 .map_err(|_| invalid("dictionary page length overflow"))?,
1297 hash: checksum(&encoded.index),
1298 });
1299 }
1300 let directory = encode_directory(&self.table)?;
1301 if directory.len() > MAX_DIRECTORY {
1302 return Err(invalid("directory exceeds the configured bound"));
1303 }
1304 let offset = self.at;
1305 self.put(&directory)?;
1306 self.file.sync_all().map_err(io)?;
1307 let slot = Slot {
1308 offset,
1309 length: u32::try_from(directory.len())
1310 .map_err(|_| invalid("directory length overflow"))?,
1311 generation: self.generation,
1312 hash: checksum(&directory),
1313 };
1314 write_at(&self.file, 16, &slot.bytes())?;
1317 self.file.sync_all().map_err(io)?;
1318 Ok(self.table)
1319 }
1320}
1321
1322#[derive(Debug, Clone)]
1324pub struct Reader {
1325 file: Arc<File>,
1326 table: Arc<Table>,
1327 dictionaries: Arc<Vec<OnceLock<Arc<Vector>>>>,
1328 loading: Arc<Vec<Mutex<()>>>,
1337 opened: Arc<AtomicUsize>,
1341 sieves: Arc<Vec<Vec<SieveSlot>>>,
1345 part_ranges: Arc<Vec<Vec<RangeSlot>>>,
1348 places: Arc<Vec<Place>>,
1350 cache: Arc<Vec<Mutex<Cached>>>,
1351 pages: Arc<AtomicUsize>,
1354 indexes: Arc<AtomicUsize>,
1357 kept: Arc<AtomicUsize>,
1360 size: u64,
1362 directory: u64,
1364 opening: Opening,
1366}
1367
1368#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1380pub struct Opening {
1381 pub reads: u32,
1384 pub bytes: u64,
1386}
1387
1388#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1390pub struct Reads {
1391 pub opening: Opening,
1393 pub pages: usize,
1395 pub indexes: usize,
1397 pub dictionaries: usize,
1400}
1401
1402#[derive(Debug, Clone, Copy)]
1404struct Place {
1405 stripe: u32,
1406 part: u32,
1407 rows: u32,
1408}
1409
1410#[derive(Debug, Clone, Copy)]
1412struct PartSpan {
1413 start: usize,
1414 length: usize,
1415 hash: u64,
1416}
1417
1418#[derive(Debug, Clone)]
1424struct CachedColumn {
1425 stripe: usize,
1426 index: Arc<Vec<PartSpan>>,
1427 page: Option<Arc<Vec<u8>>>,
1428}
1429
1430#[derive(Debug, Default)]
1450struct Cached {
1451 pages: Vec<Option<Arc<Vec<u8>>>>,
1452 order: VecDeque<usize>,
1453 loading: Vec<usize>,
1454 index: Vec<Option<Arc<Vec<PartSpan>>>>,
1455}
1456
1457const CACHED_STRIPES_PER_COLUMN: usize = 4;
1469
1470type SieveSlot = OnceLock<Arc<Vec<Option<Sieve>>>>;
1472
1473type RangeSlot = OnceLock<Arc<Vec<Range>>>;
1474
1475#[derive(Debug)]
1476struct NativeText {
1477 file: Arc<File>,
1478 values: usize,
1480 offsets: Vec<u8>,
1489 offset_bits: usize,
1492 ranks: usize,
1494 rank_at: u64,
1498 rank_ends: Vec<u64>,
1502 rank_hashes: Vec<u64>,
1503 rank_blocks: Vec<OnceLock<Result<Vec<u8>>>>,
1504 code_bits: usize,
1507 code_ranks: OnceLock<Option<Vec<u32>>>,
1514 payload: u64,
1515 ends: Vec<u64>,
1518 hashes: Vec<u64>,
1519 blocks: Vec<OnceLock<Result<Vec<u8>>>>,
1521 keep_budget: usize,
1524 payload_kept: AtomicUsize,
1532}
1533
1534const TEXT_PAYLOAD_VALUES: usize = 1024;
1550
1551const TEXT_KEEP_BUDGET: usize = 256 * 1024 * 1024;
1572
1573const TEXT_OFFSET_RUN: usize = 512;
1580
1581const DICTIONARY_HEADER: usize = 16;
1584
1585const TEXT_RANK_BLOCK: usize = 512;
1596
1597const RANK_BLOCK_HEADER: usize = size_of::<u64>() + 1;
1611
1612impl NativeText {
1613 fn payload_block(&self, block: usize) -> Result<Option<&[u8]>> {
1620 let Some(slot) = self.blocks.get(block) else { return Ok(None) };
1621 let bytes = slot.get_or_init(|| self.decode_block(block)).as_ref().map_err(Clone::clone)?;
1622 Ok(Some(bytes.as_slice()))
1623 }
1624
1625 fn decode_block(&self, block: usize) -> Result<Vec<u8>> {
1630 let start = if block == 0 { 0 } else { self.ends[block - 1] };
1631 let end = self.ends[block];
1632 let len = end
1633 .checked_sub(start)
1634 .ok_or_else(|| invalid("global dictionary block ends before it starts"))?;
1635 let mut stored = vec![
1636 0;
1637 usize::try_from(len).map_err(|_| invalid(
1638 "global dictionary block does not fit in memory"
1639 ))?
1640 ];
1641 read_at(&self.file, self.payload + start, &mut stored)?;
1642 if checksum(&stored) != self.hashes[block] {
1643 return Err(invalid("global dictionary payload checksum differs"));
1644 }
1645 let first = block * TEXT_PAYLOAD_VALUES;
1646 let last = (first + TEXT_PAYLOAD_VALUES).min(self.values);
1647 let want = self.end_within(last - 1)? as usize;
1648 let values = string::decode_flat(&stored)?;
1649 if values.len() != last - first {
1650 return Err(invalid("global dictionary block holds the wrong value count"));
1651 }
1652 let bytes = values.into_bytes();
1653 if bytes.len() != want {
1654 return Err(invalid("global dictionary block decodes to the wrong length"));
1655 }
1656 Ok(bytes)
1657 }
1658
1659 fn end_within(&self, index: usize) -> Result<u32> {
1661 let run = index / TEXT_OFFSET_RUN;
1662 let bytes = self
1663 .offsets
1664 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
1665 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1666 let end = bitpack::tail_at(bytes, self.offset_bits, index % TEXT_OFFSET_RUN)
1667 .map_err(|_| invalid("global dictionary offsets are short"))?;
1668 u32::try_from(end).map_err(|_| invalid("global dictionary offset is past the payload"))
1669 }
1670
1671 fn ends_within(&self, first: usize, last: usize) -> Result<Vec<u64>> {
1684 let mut ends = Vec::with_capacity(last.saturating_sub(first));
1685 let mut at = first;
1686 while at < last {
1687 let run = at / TEXT_OFFSET_RUN;
1688 let stop = ((run + 1) * TEXT_OFFSET_RUN).min(last);
1689 let held = self.values.saturating_sub(run * TEXT_OFFSET_RUN).min(TEXT_OFFSET_RUN);
1690 let bytes = self
1691 .offsets
1692 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
1693 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1694 let run_ends = bitpack::unpack_tail(bytes, self.offset_bits, held)
1695 .map_err(|_| invalid("global dictionary offsets are short"))?;
1696 let within = run_ends
1697 .get(at % TEXT_OFFSET_RUN..stop - run * TEXT_OFFSET_RUN)
1698 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1699 ends.extend_from_slice(within);
1700 at = stop;
1701 }
1702 Ok(ends)
1703 }
1704
1705 fn start_within(&self, index: usize) -> Result<u32> {
1708 if index % TEXT_PAYLOAD_VALUES == 0 { Ok(0) } else { self.end_within(index - 1) }
1709 }
1710
1711 fn span_within(&self, index: usize) -> Result<(u32, u32)> {
1713 let end = self.end_within(index)?;
1714 let start = self.start_within(index)?;
1715 if start > end {
1716 return Err(invalid("global dictionary value ends before it starts"));
1717 }
1718 Ok((start, end))
1719 }
1720
1721 fn rank_parts(&self, rank: usize) -> Result<(&[u8], usize)> {
1728 let slot = self
1729 .rank_blocks
1730 .get(rank / TEXT_RANK_BLOCK)
1731 .ok_or_else(|| invalid("global dictionary rank is past the order"))?;
1732 let block = slot
1733 .get_or_init(|| {
1734 let which = rank / TEXT_RANK_BLOCK;
1735 let start = if which == 0 { 0 } else { self.rank_ends[which - 1] };
1736 let end = self.rank_ends[which];
1737 let mut bytes = vec![0; (end - start) as usize];
1738 read_at(&self.file, self.rank_at + start, &mut bytes)?;
1739 if checksum(&bytes)
1740 != *self
1741 .rank_hashes
1742 .get(rank / TEXT_RANK_BLOCK)
1743 .ok_or_else(|| invalid("global dictionary rank block has no checksum"))?
1744 {
1745 return Err(invalid("global dictionary rank checksum differs"));
1746 }
1747 Ok(bytes)
1748 })
1749 .as_ref()
1750 .map_err(Clone::clone)?;
1751 Ok((block.as_slice(), rank % TEXT_RANK_BLOCK))
1752 }
1753
1754 fn head_at(&self, rank: usize) -> Result<u64> {
1756 let (block, within) = self.rank_parts(rank)?;
1757 let (base, width, packed) = rank_heads(block)?;
1758 let above = bitpack::tail_at(packed, width, within)
1759 .map_err(|_| invalid("global dictionary rank block is short of heads"))?;
1760 Ok(base.wrapping_add(above))
1761 }
1762
1763 fn rank_codes<'block>(&self, block: &'block [u8], count: usize) -> Result<&'block [u8]> {
1765 let (_, width, packed) = rank_heads(block)?;
1766 packed
1767 .get(bitpack::tail_len(count, width)..)
1768 .ok_or_else(|| invalid("global dictionary rank block is short of codes"))
1769 }
1770
1771 fn rank_block_len(&self, rank: usize) -> usize {
1773 let first = rank / TEXT_RANK_BLOCK * TEXT_RANK_BLOCK;
1774 TEXT_RANK_BLOCK.min(self.ranks - first)
1775 }
1776}
1777
1778fn rank_heads(block: &[u8]) -> Result<(u64, usize, &[u8])> {
1780 let header = block
1781 .get(..RANK_BLOCK_HEADER)
1782 .ok_or_else(|| invalid("global dictionary rank block is short"))?;
1783 let base = u64::from_le_bytes(header[..8].try_into().expect("eight bytes"));
1784 let width = header[8] as usize;
1785 if width > 64 {
1786 return Err(invalid("global dictionary rank block packs heads past a word"));
1787 }
1788 Ok((base, width, &block[RANK_BLOCK_HEADER..]))
1789}
1790
1791fn offset_width(offsets: &[u32]) -> usize {
1798 let values = offsets.len() - 1;
1799 let mut span = 0;
1800 for first in (0..values).step_by(TEXT_PAYLOAD_VALUES) {
1801 let last = (first + TEXT_PAYLOAD_VALUES).min(values);
1802 span = span.max(offsets[last] - offsets[first]);
1803 }
1804 (u32::BITS - span.leading_zeros()) as usize
1805}
1806
1807fn offset_bytes(values: usize, bits: usize) -> usize {
1810 let full = values / TEXT_OFFSET_RUN;
1811 let rest = values % TEXT_OFFSET_RUN;
1812 full * TEXT_OFFSET_RUN / 8 * bits + bitpack::tail_len(rest, bits)
1813}
1814
1815fn encode_offsets(offsets: &[u32], bits: usize, out: &mut Vec<u8>) -> Result<()> {
1817 let values = offsets.len() - 1;
1818 let mut run = Vec::with_capacity(TEXT_OFFSET_RUN);
1819 for first in (0..values).step_by(TEXT_OFFSET_RUN) {
1820 let last = (first + TEXT_OFFSET_RUN).min(values);
1821 let base = offsets[first / TEXT_PAYLOAD_VALUES * TEXT_PAYLOAD_VALUES];
1822 run.clear();
1823 run.extend((first..last).map(|value| u64::from(offsets[value + 1] - base)));
1824 bitpack::pack_tail(&run, bits, out)
1825 .map_err(|_| invalid("global dictionary offsets do not pack"))?;
1826 }
1827 Ok(())
1828}
1829
1830fn code_width(values: usize) -> usize {
1832 match u64::try_from(values).unwrap_or(u64::MAX) {
1833 0 | 1 => 0,
1834 last => (u64::BITS - (last - 1).leading_zeros()) as usize,
1835 }
1836}
1837
1838impl TextSource for NativeText {
1839 fn len(&self) -> usize {
1840 self.values
1841 }
1842
1843 fn bytes_at(&self, index: usize) -> Result<Option<&[u8]>> {
1844 if index >= self.values {
1845 return Ok(None);
1846 }
1847 let (start, end) = self.span_within(index)?;
1848 if start == end {
1849 return Ok(Some(&[]));
1850 }
1851 let block = index / TEXT_PAYLOAD_VALUES;
1854 let Some(bytes) = self.payload_block(block)? else { return Ok(None) };
1855 Ok(bytes.get(start as usize..end as usize))
1856 }
1857
1858 fn bytes_len_at(&self, index: usize) -> Result<Option<usize>> {
1859 if index >= self.values {
1860 return Ok(None);
1861 }
1862 let (start, end) = self.span_within(index)?;
1863 Ok(Some((end - start) as usize))
1864 }
1865
1866 fn sweep(
1879 &self,
1880 first: usize,
1881 limit: usize,
1882 body: &mut dyn FnMut(usize, &[u8]) -> Result<()>,
1883 ) -> Result<usize> {
1884 let limit = limit.min(self.values);
1885 if first >= limit {
1886 return Ok(first);
1887 }
1888 let block = first / TEXT_PAYLOAD_VALUES;
1889 let last = ((block + 1) * TEXT_PAYLOAD_VALUES).min(limit);
1890 let decoded;
1891 let bytes: &[u8] = match self.blocks.get(block).and_then(OnceLock::get) {
1892 Some(Ok(kept)) => kept,
1893 _ if self.payload_kept.load(Atomic::Relaxed) < self.keep_budget => {
1894 let kept = self
1895 .payload_block(block)?
1896 .ok_or_else(|| invalid("global dictionary block is past the payload"))?;
1897 self.payload_kept.fetch_add(kept.len(), Atomic::Relaxed);
1898 kept
1899 }
1900 _ => {
1901 decoded = self.decode_block(block)?;
1902 &decoded
1903 }
1904 };
1905 let ends = self.ends_within(first, last)?;
1906 if ends.len() != last - first {
1907 return Err(invalid("global dictionary offsets are short"));
1908 }
1909 let mut start = u64::from(self.start_within(first)?);
1910 for (index, &end) in (first..last).zip(&ends) {
1913 let value = usize::try_from(start)
1914 .ok()
1915 .zip(usize::try_from(end).ok())
1916 .and_then(|(from, to)| bytes.get(from..to))
1917 .ok_or_else(|| invalid("global dictionary value is past its block"))?;
1918 body(index, value)?;
1919 start = end;
1920 }
1921 Ok(last)
1922 }
1923
1924 fn ranks(&self) -> Option<usize> {
1925 (self.ranks > 0).then_some(self.ranks)
1926 }
1927
1928 fn compare_rank(&self, rank: usize, wanted: &[u8]) -> Result<Ordering> {
1929 let settled = self.head_at(rank)?.cmp(&head(wanted));
1933 if settled != Ordering::Equal {
1934 return Ok(settled);
1935 }
1936 let code = self.code_at_rank(rank)?;
1937 let bytes = self
1938 .bytes_at(code as usize)?
1939 .ok_or_else(|| invalid("global dictionary order names a code it does not have"))?;
1940 Ok(bytes.cmp(wanted))
1941 }
1942
1943 fn code_at_rank(&self, rank: usize) -> Result<u32> {
1944 let (block, within) = self.rank_parts(rank)?;
1945 let codes = self.rank_codes(block, self.rank_block_len(rank))?;
1946 let code = bitpack::tail_at(codes, self.code_bits, within)
1947 .map_err(|_| invalid("global dictionary rank block is short of codes"))?;
1948 let code = u32::try_from(code)
1949 .map_err(|_| invalid("global dictionary order names a code it does not have"))?;
1950 if code as usize >= self.len() {
1951 return Err(invalid("global dictionary order names a code it does not have"));
1952 }
1953 Ok(code)
1954 }
1955
1956 fn code_ranks(&self) -> Option<&[u32]> {
1957 if self.ranks == 0 || self.ranks != self.len() {
1961 return None;
1962 }
1963 self.code_ranks
1964 .get_or_init(|| {
1965 let mut ranks = vec![u32::MAX; self.ranks];
1966 for first in (0..self.ranks).step_by(TEXT_RANK_BLOCK) {
1969 let (block, _) = self.rank_parts(first).ok()?;
1970 let count = self.rank_block_len(first);
1971 let codes = self.rank_codes(block, count).ok()?;
1972 for (within, code) in bitpack::unpack_tail(codes, self.code_bits, count)
1973 .ok()?
1974 .into_iter()
1975 .enumerate()
1976 {
1977 let code = usize::try_from(code).ok()?;
1978 *ranks.get_mut(code)? = u32::try_from(first + within).ok()?;
1979 }
1980 }
1981 if ranks.contains(&u32::MAX) {
1982 return None;
1983 }
1984 Some(ranks)
1985 })
1986 .as_deref()
1987 }
1988
1989 fn footprint(&self) -> usize {
1990 self.offsets.capacity()
1991 + self
1992 .code_ranks
1993 .get()
1994 .and_then(Option::as_ref)
1995 .map_or(0, |ranks| ranks.capacity() * size_of::<u32>())
1996 + self.rank_hashes.capacity() * size_of::<u64>()
1997 + self.rank_ends.capacity() * size_of::<u64>()
1998 + self.rank_blocks.capacity() * size_of::<OnceLock<Result<Vec<u8>>>>()
1999 + self
2000 .rank_blocks
2001 .iter()
2002 .filter_map(OnceLock::get)
2003 .filter_map(|result| result.as_ref().ok())
2004 .map(Vec::capacity)
2005 .sum::<usize>()
2006 + self.blocks.capacity() * size_of::<OnceLock<Result<Vec<u8>>>>()
2007 + self.hashes.capacity() * size_of::<u64>()
2008 + self.ends.capacity() * size_of::<u64>()
2009 + self
2010 .blocks
2011 .iter()
2012 .filter_map(OnceLock::get)
2013 .filter_map(|result| result.as_ref().ok())
2014 .map(Vec::capacity)
2015 .sum::<usize>()
2016 }
2017}
2018
2019fn places(table: &Table) -> Result<Vec<Place>> {
2021 let mut places = Vec::with_capacity(table.stripes.len().saturating_mul(STRIPE_PARTS));
2022 for (at, stripe) in table.stripes.iter().enumerate() {
2023 let index = u32::try_from(at).map_err(|_| invalid("too many stripes"))?;
2024 for (part, &rows) in stripe.parts.iter().enumerate() {
2025 places.push(Place {
2026 stripe: index,
2027 part: u32::try_from(part).map_err(|_| invalid("too many parts in a stripe"))?,
2028 rows,
2029 });
2030 }
2031 }
2032 Ok(places)
2033}
2034
2035fn read_index(file: &File, stripe: &Stripe, column: usize) -> Result<Vec<PartSpan>> {
2040 let parts = stripe.parts.len();
2041 let section = index_section(parts)?;
2042 let at = column.checked_mul(section).ok_or_else(|| invalid("index page offset overflow"))?;
2043 let end = at.checked_add(section).ok_or_else(|| invalid("index page offset overflow"))?;
2044 if end > stripe.index.length as usize {
2045 return Err(invalid("index page is shorter than its columns"));
2046 }
2047 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
2048 let mut bytes = vec![0; section];
2049 let offset = stripe
2050 .index
2051 .offset
2052 .checked_add(at as u64)
2053 .ok_or_else(|| invalid("index page offset overflow"))?;
2054 read_at(file, offset, &mut bytes)?;
2055 let entries = section - size_of::<u64>();
2056 let stored = u64::from_le_bytes(bytes[entries..].try_into().expect("eight bytes"));
2057 if checksum(&bytes[..entries]) != stored {
2058 return Err(invalid(&format!(
2061 "index page section checksum differs, column {column} of {parts} parts at {offset}, \
2062 wanted {stored:016x} and got {:016x}",
2063 checksum(&bytes[..entries]),
2064 )));
2065 }
2066 let mut spans = Vec::with_capacity(parts);
2067 let mut start = 0_usize;
2068 for part in 0..parts {
2069 let at = part * INDEX_ENTRY;
2070 let length = u32::from_le_bytes(bytes[at..at + 4].try_into().expect("four bytes")) as usize;
2071 let hash = u64::from_le_bytes(bytes[at + 4..at + 12].try_into().expect("eight bytes"));
2072 spans.push(PartSpan { start, length, hash });
2073 start = start.checked_add(length).ok_or_else(|| invalid("column page length overflow"))?;
2074 }
2075 if start != page.length as usize {
2076 return Err(invalid("column page length differs from its index"));
2077 }
2078 Ok(spans)
2079}
2080
2081fn part_bytes(page: &[u8], span: PartSpan) -> Result<&[u8]> {
2083 let end = span.start.checked_add(span.length).ok_or_else(|| invalid("part range overflow"))?;
2084 page.get(span.start..end).ok_or_else(|| invalid("part exceeds its column page"))
2085}
2086
2087fn remember(cached: &mut Cached, held: &CachedColumn, kept: usize) {
2092 if let Some(slot) = cached.index.get_mut(held.stripe) {
2093 if slot.is_none() {
2094 *slot = Some(Arc::clone(&held.index));
2095 }
2096 }
2097 let Some(page) = held.page.clone() else { return };
2098 let Some(slot) = cached.pages.get_mut(held.stripe) else { return };
2099 if slot.is_none() {
2100 cached.order.push_back(held.stripe);
2101 }
2102 *slot = Some(page);
2103 while cached.order.len() > kept.max(1) {
2104 let Some(oldest) = cached.order.pop_front() else { break };
2105 if let Some(slot) = cached.pages.get_mut(oldest) {
2106 *slot = None;
2107 }
2108 }
2109}
2110
2111impl Reader {
2112 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
2118 let mut file = File::open(path).map_err(io)?;
2119 let size = file.metadata().map_err(io)?.len();
2120 if size < HEADER {
2121 return Err(invalid("file is shorter than its header"));
2122 }
2123 let mut header = [0; HEADER as usize];
2124 file.read_exact(&mut header).map_err(io)?;
2125 let mut opening = Opening { reads: 1, bytes: HEADER };
2126 let version = u32::from_le_bytes([header[8], header[9], header[10], header[11]]);
2127 if &header[..8] != MAGIC {
2132 return Err(invalid("the header does not begin with a rudb native magic"));
2133 }
2134 if version != FORMAT {
2135 return Err(invalid(&format!(
2136 "the file is format {version} and this build reads format {FORMAT}, so it has to \
2137 be written again"
2138 )));
2139 }
2140 let mut selected = None;
2141 for start in [16, 16 + SLOT_BYTES] {
2142 let slot = Slot::read(&header[start..start + SLOT_BYTES]);
2143 if slot.generation == 0 || slot.length == 0 || slot.length as usize > MAX_DIRECTORY {
2144 continue;
2145 }
2146 let Some(end) = slot.offset.checked_add(u64::from(slot.length)) else { continue };
2147 if slot.offset < HEADER || end > size {
2148 continue;
2149 }
2150 let mut bytes = vec![0; slot.length as usize];
2151 file.seek(SeekFrom::Start(slot.offset)).map_err(io)?;
2152 file.read_exact(&mut bytes).map_err(io)?;
2153 opening.reads += 1;
2154 opening.bytes += u64::from(slot.length);
2155 if checksum(&bytes) == slot.hash
2156 && selected
2157 .as_ref()
2158 .is_none_or(|(old, _): &(Slot, Vec<u8>)| old.generation < slot.generation)
2159 {
2160 selected = Some((slot, bytes));
2161 }
2162 }
2163 let (slot, bytes) =
2164 selected.ok_or_else(|| invalid("no committed directory slot is valid"))?;
2165 let table = decode_directory(&bytes, size)?;
2166 let places = places(&table)?;
2167 let dictionaries = (0..table.fields.len()).map(|_| OnceLock::new()).collect();
2168 let table_fields = table.fields.len();
2169 let stripes = table.stripes.len();
2170 let cache = (0..table.fields.len())
2171 .map(|_| {
2172 Mutex::new(Cached {
2173 pages: (0..stripes).map(|_| None).collect(),
2174 index: (0..stripes).map(|_| None).collect(),
2175 ..Cached::default()
2176 })
2177 })
2178 .collect::<Vec<_>>();
2179 let sieves: Vec<Vec<SieveSlot>> = (0..table.fields.len())
2180 .map(|_| table.stripes.iter().map(|_| OnceLock::new()).collect())
2181 .collect();
2182 let part_ranges: Vec<Vec<RangeSlot>> = (0..table.fields.len())
2183 .map(|_| table.stripes.iter().map(|_| OnceLock::new()).collect())
2184 .collect();
2185 Ok(Self {
2186 file: Arc::new(file),
2187 table: Arc::new(table),
2188 dictionaries: Arc::new(dictionaries),
2189 loading: Arc::new((0..table_fields).map(|_| Mutex::new(())).collect()),
2190 opened: Arc::new(AtomicUsize::new(0)),
2191 sieves: Arc::new(sieves),
2192 part_ranges: Arc::new(part_ranges),
2193 places: Arc::new(places),
2194 cache: Arc::new(cache),
2195 pages: Arc::new(AtomicUsize::new(0)),
2196 indexes: Arc::new(AtomicUsize::new(0)),
2197 kept: Arc::new(AtomicUsize::new(CACHED_STRIPES_PER_COLUMN)),
2198 size,
2199 directory: u64::from(slot.length),
2200 opening,
2201 })
2202 }
2203
2204 #[must_use]
2211 pub fn reads(&self) -> Reads {
2212 Reads {
2213 opening: self.opening,
2214 pages: self.pages.load(Atomic::Relaxed),
2215 indexes: self.indexes.load(Atomic::Relaxed),
2216 dictionaries: self.opened.load(Atomic::Relaxed),
2217 }
2218 }
2219
2220 #[must_use]
2225 pub fn layout(&self) -> Layout {
2226 let table = &self.table;
2227 let stripes = table.stripes.as_slice();
2228 let columns = table
2229 .fields
2230 .iter()
2231 .enumerate()
2232 .map(|(at, field)| ColumnLayout {
2233 name: field.name.clone(),
2234 kind: field.ty.to_string(),
2235 pages: sum(stripes.iter().map(|stripe| span_bytes(&stripe.pages, at))),
2236 memberships: sum(stripes.iter().map(|stripe| page_bytes(&stripe.memberships, at))),
2237 sieves: sum(stripes.iter().map(|stripe| page_bytes(&stripe.sieves, at))),
2238 part_ranges: sum(stripes.iter().map(|stripe| page_bytes(&stripe.part_ranges, at))),
2239 dictionary: page_bytes(&table.dictionaries, at),
2240 })
2241 .collect();
2242 Layout {
2243 file: self.size,
2244 rows: table.rows,
2245 stripes: stripes.len(),
2246 parts: self.places.len(),
2247 columns,
2248 indexes: sum(stripes.iter().map(|stripe| u64::from(stripe.index.length))),
2249 directory: self.directory,
2250 header: HEADER,
2251 }
2252 }
2253
2254 #[must_use]
2256 pub fn parts(&self) -> usize {
2257 self.places.len()
2258 }
2259
2260 #[must_use]
2267 pub fn stripe_parts(&self) -> Vec<std::ops::Range<usize>> {
2268 let mut runs = Vec::with_capacity(self.table.stripes.len());
2269 let mut start = 0;
2270 for stripe in &self.table.stripes {
2271 let end = start + stripe.parts.len();
2272 runs.push(start..end);
2273 start = end;
2274 }
2275 runs
2276 }
2277
2278 pub fn keep_stripes(&self, stripes: usize) {
2285 self.kept.fetch_max(stripes, Atomic::Relaxed);
2286 }
2287
2288 #[must_use]
2290 pub fn part_rows(&self, at: usize) -> usize {
2291 self.places.get(at).map_or(0, |place| place.rows as usize)
2292 }
2293
2294 #[must_use]
2296 pub fn table(&self) -> &Table {
2297 &self.table
2298 }
2299
2300 pub fn top_frequencies(&self, column: usize, top: usize) -> Result<Option<Vec<(Value, u64)>>> {
2309 let field = self
2310 .table
2311 .fields
2312 .get(column)
2313 .ok_or_else(|| invalid("frequency column index out of range"))?;
2314 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2315 return Ok(None);
2316 };
2317 if top == 0 || summary.entries.len() < top {
2318 return Ok(None);
2319 }
2320 let boundary = summary.entries[top - 1].count;
2321 if boundary <= summary.omitted_max {
2322 return Ok(None);
2323 }
2324 self.decode_frequencies(column, &field.ty, &summary.entries).map(Some)
2325 }
2326
2327 pub fn exact_frequencies(&self, column: usize) -> Result<Option<Vec<(Value, u64)>>> {
2347 let field = self
2348 .table
2349 .fields
2350 .get(column)
2351 .ok_or_else(|| invalid("frequency column index out of range"))?;
2352 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2353 return Ok(None);
2354 };
2355 if summary.omitted_max > 0 {
2356 return Ok(None);
2357 }
2358 self.decode_frequencies(column, &field.ty, &summary.entries).map(Some)
2359 }
2360
2361 fn decode_frequencies(
2363 &self,
2364 column: usize,
2365 ty: &LogicalType,
2366 entries: &[FrequencyEntry],
2367 ) -> Result<Vec<(Value, u64)>> {
2368 let dictionary = if *ty == LogicalType::Varchar { self.dictionary(column)? } else { None };
2369 let mut out = Vec::with_capacity(entries.len());
2370 for entry in entries {
2371 let value = match entry.value {
2372 FrequencyValue::Null => Value::Null,
2373 FrequencyValue::Integer(value) => match *ty {
2374 LogicalType::TinyInt => Value::TinyInt(
2375 i8::try_from(value)
2376 .map_err(|_| invalid("frequency TINYINT is out of range"))?,
2377 ),
2378 LogicalType::UTinyInt => Value::UTinyInt(
2379 u8::try_from(value)
2380 .map_err(|_| invalid("frequency UTINYINT is out of range"))?,
2381 ),
2382 LogicalType::USmallInt => Value::USmallInt(
2383 u16::try_from(value)
2384 .map_err(|_| invalid("frequency USMALLINT is out of range"))?,
2385 ),
2386 LogicalType::UInteger => Value::UInteger(
2387 u32::try_from(value)
2388 .map_err(|_| invalid("frequency UINTEGER is out of range"))?,
2389 ),
2390 LogicalType::UBigInt => Value::UBigInt(
2391 u64::try_from(value)
2392 .map_err(|_| invalid("frequency UBIGINT is out of range"))?,
2393 ),
2394 LogicalType::SmallInt => Value::SmallInt(
2395 i16::try_from(value)
2396 .map_err(|_| invalid("frequency SMALLINT is out of range"))?,
2397 ),
2398 LogicalType::Integer => Value::Integer(
2399 i32::try_from(value)
2400 .map_err(|_| invalid("frequency INTEGER is out of range"))?,
2401 ),
2402 LogicalType::BigInt => Value::BigInt(
2403 i64::try_from(value)
2404 .map_err(|_| invalid("frequency BIGINT is out of range"))?,
2405 ),
2406 LogicalType::Date => Value::Date(
2407 i32::try_from(value)
2408 .map_err(|_| invalid("frequency DATE is out of range"))?,
2409 ),
2410 LogicalType::Timestamp => Value::Timestamp(
2411 i64::try_from(value)
2412 .map_err(|_| invalid("frequency TIMESTAMP is out of range"))?,
2413 ),
2414 _ => return Err(invalid("integer frequency belongs to another type")),
2415 },
2416 FrequencyValue::Code(code) => dictionary
2417 .as_ref()
2418 .ok_or_else(|| invalid("frequency code has no dictionary"))?
2419 .try_value_at(code as usize)?,
2420 };
2421 out.push((value, entry.count));
2422 }
2423 Ok(out)
2424 }
2425
2426 pub fn frequency_occurrences(&self, column: usize) -> Result<Option<FrequencyOccurrences>> {
2436 self.table
2437 .fields
2438 .get(column)
2439 .ok_or_else(|| invalid("frequency column index out of range"))?;
2440 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2441 return Ok(None);
2442 };
2443 if summary.ordinals.is_empty() {
2444 return Ok(None);
2445 }
2446 Ok(Some(FrequencyOccurrences {
2447 omitted_max: summary.omitted_max,
2448 ordinals: summary.ordinals.clone(),
2449 }))
2450 }
2451
2452 pub fn distinct_values(&self, column: usize) -> Result<Option<u64>> {
2476 self.table
2477 .distincts
2478 .get(column)
2479 .copied()
2480 .ok_or_else(|| invalid("distinct column index out of range"))
2481 }
2482
2483 pub fn null_count(&self, column: usize) -> Result<u64> {
2494 if column >= self.table.fields.len() {
2495 return Err(invalid("null count column index out of range"));
2496 }
2497 let mut nulls = 0_u64;
2498 for stripe in &self.table.stripes {
2499 let range = stripe
2500 .zone
2501 .column(column)
2502 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
2503 nulls = nulls
2504 .checked_add(range.nulls as u64)
2505 .ok_or_else(|| invalid("null count overflow"))?;
2506 }
2507 Ok(nulls)
2508 }
2509
2510 pub fn text_extremes(&self, column: usize) -> Result<Option<(Value, Value)>> {
2525 if self.null_count(column)? > 0 {
2526 return Ok(None);
2527 }
2528 let Some(dictionary) = self.dictionary(column)? else { return Ok(None) };
2529 let Some(ranks) = dictionary.ranks() else { return Ok(None) };
2530 if ranks == 0 {
2531 return Ok(None);
2532 }
2533 let low = text_at_rank(&dictionary, 0)?;
2534 let high = text_at_rank(&dictionary, ranks - 1)?;
2535 Ok(Some((low, high)))
2536 }
2537
2538 pub fn exact_extremes(&self, column: usize) -> Result<Option<(Bound, Bound)>> {
2561 if column >= self.table.fields.len() {
2562 return Err(invalid("extremes column index out of range"));
2563 }
2564 let mut low: Option<Bound> = None;
2565 let mut high: Option<Bound> = None;
2566 for stripe in &self.table.stripes {
2567 let range = stripe
2568 .zone
2569 .column(column)
2570 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
2571 if !range.exact {
2572 return Ok(None);
2573 }
2574 let (Some(small), Some(large)) = (range.low.as_ref(), range.high.as_ref()) else {
2579 if stripe.rows > range.nulls {
2580 return Ok(None);
2581 }
2582 continue;
2583 };
2584 low = Some(low.map_or_else(|| small.clone(), |held| held.smaller(small.clone())));
2585 high = Some(high.map_or_else(|| large.clone(), |held| held.larger(large.clone())));
2586 }
2587 Ok(low.zip(high))
2588 }
2589
2590 pub fn exact_sum(&self, column: usize) -> Result<Option<(i128, u64)>> {
2603 if column >= self.table.fields.len() {
2604 return Err(invalid("sum column index out of range"));
2605 }
2606 let mut total = 0_i128;
2607 let mut rows = 0_u64;
2608 for stripe in &self.table.stripes {
2609 let range = stripe
2610 .zone
2611 .column(column)
2612 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
2613 let Some(part) = range.sum else { return Ok(None) };
2614 let Some(sum) = total.checked_add(part) else { return Ok(None) };
2615 total = sum;
2616 rows = rows.saturating_add(stripe.rows as u64 - range.nulls as u64);
2617 }
2618 Ok(Some((total, rows)))
2619 }
2620
2621 fn dictionary(&self, column: usize) -> Result<Option<Arc<Vector>>> {
2630 let Some(page) = self.table.dictionaries[column] else { return Ok(None) };
2631 if let Some(dictionary) = self.dictionaries[column].get() {
2632 return Ok(Some(Arc::clone(dictionary)));
2633 }
2634 let _queued = self.loading[column].lock().map_err(|_| invalid("a poisoned dictionary"))?;
2635 if let Some(dictionary) = self.dictionaries[column].get() {
2636 return Ok(Some(Arc::clone(dictionary)));
2637 }
2638 self.opened.fetch_add(1, Atomic::Relaxed);
2639 let dictionary = Arc::new(open_global_dictionary(
2640 Arc::clone(&self.file),
2641 page,
2642 &self.table.fields[column].ty,
2643 TEXT_KEEP_BUDGET,
2644 )?);
2645 let _ = self.dictionaries[column].set(Arc::clone(&dictionary));
2646 Ok(Some(dictionary))
2647 }
2648
2649 pub fn read(&self, part: usize, columns: &[usize]) -> Result<Chunk> {
2658 self.read_impl(part, columns, true)
2659 }
2660
2661 pub fn read_sparse(&self, part: usize, columns: &[usize]) -> Result<Chunk> {
2671 self.read_impl(part, columns, false)
2672 }
2673
2674 pub fn skips_codes(&self, part: usize, column: usize, candidates: &[u32]) -> Result<bool> {
2681 if candidates.is_empty() {
2682 return Ok(true);
2683 }
2684 if candidates.windows(2).any(|pair| pair[0] >= pair[1]) {
2685 return Err(Error::internal("native code candidates are not sorted and unique"));
2686 }
2687 let stripe = self.stripe_of(part)?;
2688 let Some(page) = stripe.memberships.get(column).copied().flatten() else {
2689 return Ok(false);
2690 };
2691 let mut bytes = vec![0; page.length as usize];
2692 read_at(&self.file, page.offset, &mut bytes)?;
2693 if checksum(&bytes) != page.hash {
2694 return Err(invalid("membership page checksum differs"));
2695 }
2696 let codes = decode_membership(&bytes)?;
2697 let mut left = 0;
2698 let mut right = 0;
2699 while left < codes.len() && right < candidates.len() {
2700 match codes[left].cmp(&candidates[right]) {
2701 Ordering::Less => left += 1,
2702 Ordering::Greater => right += 1,
2703 Ordering::Equal => return Ok(false),
2704 }
2705 }
2706 Ok(true)
2707 }
2708
2709 fn stripe_of(&self, part: usize) -> Result<&Stripe> {
2710 let place = self.places.get(part).ok_or_else(|| invalid("part index out of range"))?;
2711 self.table
2712 .stripes
2713 .get(place.stripe as usize)
2714 .ok_or_else(|| invalid("stripe index out of range"))
2715 }
2716
2717 fn held(&self, at: usize, stripe: &Stripe, column: usize, whole: bool) -> Result<CachedColumn> {
2734 let cache = self.cache.get(column).ok_or_else(|| invalid("column index out of range"))?;
2735 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
2736 let known = cached.index.get(at).and_then(Clone::clone);
2737 let page = cached.pages.get(at).and_then(Clone::clone);
2738 if let Some(index) = known.clone() {
2739 if !whole || page.is_some() {
2740 return Ok(CachedColumn { stripe: at, index, page });
2741 }
2742 }
2743 if cached.loading.contains(&at) {
2744 drop(cached);
2745 if let Some(index) = known {
2749 return Ok(CachedColumn { stripe: at, index, page: None });
2750 }
2751 let held = self.page_of(stripe, column, at, false, None)?;
2752 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
2753 remember(&mut cached, &held, self.kept.load(Atomic::Relaxed));
2754 return Ok(held);
2755 }
2756 cached.loading.push(at);
2757 drop(cached);
2758
2759 let read = self.page_of(stripe, column, at, whole, known);
2760
2761 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
2765 if let Some(position) = cached.loading.iter().position(|loading| *loading == at) {
2766 cached.loading.remove(position);
2767 }
2768 let held = read?;
2769 remember(&mut cached, &held, self.kept.load(Atomic::Relaxed));
2770 Ok(held)
2771 }
2772
2773 fn page_of(
2779 &self,
2780 stripe: &Stripe,
2781 column: usize,
2782 at: usize,
2783 whole: bool,
2784 known: Option<Arc<Vec<PartSpan>>>,
2785 ) -> Result<CachedColumn> {
2786 let index = match known {
2787 Some(index) => index,
2788 None => {
2789 self.indexes.fetch_add(1, Atomic::Relaxed);
2790 Arc::new(read_index(&self.file, stripe, column)?)
2791 }
2792 };
2793 let page = if whole {
2794 self.pages.fetch_add(1, Atomic::Relaxed);
2795 let span = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
2796 let mut bytes = vec![0; span.length as usize];
2797 read_at(&self.file, span.offset, &mut bytes)?;
2798 Some(Arc::new(bytes))
2799 } else {
2800 None
2801 };
2802 Ok(CachedColumn { stripe: at, index, page })
2803 }
2804
2805 fn read_impl(&self, at: usize, columns: &[usize], whole: bool) -> Result<Chunk> {
2806 let place = *self.places.get(at).ok_or_else(|| invalid("part index out of range"))?;
2807 let index = place.stripe as usize;
2808 let stripe =
2809 self.table.stripes.get(index).ok_or_else(|| invalid("stripe index out of range"))?;
2810 let rows = place.rows as usize;
2811 let mut picked = Vec::with_capacity(columns.len());
2812 for &column in columns {
2813 let field = self
2814 .table
2815 .fields
2816 .get(column)
2817 .ok_or_else(|| invalid("column index out of range"))?;
2818 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
2819 let held = self.held(index, stripe, column, whole)?;
2820 let span = *held
2821 .index
2822 .get(place.part as usize)
2823 .ok_or_else(|| invalid("part index out of range"))?;
2824 let owned;
2825 let bytes = match &held.page {
2826 Some(held) => part_bytes(held, span)?,
2827 None => {
2828 let offset = page
2829 .offset
2830 .checked_add(span.start as u64)
2831 .ok_or_else(|| invalid("part range overflow"))?;
2832 let mut bytes = vec![0; span.length];
2833 read_at(&self.file, offset, &mut bytes)?;
2834 owned = bytes;
2835 &owned
2836 }
2837 };
2838 if checksum(bytes) != span.hash {
2839 return Err(invalid(&format!(
2840 "column page checksum differs, column {column} part {} at {}+{} of {} bytes, \
2841 wanted {:016x} and got {:016x}",
2842 place.part,
2843 page.offset,
2844 span.start,
2845 span.length,
2846 span.hash,
2847 checksum(bytes),
2848 )));
2849 }
2850 let dictionary = self.dictionary(column)?;
2851 picked.push(decode(&field.ty, rows, bytes, dictionary)?);
2852 }
2853 Chunk::with_rows(picked, rows)
2854 }
2855
2856 #[must_use]
2872 pub fn skips(&self, part: usize, probes: &[Probe]) -> bool {
2873 let Some(place) = self.places.get(part).copied() else { return false };
2874 let Some(stripe) = self.table.stripes.get(place.stripe as usize) else { return false };
2875 if stripe.zone.skips(probes) {
2876 return true;
2877 }
2878 probes.iter().any(|probe| self.outside(place, probe) || self.sifted(place, probe))
2879 }
2880
2881 fn outside(&self, place: Place, probe: &Probe) -> bool {
2887 match self.stripe_part_ranges(place.stripe as usize, probe.column) {
2888 Some(ranges) => ranges
2889 .get(place.part as usize)
2890 .is_some_and(|range| range.excludes(probe.op, &probe.value)),
2891 None => false,
2892 }
2893 }
2894
2895 fn stripe_part_ranges(&self, stripe: usize, column: usize) -> Option<&[Range]> {
2901 let slot = self.part_ranges.get(column)?.get(stripe)?;
2902 if let Some(held) = slot.get() {
2903 return Some(held);
2904 }
2905 let page = self.table.stripes.get(stripe)?.part_ranges.get(column).copied().flatten()?;
2906 let mut bytes = vec![0; page.length as usize];
2907 read_at(&self.file, page.offset, &mut bytes).ok()?;
2908 if checksum(&bytes) != page.hash {
2909 return None;
2910 }
2911 let ranges = Arc::new(decode_part_ranges(&bytes).ok()?);
2912 let _ = slot.set(ranges);
2913 slot.get().map(|held| held.as_slice())
2914 }
2915
2916 #[must_use]
2928 pub fn certain(&self, part: usize, probes: &[Probe]) -> bool {
2929 let Some(place) = self.places.get(part).copied() else { return false };
2930 let Some(stripe) = self.table.stripes.get(place.stripe as usize) else { return false };
2931 stripe.zone.certain(probes)
2932 }
2933
2934 #[must_use]
2945 pub fn stripe_skips(&self, stripe: usize, probes: &[Probe]) -> bool {
2946 self.table.stripes.get(stripe).is_some_and(|held| held.zone.skips(probes))
2947 }
2948
2949 fn sifted(&self, place: Place, probe: &Probe) -> bool {
2955 if probe.op != Op::Equal {
2956 return false;
2957 }
2958 match self.stripe_sieves(place.stripe as usize, probe.column) {
2959 Some(sieves) => sieves
2960 .get(place.part as usize)
2961 .and_then(Option::as_ref)
2962 .is_some_and(|sieve| sieve.excludes(&probe.value)),
2963 None => false,
2964 }
2965 }
2966
2967 fn stripe_sieves(&self, stripe: usize, column: usize) -> Option<&[Option<Sieve>]> {
2974 let slot = self.sieves.get(column)?.get(stripe)?;
2975 if let Some(held) = slot.get() {
2976 return Some(held);
2977 }
2978 let page = self.table.stripes.get(stripe)?.sieves.get(column).copied().flatten()?;
2979 let mut bytes = vec![0; page.length as usize];
2980 read_at(&self.file, page.offset, &mut bytes).ok()?;
2981 if checksum(&bytes) != page.hash {
2982 return None;
2983 }
2984 let sieves = Arc::new(decode_sieves(&bytes).ok()?);
2985 let _ = slot.set(sieves);
2986 slot.get().map(|held| held.as_slice())
2987 }
2988}
2989
2990fn text_at_rank(dictionary: &Vector, rank: usize) -> Result<Value> {
2992 let code = dictionary.code_at_rank(rank)? as usize;
2993 let text = dictionary
2994 .try_text_at(code)?
2995 .ok_or_else(|| invalid("global dictionary order names a code it does not have"))?;
2996 Ok(Value::Varchar(text.into()))
2997}
2998
2999#[cfg(unix)]
3004fn write_at(file: &File, mut offset: u64, mut bytes: &[u8]) -> Result<()> {
3005 use std::os::unix::fs::FileExt;
3006 while !bytes.is_empty() {
3007 let written = file.write_at(bytes, offset).map_err(io)?;
3008 if written == 0 {
3009 return Err(invalid("a write to the native file wrote nothing"));
3010 }
3011 offset += written as u64;
3012 bytes = &bytes[written..];
3013 }
3014 Ok(())
3015}
3016
3017#[cfg(windows)]
3019fn write_at(file: &File, mut offset: u64, mut bytes: &[u8]) -> Result<()> {
3020 use std::os::windows::fs::FileExt;
3021 while !bytes.is_empty() {
3022 let written = file.seek_write(bytes, offset).map_err(io)?;
3023 if written == 0 {
3024 return Err(invalid("a write to the native file wrote nothing"));
3025 }
3026 offset += written as u64;
3027 bytes = &bytes[written..];
3028 }
3029 Ok(())
3030}
3031
3032#[cfg(not(any(unix, windows)))]
3034fn write_at(file: &File, offset: u64, bytes: &[u8]) -> Result<()> {
3035 use std::io::Write;
3036 let mut file = file.try_clone().map_err(io)?;
3037 file.seek(SeekFrom::Start(offset)).map_err(io)?;
3038 file.write_all(bytes).map_err(io)
3039}
3040
3041#[cfg(unix)]
3051fn read_at(file: &File, mut offset: u64, mut bytes: &mut [u8]) -> Result<()> {
3052 use std::os::unix::fs::FileExt;
3053 while !bytes.is_empty() {
3054 let read = file.read_at(bytes, offset).map_err(io)?;
3055 if read == 0 {
3056 return Err(invalid("column page ends before its declared length"));
3057 }
3058 offset += read as u64;
3059 bytes = &mut bytes[read..];
3060 }
3061 Ok(())
3062}
3063
3064#[cfg(windows)]
3070fn read_at(file: &File, mut offset: u64, mut bytes: &mut [u8]) -> Result<()> {
3071 use std::os::windows::fs::FileExt;
3072 while !bytes.is_empty() {
3073 let read = file.seek_read(bytes, offset).map_err(io)?;
3074 if read == 0 {
3075 return Err(invalid("column page ends before its declared length"));
3076 }
3077 offset += read as u64;
3078 bytes = &mut bytes[read..];
3079 }
3080 Ok(())
3081}
3082
3083#[cfg(not(any(unix, windows)))]
3088fn read_at(file: &File, offset: u64, bytes: &mut [u8]) -> Result<()> {
3089 let mut file = file.try_clone().map_err(io)?;
3090 file.seek(SeekFrom::Start(offset)).map_err(io)?;
3091 file.read_exact(bytes).map_err(io)
3092}
3093
3094fn type_tag(ty: &LogicalType) -> Result<u8> {
3095 match ty {
3096 LogicalType::SmallInt => Ok(1),
3097 LogicalType::Integer => Ok(2),
3098 LogicalType::BigInt => Ok(3),
3099 LogicalType::Varchar => Ok(4),
3100 LogicalType::Date => Ok(5),
3101 LogicalType::Timestamp => Ok(6),
3102 LogicalType::Boolean => Ok(7),
3103 LogicalType::TinyInt => Ok(8),
3104 LogicalType::UTinyInt => Ok(9),
3105 LogicalType::USmallInt => Ok(10),
3106 LogicalType::UInteger => Ok(11),
3107 LogicalType::UBigInt => Ok(12),
3108 _ => Err(Error::not_implemented(format!("native storage for {ty}"))),
3109 }
3110}
3111
3112fn tag_type(tag: u8) -> Result<LogicalType> {
3113 match tag {
3114 1 => Ok(LogicalType::SmallInt),
3115 2 => Ok(LogicalType::Integer),
3116 3 => Ok(LogicalType::BigInt),
3117 4 => Ok(LogicalType::Varchar),
3118 5 => Ok(LogicalType::Date),
3119 6 => Ok(LogicalType::Timestamp),
3120 7 => Ok(LogicalType::Boolean),
3121 8 => Ok(LogicalType::TinyInt),
3122 9 => Ok(LogicalType::UTinyInt),
3123 10 => Ok(LogicalType::USmallInt),
3124 11 => Ok(LogicalType::UInteger),
3125 12 => Ok(LogicalType::UBigInt),
3126 _ => Err(invalid("column type tag is unknown")),
3127 }
3128}
3129
3130fn put_u16(out: &mut Vec<u8>, value: u16) {
3131 out.extend_from_slice(&value.to_le_bytes());
3132}
3133fn put_u32(out: &mut Vec<u8>, value: u32) {
3134 out.extend_from_slice(&value.to_le_bytes());
3135}
3136fn put_u64(out: &mut Vec<u8>, value: u64) {
3137 out.extend_from_slice(&value.to_le_bytes());
3138}
3139fn put_var_u64(out: &mut Vec<u8>, mut value: u64) {
3140 while value >= 0x80 {
3141 out.push((value as u8 & 0x7f) | 0x80);
3142 value >>= 7;
3143 }
3144 out.push(value as u8);
3145}
3146
3147fn frequency_order(left: FrequencyValue, right: FrequencyValue) -> Ordering {
3148 match (left, right) {
3149 (FrequencyValue::Null, FrequencyValue::Null) => Ordering::Equal,
3150 (FrequencyValue::Null, _) => Ordering::Less,
3151 (_, FrequencyValue::Null) => Ordering::Greater,
3152 (FrequencyValue::Integer(left), FrequencyValue::Integer(right)) => left.cmp(&right),
3153 (FrequencyValue::Code(left), FrequencyValue::Code(right)) => left.cmp(&right),
3154 (FrequencyValue::Integer(_), FrequencyValue::Code(_)) => Ordering::Less,
3155 (FrequencyValue::Code(_), FrequencyValue::Integer(_)) => Ordering::Greater,
3156 }
3157}
3158
3159fn code_frequency(dictionary: &GlobalDictionary) -> FrequencySummary {
3160 let mut entries = dictionary
3161 .counts
3162 .iter()
3163 .enumerate()
3164 .filter(|(_, count)| **count != 0)
3165 .map(|(code, &count)| FrequencyEntry { value: FrequencyValue::Code(code as u32), count })
3166 .collect::<Vec<_>>();
3167 if dictionary.nulls != 0 {
3168 entries.push(FrequencyEntry { value: FrequencyValue::Null, count: dictionary.nulls });
3169 }
3170 entries.sort_unstable_by(|left, right| {
3171 right.count.cmp(&left.count).then_with(|| frequency_order(left.value, right.value))
3172 });
3173 let omitted_max = entries.get(FREQUENCY_ENTRIES).map_or(0, |entry| entry.count);
3174 entries.truncate(FREQUENCY_ENTRIES);
3175 FrequencySummary { entries, omitted_max, ordinals: Vec::new() }
3176}
3177
3178fn encode_directory(table: &Table) -> Result<Vec<u8>> {
3179 let mut out = DIRECTORY.to_vec();
3180 let name = table.name.as_bytes();
3181 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("table name too long"))?);
3182 out.extend_from_slice(name);
3183 put_u16(&mut out, u16::try_from(table.fields.len()).map_err(|_| invalid("too many columns"))?);
3184 for field in &table.fields {
3185 let name = field.name.as_bytes();
3186 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("column name too long"))?);
3187 out.extend_from_slice(name);
3188 out.push(type_tag(&field.ty)?);
3189 out.push(u8::from(field.not_null));
3190 }
3191 for dictionary in &table.dictionaries {
3192 match dictionary {
3193 None => out.push(0),
3194 Some(page) => {
3195 out.push(1);
3196 put_u64(&mut out, page.offset);
3197 put_u32(&mut out, page.length);
3198 put_u64(&mut out, page.hash);
3199 }
3200 }
3201 }
3202 for distinct in &table.distincts {
3203 match distinct {
3204 None => out.push(0),
3205 Some(count) => {
3206 out.push(1);
3207 put_u64(&mut out, *count);
3208 }
3209 }
3210 }
3211 put_u64(&mut out, u64::try_from(table.rows).map_err(|_| invalid("row count overflow"))?);
3212 put_u32(&mut out, u32::try_from(table.stripes.len()).map_err(|_| invalid("too many stripes"))?);
3213 for stripe in &table.stripes {
3214 put_u32(
3215 &mut out,
3216 u32::try_from(stripe.parts.len()).map_err(|_| invalid("too many parts in a stripe"))?,
3217 );
3218 for &rows in &stripe.parts {
3219 put_u32(&mut out, rows);
3220 }
3221 put_u64(&mut out, stripe.index.offset);
3222 put_u32(&mut out, stripe.index.length);
3223 for page in &stripe.pages {
3224 put_u64(&mut out, page.offset);
3225 put_u32(&mut out, page.length);
3226 }
3227 for (field, membership) in table.fields.iter().zip(&stripe.memberships) {
3228 if field.ty != LogicalType::Varchar {
3229 continue;
3230 }
3231 let page =
3232 membership.ok_or_else(|| invalid("string page has no code membership index"))?;
3233 put_u64(&mut out, page.offset);
3234 put_u32(&mut out, page.length);
3235 put_u64(&mut out, page.hash);
3236 }
3237 for sieve in &stripe.sieves {
3238 match sieve {
3239 None => out.push(0),
3240 Some(page) => {
3241 out.push(1);
3242 put_u64(&mut out, page.offset);
3243 put_u32(&mut out, page.length);
3244 put_u64(&mut out, page.hash);
3245 }
3246 }
3247 }
3248 for held in &stripe.part_ranges {
3249 match held {
3250 None => out.push(0),
3251 Some(page) => {
3252 out.push(1);
3253 put_u64(&mut out, page.offset);
3254 put_u32(&mut out, page.length);
3255 put_u64(&mut out, page.hash);
3256 }
3257 }
3258 }
3259 for range in stripe.zone.columns() {
3260 put_bound(&mut out, range.low.as_ref())?;
3261 put_bound(&mut out, range.high.as_ref())?;
3262 put_u32(
3263 &mut out,
3264 u32::try_from(range.nulls).map_err(|_| invalid("null count overflow"))?,
3265 );
3266 out.push(u8::from(range.exact));
3267 match range.sum {
3268 None => out.push(0),
3269 Some(total) => {
3270 out.push(1);
3271 out.extend_from_slice(&total.to_le_bytes());
3272 }
3273 }
3274 }
3275 }
3276 out.extend_from_slice(FREQUENCIES);
3277 put_u16(
3278 &mut out,
3279 u16::try_from(table.frequencies.len())
3280 .map_err(|_| invalid("too many frequency columns"))?,
3281 );
3282 for summary in &table.frequencies {
3283 let Some(summary) = summary else {
3284 out.push(0);
3285 continue;
3286 };
3287 out.push(1);
3288 put_u64(&mut out, summary.omitted_max);
3289 put_u32(
3290 &mut out,
3291 u32::try_from(summary.entries.len())
3292 .map_err(|_| invalid("too many frequency entries"))?,
3293 );
3294 for entry in &summary.entries {
3295 match entry.value {
3296 FrequencyValue::Null => out.push(0),
3297 FrequencyValue::Integer(value) => {
3298 out.push(1);
3299 out.extend_from_slice(&value.to_le_bytes());
3300 }
3301 FrequencyValue::Code(value) => {
3302 out.push(2);
3303 put_u32(&mut out, value);
3304 }
3305 }
3306 put_u64(&mut out, entry.count);
3307 }
3308 put_u32(
3309 &mut out,
3310 u32::try_from(summary.ordinals.len())
3311 .map_err(|_| invalid("too many frequency ordinals"))?,
3312 );
3313 let mut previous = 0_u64;
3314 for (at, &ordinal) in summary.ordinals.iter().enumerate() {
3315 let delta = if at == 0 {
3316 ordinal
3317 } else {
3318 ordinal
3319 .checked_sub(previous)
3320 .ok_or_else(|| invalid("frequency ordinals are not ordered"))?
3321 };
3322 if at != 0 && delta == 0 {
3323 return Err(invalid("frequency ordinals are not unique"));
3324 }
3325 put_var_u64(&mut out, delta);
3326 previous = ordinal;
3327 }
3328 }
3329 Ok(out)
3330}
3331
3332struct Cursor<'a> {
3333 bytes: &'a [u8],
3334 at: usize,
3335}
3336impl<'a> Cursor<'a> {
3337 fn take(&mut self, len: usize) -> Result<&'a [u8]> {
3338 let end = self.at.checked_add(len).ok_or_else(|| invalid("directory offset overflow"))?;
3339 let bytes =
3340 self.bytes.get(self.at..end).ok_or_else(|| invalid("directory is truncated"))?;
3341 self.at = end;
3342 Ok(bytes)
3343 }
3344 fn u8(&mut self) -> Result<u8> {
3345 Ok(self.take(1)?[0])
3346 }
3347 fn u16(&mut self) -> Result<u16> {
3348 Ok(u16::from_le_bytes(self.take(2)?.try_into().expect("two bytes")))
3349 }
3350 fn u32(&mut self) -> Result<u32> {
3351 Ok(u32::from_le_bytes(self.take(4)?.try_into().expect("four bytes")))
3352 }
3353 fn u64(&mut self) -> Result<u64> {
3354 Ok(u64::from_le_bytes(self.take(8)?.try_into().expect("eight bytes")))
3355 }
3356 fn var_u64(&mut self) -> Result<u64> {
3357 let mut value = 0_u64;
3358 for shift in (0..=63).step_by(7) {
3359 let byte = self.u8()?;
3360 let part = u64::from(byte & 0x7f);
3361 if shift == 63 && part > 1 {
3362 return Err(invalid("frequency ordinal varint overflows"));
3363 }
3364 value |= part << shift;
3365 if byte & 0x80 == 0 {
3366 return Ok(value);
3367 }
3368 }
3369 Err(invalid("frequency ordinal varint is too long"))
3370 }
3371 fn bound(&mut self) -> Result<Option<Bound>> {
3372 Ok(match self.u8()? {
3373 0 => None,
3374 1 => Some(Bound::Int(i128::from_le_bytes(
3375 self.take(16)?.try_into().expect("sixteen bytes"),
3376 ))),
3377 2 => Some(Bound::Real(f64::from_le_bytes(
3378 self.take(8)?.try_into().expect("eight bytes"),
3379 ))),
3380 3 => {
3381 let length = self.u32()? as usize;
3382 Some(Bound::Bytes(self.take(length)?.to_vec()))
3383 }
3384 4 => {
3385 let unscaled =
3386 i128::from_le_bytes(self.take(16)?.try_into().expect("sixteen bytes"));
3387 Some(Bound::Scaled { unscaled, scale: self.u8()? })
3388 }
3389 _ => return Err(invalid("bound tag differs")),
3390 })
3391 }
3392 fn text(&mut self) -> Result<String> {
3393 let len = self.u16()? as usize;
3394 String::from_utf8(self.take(len)?.to_vec()).map_err(|_| invalid("name is not UTF-8"))
3395 }
3396}
3397
3398fn decode_directory(bytes: &[u8], size: u64) -> Result<Table> {
3399 let mut cur = Cursor { bytes, at: 0 };
3400 if cur.take(8)? != DIRECTORY {
3401 return Err(invalid("directory magic differs"));
3402 }
3403 let name = cur.text()?;
3404 let width = cur.u16()? as usize;
3405 let mut fields = Vec::with_capacity(width);
3406 for _ in 0..width {
3407 let name = cur.text()?;
3408 let ty = tag_type(cur.u8()?)?;
3409 let not_null = match cur.u8()? {
3410 0 => false,
3411 1 => true,
3412 _ => return Err(invalid("nullability flag differs")),
3413 };
3414 fields.push(Field { name, ty, not_null });
3415 }
3416 let mut dictionaries = Vec::with_capacity(width);
3417 for _ in 0..width {
3418 dictionaries.push(match cur.u8()? {
3419 0 => None,
3420 1 => {
3421 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3422 let end = page
3423 .offset
3424 .checked_add(u64::from(page.length))
3425 .ok_or_else(|| invalid("dictionary page offset overflow"))?;
3426 if page.offset < HEADER || end > size {
3431 return Err(invalid("dictionary page range is outside the file"));
3432 }
3433 Some(page)
3434 }
3435 _ => return Err(invalid("dictionary page tag differs")),
3436 });
3437 }
3438 let mut distincts = Vec::with_capacity(width);
3439 for _ in 0..width {
3440 distincts.push(match cur.u8()? {
3441 0 => None,
3442 1 => Some(cur.u64()?),
3443 _ => return Err(invalid("distinct count tag differs")),
3444 });
3445 }
3446 let rows = usize::try_from(cur.u64()?).map_err(|_| invalid("row count does not fit"))?;
3447 let count = cur.u32()? as usize;
3448 let mut stripes = Vec::with_capacity(count);
3449 let mut total = 0_usize;
3450 for _ in 0..count {
3451 let count = cur.u32()? as usize;
3452 if count == 0 || count > STRIPE_PARTS {
3453 return Err(invalid("stripe part count is outside its bound"));
3454 }
3455 let mut parts = Vec::with_capacity(count);
3456 let mut stripe_rows = 0_usize;
3457 for _ in 0..count {
3458 let rows = cur.u32()?;
3459 if rows == 0 {
3460 return Err(invalid("empty part"));
3461 }
3462 parts.push(rows);
3463 stripe_rows = stripe_rows
3464 .checked_add(rows as usize)
3465 .ok_or_else(|| invalid("stripe row count overflow"))?;
3466 }
3467 total =
3468 total.checked_add(stripe_rows).ok_or_else(|| invalid("stripe row count overflow"))?;
3469 let index = Span { offset: cur.u64()?, length: cur.u32()? };
3470 let section = index_section(count)?;
3471 let wanted = section
3472 .checked_mul(width)
3473 .and_then(|bytes| u32::try_from(bytes).ok())
3474 .ok_or_else(|| invalid("index page length overflow"))?;
3475 let end = index
3476 .offset
3477 .checked_add(u64::from(index.length))
3478 .ok_or_else(|| invalid("index page offset overflow"))?;
3479 if index.offset < HEADER || end > size || index.length != wanted {
3480 return Err(invalid("index page range is outside the file"));
3481 }
3482 let mut pages = Vec::with_capacity(width);
3483 for _ in 0..width {
3484 let offset = cur.u64()?;
3485 let length = cur.u32()?;
3486 let end = offset
3487 .checked_add(u64::from(length))
3488 .ok_or_else(|| invalid("page offset overflow"))?;
3489 if offset < HEADER || end > size || length as usize > MAX_PAGE {
3490 return Err(invalid("page range is outside the file"));
3491 }
3492 pages.push(Span { offset, length });
3493 }
3494 let mut memberships = vec![None; width];
3495 for (column, field) in fields.iter().enumerate() {
3496 if field.ty != LogicalType::Varchar {
3497 continue;
3498 }
3499 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3500 let end = page
3501 .offset
3502 .checked_add(u64::from(page.length))
3503 .ok_or_else(|| invalid("membership page offset overflow"))?;
3504 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
3505 return Err(invalid("membership page range is outside the file"));
3506 }
3507 memberships[column] = Some(page);
3508 }
3509 let mut sieves = vec![None; width];
3510 for sieve in sieves.iter_mut().take(width) {
3511 match cur.u8()? {
3512 0 => continue,
3513 1 => {}
3514 _ => return Err(invalid("a sieve page has an unknown tag")),
3515 }
3516 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3517 let end = page
3518 .offset
3519 .checked_add(u64::from(page.length))
3520 .ok_or_else(|| invalid("sieve page offset overflow"))?;
3521 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
3522 return Err(invalid("sieve page range is outside the file"));
3523 }
3524 *sieve = Some(page);
3525 }
3526 let mut part_ranges = vec![None; width];
3527 for held in part_ranges.iter_mut().take(width) {
3528 match cur.u8()? {
3529 0 => continue,
3530 1 => {}
3531 _ => return Err(invalid("a part range page has an unknown tag")),
3532 }
3533 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3534 let end = page
3535 .offset
3536 .checked_add(u64::from(page.length))
3537 .ok_or_else(|| invalid("part range page offset overflow"))?;
3538 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
3539 return Err(invalid("part range page range is outside the file"));
3540 }
3541 *held = Some(page);
3542 }
3543 let mut ranges = Vec::with_capacity(width);
3544 for column in 0..width {
3545 let low = cur.bound()?;
3546 let high = cur.bound()?;
3547 let nulls = cur.u32()? as usize;
3548 if nulls > stripe_rows {
3549 return Err(invalid("null count exceeds stripe rows"));
3550 }
3551 let exact = cur.u8()? != 0;
3552 let sum = match cur.u8()? {
3553 0 => None,
3554 1 => Some(i128::from_le_bytes(
3555 cur.take(16)?.try_into().map_err(|_| invalid("a stripe sum is truncated"))?,
3556 )),
3557 _ => return Err(invalid("a stripe sum has an unknown tag")),
3558 };
3559 let ty = &fields.get(column).ok_or_else(|| invalid("a stripe range has no column"))?.ty;
3565 let low = low.map(|bound| scaled_as(bound, ty));
3566 let high = high.map(|bound| scaled_as(bound, ty));
3567 ranges.push(Range { low, high, nulls, exact, sum });
3568 }
3569 stripes.push(Stripe {
3570 rows: stripe_rows,
3571 parts,
3572 index,
3573 pages,
3574 memberships,
3575 sieves,
3576 part_ranges,
3577 zone: Zone::from_ranges(ranges),
3578 });
3579 }
3580 if total != rows {
3581 return Err(invalid("table row count differs from stripes"));
3582 }
3583 let frequencies = if cur.at == bytes.len() {
3584 vec![None; width]
3585 } else {
3586 if cur.take(8)? != FREQUENCIES {
3587 return Err(invalid("directory extension magic differs"));
3588 }
3589 if cur.u16()? as usize != width {
3590 return Err(invalid("frequency column count differs"));
3591 }
3592 let mut frequencies = Vec::with_capacity(width);
3593 for field in &fields {
3594 let summary = match cur.u8()? {
3595 0 => None,
3596 1 => {
3597 let omitted_max = cur.u64()?;
3598 let count = cur.u32()? as usize;
3599 if count > FREQUENCY_ENTRIES {
3600 return Err(invalid("frequency entry count exceeds its bound"));
3601 }
3602 let mut entries = Vec::with_capacity(count);
3603 for _ in 0..count {
3605 let value = match cur.u8()? {
3606 0 => FrequencyValue::Null,
3607 1 => FrequencyValue::Integer(i128::from_le_bytes(
3608 cur.take(16)?.try_into().expect("sixteen bytes"),
3609 )),
3610 2 => FrequencyValue::Code(cur.u32()?),
3611 _ => return Err(invalid("frequency value tag differs")),
3612 };
3613 let valid = matches!(
3614 (&field.ty, value),
3615 (_, FrequencyValue::Null)
3616 | (LogicalType::Varchar, FrequencyValue::Code(_))
3617 | (
3618 LogicalType::TinyInt
3619 | LogicalType::SmallInt
3620 | LogicalType::Integer
3621 | LogicalType::BigInt
3622 | LogicalType::UTinyInt
3623 | LogicalType::USmallInt
3624 | LogicalType::UInteger
3625 | LogicalType::UBigInt
3626 | LogicalType::Date
3627 | LogicalType::Timestamp,
3628 FrequencyValue::Integer(_),
3629 )
3630 );
3631 if !valid {
3632 return Err(invalid("frequency value does not match its column"));
3633 }
3634 let count = cur.u64()?;
3635 if count == 0 || count > rows as u64 {
3636 return Err(invalid("frequency count is outside the table"));
3637 }
3638 entries.push(FrequencyEntry { value, count });
3639 }
3640 if entries.windows(2).any(|pair| pair[0].count < pair[1].count) {
3641 return Err(invalid("frequency entries are not descending"));
3642 }
3643 let ordinals = {
3644 let ordinal_count = cur.u32()? as usize;
3645 if ordinal_count > FREQUENCY_ORDINALS || ordinal_count > rows {
3646 return Err(invalid("frequency ordinal count exceeds its bound"));
3647 }
3648 let mut ordinals = Vec::with_capacity(ordinal_count);
3649 let mut previous = 0_u64;
3650 for at in 0..ordinal_count {
3651 let delta = cur.var_u64()?;
3652 if at != 0 && delta == 0 {
3653 return Err(invalid("frequency ordinals are not increasing"));
3654 }
3655 let ordinal = if at == 0 {
3656 delta
3657 } else {
3658 previous
3659 .checked_add(delta)
3660 .ok_or_else(|| invalid("frequency ordinal overflows"))?
3661 };
3662 if ordinal >= rows as u64 {
3663 return Err(invalid("frequency ordinal is outside the table"));
3664 }
3665 ordinals.push(ordinal);
3666 previous = ordinal;
3667 }
3668 ordinals
3669 };
3670 Some(FrequencySummary { entries, omitted_max, ordinals })
3671 }
3672 _ => return Err(invalid("frequency summary tag differs")),
3673 };
3674 frequencies.push(summary);
3675 }
3676 frequencies
3677 };
3678 if cur.at != bytes.len() {
3679 return Err(invalid("directory has trailing bytes"));
3680 }
3681 Ok(Table { name, fields, stripes, rows, dictionaries, distincts, frequencies })
3682}
3683
3684fn put_bound(out: &mut Vec<u8>, bound: Option<&Bound>) -> Result<()> {
3685 match bound {
3686 None => out.push(0),
3687 Some(Bound::Int(value)) => {
3688 out.push(1);
3689 out.extend_from_slice(&value.to_le_bytes());
3690 }
3691 Some(Bound::Real(value)) => {
3692 out.push(2);
3693 out.extend_from_slice(&value.to_le_bytes());
3694 }
3695 Some(Bound::Bytes(value)) => {
3696 out.push(3);
3697 put_u32(out, u32::try_from(value.len()).map_err(|_| invalid("bound length overflow"))?);
3698 out.extend_from_slice(value);
3699 }
3700 Some(Bound::Scaled { unscaled, scale }) => {
3701 out.push(4);
3702 out.extend_from_slice(&unscaled.to_le_bytes());
3703 out.push(*scale);
3704 }
3705 }
3706 Ok(())
3707}
3708
3709#[derive(Debug)]
3726struct Codes;
3727
3728impl chooser::Chooser for Codes {
3729 fn name(&self) -> &'static str {
3730 "codes"
3731 }
3732
3733 fn narrow_strings(
3734 &self,
3735 _values: &[&[u8]],
3736 offered: &[string::Kind],
3737 _depth: u8,
3738 ) -> Vec<string::Kind> {
3739 offered.to_vec()
3742 }
3743
3744 fn narrow_integers(
3745 &self,
3746 _values: &[i64],
3747 offered: &[integer::Kind],
3748 depth: u8,
3749 ) -> Vec<integer::Kind> {
3750 let keep: &[integer::Kind] = if depth == 0 {
3751 &[integer::Kind::Constant, integer::Kind::Packed, integer::Kind::Rle]
3752 } else {
3753 &[integer::Kind::Constant, integer::Kind::Packed]
3754 };
3755 let narrowed: Vec<integer::Kind> =
3756 offered.iter().copied().filter(|kind| keep.contains(kind)).collect();
3757 if narrowed.is_empty() { offered.to_vec() } else { narrowed }
3760 }
3761}
3762
3763#[derive(Debug)]
3775struct Fixed;
3776
3777impl chooser::Chooser for Fixed {
3778 fn name(&self) -> &'static str {
3779 "fixed"
3780 }
3781
3782 fn narrow_strings(
3783 &self,
3784 _values: &[&[u8]],
3785 offered: &[string::Kind],
3786 _depth: u8,
3787 ) -> Vec<string::Kind> {
3788 offered.to_vec()
3789 }
3790
3791 fn narrow_integers(
3792 &self,
3793 _values: &[i64],
3794 offered: &[integer::Kind],
3795 depth: u8,
3796 ) -> Vec<integer::Kind> {
3797 let keep: &[integer::Kind] = if depth == 0 {
3798 &[
3799 integer::Kind::Constant,
3800 integer::Kind::Packed,
3801 integer::Kind::Delta,
3802 integer::Kind::Rle,
3803 integer::Kind::Sparse,
3804 integer::Kind::Strided,
3805 ]
3806 } else {
3807 &[integer::Kind::Constant, integer::Kind::Packed, integer::Kind::Delta]
3808 };
3809 let narrowed: Vec<integer::Kind> =
3810 offered.iter().copied().filter(|kind| keep.contains(kind)).collect();
3811 if narrowed.is_empty() { offered.to_vec() } else { narrowed }
3812 }
3813}
3814
3815fn widened(data: &Data) -> Option<Vec<i64>> {
3822 match data {
3823 Data::Int8(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
3824 Data::UInt8(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
3825 Data::Int16(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
3826 Data::UInt16(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
3827 Data::Int32(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
3828 Data::UInt32(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
3829 Data::Int64(values) => Some(values.to_vec()),
3830 _ => None,
3831 }
3832}
3833
3834trait Narrow: Copy {
3841 const BIASED: (u32, u64);
3846
3847 fn narrow(value: i64) -> Self;
3849}
3850
3851#[allow(clippy::cast_sign_loss, reason = "a residue is a bit pattern and not a number")]
3868fn residue<T: Narrow>(value: i64) -> u64 {
3869 let (bits, bias) = T::BIASED;
3870 (value as u64).wrapping_add(bias) >> bits
3871}
3872
3873macro_rules! narrows {
3878 ($($ty:ty => $bias:expr),* $(,)?) => {$(
3879 impl Narrow for $ty {
3880 const BIASED: (u32, u64) = (<$ty>::BITS, $bias);
3881
3882 #[allow(
3883 clippy::cast_possible_truncation,
3884 clippy::cast_sign_loss,
3885 reason = "the caller has checked the bits this truncates away"
3886 )]
3887 fn narrow(value: i64) -> Self {
3888 value as Self
3889 }
3890 }
3891 )*};
3892}
3893
3894narrows! {
3895 i8 => 1 << 7,
3896 u8 => 0,
3897 i16 => 1 << 15,
3898 u16 => 0,
3899 i32 => 1 << 31,
3900 u32 => 0,
3901}
3902
3903fn fit<T: Narrow>(values: &[i64]) -> Result<Vec<T>> {
3916 let mut spilled = 0u64;
3917 for value in values {
3918 spilled |= residue::<T>(*value);
3919 }
3920 if spilled != 0 {
3921 return Err(invalid("page value is not of its type"));
3922 }
3923 Ok(values.iter().map(|value| T::narrow(*value)).collect())
3924}
3925
3926fn narrowed(ty: &LogicalType, values: Vec<i64>) -> Result<Data> {
3931 Ok(match ty {
3932 LogicalType::TinyInt => Data::Int8(fit::<i8>(&values)?.into()),
3933 LogicalType::UTinyInt => Data::UInt8(fit::<u8>(&values)?.into()),
3934 LogicalType::SmallInt => Data::Int16(fit::<i16>(&values)?.into()),
3935 LogicalType::USmallInt => Data::UInt16(fit::<u16>(&values)?.into()),
3936 LogicalType::Integer | LogicalType::Date => Data::Int32(fit::<i32>(&values)?.into()),
3937 LogicalType::UInteger => Data::UInt32(fit::<u32>(&values)?.into()),
3938 LogicalType::BigInt | LogicalType::Timestamp => Data::Int64(values.into()),
3939 _ => return Err(invalid("cascade codec belongs to a page that is not integers")),
3940 })
3941}
3942
3943fn plain_width(ty: &LogicalType) -> Option<usize> {
3946 Some(match ty {
3947 LogicalType::TinyInt | LogicalType::UTinyInt => 1,
3948 LogicalType::SmallInt | LogicalType::USmallInt => 2,
3949 LogicalType::Integer | LogicalType::UInteger | LogicalType::Date => 4,
3950 LogicalType::BigInt | LogicalType::Timestamp => 8,
3951 _ => return None,
3952 })
3953}
3954
3955fn cascaded(
3961 flat: &Vector,
3962 ty: &LogicalType,
3963 packed: Option<&Packed<'_>>,
3964) -> Result<Option<Vec<u8>>> {
3965 let (Some(width), Some(data)) = (plain_width(ty), flat.data()) else { return Ok(None) };
3966 let Some(values) = widened(data) else { return Ok(None) };
3967 let plain = values.len().saturating_mul(width);
3968 let best = match packed {
3969 Some(packed) => plain.min(21 + size_of_val(packed.words())),
3971 None => plain,
3972 };
3973 let out = integer::encode_with(&values, &Fixed)?;
3974 Ok((out.len() < best).then_some(out))
3975}
3976
3977fn encoded_codes(codes: &[u32]) -> Result<Option<Vec<u8>>> {
3989 let wide: Vec<i64> = codes.iter().map(|code| i64::from(*code)).collect();
3990 let coded = integer::encode_with(&wide, &Codes)?;
3991 let plain = codes.len().saturating_mul(size_of::<u32>());
3992 Ok((coded.len() < plain).then_some(coded))
3993}
3994
3995fn encode(
3996 vector: &Vector,
3997 global: Option<&mut GlobalDictionary>,
3998) -> Result<(Vec<u8>, Option<Vec<u32>>)> {
3999 let ty = vector.logical_type();
4000 let flat = vector.flatten()?;
4002 let mut out = Vec::new();
4003 let mut global_codes = None;
4004 if let Some(global) = global {
4005 let mut codes = Vec::with_capacity(flat.len());
4006 for row in 0..flat.len() {
4007 let text = flat.text_at(row).unwrap_or("");
4008 let code = global.code(text)?;
4009 global.observe(code, flat.is_null_at(row))?;
4010 codes.push(code);
4011 }
4012 global_codes = Some(codes);
4013 }
4014 let membership = global_codes.as_deref().map(unique_codes);
4015 let dictionary = if global_codes.is_none() && ty == &LogicalType::Varchar {
4016 string_dictionary(&flat)?
4017 } else {
4018 None
4019 };
4020 let packed_vector = if dictionary.is_none() && global_codes.is_none() {
4021 Some(flat.bit_packed()?)
4022 } else {
4023 None
4024 };
4025 let packed = packed_vector.as_ref().and_then(Vector::packed_parts);
4026 let coded = match global_codes.as_deref() {
4027 Some(codes) => encoded_codes(codes)?,
4028 None => None,
4029 };
4030 let cascade = if dictionary.is_none() && global_codes.is_none() {
4034 cascaded(&flat, ty, packed.as_ref())?
4035 } else {
4036 None
4037 };
4038 out.push(if coded.is_some() {
4039 4
4040 } else if cascade.is_some() {
4041 5
4042 } else if global_codes.is_some() {
4043 3
4044 } else if dictionary.is_some() {
4045 1
4046 } else if packed.is_some() {
4047 2
4048 } else {
4049 0
4050 });
4051 let nulls = flat.validity();
4052 let flag = match nulls {
4053 Validity::AllValid => 0,
4054 Validity::AllInvalid => 1,
4055 Validity::Mask(_) => 2,
4056 };
4057 out.push(flag);
4058 if flag == 2 {
4059 for group in (0..vector.len()).step_by(8) {
4060 let mut bits = 0_u8;
4061 for bit in 0..8 {
4062 if group + bit < vector.len() && !flat.is_null_at(group + bit) {
4063 bits |= 1 << bit;
4064 }
4065 }
4066 out.push(bits);
4067 }
4068 }
4069 if let Some(coded) = coded {
4070 out.extend_from_slice(&coded);
4071 return Ok((out, membership));
4072 }
4073 if let Some(cascade) = cascade {
4074 out.extend_from_slice(&cascade);
4075 return Ok((out, membership));
4076 }
4077 if let Some(codes) = global_codes {
4078 for code in codes {
4079 put_u32(&mut out, code);
4080 }
4081 return Ok((out, membership));
4082 }
4083 if let Some(dictionary) = dictionary {
4084 out.extend_from_slice(&dictionary);
4085 return Ok((out, membership));
4086 }
4087 if let Some(packed) = packed {
4088 if packed.offset() != 0 {
4089 return Err(invalid("writer received a sliced packed vector"));
4090 }
4091 out.push(u8::try_from(packed.width()).map_err(|_| invalid("packed width overflow"))?);
4092 out.extend_from_slice(&packed.base().to_le_bytes());
4093 put_u32(
4094 &mut out,
4095 u32::try_from(packed.words().len()).map_err(|_| invalid("too many packed words"))?,
4096 );
4097 for word in packed.words() {
4098 put_u64(&mut out, *word);
4099 }
4100 return Ok((out, membership));
4101 }
4102 let data = flat.data().ok_or_else(|| invalid("scalar column did not flatten"))?;
4103 match (ty, data) {
4104 (LogicalType::TinyInt, Data::Int8(values)) => {
4105 for value in &**values {
4106 out.extend_from_slice(&value.to_le_bytes());
4107 }
4108 }
4109 (LogicalType::UTinyInt, Data::UInt8(values)) => {
4110 for value in &**values {
4111 out.extend_from_slice(&value.to_le_bytes());
4112 }
4113 }
4114 (LogicalType::SmallInt, Data::Int16(values)) => {
4115 for value in &**values {
4116 out.extend_from_slice(&value.to_le_bytes());
4117 }
4118 }
4119 (LogicalType::USmallInt, Data::UInt16(values)) => {
4120 for value in &**values {
4121 out.extend_from_slice(&value.to_le_bytes());
4122 }
4123 }
4124 (LogicalType::UInteger, Data::UInt32(values)) => {
4125 for value in &**values {
4126 out.extend_from_slice(&value.to_le_bytes());
4127 }
4128 }
4129 (LogicalType::UBigInt, Data::UInt64(values)) => {
4130 for value in &**values {
4131 out.extend_from_slice(&value.to_le_bytes());
4132 }
4133 }
4134 (LogicalType::Integer | LogicalType::Date, Data::Int32(values)) => {
4135 for value in &**values {
4136 out.extend_from_slice(&value.to_le_bytes());
4137 }
4138 }
4139 (LogicalType::BigInt | LogicalType::Timestamp, Data::Int64(values)) => {
4140 for value in &**values {
4141 out.extend_from_slice(&value.to_le_bytes());
4142 }
4143 }
4144 (LogicalType::Boolean, Data::Bool(values)) => {
4145 for value in &**values {
4146 out.push(u8::from(*value));
4147 }
4148 }
4149 (LogicalType::Varchar, Data::Varlen(values)) => {
4150 let mut bytes = Vec::new();
4151 put_u32(&mut out, 0);
4152 for row in 0..vector.len() {
4153 let value = values.bytes(row).ok_or_else(|| invalid("string view is invalid"))?;
4154 bytes.extend_from_slice(value);
4155 put_u32(
4156 &mut out,
4157 u32::try_from(bytes.len())
4158 .map_err(|_| invalid("string payload exceeds 4GiB"))?,
4159 );
4160 }
4161 out.extend_from_slice(&bytes);
4162 }
4163 _ => return Err(Error::not_implemented(format!("native page for {ty}"))),
4164 }
4165 Ok((out, membership))
4166}
4167
4168fn put_varint(out: &mut Vec<u8>, mut value: u32) {
4169 while value >= 0x80 {
4170 out.push((value as u8 & 0x7f) | 0x80);
4171 value >>= 7;
4172 }
4173 out.push(value as u8);
4174}
4175
4176fn unique_codes(codes: &[u32]) -> Vec<u32> {
4178 let mut unique = codes.to_vec();
4179 unique.sort_unstable();
4180 unique.dedup();
4181 unique
4182}
4183
4184fn merged_codes(lists: Vec<Vec<u32>>) -> Vec<u32> {
4190 let mut lists = lists;
4191 while lists.len() > 1 {
4192 let mut next = Vec::with_capacity(lists.len().div_ceil(2));
4193 for pair in lists.chunks(2) {
4194 match pair {
4195 [left, right] => next.push(merged_pair(left, right)),
4196 [only] => next.push(only.clone()),
4197 _ => {}
4198 }
4199 }
4200 lists = next;
4201 }
4202 lists.pop().unwrap_or_default()
4203}
4204
4205fn merged_pair(left: &[u32], right: &[u32]) -> Vec<u32> {
4206 let mut out = Vec::with_capacity(left.len().saturating_add(right.len()));
4207 let mut at = 0;
4208 let mut to = 0;
4209 while at < left.len() && to < right.len() {
4210 match left[at].cmp(&right[to]) {
4211 Ordering::Less => {
4212 out.push(left[at]);
4213 at += 1;
4214 }
4215 Ordering::Greater => {
4216 out.push(right[to]);
4217 to += 1;
4218 }
4219 Ordering::Equal => {
4220 out.push(left[at]);
4221 at += 1;
4222 to += 1;
4223 }
4224 }
4225 }
4226 out.extend_from_slice(&left[at..]);
4227 out.extend_from_slice(&right[to..]);
4228 out
4229}
4230
4231fn merged_range(ranges: impl Iterator<Item = Range>) -> Range {
4236 let mut merged = Range::default();
4237 let mut first = true;
4238 for range in ranges {
4239 merged.nulls = merged.nulls.saturating_add(range.nulls);
4240 merged.sum = match (merged.sum.take(), range.sum) {
4244 (Some(held), Some(next)) if !first => held.checked_add(next),
4245 (_, next) if first => next,
4246 _ => None,
4247 };
4248 merged.exact = if first { range.exact } else { merged.exact && range.exact };
4249 if first {
4250 merged.low = range.low;
4251 merged.high = range.high;
4252 first = false;
4253 continue;
4254 }
4255 merged.low = match (merged.low.take(), range.low) {
4256 (Some(held), Some(next)) => Some(held.smaller(next)),
4257 _ => None,
4258 };
4259 merged.high = match (merged.high.take(), range.high) {
4260 (Some(held), Some(next)) => Some(held.larger(next)),
4261 _ => None,
4262 };
4263 }
4264 merged
4265}
4266
4267fn shortened(bound: Option<Bound>, high: bool) -> Option<Bound> {
4280 match bound {
4281 Some(Bound::Bytes(mut value)) if value.len() > PART_BOUND_BYTES => {
4282 value.truncate(PART_BOUND_BYTES);
4283 if !high {
4284 return Some(Bound::Bytes(value));
4285 }
4286 while let Some(last) = value.pop() {
4287 if last < u8::MAX {
4288 value.push(last + 1);
4289 return Some(Bound::Bytes(value));
4290 }
4291 }
4292 None
4293 }
4294 other => other,
4295 }
4296}
4297
4298fn encode_part_ranges(ranges: &[Range]) -> Result<Vec<u8>> {
4306 let mut out = Vec::new();
4307 put_u32(
4308 &mut out,
4309 u32::try_from(ranges.len()).map_err(|_| invalid("too many parts in a stripe"))?,
4310 );
4311 for range in ranges {
4312 put_bound(&mut out, shortened(range.low.clone(), false).as_ref())?;
4313 put_bound(&mut out, shortened(range.high.clone(), true).as_ref())?;
4314 put_u32(&mut out, u32::try_from(range.nulls).map_err(|_| invalid("null count overflow"))?);
4315 }
4316 Ok(out)
4317}
4318
4319fn decode_part_ranges(bytes: &[u8]) -> Result<Vec<Range>> {
4321 let mut cur = Cursor { bytes, at: 0 };
4322 let parts = cur.u32()? as usize;
4323 let mut out = Vec::new();
4324 for _ in 0..parts {
4325 let low = cur.bound()?;
4326 let high = cur.bound()?;
4327 let nulls = cur.u32()? as usize;
4328 out.push(Range { low, high, nulls, exact: false, sum: None });
4329 }
4330 Ok(out)
4331}
4332
4333fn encode_sieves<'a>(sieves: impl Iterator<Item = &'a Option<Sieve>>) -> Result<Vec<u8>> {
4334 let held: Vec<&Option<Sieve>> = sieves.collect();
4335 let mut out = Vec::new();
4336 put_u32(
4337 &mut out,
4338 u32::try_from(held.len()).map_err(|_| invalid("too many parts in a stripe"))?,
4339 );
4340 for sieve in &held {
4341 let length = sieve.as_ref().map_or(0, Sieve::len);
4342 put_u32(&mut out, u32::try_from(length).map_err(|_| invalid("sieve length overflow"))?);
4343 }
4344 for sieve in held.into_iter().flatten() {
4346 out.extend_from_slice(&sieve.to_bytes());
4347 }
4348 Ok(out)
4349}
4350
4351fn decode_sieves(bytes: &[u8]) -> Result<Vec<Option<Sieve>>> {
4357 let parts = u32::from_le_bytes(
4358 bytes
4359 .get(..4)
4360 .ok_or_else(|| invalid("sieve page is truncated"))?
4361 .try_into()
4362 .map_err(|_| invalid("sieve page is truncated"))?,
4363 ) as usize;
4364 let mut lengths = Vec::with_capacity(parts);
4365 for part in 0..parts {
4366 let at = 4 + part * 4;
4367 let field = bytes.get(at..at + 4).ok_or_else(|| invalid("sieve page is truncated"))?;
4368 lengths.push(u32::from_le_bytes(
4369 field.try_into().map_err(|_| invalid("sieve page is truncated"))?,
4370 ) as usize);
4371 }
4372 let mut at = 4 + parts * 4;
4373 let mut out = Vec::with_capacity(parts);
4374 for length in lengths {
4375 if length == 0 {
4376 out.push(None);
4377 continue;
4378 }
4379 let end = at.checked_add(length).ok_or_else(|| invalid("sieve page is truncated"))?;
4380 let field = bytes.get(at..end).ok_or_else(|| invalid("sieve page is truncated"))?;
4381 out.push(Sieve::from_bytes(field));
4382 at = end;
4383 }
4384 if at != bytes.len() {
4385 return Err(invalid("sieve page has trailing bytes"));
4386 }
4387 Ok(out)
4388}
4389
4390fn encode_membership(unique: &[u32]) -> Vec<u8> {
4396 let mut out = Vec::with_capacity(unique.len().saturating_mul(2).saturating_add(5));
4397 put_varint(&mut out, u32::try_from(unique.len()).unwrap_or(u32::MAX));
4398 let mut previous = 0;
4399 for (at, &code) in unique.iter().enumerate() {
4400 put_varint(&mut out, if at == 0 { code } else { code - previous });
4401 previous = code;
4402 }
4403 out
4404}
4405
4406fn take_varint(bytes: &[u8], at: &mut usize) -> Result<u32> {
4407 let mut value = 0_u32;
4408 for shift in (0..35).step_by(7) {
4409 let byte = *bytes.get(*at).ok_or_else(|| invalid("membership varint is truncated"))?;
4410 *at += 1;
4411 let part = u32::from(byte & 0x7f);
4412 if shift == 28 && part > 0x0f {
4413 return Err(invalid("membership varint overflow"));
4414 }
4415 value = value
4416 .checked_add(
4417 part.checked_shl(shift).ok_or_else(|| invalid("membership varint overflow"))?,
4418 )
4419 .ok_or_else(|| invalid("membership varint overflow"))?;
4420 if byte & 0x80 == 0 {
4421 return Ok(value);
4422 }
4423 }
4424 Err(invalid("membership varint is too long"))
4425}
4426
4427fn decode_membership(bytes: &[u8]) -> Result<Vec<u32>> {
4428 let mut at = 0;
4429 let count = take_varint(bytes, &mut at)? as usize;
4430 let mut codes = Vec::with_capacity(count);
4431 let mut previous = 0_u32;
4432 for index in 0..count {
4433 let delta = take_varint(bytes, &mut at)?;
4434 let code = if index == 0 {
4435 delta
4436 } else {
4437 previous.checked_add(delta).ok_or_else(|| invalid("membership code overflow"))?
4438 };
4439 if index > 0 && code <= previous {
4440 return Err(invalid("membership codes are not increasing"));
4441 }
4442 codes.push(code);
4443 previous = code;
4444 }
4445 if at != bytes.len() {
4446 return Err(invalid("membership page has trailing bytes"));
4447 }
4448 Ok(codes)
4449}
4450
4451fn string_dictionary(vector: &Vector) -> Result<Option<Vec<u8>>> {
4452 let mut by_text = HashMap::new();
4453 let mut values = Vec::new();
4454 let mut codes = Vec::with_capacity(vector.len());
4455 let mut plain_bytes = 0_usize;
4456 for row in 0..vector.len() {
4457 let text = vector.text_at(row).unwrap_or("");
4458 plain_bytes = plain_bytes.saturating_add(text.len());
4459 let code = match by_text.get(text) {
4460 Some(&code) => code,
4461 None => {
4462 let code = u32::try_from(values.len())
4463 .map_err(|_| invalid("too many dictionary values"))?;
4464 by_text.insert(text, code);
4465 values.push(text);
4466 code
4467 }
4468 };
4469 codes.push(code);
4470 }
4471 let dictionary_bytes = values.iter().map(|value| value.len()).sum::<usize>();
4472 let encoded = 8_usize
4473 .saturating_add((values.len() + 1).saturating_mul(4))
4474 .saturating_add(dictionary_bytes)
4475 .saturating_add(codes.len().saturating_mul(4));
4476 let plain = (vector.len() + 1).saturating_mul(4).saturating_add(plain_bytes);
4477 if encoded >= plain {
4478 return Ok(None);
4479 }
4480 let mut out = Vec::with_capacity(encoded);
4481 put_u32(
4482 &mut out,
4483 u32::try_from(values.len()).map_err(|_| invalid("too many dictionary values"))?,
4484 );
4485 put_u32(
4486 &mut out,
4487 u32::try_from(dictionary_bytes).map_err(|_| invalid("dictionary payload exceeds 4GiB"))?,
4488 );
4489 let mut offset = 0_u32;
4490 put_u32(&mut out, offset);
4491 for value in &values {
4492 offset = offset
4493 .checked_add(
4494 u32::try_from(value.len()).map_err(|_| invalid("dictionary value is too long"))?,
4495 )
4496 .ok_or_else(|| invalid("dictionary payload exceeds 4GiB"))?;
4497 put_u32(&mut out, offset);
4498 }
4499 for value in values {
4500 out.extend_from_slice(value.as_bytes());
4501 }
4502 for code in codes {
4503 put_u32(&mut out, code);
4504 }
4505 Ok(Some(out))
4506}
4507
4508struct EncodedDictionary {
4509 index: Vec<u8>,
4510 ranks: Vec<u8>,
4511 payload: Vec<Vec<u8>>,
4514}
4515
4516fn head(bytes: &[u8]) -> u64 {
4518 let mut word = [0; 8];
4519 let take = bytes.len().min(8);
4520 word[..take].copy_from_slice(&bytes[..take]);
4521 u64::from_be_bytes(word)
4522}
4523
4524fn rankings(dictionaries: &[Option<GlobalDictionary>]) -> Result<Vec<Vec<(u64, u32)>>> {
4532 let present =
4533 dictionaries.iter().enumerate().filter(|(_, held)| held.is_some()).map(|(at, _)| at);
4534 let present = present.collect::<Vec<_>>();
4535 let mut orders = vec![Vec::new(); dictionaries.len()];
4536 let workers = std::thread::available_parallelism()
4537 .map_or(1, usize::from)
4538 .min(MAX_FREQUENCY_WORKERS)
4539 .min(present.len());
4540 if workers <= 1 {
4541 for at in present {
4542 if let Some(dictionary) = &dictionaries[at] {
4543 orders[at] = dictionary.ranked();
4544 }
4545 }
4546 return Ok(orders);
4547 }
4548 let width = present.len().div_ceil(workers);
4549 let pieces = std::thread::scope(|scope| {
4550 present
4551 .chunks(width)
4552 .map(|columns| {
4553 scope.spawn(|| {
4554 columns
4555 .iter()
4556 .filter_map(|&at| dictionaries[at].as_ref().map(|held| (at, held.ranked())))
4557 .collect::<Vec<_>>()
4558 })
4559 })
4560 .collect::<Vec<_>>()
4561 .into_iter()
4562 .map(|handle| {
4563 handle.join().map_err(|_| Error::internal("a dictionary sort worker panicked"))
4564 })
4565 .collect::<Result<Vec<_>>>()
4566 })?;
4567 for piece in pieces {
4568 for (at, order) in piece {
4569 orders[at] = order;
4570 }
4571 }
4572 Ok(orders)
4573}
4574
4575fn encode_global_dictionary(
4576 dictionary: GlobalDictionary,
4577 order: &[(u64, u32)],
4578) -> Result<EncodedDictionary> {
4579 let values = dictionary.offsets.len() - 1;
4580 if order.len() != values {
4581 return Err(invalid("global dictionary order does not cover its values"));
4582 }
4583 let blocks = values.div_ceil(TEXT_PAYLOAD_VALUES);
4584 let payload = encode_payload(&dictionary)?;
4585 if payload.len() != blocks {
4586 return Err(invalid("global dictionary payload is not the blocks it says it is"));
4587 }
4588 let (ranks, rank_ends) = encode_ranks(order, code_width(values))?;
4589 let rank_blocks = values.div_ceil(TEXT_RANK_BLOCK);
4590 let offset_bits = offset_width(&dictionary.offsets);
4591 let mut index = Vec::with_capacity(
4592 DICTIONARY_HEADER + offset_bytes(values, offset_bits) + (blocks + rank_blocks) * 16,
4593 );
4594 put_u32(
4595 &mut index,
4596 u32::try_from(values).map_err(|_| invalid("global dictionary has too many values"))?,
4597 );
4598 put_u32(&mut index, TEXT_PAYLOAD_VALUES as u32);
4599 put_u32(
4600 &mut index,
4601 u32::try_from(blocks).map_err(|_| invalid("global dictionary has too many blocks"))?,
4602 );
4603 put_u32(&mut index, offset_bits as u32);
4604 encode_offsets(&dictionary.offsets, offset_bits, &mut index)?;
4605 let mut at = 0_u64;
4609 for block in &payload {
4610 at = at
4611 .checked_add(block.len() as u64)
4612 .ok_or_else(|| invalid("global dictionary payload overflow"))?;
4613 put_u64(&mut index, at);
4614 }
4615 for block in &payload {
4616 put_u64(&mut index, checksum(block));
4617 }
4618 if rank_ends.len() != rank_blocks {
4621 return Err(invalid("global dictionary order is not the blocks it says it is"));
4622 }
4623 for end in &rank_ends {
4624 put_u64(&mut index, *end);
4625 }
4626 let mut at = 0_usize;
4627 for end in &rank_ends {
4628 let end = usize::try_from(*end).map_err(|_| invalid("global dictionary order overflow"))?;
4629 put_u64(&mut index, checksum(&ranks[at..end]));
4630 at = end;
4631 }
4632 Ok(EncodedDictionary { index, ranks, payload })
4633}
4634
4635const PAYLOAD_SAMPLE_BLOCKS: usize = 8;
4642
4643fn payload_shapes() -> Vec<chooser::Settled> {
4669 let integers = vec![integer::Kind::Packed];
4670 [
4671 vec![string::Kind::Front, string::Kind::Lz],
4672 vec![string::Kind::Lz, string::Kind::Fsst],
4673 vec![string::Kind::Lz, string::Kind::Plain],
4674 vec![string::Kind::Fsst],
4675 vec![string::Kind::Plain],
4676 ]
4677 .into_iter()
4678 .map(|strings| chooser::Settled::new(strings, integers.clone()))
4679 .collect()
4680}
4681
4682fn encode_payload(dictionary: &GlobalDictionary) -> Result<Vec<Vec<u8>>> {
4688 let values = dictionary.offsets.len() - 1;
4689 let blocks = values.div_ceil(TEXT_PAYLOAD_VALUES);
4690 let run = |block: usize| {
4691 let first = block * TEXT_PAYLOAD_VALUES;
4692 let last = (first + TEXT_PAYLOAD_VALUES).min(values);
4693 (first..last)
4694 .map(|value| {
4695 let from = dictionary.offsets[value] as usize;
4696 let to = dictionary.offsets[value + 1] as usize;
4697 &dictionary.payload[from..to]
4698 })
4699 .collect::<Vec<_>>()
4700 };
4701 let shape = (blocks > PAYLOAD_SAMPLE_BLOCKS).then(|| settle_shape(&run, blocks)).transpose()?;
4704 let one = |block: usize| match &shape {
4705 Some(shape) => string::encode_with(&run(block), shape),
4706 None => string::encode(&run(block)),
4707 };
4708 let workers = std::thread::available_parallelism()
4709 .map_or(1, usize::from)
4710 .min(MAX_FREQUENCY_WORKERS)
4711 .min(blocks);
4712 if workers <= 1 {
4713 return (0..blocks).map(one).collect();
4714 }
4715 let next = AtomicUsize::new(0);
4716 let pieces = std::thread::scope(|scope| {
4717 (0..workers)
4718 .map(|_| {
4719 scope.spawn(|| {
4720 let mut mine = Vec::new();
4721 loop {
4722 let block = next.fetch_add(1, Atomic::Relaxed);
4723 if block >= blocks {
4724 break;
4725 }
4726 mine.push((block, one(block)?));
4727 }
4728 Ok(mine)
4729 })
4730 })
4731 .collect::<Vec<_>>()
4732 .into_iter()
4733 .map(|handle| {
4734 handle.join().map_err(|_| Error::internal("a dictionary encode worker panicked"))?
4735 })
4736 .collect::<Result<Vec<_>>>()
4737 })?;
4738 let mut payload = vec![Vec::new(); blocks];
4739 for piece in pieces {
4740 for (block, bytes) in piece {
4741 payload[block] = bytes;
4742 }
4743 }
4744 Ok(payload)
4745}
4746
4747fn settle_shape<'a>(
4755 run: &dyn Fn(usize) -> Vec<&'a [u8]>,
4756 blocks: usize,
4757) -> Result<chooser::Settled> {
4758 let last = blocks - 1;
4759 let sample = (0..PAYLOAD_SAMPLE_BLOCKS)
4760 .map(|region| run(region * last / (PAYLOAD_SAMPLE_BLOCKS - 1)))
4761 .collect::<Vec<_>>();
4762 let mut best: Option<(chooser::Settled, usize)> = None;
4763 for shape in payload_shapes() {
4764 let mut size = 0;
4765 for block in &sample {
4766 size += string::encode_with(block, &shape)?.len();
4767 }
4768 if best.as_ref().is_none_or(|(_, smallest)| size < *smallest) {
4769 best = Some((shape, size));
4770 }
4771 }
4772 best.map(|(shape, _)| shape)
4773 .ok_or_else(|| invalid("no shape applies to a global dictionary payload"))
4774}
4775
4776fn encode_ranks(order: &[(u64, u32)], code_bits: usize) -> Result<(Vec<u8>, Vec<u64>)> {
4783 let mut out = Vec::with_capacity(order.len() * 4);
4784 let mut ends = Vec::with_capacity(order.len().div_ceil(TEXT_RANK_BLOCK));
4785 let mut heads = Vec::with_capacity(TEXT_RANK_BLOCK);
4786 let mut codes = Vec::with_capacity(TEXT_RANK_BLOCK);
4787 for block in order.chunks(TEXT_RANK_BLOCK) {
4788 let base = block.first().map_or(0, |&(head, _)| head);
4791 let span = block.last().map_or(0, |&(head, _)| head.wrapping_sub(base));
4792 let width = (u64::BITS - span.leading_zeros()) as usize;
4793 heads.clear();
4794 codes.clear();
4795 for &(head, code) in block {
4796 heads.push(head.wrapping_sub(base));
4797 codes.push(u64::from(code));
4798 }
4799 put_u64(&mut out, base);
4800 out.push(width as u8);
4801 bitpack::pack_tail(&heads, width, &mut out)
4802 .map_err(|_| invalid("global dictionary heads do not pack"))?;
4803 bitpack::pack_tail(&codes, code_bits, &mut out)
4804 .map_err(|_| invalid("global dictionary codes do not pack"))?;
4805 ends.push(out.len() as u64);
4806 }
4807 Ok((out, ends))
4808}
4809
4810fn open_global_dictionary(
4817 file: Arc<File>,
4818 page: Page,
4819 ty: &LogicalType,
4820 keep_budget: usize,
4821) -> Result<Vector> {
4822 if ty != &LogicalType::Varchar {
4823 return Err(invalid("global dictionary belongs to a non-string column"));
4824 }
4825 let mut header = [0; DICTIONARY_HEADER];
4826 read_at(&file, page.offset, &mut header)?;
4827 let count = u32::from_le_bytes(header[0..4].try_into().expect("four bytes")) as usize;
4828 let per_block = u32::from_le_bytes(header[4..8].try_into().expect("four bytes")) as usize;
4829 let blocks = u32::from_le_bytes(header[8..12].try_into().expect("four bytes")) as usize;
4830 let offset_bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
4831 if per_block != TEXT_PAYLOAD_VALUES {
4832 return Err(invalid("global dictionary block width differs"));
4833 }
4834 if blocks != count.div_ceil(TEXT_PAYLOAD_VALUES) {
4835 return Err(invalid("global dictionary block count differs from its value count"));
4836 }
4837 if offset_bits > u32::BITS as usize {
4838 return Err(invalid("global dictionary packs offsets past a payload"));
4839 }
4840 let offset_len = offset_bytes(count, offset_bits);
4841 let ranks = count;
4846 let rank_blocks = ranks.div_ceil(TEXT_RANK_BLOCK);
4847 let hash_len = blocks
4850 .checked_add(rank_blocks)
4851 .and_then(|words| words.checked_mul(16))
4852 .ok_or_else(|| invalid("global dictionary block count overflow"))?;
4853 let index_len = DICTIONARY_HEADER
4854 .checked_add(offset_len)
4855 .and_then(|len| len.checked_add(hash_len))
4856 .ok_or_else(|| invalid("global dictionary header overflow"))?;
4857 if index_len > page.length as usize {
4858 return Err(invalid("global dictionary offset index exceeds its page"));
4859 }
4860 let mut index = vec![0; index_len];
4861 index[..DICTIONARY_HEADER].copy_from_slice(&header);
4862 read_at(&file, page.offset + DICTIONARY_HEADER as u64, &mut index[DICTIONARY_HEADER..])?;
4863 if checksum(&index) != page.hash {
4864 return Err(invalid("global dictionary index checksum differs"));
4865 }
4866 let offsets = index[DICTIONARY_HEADER..DICTIONARY_HEADER + offset_len].to_vec();
4867 let mut words = index[DICTIONARY_HEADER + offset_len..]
4868 .chunks_exact(8)
4869 .map(|part| u64::from_le_bytes(part.try_into().expect("eight bytes")))
4870 .collect::<Vec<_>>();
4871 let mut hashes = words.split_off(blocks);
4872 let mut rank_ends = hashes.split_off(blocks);
4873 let rank_hashes = rank_ends.split_off(rank_blocks);
4874 let ends = words;
4875 if rank_ends.windows(2).any(|pair| pair[0] >= pair[1]) {
4878 return Err(invalid("global dictionary order blocks do not rise"));
4879 }
4880 let rank_len = usize::try_from(rank_ends.last().copied().unwrap_or_default())
4881 .map_err(|_| invalid("global dictionary rank overflow"))?;
4882 let body_len = index_len
4883 .checked_add(rank_len)
4884 .ok_or_else(|| invalid("global dictionary header overflow"))?;
4885 if body_len > page.length as usize {
4886 return Err(invalid("global dictionary order exceeds its page"));
4887 }
4888 let stored_len = page.length as usize - body_len;
4891 if ends.last().copied().unwrap_or_default() as usize != stored_len
4892 || ends.windows(2).any(|pair| pair[0] > pair[1])
4893 {
4894 return Err(invalid("global dictionary blocks do not bound the payload"));
4895 }
4896 Vector::external_text(
4897 LogicalType::Varchar,
4898 Arc::new(NativeText {
4899 file,
4900 values: count,
4901 offsets,
4902 offset_bits,
4903 ranks,
4904 rank_at: page.offset + index_len as u64,
4905 rank_ends,
4906 rank_hashes,
4907 rank_blocks: (0..rank_blocks).map(|_| OnceLock::new()).collect(),
4908 code_bits: code_width(count),
4909 code_ranks: OnceLock::new(),
4910 payload: page.offset + body_len as u64,
4911 ends,
4912 hashes,
4913 blocks: (0..blocks).map(|_| OnceLock::new()).collect(),
4914 keep_budget,
4915 payload_kept: AtomicUsize::new(0),
4916 }),
4917 )
4918}
4919
4920fn decode(
4921 ty: &LogicalType,
4922 rows: usize,
4923 bytes: &[u8],
4924 global: Option<Arc<Vector>>,
4925) -> Result<Vector> {
4926 let mut cur = Cursor { bytes, at: 0 };
4927 let codec = cur.u8()?;
4928 let flag = cur.u8()?;
4929 let validity = match flag {
4930 0 => Validity::AllValid,
4931 1 => Validity::AllInvalid,
4932 2 => {
4933 let mask = cur.take(rows.div_ceil(8))?;
4934 Validity::from_iter(rows, |row| mask[row / 8] >> (row % 8) & 1 == 1)
4935 }
4936 _ => return Err(invalid("page validity tag differs")),
4937 };
4938 if codec == 1 {
4939 if ty != &LogicalType::Varchar {
4940 return Err(invalid("dictionary codec belongs to a non-string page"));
4941 }
4942 let count = cur.u32()? as usize;
4943 let payload_len = cur.u32()? as usize;
4944 let offset_bytes = cur.take(
4945 (count + 1)
4946 .checked_mul(4)
4947 .ok_or_else(|| invalid("dictionary offset count overflow"))?,
4948 )?;
4949 let offsets = offset_bytes
4950 .chunks_exact(4)
4951 .map(|part| u32::from_le_bytes(part.try_into().expect("four bytes")))
4952 .collect::<Vec<_>>();
4953 let payload = cur.take(payload_len)?.to_vec();
4954 if offsets.first() != Some(&0)
4955 || offsets.last().copied().map(|last| last as usize) != Some(payload.len())
4956 || offsets.windows(2).any(|pair| pair[0] > pair[1])
4957 {
4958 return Err(invalid("dictionary offsets do not bound the payload"));
4959 }
4960 let mut strings = StringColumn::over(Buffer::from_vec(payload));
4961 for pair in offsets.windows(2) {
4962 strings.push_in_place(pair[0] as usize, (pair[1] - pair[0]) as usize)?;
4963 }
4964 let mut codes = Vec::with_capacity(rows);
4965 for _ in 0..rows {
4966 codes.push(cur.u32()?);
4967 }
4968 if codes.iter().any(|code| *code as usize >= count) {
4969 return Err(invalid("dictionary code is out of range"));
4970 }
4971 if cur.at != bytes.len() {
4972 return Err(invalid("dictionary page has trailing bytes"));
4973 }
4974 let dictionary = Vector::flat(LogicalType::Varchar, Data::Varlen(strings))?;
4975 return Ok(Vector::dictionary(codes, dictionary)?.with_validity(validity));
4976 }
4977 if codec == 3 || codec == 4 {
4978 let dictionary = global.ok_or_else(|| invalid("global code page has no dictionary"))?;
4979 let codes = if codec == 4 {
4980 let wide = integer::decode(&bytes[cur.at..])?;
4983 if wide.len() != rows {
4984 return Err(invalid("encoded code page holds the wrong number of rows"));
4985 }
4986 let mut codes = Vec::with_capacity(wide.len());
4993 let mut seen = 0_i64;
4994 for &code in &wide {
4995 seen |= code;
4996 codes.push(code as u32);
4997 }
4998 if seen < 0 || seen > i64::from(u32::MAX) {
4999 return Err(invalid("code is not a code"));
5000 }
5001 codes
5002 } else {
5003 let mut codes = Vec::with_capacity(rows);
5004 for _ in 0..rows {
5005 codes.push(cur.u32()?);
5006 }
5007 if cur.at != bytes.len() {
5008 return Err(invalid("global code page has trailing bytes"));
5009 }
5010 codes
5011 };
5012 let highest = codes.iter().copied().max();
5013 return Ok(Vector::stable_dictionary_validated(codes, dictionary, highest)?
5014 .with_validity(validity));
5015 }
5016 if codec == 5 {
5017 let values = integer::decode(&bytes[cur.at..])?;
5019 if values.len() != rows {
5020 return Err(invalid("cascade page holds the wrong number of rows"));
5021 }
5022 let data = narrowed(ty, values)?;
5023 return Ok(Vector::flat(ty.clone(), data)?.with_validity(validity));
5024 }
5025 if codec == 2 {
5026 let width = u32::from(cur.u8()?);
5027 let base = i128::from_le_bytes(cur.take(16)?.try_into().expect("sixteen bytes"));
5028 let count = cur.u32()? as usize;
5029 let mut words = Vec::with_capacity(count);
5030 for _ in 0..count {
5031 words.push(cur.u64()?);
5032 }
5033 if cur.at != bytes.len() {
5034 return Err(invalid("packed page has trailing bytes"));
5035 }
5036 return Ok(Vector::packed(ty.clone(), words, width, base, rows)?.with_validity(validity));
5037 }
5038 if codec != 0 {
5039 return Err(invalid("page codec is unknown"));
5040 }
5041 let data = match ty {
5042 LogicalType::TinyInt => {
5043 let values = cur.take(rows)?;
5044 Data::Int8(values.iter().map(|item| *item as i8).collect::<Vec<_>>().into())
5045 }
5046 LogicalType::UTinyInt => Data::UInt8(cur.take(rows)?.to_vec().into()),
5047 LogicalType::SmallInt => {
5048 let values =
5049 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5050 Data::Int16(
5051 values
5052 .chunks_exact(2)
5053 .map(|item| i16::from_le_bytes(item.try_into().expect("two bytes")))
5054 .collect::<Vec<_>>()
5055 .into(),
5056 )
5057 }
5058 LogicalType::USmallInt => {
5059 let values =
5060 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5061 Data::UInt16(
5062 values
5063 .chunks_exact(2)
5064 .map(|item| u16::from_le_bytes(item.try_into().expect("two bytes")))
5065 .collect::<Vec<_>>()
5066 .into(),
5067 )
5068 }
5069 LogicalType::UInteger => {
5070 let values =
5071 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5072 Data::UInt32(
5073 values
5074 .chunks_exact(4)
5075 .map(|item| u32::from_le_bytes(item.try_into().expect("four bytes")))
5076 .collect::<Vec<_>>()
5077 .into(),
5078 )
5079 }
5080 LogicalType::UBigInt => {
5081 let values =
5082 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5083 Data::UInt64(
5084 values
5085 .chunks_exact(8)
5086 .map(|item| u64::from_le_bytes(item.try_into().expect("eight bytes")))
5087 .collect::<Vec<_>>()
5088 .into(),
5089 )
5090 }
5091 LogicalType::Integer | LogicalType::Date => {
5092 let values =
5093 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5094 Data::Int32(
5095 values
5096 .chunks_exact(4)
5097 .map(|item| i32::from_le_bytes(item.try_into().expect("four bytes")))
5098 .collect::<Vec<_>>()
5099 .into(),
5100 )
5101 }
5102 LogicalType::BigInt | LogicalType::Timestamp => {
5103 let values =
5104 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5105 Data::Int64(
5106 values
5107 .chunks_exact(8)
5108 .map(|item| i64::from_le_bytes(item.try_into().expect("eight bytes")))
5109 .collect::<Vec<_>>()
5110 .into(),
5111 )
5112 }
5113 LogicalType::Boolean => {
5114 let values = cur.take(rows)?;
5115 if values.iter().any(|value| *value > 1) {
5116 return Err(invalid("boolean page has another value"));
5117 }
5118 Data::Bool(values.iter().map(|value| *value == 1).collect::<Vec<_>>().into())
5119 }
5120 LogicalType::Varchar => {
5121 let offset_bytes = cur
5122 .take((rows + 1).checked_mul(4).ok_or_else(|| invalid("offset count overflow"))?)?;
5123 let offsets = offset_bytes
5124 .chunks_exact(4)
5125 .map(|part| u32::from_le_bytes(part.try_into().expect("four bytes")))
5126 .collect::<Vec<_>>();
5127 let payload = cur.take(bytes.len() - cur.at)?.to_vec();
5128 if offsets.first() != Some(&0)
5129 || offsets.last().copied().map(|last| last as usize) != Some(payload.len())
5130 || offsets.windows(2).any(|pair| pair[0] > pair[1])
5131 {
5132 return Err(invalid("string offsets do not bound the payload"));
5133 }
5134 let mut values = StringColumn::over(Buffer::from_vec(payload));
5135 for pair in offsets.windows(2) {
5136 values.push_in_place(pair[0] as usize, (pair[1] - pair[0]) as usize)?;
5137 }
5138 Data::Varlen(values)
5139 }
5140 _ => return Err(Error::not_implemented(format!("native page for {ty}"))),
5141 };
5142 if cur.at != bytes.len() {
5143 return Err(invalid("page has trailing bytes"));
5144 }
5145 Ok(Vector::flat(ty.clone(), data)?.with_validity(validity))
5146}
5147
5148#[cfg(test)]
5149mod tests {
5150 use std::fs;
5151 use std::io::{Seek, SeekFrom, Write};
5152 use std::path::PathBuf;
5153 use std::time::{SystemTime, UNIX_EPOCH};
5154
5155 use rudb_common::Value;
5156 use rudb_common::bounds::Op;
5157
5158 use super::*;
5159
5160 #[test]
5161 fn checksum_matches_fixed_vectors() {
5162 assert_eq!(checksum(b""), 0xef46_db37_51d8_e999);
5163 assert_eq!(checksum(b"a"), 0xd24e_c4f1_a98c_6e5b);
5164 assert_eq!(checksum(b"abc"), 0x44bc_2cf5_ad77_0999);
5165 }
5166
5167 fn path(label: &str) -> PathBuf {
5168 let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
5169 std::env::temp_dir().join(format!("rudb-native-{label}-{}-{stamp}.rdb", std::process::id()))
5170 }
5171
5172 #[test]
5174 fn a_read_at_an_offset_ignores_where_another_thread_left_the_cursor() {
5175 const SPANS: usize = 64;
5176 const SPAN: usize = 512;
5177 let path = path("positional");
5178 let content: Vec<u8> =
5179 (0..SPANS).flat_map(|span| std::iter::repeat_n(span as u8, SPAN)).collect();
5180 fs::write(&path, &content).expect("the file is written");
5181 let file = Arc::new(File::open(&path).expect("the file opens"));
5182 std::thread::scope(|scope| {
5183 for _ in 0..8 {
5184 let file = Arc::clone(&file);
5185 scope.spawn(move || {
5186 for _ in 0..64 {
5187 for span in 0..SPANS {
5188 let mut bytes = [0_u8; SPAN];
5189 read_at(&file, (span * SPAN) as u64, &mut bytes)
5190 .expect("the span reads");
5191 assert!(
5192 bytes.iter().all(|byte| *byte == span as u8),
5193 "span {span} came back as {}",
5194 bytes[0],
5195 );
5196 }
5197 }
5198 });
5199 }
5200 });
5201 let mut past = [0_u8; SPAN];
5202 let end = (SPANS * SPAN) as u64;
5203 let error = read_at(&file, end, &mut past).expect_err("a read past the end is refused");
5204 assert!(error.message().contains("ends before its declared length"), "{error}");
5205 drop(file);
5206 let _ = fs::remove_file(&path);
5207 }
5208
5209 #[test]
5215 fn a_writer_puts_a_page_where_it_said_it_did_wherever_the_cursor_has_got_to() {
5216 let path = path("cursor");
5217 let mut writer = Writer::create(
5218 &path,
5219 "items",
5220 vec![
5221 Field::required("id", LogicalType::Integer),
5222 Field::new("text", LogicalType::Varchar),
5223 ],
5224 )
5225 .expect("new file");
5226 writer.append(&sample()).expect("first part");
5227 writer.file.seek(SeekFrom::Start(0)).expect("the cursor goes back to the header");
5228 writer.append(&sample()).expect("second part");
5229 writer.file.seek(SeekFrom::Start(1)).expect("and somewhere useless again");
5230 writer.finish().expect("commit");
5231 let reader = Reader::open(&path).expect("reopen from disk");
5232 assert_eq!(reader.table().rows(), 6);
5233 let ids = reader.read(0, &[0]).expect("the integer page reads back");
5234 assert_eq!(ids.value_at(0, 0), Value::Integer(4));
5235 assert_eq!(ids.value_at(2, 0), Value::Integer(-2));
5236 let text = reader.read(1, &[1]).expect("the text page reads back");
5237 assert_eq!(text.value_at(1, 0), Value::Null);
5238 assert_eq!(text.value_at(2, 0), Value::Varchar("long text after a slash".into()));
5239 let end = reader.table().stripes().iter().flat_map(|stripe| {
5242 stripe
5243 .pages
5244 .iter()
5245 .map(|page| page.offset + u64::from(page.length))
5246 .chain(std::iter::once(stripe.index.offset + u64::from(stripe.index.length)))
5247 });
5248 let last = end.fold(HEADER, u64::max);
5249 let directory = fs::metadata(&path).expect("the file is there").len();
5250 assert!(last <= directory, "a page runs to {last} in a file of {directory} bytes");
5251 fs::remove_file(path).expect("remove scratch file");
5252 }
5253
5254 fn dictionary_index_len(header: &[u8; DICTIONARY_HEADER]) -> u64 {
5260 let count = u64::from(u32::from_le_bytes(header[0..4].try_into().expect("four bytes")));
5261 let blocks = u64::from(u32::from_le_bytes(header[8..12].try_into().expect("four bytes")));
5262 let bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
5263 let rank_blocks = count.div_ceil(TEXT_RANK_BLOCK as u64);
5264 DICTIONARY_HEADER as u64
5265 + offset_bytes(count as usize, bits) as u64
5266 + (blocks + rank_blocks) * 16
5267 }
5268
5269 fn last_rank_end(file: &File, offset: u64, header: &[u8; DICTIONARY_HEADER]) -> u64 {
5271 let count = u64::from(u32::from_le_bytes(header[0..4].try_into().expect("four bytes")));
5272 let blocks = u64::from(u32::from_le_bytes(header[8..12].try_into().expect("four bytes")));
5273 let bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
5274 let rank_blocks = count.div_ceil(TEXT_RANK_BLOCK as u64);
5275 let at = offset
5276 + DICTIONARY_HEADER as u64
5277 + offset_bytes(count as usize, bits) as u64
5278 + blocks * 16
5279 + (rank_blocks - 1) * 8;
5280 let mut end = [0; 8];
5281 read_at(file, at, &mut end).expect("the last rank block end");
5282 u64::from_le_bytes(end)
5283 }
5284
5285 fn sample() -> Chunk {
5286 Chunk::new(vec![
5287 Vector::from_values(
5288 LogicalType::Integer,
5289 &[Value::Integer(4), Value::Integer(9), Value::Integer(-2)],
5290 )
5291 .expect("integers"),
5292 Vector::from_values(
5293 LogicalType::Varchar,
5294 &[
5295 Value::Varchar("alpha".into()),
5296 Value::Null,
5297 Value::Varchar("long text after a slash".into()),
5298 ],
5299 )
5300 .expect("strings"),
5301 ])
5302 .expect("matching rows")
5303 }
5304
5305 fn sample_ids() -> Chunk {
5306 Chunk::new(vec![
5307 Vector::flat(LogicalType::Integer, Data::Int32(vec![7, 8, 9].into()))
5308 .expect("integers"),
5309 ])
5310 .expect("one column")
5311 }
5312
5313 #[test]
5314 fn committed_file_reopens_and_reads_only_requested_columns() {
5315 let path = path("reopen");
5316 let mut writer = Writer::create(
5317 &path,
5318 "items",
5319 vec![
5320 Field::required("id", LogicalType::Integer),
5321 Field::new("text", LogicalType::Varchar),
5322 ],
5323 )
5324 .expect("new file");
5325 writer.append(&sample()).expect("first part");
5326 writer.append(&sample()).expect("second part");
5327 writer.finish().expect("commit");
5328 let reader = Reader::open(&path).expect("reopen from disk");
5329 assert_eq!(reader.table().rows(), 6);
5330 assert_eq!(reader.table().stripes().len(), 1);
5333 assert_eq!(reader.parts(), 2);
5334 assert_eq!(reader.part_rows(0), 3);
5335 assert_eq!(reader.part_rows(1), 3);
5336 let text = reader.read(1, &[1]).expect("only text page");
5337 assert_eq!(text.width(), 1);
5338 assert_eq!(text.value_at(1, 0), Value::Null);
5339 assert_eq!(text.value_at(2, 0), Value::Varchar("long text after a slash".into()));
5340 let sparse = reader.read_sparse(1, &[1]).expect("one part without its whole page");
5341 assert_eq!(sparse.width(), 1);
5342 assert_eq!(sparse.value_at(1, 0), Value::Null);
5343 assert_eq!(sparse.value_at(2, 0), Value::Varchar("long text after a slash".into()));
5344 assert!(!reader.skips_codes(0, 1, &[0]).expect("alpha is in the stripe"));
5345 assert!(!reader.skips_codes(0, 1, &[2]).expect("long text is in the stripe"));
5346 assert!(reader.skips_codes(0, 1, &[3]).expect("unknown code is absent"));
5347 let count = reader.read(0, &[]).expect("no page is needed for count");
5348 assert_eq!(count.len(), 3);
5349 assert!(reader.skips(0, &[Probe { column: 0, op: Op::Greater, value: Bound::Int(100) }]));
5350 assert!(!reader.skips(0, &[Probe { column: 0, op: Op::Greater, value: Bound::Int(0) }]));
5351 let integers = reader.top_frequencies(0, 1).expect("valid integer synopsis").expect("kept");
5352 assert_eq!(
5353 integers,
5354 vec![(Value::Integer(-2), 2), (Value::Integer(4), 2), (Value::Integer(9), 2),]
5355 );
5356 let strings = reader.top_frequencies(1, 1).expect("valid string synopsis").expect("kept");
5357 assert_eq!(strings.len(), 3);
5358 assert!(strings.contains(&(Value::Null, 2)));
5359 assert!(strings.contains(&(Value::Varchar("alpha".into()), 2)));
5360 assert!(strings.contains(&(Value::Varchar("long text after a slash".into()), 2)));
5361 fs::remove_file(path).expect("remove scratch file");
5362 }
5363
5364 #[test]
5372 fn runs_handed_over_out_of_order_still_read_back_in_source_order() {
5373 let path = path("interleaved-runs");
5374 let mut writer =
5375 Writer::create(&path, "interleaved", vec![Field::new("v", LogicalType::BigInt)])
5376 .expect("new file");
5377 for morsel in [2_u64, 0, 3, 1] {
5378 let parts = (0..4_u64)
5379 .map(|chunk| {
5380 let first = i64::try_from(morsel * 32 + chunk * 8).expect("small");
5381 let values =
5382 (0..8_i64).map(|row| Value::BigInt(first + row)).collect::<Vec<_>>();
5383 let column =
5384 Vector::from_values(LogicalType::BigInt, &values).expect("a column");
5385 ((morsel, chunk), Chunk::new(vec![column]).expect("one column"))
5386 })
5387 .collect::<Vec<_>>();
5388 writer.append_stripe(parts).expect("a stripe");
5389 }
5390 writer.finish().expect("commit");
5391
5392 let reader = Reader::open(&path).expect("valid directory");
5393 assert_eq!(reader.table().stripes().len(), 4, "a run is a stripe of its own");
5394 assert_eq!(reader.table().rows(), 128);
5395 for part in 0..16_usize {
5396 let read = reader.read(part, &[0]).expect("a part back");
5397 for row in 0..8_usize {
5398 let want = i64::try_from(part * 8 + row).expect("small");
5399 assert_eq!(read.value_at(row, 0), Value::BigInt(want), "part {part} row {row}");
5400 }
5401 }
5402 fs::remove_file(path).expect("remove scratch file");
5403 }
5404
5405 #[test]
5408 fn runs_that_overlap_each_other_are_refused_at_commit() {
5409 let path = path("overlapping-runs");
5410 let mut writer =
5411 Writer::create(&path, "overlapping", vec![Field::new("v", LogicalType::BigInt)])
5412 .expect("new file");
5413 let one = |order: (u64, u64)| {
5414 let column =
5415 Vector::from_values(LogicalType::BigInt, &[Value::BigInt(1)]).expect("a column");
5416 (order, Chunk::new(vec![column]).expect("one column"))
5417 };
5418 writer.append_stripe(vec![one((0, 0)), one((0, 2))]).expect("a stripe");
5421 writer.append_stripe(vec![one((0, 1))]).expect("a stripe");
5422 let error = writer.finish().expect_err("the runs overlap");
5423 assert!(error.message().contains("source order"), "{error}");
5424 fs::remove_file(path).expect("remove scratch file");
5425 }
5426
5427 #[test]
5430 fn a_run_longer_than_a_stripe_is_refused() {
5431 let path = path("overlong-run");
5432 let mut writer =
5433 Writer::create(&path, "overlong", vec![Field::new("v", LogicalType::BigInt)])
5434 .expect("new file");
5435 let parts = (0..=STRIPE_PARTS)
5436 .map(|at| {
5437 let column = Vector::from_values(LogicalType::BigInt, &[Value::BigInt(1)])
5438 .expect("a column");
5439 let chunk = Chunk::new(vec![column]).expect("one column");
5440 ((0, u64::try_from(at).expect("small")), chunk)
5441 })
5442 .collect::<Vec<_>>();
5443 let error = writer.append_stripe(parts).expect_err("one part too many");
5444 assert!(error.message().contains("more parts than it holds"), "{error}");
5445 fs::remove_file(path).expect("remove scratch file");
5446 }
5447
5448 #[test]
5454 fn parts_past_the_stripe_bound_start_a_new_stripe() {
5455 let path = path("stripe-bound");
5456 let mut writer = Writer::create(
5457 &path,
5458 "items",
5459 vec![
5460 Field::required("id", LogicalType::Integer),
5461 Field::new("text", LogicalType::Varchar),
5462 ],
5463 )
5464 .expect("new file");
5465 let parts = STRIPE_PARTS * 2 + 3;
5466 for part in 0..parts {
5467 let id = part as i32;
5468 let chunk = Chunk::new(vec![
5469 Vector::from_values(
5470 LogicalType::Integer,
5471 &[Value::Integer(id), Value::Integer(-id)],
5472 )
5473 .expect("integers"),
5474 Vector::from_values(
5475 LogicalType::Varchar,
5476 &[Value::Varchar(format!("value {part}")), Value::Null],
5477 )
5478 .expect("strings"),
5479 ])
5480 .expect("matching rows");
5481 writer.append(&chunk).expect("one part");
5482 }
5483 writer.finish().expect("commit");
5484
5485 let reader = Reader::open(&path).expect("reopen from disk");
5486 assert_eq!(reader.parts(), parts);
5487 assert_eq!(reader.table().rows(), parts * 2);
5488 assert_eq!(reader.table().stripes().len(), parts.div_ceil(STRIPE_PARTS));
5489 assert_eq!(reader.table().stripes()[0].parts(), STRIPE_PARTS);
5490 assert_eq!(reader.table().stripes()[0].rows(), STRIPE_PARTS * 2);
5491 assert_eq!(reader.table().stripes()[2].parts(), 3);
5492 for part in (0..parts).rev() {
5495 let dense = reader.read(part, &[0, 1]).expect("a whole page read");
5496 let sparse = reader.read_sparse(part, &[0, 1]).expect("one part read");
5497 for chunk in [&dense, &sparse] {
5498 assert_eq!(chunk.len(), 2, "part {part} has its own row count");
5499 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
5500 assert_eq!(chunk.value_at(1, 0), Value::Integer(-(part as i32)));
5501 assert_eq!(chunk.value_at(0, 1), Value::Varchar(format!("value {part}")));
5502 assert_eq!(chunk.value_at(1, 1), Value::Null);
5503 }
5504 }
5505 let above = [Probe { column: 0, op: Op::Greater, value: Bound::Int(100) }];
5508 assert!(reader.skips(0, &above), "the first stripe stops at 63");
5509 assert!(!reader.skips(STRIPE_PARTS * 2, &above), "the third stripe reaches 130");
5510 fs::remove_file(path).expect("remove scratch file");
5511 }
5512
5513 fn scattered(n: i64) -> i64 {
5515 n.wrapping_mul(-7_046_029_254_386_353_131)
5516 }
5517
5518 #[test]
5524 fn a_part_is_skipped_when_its_sieve_does_not_hold_the_constant() {
5525 let path = path("sieve-skip");
5526 let mut writer =
5527 Writer::create(&path, "hits", vec![Field::required("id", LogicalType::BigInt)])
5528 .expect("new file");
5529 let parts = STRIPE_PARTS + 3;
5530 let per_part = 128;
5534 for part in 0..parts {
5535 let held: Vec<Value> = (0..per_part)
5536 .map(|row| Value::BigInt(scattered((part * per_part + row) as i64)))
5537 .collect();
5538 let chunk =
5539 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
5540 .expect("one column");
5541 writer.append(&chunk).expect("one part");
5542 }
5543 writer.finish().expect("commit");
5544
5545 let reader = Reader::open(&path).expect("reopen from disk");
5546 let probe = |value: i64| Probe {
5547 column: 0,
5548 op: Op::Equal,
5549 value: Bound::Int(i128::from(scattered(value))),
5550 };
5551 for wanted in [0_i64, (per_part + 1) as i64, (parts * per_part - 1) as i64] {
5552 let tests = [probe(wanted)];
5553 let kept: Vec<usize> = (0..parts).filter(|&part| !reader.skips(part, &tests)).collect();
5554 let home = wanted as usize / per_part;
5555 assert!(kept.contains(&home), "the part holding {wanted} is read");
5556 assert!(kept.len() <= 2, "{wanted} keeps {kept:?}, which is more than one stray part");
5560 }
5561 let absent = [probe((parts * per_part) as i64 + 1)];
5562 let kept = (0..parts).filter(|&part| !reader.skips(part, &absent)).count();
5563 assert!(kept <= 1, "{kept} parts of {parts} kept a value no part holds");
5564 let tests = [probe(0)];
5567 assert!(
5568 reader.table().stripes().iter().all(|stripe| !stripe.zone.skips(&tests)),
5569 "the bounds rule out no stripe at all"
5570 );
5571 fs::remove_file(path).expect("remove scratch file");
5572 }
5573
5574 #[test]
5580 fn a_part_is_skipped_when_its_own_bounds_rule_out_a_comparison_the_stripe_keeps() {
5581 let path = path("part-range-skip");
5582 let mut writer =
5583 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
5584 .expect("new file");
5585 let parts = STRIPE_PARTS + 3;
5586 let per_part = 128;
5587 for part in 0..parts {
5588 let held: Vec<Value> = (0..per_part)
5592 .map(|row| {
5593 Value::BigInt((part * 1_000) as i64 + (scattered(row as i64).rem_euclid(900)))
5594 })
5595 .collect();
5596 let chunk =
5597 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
5598 .expect("one column");
5599 writer.append(&chunk).expect("one part");
5600 }
5601 writer.finish().expect("commit");
5602
5603 let reader = Reader::open(&path).expect("reopen from disk");
5604 let under = [Probe { column: 0, op: Op::Less, value: Bound::Int(3_000) }];
5605 let kept: Vec<usize> = (0..parts).filter(|&part| !reader.skips(part, &under)).collect();
5606 assert_eq!(kept, vec![0, 1, 2], "only the three parts that start under three thousand");
5607 assert!(!reader.stripe_skips(0, &under), "the stripe reaches from zero and keeps itself");
5609 fs::remove_file(path).expect("remove scratch file");
5610 }
5611
5612 #[test]
5615 fn a_stripe_of_one_part_writes_no_range_page_and_a_stripe_of_many_does() {
5616 for (parts, wanted) in [(1_usize, false), (STRIPE_PARTS, true)] {
5617 let path = path("part-range-page");
5618 let mut writer =
5619 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
5620 .expect("new file");
5621 for part in 0..parts {
5622 let held: Vec<Value> = (0..128)
5623 .map(|row| {
5624 Value::BigInt((part * 1_000) as i64 + scattered(row as i64).rem_euclid(900))
5625 })
5626 .collect();
5627 let chunk = Chunk::new(vec![
5628 Vector::from_values(LogicalType::BigInt, &held).expect("numbers"),
5629 ])
5630 .expect("one column");
5631 writer.append(&chunk).expect("one part");
5632 }
5633 writer.finish().expect("commit");
5634 let reader = Reader::open(&path).expect("reopen from disk");
5635 let bytes = reader.layout().columns[0].part_ranges;
5636 assert_eq!(bytes > 0, wanted, "{parts} parts wrote {bytes} bytes of ranges");
5637 fs::remove_file(path).expect("remove scratch file");
5638 }
5639 }
5640
5641 #[test]
5644 fn a_string_end_that_is_cut_down_still_covers_the_value_it_came_from() {
5645 let long = vec![b'a'; PART_BOUND_BYTES * 2];
5646 let low = shortened(Some(Bound::Bytes(long.clone())), false).expect("a low end");
5647 let high = shortened(Some(Bound::Bytes(long.clone())), true).expect("a high end");
5648 let Bound::Bytes(low) = low else { panic!("a string stays a string") };
5649 let Bound::Bytes(high) = high else { panic!("a string stays a string") };
5650 assert!(low.len() <= PART_BOUND_BYTES && high.len() <= PART_BOUND_BYTES);
5651 assert!(low.as_slice() <= long.as_slice(), "the low end is at or under the value");
5652 assert!(high.as_slice() >= long.as_slice(), "the high end is at or over the value");
5653 }
5654
5655 #[test]
5658 fn a_string_end_with_no_room_to_step_up_gives_up_the_bound() {
5659 let long = vec![u8::MAX; PART_BOUND_BYTES * 2];
5660 assert_eq!(shortened(Some(Bound::Bytes(long.clone())), true), None);
5661 let low = shortened(Some(Bound::Bytes(long)), false).expect("a low end is still a prefix");
5662 assert_eq!(low, Bound::Bytes(vec![u8::MAX; PART_BOUND_BYTES]));
5663 }
5664
5665 #[test]
5675 fn a_sieve_larger_than_the_part_it_indexes_is_not_written() {
5676 let path = path("sieve-pays");
5677 let fields = vec![
5678 Field::required("spread", LogicalType::BigInt),
5679 Field::required("repeated", LogicalType::BigInt),
5680 ];
5681 let mut writer = Writer::create(&path, "hits", fields).expect("new file");
5682 let parts = 3;
5683 let per_part = 1024;
5684 for part in 0..parts {
5685 let base = (part * per_part) as i64;
5686 let spread: Vec<Value> =
5687 (0..per_part).map(|row| Value::BigInt(scattered(base + row as i64))).collect();
5688 let repeated: Vec<Value> =
5689 (0..per_part).map(|row| Value::BigInt(scattered((row / 256) as i64))).collect();
5690 let chunk = Chunk::new(vec![
5691 Vector::from_values(LogicalType::BigInt, &spread).expect("numbers"),
5692 Vector::from_values(LogicalType::BigInt, &repeated).expect("numbers"),
5693 ])
5694 .expect("two columns");
5695 writer.append(&chunk).expect("one part");
5696 }
5697 writer.finish().expect("commit");
5698
5699 let reader = Reader::open(&path).expect("reopen from disk");
5700 let layout = reader.layout();
5701 let spread = &layout.columns[0];
5702 let repeated = &layout.columns[1];
5703 assert!(spread.sieves > 0, "a column whose parts are worth a filter keeps one");
5704 assert_eq!(
5705 repeated.sieves, 0,
5706 "a column whose filter costs more than its parts keeps none"
5707 );
5708 for column in &layout.columns {
5711 assert!(
5712 column.sieves < column.pages,
5713 "{} spends {} on sieves over {} of data",
5714 column.name,
5715 column.sieves,
5716 column.pages
5717 );
5718 }
5719 let absent = [Probe {
5721 column: 0,
5722 op: Op::Equal,
5723 value: Bound::Int(i128::from(scattered((parts * per_part) as i64 + 1))),
5724 }];
5725 assert!((0..parts).all(|part| reader.skips(part, &absent)), "no part holds it");
5726 fs::remove_file(path).expect("remove scratch file");
5727 }
5728
5729 #[test]
5735 fn a_damaged_sieve_page_is_read_through_rather_than_refused() {
5736 let path = path("sieve-damaged");
5737 let mut writer =
5738 Writer::create(&path, "hits", vec![Field::required("id", LogicalType::BigInt)])
5739 .expect("new file");
5740 let rows = 128;
5741 let held: Vec<Value> = (0..rows).map(|row| Value::BigInt(scattered(row))).collect();
5742 let chunk =
5743 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
5744 .expect("one column");
5745 writer.append(&chunk).expect("one part");
5746 writer.finish().expect("commit");
5747
5748 let page =
5749 Reader::open(&path).expect("reopen").table.stripes[0].sieves[0].expect("a sieve page");
5750 let mut file = OpenOptions::new().write(true).open(&path).expect("open the sieve page");
5751 file.seek(SeekFrom::Start(page.offset + u64::from(page.length) - 1)).expect("seek");
5752 file.write_all(&[0xff]).expect("damage one byte");
5753 drop(file);
5754
5755 let reader = Reader::open(&path).expect("reopen the damaged file");
5756 let absent =
5757 [Probe { column: 0, op: Op::Equal, value: Bound::Int(i128::from(scattered(99))) }];
5758 assert!(!reader.skips(0, &absent), "a sieve that cannot be read skips nothing");
5759 assert_eq!(
5760 reader.read(0, &[0]).expect("the rows are untouched").len(),
5761 usize::try_from(rows).expect("a small count")
5762 );
5763 fs::remove_file(path).expect("remove scratch file");
5764 }
5765
5766 #[test]
5777 fn workers_that_want_the_same_stripe_read_it_once() {
5778 let path = path("single-flight");
5779 let mut writer =
5780 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
5781 .expect("new file");
5782 for part in 0..STRIPE_PARTS {
5783 let id = part as i32;
5784 let chunk = Chunk::new(vec![
5785 Vector::from_values(
5786 LogicalType::Integer,
5787 &[Value::Integer(id), Value::Integer(-id)],
5788 )
5789 .expect("integers"),
5790 ])
5791 .expect("matching rows");
5792 writer.append(&chunk).expect("one part");
5793 }
5794 writer.finish().expect("commit");
5795
5796 let reader = Reader::open(&path).expect("reopen from disk");
5797 assert_eq!(reader.table().stripes().len(), 1, "one stripe is the point of the test");
5798 let barrier = std::sync::Barrier::new(8);
5799 std::thread::scope(|scope| {
5800 for worker in 0..8 {
5801 let reader = &reader;
5802 let barrier = &barrier;
5803 scope.spawn(move || {
5804 barrier.wait();
5805 for part in (worker..STRIPE_PARTS).step_by(8) {
5806 let chunk = reader.read(part, &[0]).expect("a whole page read");
5807 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
5808 assert_eq!(chunk.value_at(1, 0), Value::Integer(-(part as i32)));
5809 }
5810 });
5811 }
5812 });
5813 assert_eq!(reader.pages.load(Atomic::Relaxed), 1, "one stripe, one page read, whoever won");
5814 fs::remove_file(path).expect("remove scratch file");
5815 }
5816
5817 #[test]
5830 fn opening_costs_the_same_over_a_thousand_times_the_rows() {
5831 let opened = |label: &str, rows_per_part: i32| {
5832 let path = path(label);
5833 let mut writer =
5834 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
5835 .expect("new file");
5836 for part in 0..STRIPE_PARTS * 3 {
5837 let values = (0..rows_per_part)
5841 .map(|row| {
5842 Value::Integer((part as i32 * rows_per_part + row).wrapping_mul(2_654_435))
5843 })
5844 .collect::<Vec<_>>();
5845 let chunk = Chunk::new(vec![
5846 Vector::from_values(LogicalType::Integer, &values).expect("integers"),
5847 ])
5848 .expect("matching rows");
5849 writer.append(&chunk).expect("one part");
5850 }
5851 writer.finish().expect("commit");
5852 let reader = Reader::open(&path).expect("reopen from disk");
5853 let size = fs::metadata(&path).expect("the file is there").len();
5854 let out = (reader.reads(), reader.table().stripes().len(), size);
5855 fs::remove_file(path).expect("remove scratch file");
5856 out
5857 };
5858
5859 let (thin, thin_stripes, thin_size) = opened("open-thin", 1);
5860 let (fat, fat_stripes, fat_size) = opened("open-fat", 1000);
5861 assert_eq!(
5862 thin_stripes, fat_stripes,
5863 "the same stripe count is what makes this a fair ask"
5864 );
5865 assert!(
5866 fat_size > thin_size * 50,
5867 "the fat file has to actually be larger, and it is {fat_size} against {thin_size}"
5868 );
5869
5870 assert_eq!(thin.opening.reads, fat.opening.reads, "the same reads either way");
5871 assert_eq!(thin.pages, 0, "opening read a page");
5872 assert_eq!(fat.pages, 0, "opening read a page");
5873 assert_eq!(thin.indexes, 0, "opening read an index");
5874 assert_eq!(fat.indexes, 0, "opening read an index");
5875 assert!(
5878 fat.opening.bytes < thin.opening.bytes * 2,
5879 "opening the thin file read {} bytes and the fat one read {}",
5880 thin.opening.bytes,
5881 fat.opening.bytes
5882 );
5883 }
5884
5885 #[test]
5893 fn two_opens_of_one_file_cost_the_same_and_the_second_is_not_cheaper() {
5894 let path = path("open-twice");
5895 let mut writer =
5896 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
5897 .expect("new file");
5898 for part in 0..STRIPE_PARTS * 3 {
5899 let chunk = Chunk::new(vec![
5900 Vector::from_values(LogicalType::Integer, &[Value::Integer(part as i32)])
5901 .expect("integers"),
5902 ])
5903 .expect("matching rows");
5904 writer.append(&chunk).expect("one part");
5905 }
5906 writer.finish().expect("commit");
5907
5908 let first = Reader::open(&path).expect("open");
5909 for part in 0..first.parts() {
5912 first.read(part, &[0]).expect("a part");
5913 }
5914 assert!(first.reads().pages > 0, "the scan has to have read something");
5915 let second = Reader::open(&path).expect("open again");
5916
5917 assert_eq!(first.reads().opening, second.reads().opening);
5918 assert_eq!(
5919 second.reads().pages,
5920 0,
5921 "the second open read a page off the back of the first"
5922 );
5923 assert_eq!(second.reads().indexes, 0, "the second open read an index it inherited");
5924 fs::remove_file(path).expect("remove scratch file");
5925 }
5926
5927 #[test]
5935 fn an_index_is_read_once_per_stripe_however_often_the_page_is_evicted() {
5936 let path = path("index-cache");
5937 let mut writer =
5938 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
5939 .expect("new file");
5940 let parts = STRIPE_PARTS * (CACHED_STRIPES_PER_COLUMN + 2);
5941 for part in 0..parts {
5942 let id = part as i32;
5943 let chunk = Chunk::new(vec![
5944 Vector::from_values(LogicalType::Integer, &[Value::Integer(id)]).expect("integers"),
5945 ])
5946 .expect("matching rows");
5947 writer.append(&chunk).expect("one part");
5948 }
5949 writer.finish().expect("commit");
5950
5951 let reader = Reader::open(&path).expect("reopen from disk");
5952 let stripes = reader.table().stripes().len();
5953 assert!(stripes > CACHED_STRIPES_PER_COLUMN, "the page cache has to be too small for this");
5954 for _ in 0..2 {
5956 for part in 0..parts {
5957 let chunk = reader.read(part, &[0]).expect("a part");
5958 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
5959 }
5960 }
5961 assert_eq!(reader.indexes.load(Atomic::Relaxed), stripes, "one index read per stripe");
5962 assert!(
5963 reader.pages.load(Atomic::Relaxed) > stripes,
5964 "the pages are the ones that get read again, which is what makes the index count mean \
5965 something"
5966 );
5967 fs::remove_file(path).expect("remove scratch file");
5968 }
5969
5970 #[test]
5979 fn a_worker_per_stripe_reads_its_page_once_when_the_cache_was_told_to_expect_it() {
5980 let workers = CACHED_STRIPES_PER_COLUMN + 4;
5981 let path = path("stripe-per-worker");
5982 let mut writer =
5983 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
5984 .expect("new file");
5985 for part in 0..STRIPE_PARTS * workers {
5986 let chunk = Chunk::new(vec![
5987 Vector::from_values(LogicalType::Integer, &[Value::Integer(part as i32)])
5988 .expect("integers"),
5989 ])
5990 .expect("matching rows");
5991 writer.append(&chunk).expect("one part");
5992 }
5993 writer.finish().expect("commit");
5994
5995 let read = |told: bool| {
5996 let reader = Reader::open(&path).expect("reopen from disk");
5997 assert_eq!(reader.table().stripes().len(), workers, "a stripe per worker");
5998 if told {
5999 reader.keep_stripes(workers);
6000 }
6001 let barrier = std::sync::Barrier::new(workers);
6002 std::thread::scope(|scope| {
6003 for (worker, run) in reader.stripe_parts().into_iter().enumerate() {
6004 let reader = &reader;
6005 let barrier = &barrier;
6006 scope.spawn(move || {
6007 for part in run {
6008 barrier.wait();
6009 let chunk = reader.read(part, &[0]).expect("a part of my own stripe");
6010 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6011 }
6012 assert!(worker < workers);
6013 });
6014 }
6015 });
6016 reader.pages.load(Atomic::Relaxed)
6017 };
6018
6019 assert_eq!(read(true), workers, "one page read per stripe and no more");
6020 assert!(read(false) > workers, "a cache that small is read again on every part");
6021 fs::remove_file(path).expect("remove scratch file");
6022 }
6023
6024 #[test]
6029 fn a_damaged_index_page_is_an_error() {
6030 let path = path("damaged-index");
6031 let mut writer =
6032 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6033 .expect("new file");
6034 writer.append(&sample_ids()).expect("first part");
6035 writer.append(&sample_ids()).expect("second part");
6036 writer.finish().expect("commit");
6037
6038 let reader = Reader::open(&path).expect("valid directory");
6039 let index = reader.table.stripes[0].index;
6040 let mut byte = [0; 1];
6041 read_at(&reader.file, index.offset, &mut byte).expect("the first part length");
6042 let mut file = OpenOptions::new().write(true).open(&path).expect("open index page");
6043 file.seek(SeekFrom::Start(index.offset)).expect("index start");
6044 file.write_all(&[!byte[0]]).expect("damage the first part length");
6045 let error = reader.read(1, &[0]).expect_err("a damaged index must not be used");
6046 assert!(error.message().contains("index page section checksum differs"), "{error}");
6047 fs::remove_file(path).expect("remove scratch file");
6048 }
6049
6050 #[test]
6057 fn every_integer_width_round_trips_through_a_page() {
6058 let path = path("integer-widths");
6059 let columns = [
6060 (LogicalType::TinyInt, vec![Value::TinyInt(i8::MIN), Value::TinyInt(i8::MAX)]),
6061 (LogicalType::UTinyInt, vec![Value::UTinyInt(0), Value::UTinyInt(u8::MAX)]),
6062 (LogicalType::SmallInt, vec![Value::SmallInt(i16::MIN), Value::SmallInt(i16::MAX)]),
6063 (LogicalType::USmallInt, vec![Value::USmallInt(0), Value::USmallInt(u16::MAX)]),
6064 (LogicalType::Integer, vec![Value::Integer(i32::MIN), Value::Integer(i32::MAX)]),
6065 (LogicalType::UInteger, vec![Value::UInteger(0), Value::UInteger(u32::MAX)]),
6066 (LogicalType::BigInt, vec![Value::BigInt(i64::MIN), Value::BigInt(i64::MAX)]),
6067 (LogicalType::UBigInt, vec![Value::UBigInt(0), Value::UBigInt(u64::MAX)]),
6068 ];
6069 let fields = columns
6070 .iter()
6071 .enumerate()
6072 .map(|(at, (ty, _))| Field::required(format!("c{at}"), ty.clone()))
6073 .collect::<Vec<_>>();
6074 let vectors = columns
6075 .iter()
6076 .map(|(ty, values)| Vector::from_values(ty.clone(), values).expect("a vector"))
6077 .collect::<Vec<_>>();
6078 let mut writer = Writer::create(&path, "widths", fields).expect("new file");
6079 writer.append(&Chunk::new(vectors).expect("matching rows")).expect("one stripe");
6080 writer.finish().expect("commit");
6081
6082 let reader = Reader::open(&path).expect("reopen from disk");
6083 let wanted = (0..columns.len()).collect::<Vec<_>>();
6084 let read = reader.read(0, &wanted).expect("every column");
6085 assert_eq!(read.len(), 2);
6086 for (at, (ty, values)) in columns.iter().enumerate() {
6088 assert_eq!(read.value_at(0, at), values[0], "the low end of {ty}");
6089 assert_eq!(read.value_at(1, at), values[1], "the high end of {ty}");
6090 }
6091 fs::remove_file(path).expect("remove scratch file");
6092 }
6093
6094 #[test]
6095 fn numeric_frequency_candidates_keep_bounded_row_ordinals() {
6096 let path = path("frequency-ordinals");
6097 let mut writer =
6098 Writer::create(&path, "items", vec![Field::required("id", LogicalType::BigInt)])
6099 .expect("new file");
6100 let mut values = Vec::new();
6101 for leader in 0..10_i64 {
6102 values.extend(std::iter::repeat_n(leader, 100));
6103 }
6104 values.extend(1_000_i64..41_000);
6105 for part in values.chunks(1_024) {
6106 let vector = Vector::flat(LogicalType::BigInt, Data::Int64(part.to_vec().into()))
6107 .expect("big integers");
6108 writer.append(&Chunk::new(vec![vector]).expect("one column")).expect("one stripe");
6109 }
6110 writer.finish().expect("commit");
6111
6112 let reader = Reader::open(&path).expect("reopen from disk");
6113 let occurrences =
6114 reader.frequency_occurrences(0).expect("valid metadata").expect("bounded ordinals");
6115 assert!(occurrences.omitted_max < 100);
6116 assert!(occurrences.ordinals.len() <= FREQUENCY_ORDINALS);
6117 assert!(occurrences.ordinals.windows(2).all(|pair| pair[0] < pair[1]));
6118 assert_eq!(&occurrences.ordinals[..1_000], &(0_u64..1_000).collect::<Vec<_>>());
6119 fs::remove_file(path).expect("remove scratch file");
6120 }
6121
6122 #[test]
6128 fn a_file_from_another_format_says_which_format_it_is() {
6129 let older = path("older-format");
6130 let mut writer =
6131 Writer::create(&older, "items", vec![Field::new("id", LogicalType::Integer)])
6132 .expect("new file");
6133 let chunk = Chunk::new(vec![
6134 Vector::flat(LogicalType::Integer, Data::Int32(vec![1, 2, 3].into()))
6135 .expect("integers"),
6136 ])
6137 .expect("chunk");
6138 writer.append(&chunk).expect("page written");
6139 writer.finish().expect("commit");
6140
6141 let mut file = OpenOptions::new().write(true).open(&older).expect("open for the header");
6142 file.seek(SeekFrom::Start(8)).expect("the version follows the magic");
6143 file.write_all(&(FORMAT - 1).to_le_bytes()).expect("write an older version");
6144 drop(file);
6145 let complaint = Reader::open(&older).expect_err("an older format is refused").to_string();
6146 assert!(complaint.contains(&format!("format {}", FORMAT - 1)), "{complaint}");
6147 assert!(complaint.contains(&format!("format {FORMAT}")), "{complaint}");
6148
6149 let mut file = OpenOptions::new().write(true).open(&older).expect("open for the header");
6150 file.seek(SeekFrom::Start(0)).expect("the magic is first");
6151 file.write_all(b"NOTRUDB!").expect("write another engine's magic");
6152 drop(file);
6153 let complaint = Reader::open(&older).expect_err("a foreign file is refused").to_string();
6154 assert!(complaint.contains("magic"), "{complaint}");
6155 assert!(!complaint.contains("format"), "a version has nothing to do with it: {complaint}");
6156 fs::remove_file(older).expect("remove scratch file");
6157 }
6158
6159 #[test]
6160 fn an_unfinished_or_damaged_file_does_not_answer_with_partial_rows() {
6161 let unfinished = path("unfinished");
6162 let mut writer =
6163 Writer::create(&unfinished, "items", vec![Field::new("id", LogicalType::Integer)])
6164 .expect("new file");
6165 let chunk = Chunk::new(vec![
6166 Vector::flat(LogicalType::Integer, Data::Int32(vec![1, 2, 3].into()))
6167 .expect("integers"),
6168 ])
6169 .expect("chunk");
6170 writer.append(&chunk).expect("page written");
6171 drop(writer);
6172 assert!(Reader::open(&unfinished).is_err(), "no directory was committed");
6173 fs::remove_file(unfinished).expect("remove scratch file");
6174
6175 let damaged = path("damaged");
6176 let mut writer =
6177 Writer::create(&damaged, "items", vec![Field::new("id", LogicalType::Integer)])
6178 .expect("new file");
6179 writer.append(&chunk).expect("page written");
6180 writer.finish().expect("commit");
6181 let reader = Reader::open(&damaged).expect("valid directory");
6182 let mut file =
6183 OpenOptions::new().write(true).open(&damaged).expect("open for a damaged page");
6184 file.seek(SeekFrom::Start(HEADER + 1)).expect("inside first page");
6185 file.write_all(&[255]).expect("damage one byte");
6186 assert!(reader.read(0, &[0]).is_err(), "page checksum rejects corruption");
6187 fs::remove_file(damaged).expect("remove scratch file");
6188 }
6189
6190 #[test]
6191 fn damaged_lazy_dictionary_payload_is_an_error() {
6192 let path = path("damaged-dictionary");
6193 let mut writer = Writer::create(
6194 &path,
6195 "items",
6196 vec![
6197 Field::required("id", LogicalType::Integer),
6198 Field::new("text", LogicalType::Varchar),
6199 ],
6200 )
6201 .expect("new file");
6202 writer.append(&sample()).expect("stripe written");
6203 writer.finish().expect("commit");
6204
6205 let reader = Reader::open(&path).expect("valid directory");
6206 let dictionary = reader.table.dictionaries[1].expect("string dictionary page");
6207 let mut header = [0; DICTIONARY_HEADER];
6210 read_at(&reader.file, dictionary.offset, &mut header).expect("dictionary header");
6211 let index_len = dictionary_index_len(&header);
6212 let rank_len = last_rank_end(&reader.file, dictionary.offset, &header);
6213 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
6214 file.seek(SeekFrom::Start(dictionary.offset + index_len + rank_len))
6215 .expect("inside dictionary payload");
6216 file.write_all(&[255]).expect("damage dictionary payload");
6217
6218 let chunk = reader.read(0, &[1]).expect("code page and dictionary index remain valid");
6219 let error =
6220 chunk.validate_external().expect_err("payload corruption must reach the caller");
6221 assert!(error.message().contains("payload checksum differs"), "{error}");
6222 fs::remove_file(path).expect("remove scratch file");
6223 }
6224
6225 #[test]
6232 fn a_dictionary_over_many_blocks_checks_every_block_of_it() {
6233 let path = path("dictionary-blocks");
6234 let value =
6235 |row: usize| format!("{row:07} a value long enough to be worth a payload block");
6236 let parts = 30;
6237 let per_part = 1000;
6238 let mut writer =
6239 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
6240 .expect("new file");
6241 for part in 0..parts {
6242 let values = (0..per_part)
6243 .map(|row| Value::Varchar(value(part * per_part + row)))
6244 .collect::<Vec<_>>();
6245 let chunk = Chunk::new(vec![
6246 Vector::from_values(LogicalType::Varchar, &values).expect("strings"),
6247 ])
6248 .expect("matching rows");
6249 writer.append(&chunk).expect("a part");
6250 }
6251 writer.finish().expect("commit");
6252
6253 let reader = Reader::open(&path).expect("reopen from disk");
6254 let dictionary = reader.table.dictionaries[0].expect("string dictionary page");
6255 assert!(
6256 parts * per_part > TEXT_PAYLOAD_VALUES * 4,
6257 "the dictionary has to be several blocks for this to be testing anything"
6258 );
6259 for part in [0, parts - 1] {
6260 let chunk = reader.read(part, &[0]).expect("a part");
6261 chunk.validate_external().expect("every payload block checks out");
6262 assert_eq!(chunk.value_at(0, 0), Value::Varchar(value(part * per_part)));
6263 }
6264
6265 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
6266 file.seek(SeekFrom::Start(dictionary.offset + u64::from(dictionary.length) - 4))
6267 .expect("the last bytes of the page are payload");
6268 file.write_all(&[255]).expect("damage the last payload block");
6269 let reader = Reader::open(&path).expect("the directory and the index are untouched");
6270 let chunk = reader.read(parts - 1, &[0]).expect("the code page remains valid");
6271 let error = chunk.validate_external().expect_err("the damage must reach the caller");
6272 assert!(error.message().contains("payload checksum differs"), "{error}");
6273 fs::remove_file(path).expect("remove scratch file");
6274 }
6275
6276 #[test]
6286 fn values_of_different_lengths_read_back_out_of_packed_offsets() {
6287 let path = path("dictionary-offsets");
6288 let value = |row: usize| {
6289 if row % 511 == 3 { String::new() } else { "x".repeat(row % 97) + &format!("{row:05}") }
6290 };
6291 let rows = 5_000;
6292 let mut writer =
6293 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
6294 .expect("new file");
6295 let values = (0..rows).map(|row| Value::Varchar(value(row))).collect::<Vec<_>>();
6296 for part in values.chunks(1_000) {
6297 let chunk =
6298 Chunk::new(vec![Vector::from_values(LogicalType::Varchar, part).expect("strings")])
6299 .expect("matching rows");
6300 writer.append(&chunk).expect("a part");
6301 }
6302 writer.finish().expect("commit");
6303
6304 let reader = Reader::open(&path).expect("reopen from disk");
6305 assert!(
6306 rows > TEXT_PAYLOAD_VALUES * 4,
6307 "the dictionary has to be several blocks for this to be testing anything"
6308 );
6309 for part in 0..rows / 1_000 {
6310 let chunk = reader.read(part, &[0]).expect("a part");
6311 for row in 0..1_000 {
6312 let row = part * 1_000 + row;
6313 assert_eq!(
6314 chunk.value_at(row % 1_000, 0),
6315 Value::Varchar(value(row)),
6316 "value {row}"
6317 );
6318 }
6319 }
6320 fs::remove_file(path).expect("remove scratch file");
6321 }
6322
6323 #[test]
6335 fn a_global_dictionary_is_opened_once_however_many_workers_ask_at_once() {
6336 let path = path("dictionary-once");
6337 let parts = 8;
6338 let per_part = 500;
6339 let value =
6340 |row: usize| format!("{row:07} a value long enough to be worth a payload block");
6341 let mut writer =
6342 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
6343 .expect("new file");
6344 for part in 0..parts {
6345 let values = (0..per_part)
6346 .map(|row| Value::Varchar(value(part * per_part + row)))
6347 .collect::<Vec<_>>();
6348 let chunk = Chunk::new(vec![
6349 Vector::from_values(LogicalType::Varchar, &values).expect("strings"),
6350 ])
6351 .expect("matching rows");
6352 writer.append(&chunk).expect("a part");
6353 }
6354 writer.finish().expect("commit");
6355
6356 let reader = Reader::open(&path).expect("reopen from disk");
6357 assert!(reader.table.dictionaries[0].is_some(), "the column has to have one to share");
6358 assert_eq!(reader.reads().dictionaries, 0, "opening the file does not open a dictionary");
6359
6360 let workers = 16;
6361 let gate = std::sync::Barrier::new(workers);
6362 std::thread::scope(|scope| {
6363 for worker in 0..workers {
6364 let reader = reader.clone();
6365 let gate = &gate;
6366 scope.spawn(move || {
6367 gate.wait();
6368 let chunk = reader.read(worker % parts, &[0]).expect("a part");
6369 assert_eq!(
6370 chunk.value_at(0, 0),
6371 Value::Varchar(value((worker % parts) * per_part))
6372 );
6373 });
6374 }
6375 });
6376
6377 assert_eq!(reader.reads().dictionaries, 1, "sixteen workers, one dictionary, one open");
6378 fs::remove_file(path).expect("remove scratch file");
6379 }
6380
6381 #[test]
6386 fn a_damaged_sorted_order_is_an_error() {
6387 let path = path("damaged-order");
6388 let mut writer = Writer::create(
6389 &path,
6390 "items",
6391 vec![
6392 Field::required("id", LogicalType::Integer),
6393 Field::new("text", LogicalType::Varchar),
6394 ],
6395 )
6396 .expect("new file");
6397 writer.append(&sample()).expect("stripe written");
6398 writer.finish().expect("commit");
6399
6400 let reader = Reader::open(&path).expect("valid directory");
6401 let page = reader.table.dictionaries[1].expect("string dictionary page");
6402 let mut header = [0; DICTIONARY_HEADER];
6403 read_at(&reader.file, page.offset, &mut header).expect("dictionary header");
6404 let index_len = dictionary_index_len(&header);
6405 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
6406 file.seek(SeekFrom::Start(page.offset + index_len)).expect("the first head");
6407 file.write_all(&[255]).expect("damage the order");
6408
6409 let dictionary = reader.dictionary(1).expect("read").expect("a string column has one");
6410 let error = dictionary.compare_rank(0, b"anything").expect_err("a damaged order is caught");
6411 assert!(error.message().contains("rank checksum differs"), "{error}");
6412 fs::remove_file(path).expect("remove scratch file");
6413 }
6414
6415 #[test]
6419 fn a_global_dictionary_carries_the_sorted_order_of_its_values() {
6420 let spellings = ["overlong1z", "b", "", "overlong1a", "overlong", "ab", "a", "overlong1"];
6423 let path = path("dictionary-order");
6424 let mut writer =
6425 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
6426 .expect("new file");
6427 writer
6428 .append(
6429 &Chunk::new(vec![
6430 Vector::from_values(
6431 LogicalType::Varchar,
6432 &spellings.map(|text| Value::Varchar(text.into())),
6433 )
6434 .expect("strings"),
6435 ])
6436 .expect("one column"),
6437 )
6438 .expect("stripe written");
6439 writer.finish().expect("commit");
6440
6441 let reader = Reader::open(&path).expect("valid directory");
6442 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
6443 let count = dictionary.ranks().expect("a v10 file stores one");
6444 assert_eq!(count, spellings.len(), "every distinct value has a rank");
6445 let order = (0..count)
6446 .map(|rank| dictionary.code_at_rank(rank).expect("a code"))
6447 .collect::<Vec<_>>();
6448 let mut seen = order.clone();
6449 seen.sort_unstable();
6450 assert_eq!(seen, (0..spellings.len() as u32).collect::<Vec<_>>(), "a permutation of codes");
6451
6452 let ranked = order
6453 .iter()
6454 .map(|&code| {
6455 dictionary.try_bytes_at(code as usize).expect("read").expect("a value").to_vec()
6456 })
6457 .collect::<Vec<_>>();
6458 let mut expected = spellings.map(|text| text.as_bytes().to_vec()).to_vec();
6459 expected.sort();
6460 assert_eq!(ranked, expected, "rank order is value order");
6461
6462 for (rank, value) in expected.iter().enumerate() {
6465 assert_eq!(
6466 dictionary.compare_rank(rank, value).expect("compare"),
6467 Ordering::Equal,
6468 "rank {rank} is its own value"
6469 );
6470 if rank > 0 {
6471 assert_eq!(
6472 dictionary.compare_rank(rank - 1, value).expect("compare"),
6473 Ordering::Less,
6474 "rank {rank} follows the one before it"
6475 );
6476 }
6477 }
6478 fs::remove_file(path).expect("remove scratch file");
6479 }
6480
6481 #[test]
6489 fn a_dictionary_sweep_reads_every_value_and_keeps_it_under_the_budget() {
6490 let path = path("dictionary-sweep");
6491 let spellings = (0..2_500)
6494 .map(|index| Value::Varchar(format!("value {index:08} {}", "x".repeat(index % 40))))
6495 .collect::<Vec<_>>();
6496 let mut writer =
6497 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
6498 .expect("new file");
6499 for part in spellings.chunks(1_024) {
6502 writer
6503 .append(
6504 &Chunk::new(vec![
6505 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
6506 ])
6507 .expect("one column"),
6508 )
6509 .expect("stripe written");
6510 }
6511 writer.finish().expect("commit");
6512
6513 let reader = Reader::open(&path).expect("valid directory");
6514 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
6515 assert_eq!(dictionary.len(), spellings.len(), "every value is distinct");
6516
6517 let resting = dictionary.footprint();
6518 let mut swept: Vec<Vec<u8>> = Vec::new();
6519 let mut at = 0;
6520 let mut calls = 0;
6521 while at < dictionary.len() {
6522 let stopped = dictionary
6523 .sweep_text(at, dictionary.len(), &mut |index: usize, text: &[u8]| {
6524 assert_eq!(index, swept.len(), "a sweep hands its values over in order");
6525 swept.push(text.to_vec());
6526 Ok(())
6527 })
6528 .expect("a sweep reads");
6529 assert!(stopped > at, "a sweep moves");
6530 at = stopped;
6531 calls += 1;
6532 }
6533 assert_eq!(calls, 3, "a sweep hands over one block at a time");
6534 let after = dictionary.footprint();
6535 assert!(after > resting, "a sweep under the budget keeps what it decoded");
6536
6537 let read = (0..dictionary.len())
6538 .map(|code| dictionary.try_bytes_at(code).expect("read").expect("a value").to_vec())
6539 .collect::<Vec<_>>();
6540 assert_eq!(swept, read, "a sweep answers what a point read answers");
6541 assert_eq!(dictionary.footprint(), after, "a point read of a kept block decodes nothing");
6542 fs::remove_file(path).expect("remove scratch file");
6543 }
6544
6545 #[test]
6556 fn a_sweep_over_a_block_with_a_short_second_run_reads_what_a_point_read_reads() {
6557 let path = path("dictionary-sweep-short-run");
6558 let spellings = (0..2_800)
6559 .map(|index| Value::Varchar(format!("value {index:08} {}", "x".repeat(index % 40))))
6560 .collect::<Vec<_>>();
6561 let mut writer =
6562 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
6563 .expect("new file");
6564 for part in spellings.chunks(1_024) {
6565 writer
6566 .append(
6567 &Chunk::new(vec![
6568 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
6569 ])
6570 .expect("one column"),
6571 )
6572 .expect("stripe written");
6573 }
6574 writer.finish().expect("commit");
6575
6576 let reader = Reader::open(&path).expect("valid directory");
6577 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
6578 assert_eq!(dictionary.len(), spellings.len(), "every value is distinct");
6579 let last = dictionary.len() % TEXT_PAYLOAD_VALUES;
6580 assert!(last > TEXT_OFFSET_RUN, "the last block has to reach into a second run of offsets");
6581 assert!(last < TEXT_PAYLOAD_VALUES, "and that second run has to be short of a whole one");
6582
6583 let mut swept: Vec<Vec<u8>> = Vec::new();
6584 let mut at = 0;
6585 while at < dictionary.len() {
6586 let stopped = dictionary
6587 .sweep_text(at, dictionary.len(), &mut |index: usize, text: &[u8]| {
6588 assert_eq!(index, swept.len(), "a sweep hands its values over in order");
6589 swept.push(text.to_vec());
6590 Ok(())
6591 })
6592 .expect("a sweep reads");
6593 assert!(stopped > at, "a sweep moves");
6594 at = stopped;
6595 }
6596 let read = (0..dictionary.len())
6597 .map(|code| dictionary.try_bytes_at(code).expect("read").expect("a value").to_vec())
6598 .collect::<Vec<_>>();
6599 assert_eq!(swept, read, "a sweep answers what a point read answers");
6600 fs::remove_file(path).expect("remove scratch file");
6601 }
6602
6603 #[test]
6613 fn narrowing_a_page_takes_what_fits_and_refuses_what_does_not() {
6614 assert_eq!(fit::<i8>(&[]).expect("an empty page fits anything"), Vec::<i8>::new());
6615 assert_eq!(fit::<i8>(&[-128, 0, 127]).expect("the edges fit"), vec![-128_i8, 0, 127]);
6616 fit::<i8>(&[128]).expect_err("one past the top does not fit");
6617 fit::<i8>(&[-129]).expect_err("one past the bottom does not fit");
6618 assert_eq!(fit::<u8>(&[0, 255]).expect("the edges fit"), vec![0_u8, 255]);
6619 fit::<u8>(&[256]).expect_err("one past the top does not fit");
6620 fit::<u8>(&[-1]).expect_err("a negative does not fit an unsigned page");
6621 assert_eq!(
6622 fit::<i16>(&[-32_768, 0, 32_767]).expect("the edges fit"),
6623 vec![-32_768_i16, 0, 32_767]
6624 );
6625 fit::<i16>(&[32_768]).expect_err("one past the top does not fit");
6626 fit::<i16>(&[-32_769]).expect_err("one past the bottom does not fit");
6627 assert_eq!(fit::<u16>(&[0, 65_535]).expect("the edges fit"), vec![0_u16, 65_535]);
6628 fit::<u16>(&[65_536]).expect_err("one past the top does not fit");
6629 fit::<u16>(&[-1]).expect_err("a negative does not fit an unsigned page");
6630 assert_eq!(
6631 fit::<i32>(&[i64::from(i32::MIN), 0, i64::from(i32::MAX)]).expect("the edges fit"),
6632 vec![i32::MIN, 0, i32::MAX]
6633 );
6634 fit::<i32>(&[i64::from(i32::MAX) + 1]).expect_err("one past the top does not fit");
6635 fit::<i32>(&[i64::from(i32::MIN) - 1]).expect_err("one past the bottom does not fit");
6636 assert_eq!(
6637 fit::<u32>(&[0, 4_294_967_295]).expect("the edges fit"),
6638 vec![0_u32, 4_294_967_295]
6639 );
6640 fit::<u32>(&[4_294_967_296]).expect_err("one past the top does not fit");
6641 fit::<u32>(&[-1]).expect_err("a negative does not fit an unsigned page");
6642
6643 fit::<i8>(&[0, 1, 2, 128, 3]).expect_err("one bad value spoils the page");
6646 }
6647
6648 #[test]
6655 fn the_residue_agrees_with_a_checked_conversion_everywhere() {
6656 for value in -70_000_i64..70_000 {
6657 assert_eq!(fit::<i8>(&[value]).is_ok(), i8::try_from(value).is_ok(), "{value} as i8");
6658 assert_eq!(fit::<u8>(&[value]).is_ok(), u8::try_from(value).is_ok(), "{value} as u8");
6659 assert_eq!(fit::<i16>(&[value]).is_ok(), i16::try_from(value).is_ok(), "{value} i16");
6660 assert_eq!(fit::<u16>(&[value]).is_ok(), u16::try_from(value).is_ok(), "{value} u16");
6661 }
6662 let wide = [i64::MIN, i64::MIN + 1, i64::from(i32::MIN), 0, i64::from(u32::MAX), i64::MAX];
6663 for edge in wide {
6664 for step in -2_i64..=2 {
6665 let value = edge.saturating_add(step);
6666 assert_eq!(
6667 fit::<i32>(&[value]).is_ok(),
6668 i32::try_from(value).is_ok(),
6669 "{value} as i32"
6670 );
6671 assert_eq!(
6672 fit::<u32>(&[value]).is_ok(),
6673 u32::try_from(value).is_ok(),
6674 "{value} as u32"
6675 );
6676 }
6677 }
6678 }
6679
6680 #[test]
6688 fn a_dictionary_at_its_budget_sweeps_without_keeping() {
6689 let path = path("dictionary-budget");
6690 let spellings = (0..2_500)
6691 .map(|index| Value::Varchar(format!("value {index:08} {}", "y".repeat(index % 40))))
6692 .collect::<Vec<_>>();
6693 let mut writer =
6694 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
6695 .expect("new file");
6696 for part in spellings.chunks(1_024) {
6697 writer
6698 .append(
6699 &Chunk::new(vec![
6700 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
6701 ])
6702 .expect("one column"),
6703 )
6704 .expect("stripe written");
6705 }
6706 writer.finish().expect("commit");
6707
6708 let reader = Reader::open(&path).expect("valid directory");
6709 let page = reader.table.dictionaries[0].expect("a string column has one");
6710 let file = Arc::clone(&reader.file);
6711 let starved = open_global_dictionary(file, page, &LogicalType::Varchar, 0)
6712 .expect("a dictionary opens whatever it may keep");
6713
6714 let resting = starved.footprint();
6715 let mut swept: Vec<Vec<u8>> = Vec::new();
6716 let mut at = 0;
6717 while at < starved.len() {
6718 at = starved
6719 .sweep_text(at, starved.len(), &mut |_index: usize, text: &[u8]| {
6720 swept.push(text.to_vec());
6721 Ok(())
6722 })
6723 .expect("a sweep reads");
6724 }
6725 assert_eq!(swept.len(), spellings.len(), "a starved sweep still reads every value");
6726 assert_eq!(starved.footprint(), resting, "and keeps no block it decoded");
6727
6728 let generous = reader.dictionary(0).expect("read").expect("a string column has one");
6729 let read = (0..generous.len())
6730 .map(|code| generous.try_bytes_at(code).expect("read").expect("a value").to_vec())
6731 .collect::<Vec<_>>();
6732 assert_eq!(swept, read, "a starved sweep answers what a point read answers");
6733 fs::remove_file(path).expect("remove scratch file");
6734 }
6735
6736 #[test]
6737 fn damaged_membership_cannot_skip_a_string_page() {
6738 let path = path("damaged-membership");
6739 let mut writer = Writer::create(
6740 &path,
6741 "items",
6742 vec![
6743 Field::required("id", LogicalType::Integer),
6744 Field::new("text", LogicalType::Varchar),
6745 ],
6746 )
6747 .expect("new file");
6748 writer.append(&sample()).expect("stripe written");
6749 writer.finish().expect("commit");
6750
6751 let reader = Reader::open(&path).expect("valid directory");
6752 let membership = reader.table.stripes[0].memberships[1].expect("string membership");
6753 let mut file = OpenOptions::new().write(true).open(&path).expect("open membership page");
6754 file.seek(SeekFrom::Start(membership.offset)).expect("membership start");
6755 file.write_all(&[255]).expect("damage membership");
6756 let error = reader.skips_codes(0, 1, &[3]).expect_err("corruption must not skip rows");
6757 assert!(error.message().contains("membership page checksum differs"), "{error}");
6758 fs::remove_file(path).expect("remove scratch file");
6759 }
6760
6761 #[test]
6762 fn membership_delta_stream_is_sorted_exact_and_bounded() {
6763 let unique = unique_codes(&[900, 4, 4, 72, 9, u32::MAX]);
6764 assert_eq!(unique, [4, 9, 72, 900, u32::MAX]);
6765 let encoded = encode_membership(&unique);
6766 assert_eq!(
6767 decode_membership(&encoded).expect("valid membership"),
6768 [4, 9, 72, 900, u32::MAX]
6769 );
6770 let merged = merged_codes(vec![vec![4, 900], vec![9, 900, u32::MAX], vec![72]]);
6773 assert_eq!(merged, [4, 9, 72, 900, u32::MAX]);
6774 assert_eq!(
6775 decode_membership(&encode_membership(&merged)).expect("valid membership"),
6776 unique
6777 );
6778 assert!(decode_membership(&[1, 0x80]).is_err(), "a truncated varint is invalid");
6779 assert!(
6780 decode_membership(&[1, 0xff, 0xff, 0xff, 0xff, 0x10]).is_err(),
6781 "a value past u32 is invalid"
6782 );
6783 }
6784
6785 #[test]
6786 fn a_global_dictionary_may_be_larger_than_one_column_page() {
6787 let dictionary = Page {
6788 offset: HEADER,
6789 length: u32::try_from(MAX_PAGE + 1).expect("the page bound fits on disk"),
6790 hash: 0,
6791 };
6792 let table = Table {
6793 name: "items".to_owned(),
6794 fields: vec![Field::new("text", LogicalType::Varchar)],
6795 stripes: Vec::new(),
6796 rows: 0,
6797 dictionaries: vec![Some(dictionary)],
6798 distincts: vec![None],
6799 frequencies: vec![None],
6800 };
6801 let directory = encode_directory(&table).expect("directory");
6802 let file_size = dictionary.offset + u64::from(dictionary.length) + 1;
6803
6804 let decoded = decode_directory(&directory, file_size).expect("large lazy dictionary");
6805 assert_eq!(decoded.dictionaries[0].expect("dictionary").length, dictionary.length);
6806 }
6807
6808 #[test]
6809 fn a_column_with_one_value_everywhere_costs_almost_nothing_a_row() {
6810 let path = path("constant-codes");
6811 let mut writer =
6812 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
6813 .expect("new file");
6814 let empty = vec![Value::Varchar(String::new()); 1024];
6815 for _ in 0..4 {
6816 let column = Vector::from_values(LogicalType::Varchar, &empty).expect("strings");
6817 writer.append(&Chunk::new(vec![column]).expect("one column")).expect("a part");
6818 }
6819 writer.finish().expect("commit");
6820
6821 let reader = Reader::open(&path).expect("valid directory");
6822 let pages = reader.layout().columns.first().expect("one column").pages;
6823 assert!(pages < 256, "{pages} bytes of pages for 4,096 rows of one value");
6827 let read = reader.read(3, &[0]).expect("the last part back");
6828 assert_eq!(read.value_at(0, 0), Value::Varchar(String::new()));
6829 assert_eq!(read.value_at(1023, 0), Value::Varchar(String::new()));
6830 fs::remove_file(path).expect("remove scratch file");
6831 }
6832
6833 #[test]
6834 fn a_cascade_value_too_wide_for_its_column_is_refused_rather_than_cut() {
6835 let over = vec![i64::from(i32::MAX) + 1];
6838 let error = narrowed(&LogicalType::Integer, over).expect_err("a page that disagrees");
6839 assert!(format!("{error}").contains("not of its type"), "{error}");
6840 assert!(narrowed(&LogicalType::BigInt, vec![i64::MIN]).is_ok(), "bigint holds all of i64");
6841 assert!(narrowed(&LogicalType::Varchar, vec![0]).is_err(), "strings are not integers");
6842 }
6843
6844 #[test]
6845 fn a_code_stream_the_cascade_cannot_shrink_is_left_alone() {
6846 let mut state: u32 = 0x9e37_79b9;
6850 let spread: Vec<u32> = (0..1024)
6851 .map(|_| {
6852 state ^= state << 13;
6853 state ^= state >> 17;
6854 state ^= state << 5;
6855 state
6856 })
6857 .collect();
6858 assert_eq!(encoded_codes(&spread).expect("no failure"), None);
6859 let near: Vec<u32> = (0..1024).collect();
6860 let coded = encoded_codes(&near).expect("no failure").expect("counting up is packable");
6861 assert!(coded.len() < near.len() * 4, "{} bytes for a run of 1,024", coded.len());
6862 }
6863
6864 #[test]
6870 fn two_writes_of_the_same_rows_give_the_same_bytes() {
6871 fn written(path: &PathBuf) {
6872 let fields = (0..40)
6873 .map(|column| {
6874 let ty =
6875 if column % 4 == 0 { LogicalType::Varchar } else { LogicalType::BigInt };
6876 Field::new(format!("c{column}"), ty)
6877 })
6878 .collect::<Vec<_>>();
6879 let mut writer = Writer::create(path, "wide", fields).expect("new file");
6880 for part in 0..70_u64 {
6881 let columns = (0..40)
6882 .map(|column| {
6883 let values = (0..64_u64)
6884 .map(|row| {
6885 let seed = part.wrapping_mul(31).wrapping_add(row);
6886 if column % 4 == 0 {
6887 Value::Varchar(format!("v{}", seed % 17))
6888 } else {
6889 Value::BigInt(i64::try_from(seed % 97).expect("small"))
6890 }
6891 })
6892 .collect::<Vec<_>>();
6893 let ty = if column % 4 == 0 {
6894 LogicalType::Varchar
6895 } else {
6896 LogicalType::BigInt
6897 };
6898 Vector::from_values(ty, &values).expect("a column")
6899 })
6900 .collect::<Vec<_>>();
6901 writer.append(&Chunk::new(columns).expect("forty columns")).expect("a part");
6902 }
6903 writer.finish().expect("commit");
6904 }
6905
6906 let first = path("repeatable-one");
6907 let second = path("repeatable-two");
6908 written(&first);
6909 written(&second);
6910 let left = fs::read(&first).expect("the first file");
6911 let right = fs::read(&second).expect("the second file");
6912 assert_eq!(left.len(), right.len(), "two writes of the same rows differ in length");
6913 assert!(left == right, "two writes of the same rows differ in their bytes");
6914
6915 let reader = Reader::open(&first).expect("valid directory");
6918 assert_eq!(reader.table().rows(), 70 * 64);
6919 let read = reader.read(0, &[0, 1]).expect("the first part back");
6920 assert_eq!(read.value_at(0, 0), Value::Varchar("v0".to_owned()));
6921 assert_eq!(read.value_at(0, 1), Value::BigInt(0));
6922 fs::remove_file(first).expect("remove scratch file");
6923 fs::remove_file(second).expect("remove scratch file");
6924 }
6925}