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::slice;
37use std::sync::atomic::{AtomicUsize, Ordering as Atomic};
38use std::sync::{Arc, Mutex, OnceLock};
39
40use rudb_common::bounds::{Bound, Op, scaled_as};
41use rudb_common::{Error, Field, LogicalType, Result, Value};
42use rudb_encoding::{bitpack, chooser, integer, string};
43use rudb_storage::sieve::Sieve;
44use rudb_storage::{Probe, Range, Zone};
45use rudb_vector::string::StringColumn;
46use rudb_vector::validity::Validity;
47use rudb_vector::{Buffer, Chunk, Data, Packed, TextSource, Vector};
48
49const MAGIC: &[u8; 8] = b"RUDBNV10";
50const DIRECTORY: &[u8; 8] = b"RUDBDI10";
51const FORMAT: u32 = 21;
52const HEADER: u64 = 80;
53const SLOT_BYTES: usize = 28;
54const MAX_PAGE: usize = 256 * 1024 * 1024;
55const MAX_DIRECTORY: usize = 128 * 1024 * 1024;
56const FREQUENCIES: &[u8; 8] = b"RUDBFQ2\0";
57const FREQUENCY_CANDIDATES: usize = 32_768;
58const FREQUENCY_ENTRIES: usize = 512;
59const FREQUENCY_BUILD_RANK: usize = 10;
60const FREQUENCY_ORDINALS: usize = 65_536;
61const MAX_FREQUENCY_WORKERS: usize = 32;
68
69const MAX_ENCODE_WORKERS: usize = 32;
76
77const SIEVE_BUDGET: usize = 8 * 1024;
85
86const PART_BOUND_BYTES: usize = 24;
95
96fn io(error: std::io::Error) -> Error {
97 Error::io(error.to_string())
98}
99
100fn invalid(message: &str) -> Error {
101 Error::invalid_input(format!("invalid rudb native file: {message}"))
102}
103
104fn sum(counts: impl Iterator<Item = u64>) -> u64 {
106 counts.fold(0, u64::saturating_add)
107}
108
109fn span_bytes(spans: &[Span], at: usize) -> u64 {
111 spans.get(at).map_or(0, |span| u64::from(span.length))
112}
113
114fn page_bytes(pages: &[Option<Page>], at: usize) -> u64 {
116 pages.get(at).and_then(Option::as_ref).map_or(0, Page::bytes)
117}
118
119fn checksum(bytes: &[u8]) -> u64 {
120 const P1: u64 = 11_400_714_785_074_694_791;
121 const P2: u64 = 14_029_467_366_897_019_727;
122 const P3: u64 = 1_609_587_929_392_839_161;
123 const P4: u64 = 9_650_029_242_287_828_579;
124 const P5: u64 = 2_870_177_450_012_600_261;
125 let round = |state: u64, word: u64| {
126 state.wrapping_add(word.wrapping_mul(P2)).rotate_left(31).wrapping_mul(P1)
127 };
128 let merge = |state: u64, lane: u64| (state ^ round(0, lane)).wrapping_mul(P1).wrapping_add(P4);
129 let word =
130 |at: usize| u64::from_le_bytes(bytes[at..at + 8].try_into().expect("eight checksum bytes"));
131
132 let mut at = 0;
133 let mut hash = if bytes.len() >= 32 {
134 let mut one = P1.wrapping_add(P2);
135 let mut two = P2;
136 let mut three = 0;
137 let mut four = 0_u64.wrapping_sub(P1);
138 while at + 32 <= bytes.len() {
139 one = round(one, word(at));
140 two = round(two, word(at + 8));
141 three = round(three, word(at + 16));
142 four = round(four, word(at + 24));
143 at += 32;
144 }
145 let combined = one
146 .rotate_left(1)
147 .wrapping_add(two.rotate_left(7))
148 .wrapping_add(three.rotate_left(12))
149 .wrapping_add(four.rotate_left(18));
150 merge(merge(merge(merge(combined, one), two), three), four)
151 } else {
152 P5
153 };
154 hash = hash.wrapping_add(bytes.len() as u64);
155 while at + 8 <= bytes.len() {
156 hash ^= round(0, word(at));
157 hash = hash.rotate_left(27).wrapping_mul(P1).wrapping_add(P4);
158 at += 8;
159 }
160 if at + 4 <= bytes.len() {
161 let tail = u32::from_le_bytes(bytes[at..at + 4].try_into().expect("four checksum bytes"));
162 hash ^= u64::from(tail).wrapping_mul(P1);
163 hash = hash.rotate_left(23).wrapping_mul(P2).wrapping_add(P3);
164 at += 4;
165 }
166 while at < bytes.len() {
167 hash ^= u64::from(bytes[at]).wrapping_mul(P5);
168 hash = hash.rotate_left(11).wrapping_mul(P1);
169 at += 1;
170 }
171 hash ^= hash >> 33;
172 hash = hash.wrapping_mul(P2);
173 hash ^= hash >> 29;
174 hash = hash.wrapping_mul(P3);
175 hash ^ (hash >> 32)
176}
177
178#[derive(Debug, Clone, Copy)]
179struct Slot {
180 offset: u64,
181 length: u32,
182 generation: u64,
183 hash: u64,
184}
185
186impl Slot {
187 fn bytes(self) -> [u8; SLOT_BYTES] {
188 let mut result = [0; SLOT_BYTES];
189 result[..8].copy_from_slice(&self.offset.to_le_bytes());
190 result[8..12].copy_from_slice(&self.length.to_le_bytes());
191 result[12..20].copy_from_slice(&self.generation.to_le_bytes());
192 result[20..28].copy_from_slice(&self.hash.to_le_bytes());
193 result
194 }
195
196 fn read(bytes: &[u8]) -> Self {
197 Self {
198 offset: u64::from_le_bytes(bytes[..8].try_into().expect("eight bytes")),
199 length: u32::from_le_bytes(bytes[8..12].try_into().expect("four bytes")),
200 generation: u64::from_le_bytes(bytes[12..20].try_into().expect("eight bytes")),
201 hash: u64::from_le_bytes(bytes[20..28].try_into().expect("eight bytes")),
202 }
203 }
204}
205
206#[derive(Debug, Clone, Copy)]
207struct Page {
208 offset: u64,
209 length: u32,
210 hash: u64,
211}
212
213impl Page {
214 fn bytes(&self) -> u64 {
216 u64::from(self.length)
217 }
218}
219
220#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
221enum FrequencyValue {
222 Null,
223 Integer(i128),
224 Code(u32),
225}
226
227#[derive(Debug, Clone)]
228struct FrequencyEntry {
229 value: FrequencyValue,
230 count: u64,
231}
232
233#[derive(Debug, Clone)]
238struct FrequencySummary {
239 entries: Vec<FrequencyEntry>,
240 omitted_max: u64,
241 ordinals: Vec<u64>,
242}
243
244#[derive(Debug, Clone, PartialEq, Eq)]
246pub struct FrequencyOccurrences {
247 pub omitted_max: u64,
249 pub ordinals: Vec<u64>,
251}
252
253#[derive(Debug, Clone, Copy, Default)]
260struct Span {
261 offset: u64,
262 length: u32,
263}
264
265#[derive(Debug, Clone)]
267pub struct Stripe {
268 rows: usize,
269 parts: Vec<u32>,
272 index: Span,
276 pages: Vec<Span>,
277 memberships: Vec<Option<Page>>,
278 sieves: Vec<Option<Page>>,
281 part_ranges: Vec<Option<Page>>,
292 zone: Zone,
293}
294
295impl Stripe {
296 #[must_use]
298 pub fn rows(&self) -> usize {
299 self.rows
300 }
301
302 #[must_use]
304 pub fn parts(&self) -> usize {
305 self.parts.len()
306 }
307}
308
309#[derive(Debug, Clone)]
311pub struct Table {
312 name: String,
313 fields: Vec<Field>,
314 stripes: Vec<Stripe>,
315 rows: usize,
316 dictionaries: Vec<Option<Page>>,
317 frequencies: Vec<Option<FrequencySummary>>,
318 distincts: Vec<Option<u64>>,
328}
329
330impl Table {
331 #[must_use]
333 pub fn name(&self) -> &str {
334 &self.name
335 }
336
337 #[must_use]
339 pub fn fields(&self) -> &[Field] {
340 &self.fields
341 }
342
343 #[must_use]
345 pub fn rows(&self) -> usize {
346 self.rows
347 }
348
349 #[must_use]
351 pub fn stripes(&self) -> &[Stripe] {
352 &self.stripes
353 }
354}
355
356#[derive(Debug, Clone)]
358pub struct ColumnLayout {
359 pub name: String,
361 pub kind: String,
363 pub pages: u64,
365 pub memberships: u64,
367 pub sieves: u64,
369 pub part_ranges: u64,
371 pub dictionary: u64,
373}
374
375impl ColumnLayout {
376 #[must_use]
378 pub fn total(&self) -> u64 {
379 self.pages
380 .saturating_add(self.memberships)
381 .saturating_add(self.sieves)
382 .saturating_add(self.part_ranges)
383 .saturating_add(self.dictionary)
384 }
385}
386
387#[derive(Debug, Clone)]
398pub struct Layout {
399 pub file: u64,
401 pub rows: usize,
403 pub stripes: usize,
405 pub parts: usize,
407 pub columns: Vec<ColumnLayout>,
409 pub indexes: u64,
412 pub directory: u64,
414 pub header: u64,
416}
417
418impl Layout {
419 #[must_use]
421 pub fn columns_total(&self) -> u64 {
422 self.columns.iter().map(ColumnLayout::total).fold(0, u64::saturating_add)
423 }
424
425 #[must_use]
431 pub fn unaccounted(&self) -> u64 {
432 self.file
433 .saturating_sub(self.columns_total())
434 .saturating_sub(self.indexes)
435 .saturating_sub(self.directory)
436 .saturating_sub(self.header)
437 }
438}
439
440#[derive(Debug)]
442struct GlobalDictionary {
443 primary: HashMap<u64, u32>,
444 collisions: HashMap<u64, Vec<u32>>,
445 offsets: Vec<u32>,
446 payload: Vec<u8>,
447 counts: Vec<u64>,
448 nulls: u64,
449}
450
451impl GlobalDictionary {
452 fn new() -> Self {
453 Self {
454 primary: HashMap::new(),
455 collisions: HashMap::new(),
456 offsets: vec![0],
457 payload: Vec::new(),
458 counts: Vec::new(),
459 nulls: 0,
460 }
461 }
462
463 fn bytes(&self, code: u32) -> Option<&[u8]> {
464 let start = *self.offsets.get(code as usize)? as usize;
465 let end = *self.offsets.get(code as usize + 1)? as usize;
466 self.payload.get(start..end)
467 }
468
469 fn code(&mut self, text: &str) -> Result<u32> {
470 let hash = checksum(text.as_bytes());
471 if let Some(&code) = self.primary.get(&hash) {
472 if self.bytes(code) == Some(text.as_bytes()) {
473 return Ok(code);
474 }
475 if let Some(codes) = self.collisions.get(&hash) {
476 if let Some(code) =
477 codes.iter().copied().find(|&code| self.bytes(code) == Some(text.as_bytes()))
478 {
479 return Ok(code);
480 }
481 }
482 let code = self.insert(text)?;
483 self.collisions.entry(hash).or_default().push(code);
484 return Ok(code);
485 }
486 let code = self.insert(text)?;
487 self.primary.insert(hash, code);
488 Ok(code)
489 }
490
491 fn insert(&mut self, text: &str) -> Result<u32> {
492 let code = u32::try_from(self.offsets.len() - 1)
493 .map_err(|_| invalid("global dictionary has too many values"))?;
494 self.payload.extend_from_slice(text.as_bytes());
495 self.offsets.push(
496 u32::try_from(self.payload.len())
497 .map_err(|_| invalid("global dictionary payload exceeds 4 GiB"))?,
498 );
499 self.counts.push(0);
500 Ok(code)
501 }
502
503 fn ranked(&self) -> Vec<(u64, u32)> {
523 let count = self.offsets.len() - 1;
524 let mut ranked = (0..count)
525 .map(|code| {
526 let code = code as u32;
527 (head(self.bytes(code).unwrap_or_default()), code)
528 })
529 .collect::<Vec<_>>();
530 ranked.sort_unstable_by(|left, right| {
531 left.0.cmp(&right.0).then_with(|| self.bytes(left.1).cmp(&self.bytes(right.1)))
532 });
533 ranked
534 }
535
536 fn observe(&mut self, code: u32, null: bool) -> Result<()> {
537 if null {
538 self.nulls = self.nulls.saturating_add(1);
539 return Ok(());
540 }
541 let count = self
542 .counts
543 .get_mut(code as usize)
544 .ok_or_else(|| invalid("global dictionary count code is out of range"))?;
545 *count = count.saturating_add(1);
546 Ok(())
547 }
548}
549
550#[derive(Debug)]
552pub struct Writer {
553 file: File,
554 at: u64,
562 table: Table,
563 generation: u64,
564 order: Vec<((u64, u64), (u64, u64))>,
567 next_order: u64,
568 dictionaries: Vec<Option<GlobalDictionary>>,
569 pending: Vec<PendingChunk>,
570}
571
572#[derive(Debug)]
580struct PendingChunk {
581 order: (u64, u64),
582 chunk: Chunk,
583}
584
585#[derive(Debug)]
591struct ColumnStripe {
592 pages: Vec<Vec<u8>>,
593 codes: Vec<Option<Vec<u32>>>,
594 sieves: Vec<Option<Sieve>>,
595 ranges: Vec<Range>,
596}
597
598fn weight(ty: &LogicalType) -> usize {
606 match ty {
607 LogicalType::Varchar | LogicalType::Blob => 64,
608 LogicalType::BigInt
609 | LogicalType::UBigInt
610 | LogicalType::Timestamp
611 | LogicalType::Double
612 | LogicalType::Decimal { .. } => 8,
613 LogicalType::Integer | LogicalType::UInteger | LogicalType::Date | LogicalType::Float => 4,
614 LogicalType::SmallInt | LogicalType::USmallInt => 2,
615 _ => 1,
616 }
617}
618
619pub const STRIPE_PARTS: usize = 64;
626
627const INDEX_ENTRY: usize = size_of::<u32>() + size_of::<u64>();
629
630fn index_section(parts: usize) -> Result<usize> {
632 parts
633 .checked_mul(INDEX_ENTRY)
634 .and_then(|bytes| bytes.checked_add(size_of::<u64>()))
635 .ok_or_else(|| invalid("index page length overflow"))
636}
637
638impl Writer {
639 pub fn create(
645 path: impl AsRef<Path>,
646 name: impl Into<String>,
647 fields: Vec<Field>,
648 ) -> Result<Self> {
649 for field in &fields {
650 type_tag(&field.ty)?;
651 }
652 let file =
653 OpenOptions::new().write(true).read(true).create_new(true).open(path).map_err(io)?;
654 let mut header = [0; HEADER as usize];
655 header[..8].copy_from_slice(MAGIC);
656 header[8..12].copy_from_slice(&FORMAT.to_le_bytes());
657 write_at(&file, 0, &header)?;
658 Ok(Self {
659 file,
660 at: HEADER,
661 dictionaries: fields
662 .iter()
663 .map(|field| (field.ty == LogicalType::Varchar).then(GlobalDictionary::new))
664 .collect(),
665 table: Table {
666 name: name.into(),
667 dictionaries: vec![None; fields.len()],
668 distincts: vec![None; fields.len()],
669 fields,
670 stripes: Vec::new(),
671 rows: 0,
672 frequencies: Vec::new(),
673 },
674 generation: 1,
675 order: Vec::new(),
676 next_order: 0,
677 pending: Vec::with_capacity(STRIPE_PARTS),
678 })
679 }
680
681 fn put(&mut self, bytes: &[u8]) -> Result<()> {
686 write_at(&self.file, self.at, bytes)?;
687 self.at = self
688 .at
689 .checked_add(bytes.len() as u64)
690 .ok_or_else(|| invalid("native file length overflow"))?;
691 Ok(())
692 }
693
694 pub fn append(&mut self, chunk: &Chunk) -> Result<()> {
700 let order = (self.next_order, 0);
701 self.next_order = self.next_order.saturating_add(1);
702 self.append_at(order, chunk)
703 }
704
705 pub fn append_at(&mut self, order: (u64, u64), chunk: &Chunk) -> Result<()> {
716 if chunk.is_empty() {
717 return Ok(());
718 }
719 self.admit(chunk)?;
720 if self.pending.last().is_some_and(|last| last.order > order) {
721 self.flush_pending()?;
722 }
723 self.pending.push(PendingChunk { order, chunk: chunk.clone() });
728 if self.pending.len() == STRIPE_PARTS {
729 self.flush_pending()?;
730 }
731 Ok(())
732 }
733
734 pub fn append_stripe(&mut self, parts: Vec<((u64, u64), Chunk)>) -> Result<()> {
750 if parts.len() > STRIPE_PARTS {
751 return Err(invalid("a stripe was handed more parts than it holds"));
752 }
753 self.flush_pending()?;
756 for (order, chunk) in parts {
757 if chunk.is_empty() {
758 continue;
759 }
760 self.admit(&chunk)?;
761 self.pending.push(PendingChunk { order, chunk });
762 }
763 self.flush_pending()
764 }
765
766 fn admit(&mut self, chunk: &Chunk) -> Result<()> {
768 if chunk.width() != self.table.fields.len() {
769 return Err(invalid("chunk width differs from table schema"));
770 }
771 for (index, field) in self.table.fields.iter().enumerate() {
772 if chunk.column(index)?.logical_type() != &field.ty {
773 return Err(invalid("chunk type differs from table schema"));
774 }
775 }
776 self.table.rows = self
777 .table
778 .rows
779 .checked_add(chunk.len())
780 .ok_or_else(|| invalid("row count overflow"))?;
781 Ok(())
782 }
783
784 fn encode_column(
792 index: usize,
793 held: &[PendingChunk],
794 mut dictionary: Option<&mut GlobalDictionary>,
795 ) -> Result<ColumnStripe> {
796 let mut stripe = ColumnStripe {
797 pages: Vec::with_capacity(held.len()),
798 codes: Vec::with_capacity(held.len()),
799 sieves: Vec::with_capacity(held.len()),
800 ranges: Vec::with_capacity(held.len()),
801 };
802 for pending in held {
803 let column = pending.chunk.column(index)?;
804 let (bytes, unique) = encode(column, dictionary.as_deref_mut())?;
805 if bytes.len() > MAX_PAGE {
806 return Err(invalid("column page exceeds the configured bound"));
807 }
808 let range = Range::of(column);
811 let sieve = match dictionary {
824 Some(_) => None,
825 None => Sieve::of(column, &range, SIEVE_BUDGET)
826 .filter(|sieve| sieve.len() < bytes.len()),
827 };
828 stripe.pages.push(bytes);
829 stripe.codes.push(unique);
830 stripe.sieves.push(sieve);
831 stripe.ranges.push(range);
832 }
833 Ok(stripe)
834 }
835
836 fn encode_columns(&mut self, held: &[PendingChunk]) -> Result<Vec<ColumnStripe>> {
845 let width = self.table.fields.len();
846 let workers = std::thread::available_parallelism()
847 .map_or(1, usize::from)
848 .min(MAX_ENCODE_WORKERS)
849 .min(width);
850 if workers <= 1 || held.len() <= 1 {
851 return self
852 .dictionaries
853 .iter_mut()
854 .enumerate()
855 .map(|(index, dictionary)| Self::encode_column(index, held, dictionary.as_mut()))
856 .collect();
857 }
858 let mut jobs: Vec<(usize, Option<GlobalDictionary>)> =
861 std::mem::take(&mut self.dictionaries).into_iter().enumerate().collect();
862 jobs.sort_by_key(|(index, _)| weight(&self.table.fields[*index].ty));
864 let queue = Mutex::new(jobs);
865 let pieces = std::thread::scope(|scope| {
866 (0..workers)
867 .map(|_| {
868 scope.spawn(|| {
869 let mut mine = Vec::new();
870 loop {
871 let taken = queue
872 .lock()
873 .map_err(|_| Error::internal("a native encode worker panicked"))?
874 .pop();
875 let Some((index, mut dictionary)) = taken else { break };
876 let encoded = Self::encode_column(index, held, dictionary.as_mut())?;
877 mine.push((index, dictionary, encoded));
878 }
879 Ok(mine)
880 })
881 })
882 .collect::<Vec<_>>()
883 .into_iter()
884 .map(|handle| {
885 handle.join().map_err(|_| Error::internal("a native encode worker panicked"))?
886 })
887 .collect::<Result<Vec<_>>>()
888 })?;
889 let mut dictionaries: Vec<Option<GlobalDictionary>> = (0..width).map(|_| None).collect();
890 let mut encoded: Vec<Option<ColumnStripe>> = (0..width).map(|_| None).collect();
891 for piece in pieces {
892 for (index, dictionary, stripe) in piece {
893 dictionaries[index] = dictionary;
894 encoded[index] = Some(stripe);
895 }
896 }
897 self.dictionaries = dictionaries;
898 encoded
899 .into_iter()
900 .map(|stripe| stripe.ok_or_else(|| Error::internal("a column was never encoded")))
901 .collect()
902 }
903
904 fn flush_pending(&mut self) -> Result<()> {
906 if self.pending.is_empty() {
907 return Ok(());
908 }
909 let width = self.table.fields.len();
910 let mut held = std::mem::take(&mut self.pending);
913 let parts = held.len();
914 let encoded = self.encode_columns(&held)?;
915 let mut pages = Vec::with_capacity(width);
916 let mut memberships = vec![None; width];
917 let mut ranges = Vec::with_capacity(width);
918 let mut index = Vec::with_capacity(width.saturating_mul(index_section(parts)?));
919 for stripe in &encoded {
920 let offset = self.at;
921 let section = index.len();
922 let mut length = 0_usize;
923 for bytes in &stripe.pages {
924 write_at(&self.file, self.at + length as u64, bytes)?;
925 put_u32(
926 &mut index,
927 u32::try_from(bytes.len()).map_err(|_| invalid("part length overflow"))?,
928 );
929 put_u64(&mut index, checksum(bytes));
930 length = length
931 .checked_add(bytes.len())
932 .ok_or_else(|| invalid("column page length overflow"))?;
933 }
934 let hash = checksum(&index[section..]);
935 put_u64(&mut index, hash);
936 if length > MAX_PAGE {
937 return Err(invalid("column page exceeds the configured bound"));
938 }
939 self.at = self
940 .at
941 .checked_add(length as u64)
942 .ok_or_else(|| invalid("native file length overflow"))?;
943 pages.push(Span {
944 offset,
945 length: u32::try_from(length).map_err(|_| invalid("page length overflow"))?,
946 });
947 ranges.push(merged_range(stripe.ranges.iter().cloned()));
948 }
949 for (membership, stripe) in memberships.iter_mut().zip(&encoded) {
950 if stripe.codes.iter().all(Option::is_none) {
951 continue;
952 }
953 let lists = stripe
954 .codes
955 .iter()
956 .map(|codes| codes.clone().unwrap_or_default())
957 .collect::<Vec<_>>();
958 let bytes = encode_membership(&merged_codes(lists));
959 let offset = self.at;
960 self.put(&bytes)?;
961 *membership = Some(Page {
962 offset,
963 length: u32::try_from(bytes.len())
964 .map_err(|_| invalid("membership page length overflow"))?,
965 hash: checksum(&bytes),
966 });
967 }
968 let mut sieves = vec![None; width];
969 for (page, stripe) in sieves.iter_mut().zip(&encoded) {
970 if stripe.sieves.iter().all(Option::is_none) {
971 continue;
972 }
973 let bytes = encode_sieves(stripe.sieves.iter())?;
974 let offset = self.at;
975 self.put(&bytes)?;
976 *page = Some(Page {
977 offset,
978 length: u32::try_from(bytes.len())
979 .map_err(|_| invalid("sieve page length overflow"))?,
980 hash: checksum(&bytes),
981 });
982 }
983 let mut part_ranges = vec![None; width];
989 if parts > 1 {
990 for ((page, stripe), span) in part_ranges.iter_mut().zip(&encoded).zip(&pages) {
991 let bytes = encode_part_ranges(&stripe.ranges)?;
992 if bytes.len() >= span.length as usize {
993 continue;
994 }
995 let offset = self.at;
996 self.put(&bytes)?;
997 *page = Some(Page {
998 offset,
999 length: u32::try_from(bytes.len())
1000 .map_err(|_| invalid("part range page length overflow"))?,
1001 hash: checksum(&bytes),
1002 });
1003 }
1004 }
1005 let offset = self.at;
1006 self.put(&index)?;
1007 let index = Span {
1008 offset,
1009 length: u32::try_from(index.len())
1010 .map_err(|_| invalid("index page length overflow"))?,
1011 };
1012 let mut rows = 0_usize;
1013 let mut lengths = Vec::with_capacity(parts);
1014 let mut span = None;
1015 for pending in held.drain(..) {
1016 let part = pending.chunk.len();
1017 rows = rows.checked_add(part).ok_or_else(|| invalid("row count overflow"))?;
1018 lengths.push(u32::try_from(part).map_err(|_| invalid("part row count overflow"))?);
1019 span = Some(
1020 span.map_or((pending.order, pending.order), |(first, _)| (first, pending.order)),
1021 );
1022 }
1023 self.order.push(span.ok_or_else(|| invalid("a stripe was flushed with no parts"))?);
1024 self.table.stripes.push(Stripe {
1025 rows,
1026 parts: lengths,
1027 index,
1028 pages,
1029 memberships,
1030 sieves,
1031 part_ranges,
1032 zone: Zone::from_ranges(ranges),
1033 });
1034 self.pending = held;
1036 Ok(())
1037 }
1038
1039 fn numeric_frequency(&self, column: usize) -> Result<Option<FrequencySummary>> {
1043 let ty = &self.table.fields[column].ty;
1044 if !matches!(
1045 ty,
1046 LogicalType::TinyInt
1047 | LogicalType::SmallInt
1048 | LogicalType::Integer
1049 | LogicalType::BigInt
1050 | LogicalType::UTinyInt
1051 | LogicalType::USmallInt
1052 | LogicalType::UInteger
1053 | LogicalType::UBigInt
1054 | LogicalType::Date
1055 | LogicalType::Timestamp
1056 ) {
1057 return Ok(None);
1058 }
1059 let mut candidates: HashMap<FrequencyValue, u32> = HashMap::new();
1060 let mut decrements = 0_u64;
1061 self.visit_numeric(column, |_, value| {
1062 if let Some(count) = candidates.get_mut(&value) {
1063 *count = count.saturating_add(1);
1064 } else if candidates.len() < FREQUENCY_CANDIDATES {
1065 candidates.insert(value, 1);
1066 } else {
1067 candidates.retain(|_, count| {
1068 *count -= 1;
1069 *count != 0
1070 });
1071 decrements = decrements.saturating_add(1);
1072 }
1073 })?;
1074 let (exact, ordinals) = if decrements == 0 {
1075 (
1076 candidates
1077 .into_iter()
1078 .map(|(value, count)| (value, u64::from(count)))
1079 .collect::<HashMap<_, _>>(),
1080 Vec::new(),
1081 )
1082 } else {
1083 let mut lower = candidates.values().copied().collect::<Vec<_>>();
1084 lower.sort_unstable_by(|left, right| right.cmp(left));
1085 if lower.len() < FREQUENCY_BUILD_RANK
1086 || u64::from(lower[FREQUENCY_BUILD_RANK - 1]) <= decrements
1087 {
1088 return Ok(None);
1089 }
1090 let mut exact =
1091 candidates.into_keys().map(|value| (value, 0_u64)).collect::<HashMap<_, _>>();
1092 let mut ordinals = Vec::new();
1093 let mut exceeded = false;
1094 self.visit_numeric(column, |ordinal, value| {
1095 if let Some(count) = exact.get_mut(&value) {
1096 *count = count.saturating_add(1);
1097 if !exceeded {
1098 if ordinals.len() < FREQUENCY_ORDINALS {
1099 ordinals.push(ordinal);
1100 } else {
1101 ordinals.clear();
1102 exceeded = true;
1103 }
1104 }
1105 }
1106 })?;
1107 (exact, ordinals)
1108 };
1109 let mut entries = exact
1110 .into_iter()
1111 .map(|(value, count)| FrequencyEntry { value, count })
1112 .collect::<Vec<_>>();
1113 entries.sort_unstable_by(|left, right| {
1114 right.count.cmp(&left.count).then_with(|| frequency_order(left.value, right.value))
1115 });
1116 let omitted_max =
1117 entries.get(FREQUENCY_ENTRIES).map_or(decrements, |entry| decrements.max(entry.count));
1118 entries.truncate(FREQUENCY_ENTRIES);
1119 Ok(Some(FrequencySummary { entries, omitted_max, ordinals }))
1120 }
1121
1122 fn visit_numeric(
1123 &self,
1124 column: usize,
1125 mut visit: impl FnMut(u64, FrequencyValue),
1126 ) -> Result<()> {
1127 let ty = &self.table.fields[column].ty;
1128 let mut start = 0_u64;
1129 for stripe in &self.table.stripes {
1130 let spans = read_index(&self.file, stripe, column)?;
1131 let page = stripe.pages[column];
1132 let mut bytes = vec![0; page.length as usize];
1133 read_at(&self.file, page.offset, &mut bytes)?;
1134 for (span, &rows) in spans.iter().zip(&stripe.parts) {
1135 let part = part_bytes(&bytes, *span)?;
1136 if checksum(part) != span.hash {
1137 return Err(invalid("column page checksum differs while building frequencies"));
1138 }
1139 let rows = rows as usize;
1140 let vector = decode(ty, rows, part, None)?;
1141 for row in 0..rows {
1143 let value = if vector.is_null_at(row) {
1144 FrequencyValue::Null
1145 } else {
1146 let widened = match vector.signed_at(row) {
1150 Some(value) => Some(value),
1151 None => match vector.value_at(row) {
1152 Value::UTinyInt(value) => Some(i128::from(value)),
1153 Value::USmallInt(value) => Some(i128::from(value)),
1154 Value::UInteger(value) => Some(i128::from(value)),
1155 Value::UBigInt(value) => Some(i128::from(value)),
1156 _ => None,
1157 },
1158 };
1159 FrequencyValue::Integer(widened.ok_or_else(|| {
1160 invalid("numeric frequency page did not contain an integer value")
1161 })?)
1162 };
1163 visit(start.saturating_add(row as u64), value);
1164 }
1165 start = start.saturating_add(rows as u64);
1166 }
1167 }
1168 Ok(())
1169 }
1170
1171 fn numeric_frequencies(&self) -> Result<Vec<Option<FrequencySummary>>> {
1179 let mut columns = self
1180 .table
1181 .fields
1182 .iter()
1183 .enumerate()
1184 .filter_map(|(column, field)| {
1185 matches!(
1186 field.ty,
1187 LogicalType::TinyInt
1188 | LogicalType::SmallInt
1189 | LogicalType::Integer
1190 | LogicalType::BigInt
1191 | LogicalType::UTinyInt
1192 | LogicalType::USmallInt
1193 | LogicalType::UInteger
1194 | LogicalType::UBigInt
1195 | LogicalType::Date
1196 | LogicalType::Timestamp
1197 )
1198 .then_some(column)
1199 })
1200 .collect::<Vec<_>>();
1201 let workers = std::thread::available_parallelism()
1202 .map_or(1, usize::from)
1203 .min(MAX_FREQUENCY_WORKERS)
1204 .min(columns.len());
1205 if workers <= 1 {
1206 let mut frequencies = vec![None; self.table.fields.len()];
1207 for column in columns {
1208 frequencies[column] = self.numeric_frequency(column)?;
1209 }
1210 return Ok(frequencies);
1211 }
1212 columns.sort_by_key(|&column| weight(&self.table.fields[column].ty));
1215 let queue = Mutex::new(columns);
1216 let pieces = std::thread::scope(|scope| {
1217 (0..workers)
1218 .map(|_| {
1219 scope.spawn(|| {
1220 let mut mine = Vec::new();
1221 loop {
1222 let taken = queue
1223 .lock()
1224 .map_err(|_| Error::internal("a native frequency worker panicked"))?
1225 .pop();
1226 let Some(column) = taken else { break };
1227 mine.push((column, self.numeric_frequency(column)?));
1228 }
1229 Ok(mine)
1230 })
1231 })
1232 .collect::<Vec<_>>()
1233 .into_iter()
1234 .map(|handle| {
1235 handle
1236 .join()
1237 .map_err(|_| Error::internal("a native frequency worker panicked"))?
1238 })
1239 .collect::<Result<Vec<_>>>()
1240 })?;
1241 let mut frequencies = vec![None; self.table.fields.len()];
1242 for piece in pieces {
1243 for (column, summary) in piece {
1244 frequencies[column] = summary;
1245 }
1246 }
1247 Ok(frequencies)
1248 }
1249
1250 pub fn finish(mut self) -> Result<Table> {
1256 self.flush_pending()?;
1257 let mut stripes = std::mem::take(&mut self.order)
1258 .into_iter()
1259 .zip(std::mem::take(&mut self.table.stripes))
1260 .collect::<Vec<_>>();
1261 stripes.sort_by_key(|(order, _)| order.0);
1262 let mut previous: Option<(u64, u64)> = None;
1263 for ((first, last), _) in &stripes {
1264 if previous.is_some_and(|previous| previous >= *first) {
1265 return Err(invalid("chunks did not arrive in source order"));
1266 }
1267 previous = Some(*last);
1268 }
1269 self.table.stripes = stripes.into_iter().map(|(_, stripe)| stripe).collect();
1270 self.table.frequencies = self.numeric_frequencies()?;
1271 let dictionaries = std::mem::take(&mut self.dictionaries);
1272 let orders = rankings(&dictionaries)?;
1273 for (index, (dictionary, order)) in dictionaries.into_iter().zip(orders).enumerate() {
1274 let Some(dictionary) = dictionary else { continue };
1275 self.table.distincts[index] =
1279 Some(dictionary.counts.iter().filter(|count| **count != 0).count() as u64);
1280 self.table.frequencies[index] = Some(code_frequency(&dictionary));
1281 let encoded = encode_global_dictionary(dictionary, &order)?;
1282 let offset = self.at;
1283 self.put(&encoded.index)?;
1284 self.put(&encoded.ranks)?;
1285 for block in &encoded.payload {
1286 self.put(block)?;
1287 }
1288 let payload_len =
1289 encoded.payload.iter().try_fold(0_usize, |len, block| len.checked_add(block.len()));
1290 let length = payload_len
1291 .and_then(|len| len.checked_add(encoded.index.len()))
1292 .and_then(|len| len.checked_add(encoded.ranks.len()))
1293 .ok_or_else(|| invalid("dictionary page length overflow"))?;
1294 self.table.dictionaries[index] = Some(Page {
1295 offset,
1296 length: u32::try_from(length)
1297 .map_err(|_| invalid("dictionary page length overflow"))?,
1298 hash: checksum(&encoded.index),
1299 });
1300 }
1301 let directory = encode_directory(&self.table)?;
1302 if directory.len() > MAX_DIRECTORY {
1303 return Err(invalid("directory exceeds the configured bound"));
1304 }
1305 let offset = self.at;
1306 self.put(&directory)?;
1307 self.file.sync_all().map_err(io)?;
1308 let slot = Slot {
1309 offset,
1310 length: u32::try_from(directory.len())
1311 .map_err(|_| invalid("directory length overflow"))?,
1312 generation: self.generation,
1313 hash: checksum(&directory),
1314 };
1315 write_at(&self.file, 16, &slot.bytes())?;
1318 self.file.sync_all().map_err(io)?;
1319 Ok(self.table)
1320 }
1321}
1322
1323#[derive(Debug, Clone)]
1325pub struct Reader {
1326 file: Arc<File>,
1327 table: Arc<Table>,
1328 dictionaries: Arc<Vec<OnceLock<Arc<Vector>>>>,
1329 loading: Arc<Vec<Mutex<()>>>,
1338 opened: Arc<AtomicUsize>,
1342 sieves: Arc<Vec<Vec<SieveSlot>>>,
1346 part_ranges: Arc<Vec<Vec<RangeSlot>>>,
1349 places: Arc<Vec<Place>>,
1351 cache: Arc<Vec<Mutex<Cached>>>,
1352 pages: Arc<AtomicUsize>,
1355 indexes: Arc<AtomicUsize>,
1358 kept: Arc<AtomicUsize>,
1361 size: u64,
1363 directory: u64,
1365 opening: Opening,
1367}
1368
1369#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1381pub struct Opening {
1382 pub reads: u32,
1385 pub bytes: u64,
1387}
1388
1389#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1391pub struct Reads {
1392 pub opening: Opening,
1394 pub pages: usize,
1396 pub indexes: usize,
1398 pub dictionaries: usize,
1401}
1402
1403#[derive(Debug, Clone, Copy)]
1405struct Place {
1406 stripe: u32,
1407 part: u32,
1408 rows: u32,
1409}
1410
1411#[derive(Debug, Clone, Copy)]
1413struct PartSpan {
1414 start: usize,
1415 length: usize,
1416 hash: u64,
1417}
1418
1419#[derive(Debug, Clone)]
1425struct CachedColumn {
1426 stripe: usize,
1427 index: Arc<Vec<PartSpan>>,
1428 page: Option<Arc<Vec<u8>>>,
1429}
1430
1431#[derive(Debug, Default)]
1451struct Cached {
1452 pages: Vec<Option<Arc<Vec<u8>>>>,
1453 order: VecDeque<usize>,
1454 loading: Vec<usize>,
1455 index: Vec<Option<Arc<Vec<PartSpan>>>>,
1456}
1457
1458const CACHED_STRIPES_PER_COLUMN: usize = 4;
1470
1471type SieveSlot = OnceLock<Arc<Vec<Option<Sieve>>>>;
1473
1474type RangeSlot = OnceLock<Arc<Vec<Range>>>;
1475
1476#[derive(Debug)]
1477struct NativeText {
1478 file: Arc<File>,
1479 values: usize,
1481 offsets: Vec<u8>,
1490 offset_bits: usize,
1493 ranks: usize,
1495 rank_at: u64,
1499 rank_ends: Vec<u64>,
1503 rank_hashes: Vec<u64>,
1504 rank_blocks: Vec<OnceLock<Result<Vec<u8>>>>,
1505 code_bits: usize,
1508 code_ranks: OnceLock<Option<Vec<u32>>>,
1515 payload: u64,
1516 ends: Vec<u64>,
1519 hashes: Vec<u64>,
1520 blocks: Vec<OnceLock<Result<Vec<u8>>>>,
1522 keep_budget: usize,
1525 payload_kept: AtomicUsize,
1533}
1534
1535const TEXT_PAYLOAD_VALUES: usize = 1024;
1551
1552const TEXT_KEEP_BUDGET: usize = 256 * 1024 * 1024;
1573
1574const TEXT_OFFSET_RUN: usize = 512;
1581
1582const DICTIONARY_HEADER: usize = 16;
1585
1586const TEXT_RANK_BLOCK: usize = 512;
1597
1598const RANK_BLOCK_HEADER: usize = size_of::<u64>() + 1;
1612
1613impl NativeText {
1614 fn payload_block(&self, block: usize) -> Result<Option<&[u8]>> {
1621 let Some(slot) = self.blocks.get(block) else { return Ok(None) };
1622 let bytes = slot.get_or_init(|| self.decode_block(block)).as_ref().map_err(Clone::clone)?;
1623 Ok(Some(bytes.as_slice()))
1624 }
1625
1626 fn decode_block(&self, block: usize) -> Result<Vec<u8>> {
1631 let start = if block == 0 { 0 } else { self.ends[block - 1] };
1632 let end = self.ends[block];
1633 let len = end
1634 .checked_sub(start)
1635 .ok_or_else(|| invalid("global dictionary block ends before it starts"))?;
1636 let mut stored = vec![
1637 0;
1638 usize::try_from(len).map_err(|_| invalid(
1639 "global dictionary block does not fit in memory"
1640 ))?
1641 ];
1642 read_at(&self.file, self.payload + start, &mut stored)?;
1643 if checksum(&stored) != self.hashes[block] {
1644 return Err(invalid("global dictionary payload checksum differs"));
1645 }
1646 let first = block * TEXT_PAYLOAD_VALUES;
1647 let last = (first + TEXT_PAYLOAD_VALUES).min(self.values);
1648 let want = self.end_within(last - 1)? as usize;
1649 let values = string::decode_flat(&stored)?;
1650 if values.len() != last - first {
1651 return Err(invalid("global dictionary block holds the wrong value count"));
1652 }
1653 let bytes = values.into_bytes();
1654 if bytes.len() != want {
1655 return Err(invalid("global dictionary block decodes to the wrong length"));
1656 }
1657 Ok(bytes)
1658 }
1659
1660 fn end_within(&self, index: usize) -> Result<u32> {
1662 let run = index / TEXT_OFFSET_RUN;
1663 let bytes = self
1664 .offsets
1665 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
1666 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1667 let end = bitpack::tail_at(bytes, self.offset_bits, index % TEXT_OFFSET_RUN)
1668 .map_err(|_| invalid("global dictionary offsets are short"))?;
1669 u32::try_from(end).map_err(|_| invalid("global dictionary offset is past the payload"))
1670 }
1671
1672 fn ends_within(&self, first: usize, last: usize) -> Result<Vec<u64>> {
1685 let mut ends = Vec::with_capacity(last.saturating_sub(first));
1686 let mut at = first;
1687 while at < last {
1688 let run = at / TEXT_OFFSET_RUN;
1689 let stop = ((run + 1) * TEXT_OFFSET_RUN).min(last);
1690 let held = self.values.saturating_sub(run * TEXT_OFFSET_RUN).min(TEXT_OFFSET_RUN);
1691 let bytes = self
1692 .offsets
1693 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
1694 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1695 let run_ends = bitpack::unpack_tail(bytes, self.offset_bits, held)
1696 .map_err(|_| invalid("global dictionary offsets are short"))?;
1697 let within = run_ends
1698 .get(at % TEXT_OFFSET_RUN..stop - run * TEXT_OFFSET_RUN)
1699 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1700 ends.extend_from_slice(within);
1701 at = stop;
1702 }
1703 Ok(ends)
1704 }
1705
1706 fn start_within(&self, index: usize) -> Result<u32> {
1709 if index % TEXT_PAYLOAD_VALUES == 0 { Ok(0) } else { self.end_within(index - 1) }
1710 }
1711
1712 fn span_within(&self, index: usize) -> Result<(u32, u32)> {
1720 let within = index % TEXT_OFFSET_RUN;
1721 let (start, end) = if within == 0 {
1722 (self.start_within(index)?, self.end_within(index)?)
1723 } else {
1724 let run = index / TEXT_OFFSET_RUN;
1725 let bytes = self
1726 .offsets
1727 .get(run * TEXT_OFFSET_RUN / 8 * self.offset_bits..)
1728 .ok_or_else(|| invalid("global dictionary offsets are short"))?;
1729 let (start, end) = bitpack::tail_pair(bytes, self.offset_bits, within)
1730 .map_err(|_| invalid("global dictionary offsets are short"))?;
1731 let ends = u32::try_from(end)
1732 .map_err(|_| invalid("global dictionary offset is past the payload"))?;
1733 let starts = u32::try_from(start)
1734 .map_err(|_| invalid("global dictionary offset is past the payload"))?;
1735 (starts, ends)
1736 };
1737 if start > end {
1738 return Err(invalid("global dictionary value ends before it starts"));
1739 }
1740 Ok((start, end))
1741 }
1742
1743 fn rank_parts(&self, rank: usize) -> Result<(&[u8], usize)> {
1750 let slot = self
1751 .rank_blocks
1752 .get(rank / TEXT_RANK_BLOCK)
1753 .ok_or_else(|| invalid("global dictionary rank is past the order"))?;
1754 let block = slot
1755 .get_or_init(|| {
1756 let which = rank / TEXT_RANK_BLOCK;
1757 let start = if which == 0 { 0 } else { self.rank_ends[which - 1] };
1758 let end = self.rank_ends[which];
1759 let mut bytes = vec![0; (end - start) as usize];
1760 read_at(&self.file, self.rank_at + start, &mut bytes)?;
1761 if checksum(&bytes)
1762 != *self
1763 .rank_hashes
1764 .get(rank / TEXT_RANK_BLOCK)
1765 .ok_or_else(|| invalid("global dictionary rank block has no checksum"))?
1766 {
1767 return Err(invalid("global dictionary rank checksum differs"));
1768 }
1769 Ok(bytes)
1770 })
1771 .as_ref()
1772 .map_err(Clone::clone)?;
1773 Ok((block.as_slice(), rank % TEXT_RANK_BLOCK))
1774 }
1775
1776 fn head_at(&self, rank: usize) -> Result<u64> {
1778 let (block, within) = self.rank_parts(rank)?;
1779 let (base, width, packed) = rank_heads(block)?;
1780 let above = bitpack::tail_at(packed, width, within)
1781 .map_err(|_| invalid("global dictionary rank block is short of heads"))?;
1782 Ok(base.wrapping_add(above))
1783 }
1784
1785 fn rank_codes<'block>(&self, block: &'block [u8], count: usize) -> Result<&'block [u8]> {
1787 let (_, width, packed) = rank_heads(block)?;
1788 packed
1789 .get(bitpack::tail_len(count, width)..)
1790 .ok_or_else(|| invalid("global dictionary rank block is short of codes"))
1791 }
1792
1793 fn rank_block_len(&self, rank: usize) -> usize {
1795 let first = rank / TEXT_RANK_BLOCK * TEXT_RANK_BLOCK;
1796 TEXT_RANK_BLOCK.min(self.ranks - first)
1797 }
1798}
1799
1800fn rank_heads(block: &[u8]) -> Result<(u64, usize, &[u8])> {
1802 let header = block
1803 .get(..RANK_BLOCK_HEADER)
1804 .ok_or_else(|| invalid("global dictionary rank block is short"))?;
1805 let base = u64::from_le_bytes(header[..8].try_into().expect("eight bytes"));
1806 let width = header[8] as usize;
1807 if width > 64 {
1808 return Err(invalid("global dictionary rank block packs heads past a word"));
1809 }
1810 Ok((base, width, &block[RANK_BLOCK_HEADER..]))
1811}
1812
1813fn offset_width(offsets: &[u32]) -> usize {
1820 let values = offsets.len() - 1;
1821 let mut span = 0;
1822 for first in (0..values).step_by(TEXT_PAYLOAD_VALUES) {
1823 let last = (first + TEXT_PAYLOAD_VALUES).min(values);
1824 span = span.max(offsets[last] - offsets[first]);
1825 }
1826 (u32::BITS - span.leading_zeros()) as usize
1827}
1828
1829fn offset_bytes(values: usize, bits: usize) -> usize {
1832 let full = values / TEXT_OFFSET_RUN;
1833 let rest = values % TEXT_OFFSET_RUN;
1834 full * TEXT_OFFSET_RUN / 8 * bits + bitpack::tail_len(rest, bits)
1835}
1836
1837fn encode_offsets(offsets: &[u32], bits: usize, out: &mut Vec<u8>) -> Result<()> {
1839 let values = offsets.len() - 1;
1840 let mut run = Vec::with_capacity(TEXT_OFFSET_RUN);
1841 for first in (0..values).step_by(TEXT_OFFSET_RUN) {
1842 let last = (first + TEXT_OFFSET_RUN).min(values);
1843 let base = offsets[first / TEXT_PAYLOAD_VALUES * TEXT_PAYLOAD_VALUES];
1844 run.clear();
1845 run.extend((first..last).map(|value| u64::from(offsets[value + 1] - base)));
1846 bitpack::pack_tail(&run, bits, out)
1847 .map_err(|_| invalid("global dictionary offsets do not pack"))?;
1848 }
1849 Ok(())
1850}
1851
1852fn code_width(values: usize) -> usize {
1854 match u64::try_from(values).unwrap_or(u64::MAX) {
1855 0 | 1 => 0,
1856 last => (u64::BITS - (last - 1).leading_zeros()) as usize,
1857 }
1858}
1859
1860impl TextSource for NativeText {
1861 fn len(&self) -> usize {
1862 self.values
1863 }
1864
1865 fn bytes_at(&self, index: usize) -> Result<Option<&[u8]>> {
1866 if index >= self.values {
1867 return Ok(None);
1868 }
1869 let (start, end) = self.span_within(index)?;
1870 if start == end {
1871 return Ok(Some(&[]));
1872 }
1873 let block = index / TEXT_PAYLOAD_VALUES;
1876 let Some(bytes) = self.payload_block(block)? else { return Ok(None) };
1877 Ok(bytes.get(start as usize..end as usize))
1878 }
1879
1880 fn bytes_len_at(&self, index: usize) -> Result<Option<usize>> {
1881 if index >= self.values {
1882 return Ok(None);
1883 }
1884 let (start, end) = self.span_within(index)?;
1885 Ok(Some((end - start) as usize))
1886 }
1887
1888 fn sweep(
1901 &self,
1902 first: usize,
1903 limit: usize,
1904 body: &mut dyn FnMut(usize, &[u8]) -> Result<()>,
1905 ) -> Result<usize> {
1906 let limit = limit.min(self.values);
1907 if first >= limit {
1908 return Ok(first);
1909 }
1910 let block = first / TEXT_PAYLOAD_VALUES;
1911 let last = ((block + 1) * TEXT_PAYLOAD_VALUES).min(limit);
1912 let decoded;
1913 let bytes: &[u8] = match self.blocks.get(block).and_then(OnceLock::get) {
1914 Some(Ok(kept)) => kept,
1915 _ if self.payload_kept.load(Atomic::Relaxed) < self.keep_budget => {
1916 let kept = self
1917 .payload_block(block)?
1918 .ok_or_else(|| invalid("global dictionary block is past the payload"))?;
1919 self.payload_kept.fetch_add(kept.len(), Atomic::Relaxed);
1920 kept
1921 }
1922 _ => {
1923 decoded = self.decode_block(block)?;
1924 &decoded
1925 }
1926 };
1927 let ends = self.ends_within(first, last)?;
1928 if ends.len() != last - first {
1929 return Err(invalid("global dictionary offsets are short"));
1930 }
1931 let mut start = u64::from(self.start_within(first)?);
1932 for (index, &end) in (first..last).zip(&ends) {
1935 let value = usize::try_from(start)
1936 .ok()
1937 .zip(usize::try_from(end).ok())
1938 .and_then(|(from, to)| bytes.get(from..to))
1939 .ok_or_else(|| invalid("global dictionary value is past its block"))?;
1940 body(index, value)?;
1941 start = end;
1942 }
1943 Ok(last)
1944 }
1945
1946 fn ranks(&self) -> Option<usize> {
1947 (self.ranks > 0).then_some(self.ranks)
1948 }
1949
1950 fn compare_rank(&self, rank: usize, wanted: &[u8]) -> Result<Ordering> {
1951 let settled = self.head_at(rank)?.cmp(&head(wanted));
1955 if settled != Ordering::Equal {
1956 return Ok(settled);
1957 }
1958 let code = self.code_at_rank(rank)?;
1959 let bytes = self
1960 .bytes_at(code as usize)?
1961 .ok_or_else(|| invalid("global dictionary order names a code it does not have"))?;
1962 Ok(bytes.cmp(wanted))
1963 }
1964
1965 fn code_at_rank(&self, rank: usize) -> Result<u32> {
1966 let (block, within) = self.rank_parts(rank)?;
1967 let codes = self.rank_codes(block, self.rank_block_len(rank))?;
1968 let code = bitpack::tail_at(codes, self.code_bits, within)
1969 .map_err(|_| invalid("global dictionary rank block is short of codes"))?;
1970 let code = u32::try_from(code)
1971 .map_err(|_| invalid("global dictionary order names a code it does not have"))?;
1972 if code as usize >= self.len() {
1973 return Err(invalid("global dictionary order names a code it does not have"));
1974 }
1975 Ok(code)
1976 }
1977
1978 fn code_ranks(&self) -> Option<&[u32]> {
1979 if self.ranks == 0 || self.ranks != self.len() {
1983 return None;
1984 }
1985 self.code_ranks
1986 .get_or_init(|| {
1987 let mut ranks = vec![u32::MAX; self.ranks];
1988 for first in (0..self.ranks).step_by(TEXT_RANK_BLOCK) {
1991 let (block, _) = self.rank_parts(first).ok()?;
1992 let count = self.rank_block_len(first);
1993 let codes = self.rank_codes(block, count).ok()?;
1994 for (within, code) in bitpack::unpack_tail(codes, self.code_bits, count)
1995 .ok()?
1996 .into_iter()
1997 .enumerate()
1998 {
1999 let code = usize::try_from(code).ok()?;
2000 *ranks.get_mut(code)? = u32::try_from(first + within).ok()?;
2001 }
2002 }
2003 if ranks.contains(&u32::MAX) {
2004 return None;
2005 }
2006 Some(ranks)
2007 })
2008 .as_deref()
2009 }
2010
2011 fn footprint(&self) -> usize {
2012 self.offsets.capacity()
2013 + self
2014 .code_ranks
2015 .get()
2016 .and_then(Option::as_ref)
2017 .map_or(0, |ranks| ranks.capacity() * size_of::<u32>())
2018 + self.rank_hashes.capacity() * size_of::<u64>()
2019 + self.rank_ends.capacity() * size_of::<u64>()
2020 + self.rank_blocks.capacity() * size_of::<OnceLock<Result<Vec<u8>>>>()
2021 + self
2022 .rank_blocks
2023 .iter()
2024 .filter_map(OnceLock::get)
2025 .filter_map(|result| result.as_ref().ok())
2026 .map(Vec::capacity)
2027 .sum::<usize>()
2028 + self.blocks.capacity() * size_of::<OnceLock<Result<Vec<u8>>>>()
2029 + self.hashes.capacity() * size_of::<u64>()
2030 + self.ends.capacity() * size_of::<u64>()
2031 + self
2032 .blocks
2033 .iter()
2034 .filter_map(OnceLock::get)
2035 .filter_map(|result| result.as_ref().ok())
2036 .map(Vec::capacity)
2037 .sum::<usize>()
2038 }
2039}
2040
2041fn places(table: &Table) -> Result<Vec<Place>> {
2043 let mut places = Vec::with_capacity(table.stripes.len().saturating_mul(STRIPE_PARTS));
2044 for (at, stripe) in table.stripes.iter().enumerate() {
2045 let index = u32::try_from(at).map_err(|_| invalid("too many stripes"))?;
2046 for (part, &rows) in stripe.parts.iter().enumerate() {
2047 places.push(Place {
2048 stripe: index,
2049 part: u32::try_from(part).map_err(|_| invalid("too many parts in a stripe"))?,
2050 rows,
2051 });
2052 }
2053 }
2054 Ok(places)
2055}
2056
2057fn read_index(file: &File, stripe: &Stripe, column: usize) -> Result<Vec<PartSpan>> {
2062 let parts = stripe.parts.len();
2063 let section = index_section(parts)?;
2064 let at = column.checked_mul(section).ok_or_else(|| invalid("index page offset overflow"))?;
2065 let end = at.checked_add(section).ok_or_else(|| invalid("index page offset overflow"))?;
2066 if end > stripe.index.length as usize {
2067 return Err(invalid("index page is shorter than its columns"));
2068 }
2069 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
2070 let mut bytes = vec![0; section];
2071 let offset = stripe
2072 .index
2073 .offset
2074 .checked_add(at as u64)
2075 .ok_or_else(|| invalid("index page offset overflow"))?;
2076 read_at(file, offset, &mut bytes)?;
2077 let entries = section - size_of::<u64>();
2078 let stored = u64::from_le_bytes(bytes[entries..].try_into().expect("eight bytes"));
2079 if checksum(&bytes[..entries]) != stored {
2080 return Err(invalid(&format!(
2083 "index page section checksum differs, column {column} of {parts} parts at {offset}, \
2084 wanted {stored:016x} and got {:016x}",
2085 checksum(&bytes[..entries]),
2086 )));
2087 }
2088 let mut spans = Vec::with_capacity(parts);
2089 let mut start = 0_usize;
2090 for part in 0..parts {
2091 let at = part * INDEX_ENTRY;
2092 let length = u32::from_le_bytes(bytes[at..at + 4].try_into().expect("four bytes")) as usize;
2093 let hash = u64::from_le_bytes(bytes[at + 4..at + 12].try_into().expect("eight bytes"));
2094 spans.push(PartSpan { start, length, hash });
2095 start = start.checked_add(length).ok_or_else(|| invalid("column page length overflow"))?;
2096 }
2097 if start != page.length as usize {
2098 return Err(invalid("column page length differs from its index"));
2099 }
2100 Ok(spans)
2101}
2102
2103fn part_bytes(page: &[u8], span: PartSpan) -> Result<&[u8]> {
2105 let end = span.start.checked_add(span.length).ok_or_else(|| invalid("part range overflow"))?;
2106 page.get(span.start..end).ok_or_else(|| invalid("part exceeds its column page"))
2107}
2108
2109fn remember(cached: &mut Cached, held: &CachedColumn, kept: usize) {
2114 if let Some(slot) = cached.index.get_mut(held.stripe) {
2115 if slot.is_none() {
2116 *slot = Some(Arc::clone(&held.index));
2117 }
2118 }
2119 let Some(page) = held.page.clone() else { return };
2120 let Some(slot) = cached.pages.get_mut(held.stripe) else { return };
2121 if slot.is_none() {
2122 cached.order.push_back(held.stripe);
2123 }
2124 *slot = Some(page);
2125 while cached.order.len() > kept.max(1) {
2126 let Some(oldest) = cached.order.pop_front() else { break };
2127 if let Some(slot) = cached.pages.get_mut(oldest) {
2128 *slot = None;
2129 }
2130 }
2131}
2132
2133impl Reader {
2134 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
2140 let mut file = File::open(path).map_err(io)?;
2141 let size = file.metadata().map_err(io)?.len();
2142 if size < HEADER {
2143 return Err(invalid("file is shorter than its header"));
2144 }
2145 let mut header = [0; HEADER as usize];
2146 file.read_exact(&mut header).map_err(io)?;
2147 let mut opening = Opening { reads: 1, bytes: HEADER };
2148 let version = u32::from_le_bytes([header[8], header[9], header[10], header[11]]);
2149 if &header[..8] != MAGIC {
2154 return Err(invalid("the header does not begin with a rudb native magic"));
2155 }
2156 if version != FORMAT {
2157 return Err(invalid(&format!(
2158 "the file is format {version} and this build reads format {FORMAT}, so it has to \
2159 be written again"
2160 )));
2161 }
2162 let mut selected = None;
2163 for start in [16, 16 + SLOT_BYTES] {
2164 let slot = Slot::read(&header[start..start + SLOT_BYTES]);
2165 if slot.generation == 0 || slot.length == 0 || slot.length as usize > MAX_DIRECTORY {
2166 continue;
2167 }
2168 let Some(end) = slot.offset.checked_add(u64::from(slot.length)) else { continue };
2169 if slot.offset < HEADER || end > size {
2170 continue;
2171 }
2172 let mut bytes = vec![0; slot.length as usize];
2173 file.seek(SeekFrom::Start(slot.offset)).map_err(io)?;
2174 file.read_exact(&mut bytes).map_err(io)?;
2175 opening.reads += 1;
2176 opening.bytes += u64::from(slot.length);
2177 if checksum(&bytes) == slot.hash
2178 && selected
2179 .as_ref()
2180 .is_none_or(|(old, _): &(Slot, Vec<u8>)| old.generation < slot.generation)
2181 {
2182 selected = Some((slot, bytes));
2183 }
2184 }
2185 let (slot, bytes) =
2186 selected.ok_or_else(|| invalid("no committed directory slot is valid"))?;
2187 let table = decode_directory(&bytes, size)?;
2188 let places = places(&table)?;
2189 let dictionaries = (0..table.fields.len()).map(|_| OnceLock::new()).collect();
2190 let table_fields = table.fields.len();
2191 let stripes = table.stripes.len();
2192 let cache = (0..table.fields.len())
2193 .map(|_| {
2194 Mutex::new(Cached {
2195 pages: (0..stripes).map(|_| None).collect(),
2196 index: (0..stripes).map(|_| None).collect(),
2197 ..Cached::default()
2198 })
2199 })
2200 .collect::<Vec<_>>();
2201 let sieves: Vec<Vec<SieveSlot>> = (0..table.fields.len())
2202 .map(|_| table.stripes.iter().map(|_| OnceLock::new()).collect())
2203 .collect();
2204 let part_ranges: Vec<Vec<RangeSlot>> = (0..table.fields.len())
2205 .map(|_| table.stripes.iter().map(|_| OnceLock::new()).collect())
2206 .collect();
2207 Ok(Self {
2208 file: Arc::new(file),
2209 table: Arc::new(table),
2210 dictionaries: Arc::new(dictionaries),
2211 loading: Arc::new((0..table_fields).map(|_| Mutex::new(())).collect()),
2212 opened: Arc::new(AtomicUsize::new(0)),
2213 sieves: Arc::new(sieves),
2214 part_ranges: Arc::new(part_ranges),
2215 places: Arc::new(places),
2216 cache: Arc::new(cache),
2217 pages: Arc::new(AtomicUsize::new(0)),
2218 indexes: Arc::new(AtomicUsize::new(0)),
2219 kept: Arc::new(AtomicUsize::new(CACHED_STRIPES_PER_COLUMN)),
2220 size,
2221 directory: u64::from(slot.length),
2222 opening,
2223 })
2224 }
2225
2226 #[must_use]
2233 pub fn reads(&self) -> Reads {
2234 Reads {
2235 opening: self.opening,
2236 pages: self.pages.load(Atomic::Relaxed),
2237 indexes: self.indexes.load(Atomic::Relaxed),
2238 dictionaries: self.opened.load(Atomic::Relaxed),
2239 }
2240 }
2241
2242 #[must_use]
2247 pub fn layout(&self) -> Layout {
2248 let table = &self.table;
2249 let stripes = table.stripes.as_slice();
2250 let columns = table
2251 .fields
2252 .iter()
2253 .enumerate()
2254 .map(|(at, field)| ColumnLayout {
2255 name: field.name.clone(),
2256 kind: field.ty.to_string(),
2257 pages: sum(stripes.iter().map(|stripe| span_bytes(&stripe.pages, at))),
2258 memberships: sum(stripes.iter().map(|stripe| page_bytes(&stripe.memberships, at))),
2259 sieves: sum(stripes.iter().map(|stripe| page_bytes(&stripe.sieves, at))),
2260 part_ranges: sum(stripes.iter().map(|stripe| page_bytes(&stripe.part_ranges, at))),
2261 dictionary: page_bytes(&table.dictionaries, at),
2262 })
2263 .collect();
2264 Layout {
2265 file: self.size,
2266 rows: table.rows,
2267 stripes: stripes.len(),
2268 parts: self.places.len(),
2269 columns,
2270 indexes: sum(stripes.iter().map(|stripe| u64::from(stripe.index.length))),
2271 directory: self.directory,
2272 header: HEADER,
2273 }
2274 }
2275
2276 #[must_use]
2278 pub fn parts(&self) -> usize {
2279 self.places.len()
2280 }
2281
2282 #[must_use]
2289 pub fn stripe_parts(&self) -> Vec<std::ops::Range<usize>> {
2290 let mut runs = Vec::with_capacity(self.table.stripes.len());
2291 let mut start = 0;
2292 for stripe in &self.table.stripes {
2293 let end = start + stripe.parts.len();
2294 runs.push(start..end);
2295 start = end;
2296 }
2297 runs
2298 }
2299
2300 pub fn keep_stripes(&self, stripes: usize) {
2307 self.kept.fetch_max(stripes, Atomic::Relaxed);
2308 }
2309
2310 #[must_use]
2312 pub fn part_rows(&self, at: usize) -> usize {
2313 self.places.get(at).map_or(0, |place| place.rows as usize)
2314 }
2315
2316 #[must_use]
2318 pub fn table(&self) -> &Table {
2319 &self.table
2320 }
2321
2322 pub fn top_frequencies(&self, column: usize, top: usize) -> Result<Option<Vec<(Value, u64)>>> {
2331 let field = self
2332 .table
2333 .fields
2334 .get(column)
2335 .ok_or_else(|| invalid("frequency column index out of range"))?;
2336 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2337 return Ok(None);
2338 };
2339 if top == 0 || summary.entries.len() < top {
2340 return Ok(None);
2341 }
2342 let boundary = summary.entries[top - 1].count;
2343 if boundary <= summary.omitted_max {
2344 return Ok(None);
2345 }
2346 self.decode_frequencies(column, &field.ty, &summary.entries).map(Some)
2347 }
2348
2349 pub fn exact_frequencies(&self, column: usize) -> Result<Option<Vec<(Value, u64)>>> {
2369 let field = self
2370 .table
2371 .fields
2372 .get(column)
2373 .ok_or_else(|| invalid("frequency column index out of range"))?;
2374 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2375 return Ok(None);
2376 };
2377 if summary.omitted_max > 0 {
2378 return Ok(None);
2379 }
2380 self.decode_frequencies(column, &field.ty, &summary.entries).map(Some)
2381 }
2382
2383 fn decode_frequencies(
2385 &self,
2386 column: usize,
2387 ty: &LogicalType,
2388 entries: &[FrequencyEntry],
2389 ) -> Result<Vec<(Value, u64)>> {
2390 let dictionary = if *ty == LogicalType::Varchar { self.dictionary(column)? } else { None };
2391 let mut out = Vec::with_capacity(entries.len());
2392 for entry in entries {
2393 let value = match entry.value {
2394 FrequencyValue::Null => Value::Null,
2395 FrequencyValue::Integer(value) => match *ty {
2396 LogicalType::TinyInt => Value::TinyInt(
2397 i8::try_from(value)
2398 .map_err(|_| invalid("frequency TINYINT is out of range"))?,
2399 ),
2400 LogicalType::UTinyInt => Value::UTinyInt(
2401 u8::try_from(value)
2402 .map_err(|_| invalid("frequency UTINYINT is out of range"))?,
2403 ),
2404 LogicalType::USmallInt => Value::USmallInt(
2405 u16::try_from(value)
2406 .map_err(|_| invalid("frequency USMALLINT is out of range"))?,
2407 ),
2408 LogicalType::UInteger => Value::UInteger(
2409 u32::try_from(value)
2410 .map_err(|_| invalid("frequency UINTEGER is out of range"))?,
2411 ),
2412 LogicalType::UBigInt => Value::UBigInt(
2413 u64::try_from(value)
2414 .map_err(|_| invalid("frequency UBIGINT is out of range"))?,
2415 ),
2416 LogicalType::SmallInt => Value::SmallInt(
2417 i16::try_from(value)
2418 .map_err(|_| invalid("frequency SMALLINT is out of range"))?,
2419 ),
2420 LogicalType::Integer => Value::Integer(
2421 i32::try_from(value)
2422 .map_err(|_| invalid("frequency INTEGER is out of range"))?,
2423 ),
2424 LogicalType::BigInt => Value::BigInt(
2425 i64::try_from(value)
2426 .map_err(|_| invalid("frequency BIGINT is out of range"))?,
2427 ),
2428 LogicalType::Date => Value::Date(
2429 i32::try_from(value)
2430 .map_err(|_| invalid("frequency DATE is out of range"))?,
2431 ),
2432 LogicalType::Timestamp => Value::Timestamp(
2433 i64::try_from(value)
2434 .map_err(|_| invalid("frequency TIMESTAMP is out of range"))?,
2435 ),
2436 _ => return Err(invalid("integer frequency belongs to another type")),
2437 },
2438 FrequencyValue::Code(code) => dictionary
2439 .as_ref()
2440 .ok_or_else(|| invalid("frequency code has no dictionary"))?
2441 .try_value_at(code as usize)?,
2442 };
2443 out.push((value, entry.count));
2444 }
2445 Ok(out)
2446 }
2447
2448 pub fn frequency_occurrences(&self, column: usize) -> Result<Option<FrequencyOccurrences>> {
2458 self.table
2459 .fields
2460 .get(column)
2461 .ok_or_else(|| invalid("frequency column index out of range"))?;
2462 let Some(summary) = self.table.frequencies.get(column).and_then(Option::as_ref) else {
2463 return Ok(None);
2464 };
2465 if summary.ordinals.is_empty() {
2466 return Ok(None);
2467 }
2468 Ok(Some(FrequencyOccurrences {
2469 omitted_max: summary.omitted_max,
2470 ordinals: summary.ordinals.clone(),
2471 }))
2472 }
2473
2474 pub fn distinct_values(&self, column: usize) -> Result<Option<u64>> {
2498 self.table
2499 .distincts
2500 .get(column)
2501 .copied()
2502 .ok_or_else(|| invalid("distinct column index out of range"))
2503 }
2504
2505 pub fn null_count(&self, column: usize) -> Result<u64> {
2516 if column >= self.table.fields.len() {
2517 return Err(invalid("null count column index out of range"));
2518 }
2519 let mut nulls = 0_u64;
2520 for stripe in &self.table.stripes {
2521 let range = stripe
2522 .zone
2523 .column(column)
2524 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
2525 nulls = nulls
2526 .checked_add(range.nulls as u64)
2527 .ok_or_else(|| invalid("null count overflow"))?;
2528 }
2529 Ok(nulls)
2530 }
2531
2532 pub fn text_extremes(&self, column: usize) -> Result<Option<(Value, Value)>> {
2547 if self.null_count(column)? > 0 {
2548 return Ok(None);
2549 }
2550 let Some(dictionary) = self.dictionary(column)? else { return Ok(None) };
2551 let Some(ranks) = dictionary.ranks() else { return Ok(None) };
2552 if ranks == 0 {
2553 return Ok(None);
2554 }
2555 let low = text_at_rank(&dictionary, 0)?;
2556 let high = text_at_rank(&dictionary, ranks - 1)?;
2557 Ok(Some((low, high)))
2558 }
2559
2560 pub fn exact_extremes(&self, column: usize) -> Result<Option<(Bound, Bound)>> {
2583 if column >= self.table.fields.len() {
2584 return Err(invalid("extremes column index out of range"));
2585 }
2586 let mut low: Option<Bound> = None;
2587 let mut high: Option<Bound> = None;
2588 for stripe in &self.table.stripes {
2589 let range = stripe
2590 .zone
2591 .column(column)
2592 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
2593 if !range.exact {
2594 return Ok(None);
2595 }
2596 let (Some(small), Some(large)) = (range.low.as_ref(), range.high.as_ref()) else {
2601 if stripe.rows > range.nulls {
2602 return Ok(None);
2603 }
2604 continue;
2605 };
2606 low = Some(low.map_or_else(|| small.clone(), |held| held.smaller(small.clone())));
2607 high = Some(high.map_or_else(|| large.clone(), |held| held.larger(large.clone())));
2608 }
2609 Ok(low.zip(high))
2610 }
2611
2612 pub fn exact_sum(&self, column: usize) -> Result<Option<(i128, u64)>> {
2625 if column >= self.table.fields.len() {
2626 return Err(invalid("sum column index out of range"));
2627 }
2628 let mut total = 0_i128;
2629 let mut rows = 0_u64;
2630 for stripe in &self.table.stripes {
2631 let range = stripe
2632 .zone
2633 .column(column)
2634 .ok_or_else(|| invalid("stripe zone is narrower than the schema"))?;
2635 let Some(part) = range.sum else { return Ok(None) };
2636 let Some(sum) = total.checked_add(part) else { return Ok(None) };
2637 total = sum;
2638 rows = rows.saturating_add(stripe.rows as u64 - range.nulls as u64);
2639 }
2640 Ok(Some((total, rows)))
2641 }
2642
2643 fn dictionary(&self, column: usize) -> Result<Option<Arc<Vector>>> {
2652 let Some(page) = self.table.dictionaries[column] else { return Ok(None) };
2653 if let Some(dictionary) = self.dictionaries[column].get() {
2654 return Ok(Some(Arc::clone(dictionary)));
2655 }
2656 let _queued = self.loading[column].lock().map_err(|_| invalid("a poisoned dictionary"))?;
2657 if let Some(dictionary) = self.dictionaries[column].get() {
2658 return Ok(Some(Arc::clone(dictionary)));
2659 }
2660 self.opened.fetch_add(1, Atomic::Relaxed);
2661 let dictionary = Arc::new(open_global_dictionary(
2662 Arc::clone(&self.file),
2663 page,
2664 &self.table.fields[column].ty,
2665 TEXT_KEEP_BUDGET,
2666 )?);
2667 let _ = self.dictionaries[column].set(Arc::clone(&dictionary));
2668 Ok(Some(dictionary))
2669 }
2670
2671 pub fn read(&self, part: usize, columns: &[usize]) -> Result<Chunk> {
2680 self.read_impl(part, columns, true)
2681 }
2682
2683 pub fn read_sparse(&self, part: usize, columns: &[usize]) -> Result<Chunk> {
2693 self.read_impl(part, columns, false)
2694 }
2695
2696 pub fn skips_codes(&self, part: usize, column: usize, candidates: &[u32]) -> Result<bool> {
2703 if candidates.is_empty() {
2704 return Ok(true);
2705 }
2706 if candidates.windows(2).any(|pair| pair[0] >= pair[1]) {
2707 return Err(Error::internal("native code candidates are not sorted and unique"));
2708 }
2709 let stripe = self.stripe_of(part)?;
2710 let Some(page) = stripe.memberships.get(column).copied().flatten() else {
2711 return Ok(false);
2712 };
2713 let mut bytes = vec![0; page.length as usize];
2714 read_at(&self.file, page.offset, &mut bytes)?;
2715 if checksum(&bytes) != page.hash {
2716 return Err(invalid("membership page checksum differs"));
2717 }
2718 let codes = decode_membership(&bytes)?;
2719 let mut left = 0;
2720 let mut right = 0;
2721 while left < codes.len() && right < candidates.len() {
2722 match codes[left].cmp(&candidates[right]) {
2723 Ordering::Less => left += 1,
2724 Ordering::Greater => right += 1,
2725 Ordering::Equal => return Ok(false),
2726 }
2727 }
2728 Ok(true)
2729 }
2730
2731 fn stripe_of(&self, part: usize) -> Result<&Stripe> {
2732 let place = self.places.get(part).ok_or_else(|| invalid("part index out of range"))?;
2733 self.table
2734 .stripes
2735 .get(place.stripe as usize)
2736 .ok_or_else(|| invalid("stripe index out of range"))
2737 }
2738
2739 fn held(&self, at: usize, stripe: &Stripe, column: usize, whole: bool) -> Result<CachedColumn> {
2756 let cache = self.cache.get(column).ok_or_else(|| invalid("column index out of range"))?;
2757 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
2758 let known = cached.index.get(at).and_then(Clone::clone);
2759 let page = cached.pages.get(at).and_then(Clone::clone);
2760 if let Some(index) = known.clone() {
2761 if !whole || page.is_some() {
2762 return Ok(CachedColumn { stripe: at, index, page });
2763 }
2764 }
2765 if cached.loading.contains(&at) {
2766 drop(cached);
2767 if let Some(index) = known {
2771 return Ok(CachedColumn { stripe: at, index, page: None });
2772 }
2773 let held = self.page_of(stripe, column, at, false, None)?;
2774 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
2775 remember(&mut cached, &held, self.kept.load(Atomic::Relaxed));
2776 return Ok(held);
2777 }
2778 cached.loading.push(at);
2779 drop(cached);
2780
2781 let read = self.page_of(stripe, column, at, whole, known);
2782
2783 let mut cached = cache.lock().map_err(|_| invalid("column page cache is poisoned"))?;
2787 if let Some(position) = cached.loading.iter().position(|loading| *loading == at) {
2788 cached.loading.remove(position);
2789 }
2790 let held = read?;
2791 remember(&mut cached, &held, self.kept.load(Atomic::Relaxed));
2792 Ok(held)
2793 }
2794
2795 fn page_of(
2801 &self,
2802 stripe: &Stripe,
2803 column: usize,
2804 at: usize,
2805 whole: bool,
2806 known: Option<Arc<Vec<PartSpan>>>,
2807 ) -> Result<CachedColumn> {
2808 let index = match known {
2809 Some(index) => index,
2810 None => {
2811 self.indexes.fetch_add(1, Atomic::Relaxed);
2812 Arc::new(read_index(&self.file, stripe, column)?)
2813 }
2814 };
2815 let page = if whole {
2816 self.pages.fetch_add(1, Atomic::Relaxed);
2817 let span = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
2818 let mut bytes = vec![0; span.length as usize];
2819 read_at(&self.file, span.offset, &mut bytes)?;
2820 Some(Arc::new(bytes))
2821 } else {
2822 None
2823 };
2824 Ok(CachedColumn { stripe: at, index, page })
2825 }
2826
2827 fn read_impl(&self, at: usize, columns: &[usize], whole: bool) -> Result<Chunk> {
2828 let place = *self.places.get(at).ok_or_else(|| invalid("part index out of range"))?;
2829 let index = place.stripe as usize;
2830 let stripe =
2831 self.table.stripes.get(index).ok_or_else(|| invalid("stripe index out of range"))?;
2832 let rows = place.rows as usize;
2833 let mut picked = Vec::with_capacity(columns.len());
2834 for &column in columns {
2835 let field = self
2836 .table
2837 .fields
2838 .get(column)
2839 .ok_or_else(|| invalid("column index out of range"))?;
2840 let page = stripe.pages.get(column).ok_or_else(|| invalid("stripe page is missing"))?;
2841 let held = self.held(index, stripe, column, whole)?;
2842 let span = *held
2843 .index
2844 .get(place.part as usize)
2845 .ok_or_else(|| invalid("part index out of range"))?;
2846 let owned;
2847 let bytes = match &held.page {
2848 Some(held) => part_bytes(held, span)?,
2849 None => {
2850 let offset = page
2851 .offset
2852 .checked_add(span.start as u64)
2853 .ok_or_else(|| invalid("part range overflow"))?;
2854 let mut bytes = vec![0; span.length];
2855 read_at(&self.file, offset, &mut bytes)?;
2856 owned = bytes;
2857 &owned
2858 }
2859 };
2860 if checksum(bytes) != span.hash {
2861 return Err(invalid(&format!(
2862 "column page checksum differs, column {column} part {} at {}+{} of {} bytes, \
2863 wanted {:016x} and got {:016x}",
2864 place.part,
2865 page.offset,
2866 span.start,
2867 span.length,
2868 span.hash,
2869 checksum(bytes),
2870 )));
2871 }
2872 let dictionary = self.dictionary(column)?;
2873 picked.push(decode(&field.ty, rows, bytes, dictionary)?);
2874 }
2875 Chunk::with_rows(picked, rows)
2876 }
2877
2878 #[must_use]
2894 pub fn skips(&self, part: usize, probes: &[Probe]) -> bool {
2895 let Some(place) = self.places.get(part).copied() else { return false };
2896 let Some(stripe) = self.table.stripes.get(place.stripe as usize) else { return false };
2897 if stripe.zone.skips(probes) {
2898 return true;
2899 }
2900 probes.iter().any(|probe| self.outside(place, probe) || self.sifted(place, probe))
2901 }
2902
2903 fn outside(&self, place: Place, probe: &Probe) -> bool {
2909 match self.stripe_part_ranges(place.stripe as usize, probe.column) {
2910 Some(ranges) => ranges
2911 .get(place.part as usize)
2912 .is_some_and(|range| range.excludes(probe.op, &probe.value)),
2913 None => false,
2914 }
2915 }
2916
2917 fn stripe_part_ranges(&self, stripe: usize, column: usize) -> Option<&[Range]> {
2923 let slot = self.part_ranges.get(column)?.get(stripe)?;
2924 if let Some(held) = slot.get() {
2925 return Some(held);
2926 }
2927 let page = self.table.stripes.get(stripe)?.part_ranges.get(column).copied().flatten()?;
2928 let mut bytes = vec![0; page.length as usize];
2929 read_at(&self.file, page.offset, &mut bytes).ok()?;
2930 if checksum(&bytes) != page.hash {
2931 return None;
2932 }
2933 let ranges = Arc::new(decode_part_ranges(&bytes).ok()?);
2934 let _ = slot.set(ranges);
2935 slot.get().map(|held| held.as_slice())
2936 }
2937
2938 #[must_use]
2955 pub fn certain(&self, part: usize, probes: &[Probe]) -> bool {
2956 let Some(place) = self.places.get(part).copied() else { return false };
2957 let Some(stripe) = self.table.stripes.get(place.stripe as usize) else { return false };
2958 if stripe.zone.certain(probes) {
2959 return true;
2960 }
2961 probes
2962 .iter()
2963 .all(|probe| stripe.zone.certain(slice::from_ref(probe)) || self.inside(place, probe))
2964 }
2965
2966 fn inside(&self, place: Place, probe: &Probe) -> bool {
2972 match self.stripe_part_ranges(place.stripe as usize, probe.column) {
2973 Some(ranges) => ranges
2974 .get(place.part as usize)
2975 .is_some_and(|range| range.certain(probe.op, &probe.value)),
2976 None => false,
2977 }
2978 }
2979
2980 #[must_use]
2991 pub fn stripe_skips(&self, stripe: usize, probes: &[Probe]) -> bool {
2992 self.table.stripes.get(stripe).is_some_and(|held| held.zone.skips(probes))
2993 }
2994
2995 fn sifted(&self, place: Place, probe: &Probe) -> bool {
3001 if probe.op != Op::Equal {
3002 return false;
3003 }
3004 match self.stripe_sieves(place.stripe as usize, probe.column) {
3005 Some(sieves) => sieves
3006 .get(place.part as usize)
3007 .and_then(Option::as_ref)
3008 .is_some_and(|sieve| sieve.excludes(&probe.value)),
3009 None => false,
3010 }
3011 }
3012
3013 fn stripe_sieves(&self, stripe: usize, column: usize) -> Option<&[Option<Sieve>]> {
3020 let slot = self.sieves.get(column)?.get(stripe)?;
3021 if let Some(held) = slot.get() {
3022 return Some(held);
3023 }
3024 let page = self.table.stripes.get(stripe)?.sieves.get(column).copied().flatten()?;
3025 let mut bytes = vec![0; page.length as usize];
3026 read_at(&self.file, page.offset, &mut bytes).ok()?;
3027 if checksum(&bytes) != page.hash {
3028 return None;
3029 }
3030 let sieves = Arc::new(decode_sieves(&bytes).ok()?);
3031 let _ = slot.set(sieves);
3032 slot.get().map(|held| held.as_slice())
3033 }
3034}
3035
3036fn text_at_rank(dictionary: &Vector, rank: usize) -> Result<Value> {
3038 let code = dictionary.code_at_rank(rank)? as usize;
3039 let text = dictionary
3040 .try_text_at(code)?
3041 .ok_or_else(|| invalid("global dictionary order names a code it does not have"))?;
3042 Ok(Value::Varchar(text.into()))
3043}
3044
3045#[cfg(unix)]
3050fn write_at(file: &File, mut offset: u64, mut bytes: &[u8]) -> Result<()> {
3051 use std::os::unix::fs::FileExt;
3052 while !bytes.is_empty() {
3053 let written = file.write_at(bytes, offset).map_err(io)?;
3054 if written == 0 {
3055 return Err(invalid("a write to the native file wrote nothing"));
3056 }
3057 offset += written as u64;
3058 bytes = &bytes[written..];
3059 }
3060 Ok(())
3061}
3062
3063#[cfg(windows)]
3065fn write_at(file: &File, mut offset: u64, mut bytes: &[u8]) -> Result<()> {
3066 use std::os::windows::fs::FileExt;
3067 while !bytes.is_empty() {
3068 let written = file.seek_write(bytes, offset).map_err(io)?;
3069 if written == 0 {
3070 return Err(invalid("a write to the native file wrote nothing"));
3071 }
3072 offset += written as u64;
3073 bytes = &bytes[written..];
3074 }
3075 Ok(())
3076}
3077
3078#[cfg(not(any(unix, windows)))]
3080fn write_at(file: &File, offset: u64, bytes: &[u8]) -> Result<()> {
3081 use std::io::Write;
3082 let mut file = file.try_clone().map_err(io)?;
3083 file.seek(SeekFrom::Start(offset)).map_err(io)?;
3084 file.write_all(bytes).map_err(io)
3085}
3086
3087#[cfg(unix)]
3097fn read_at(file: &File, mut offset: u64, mut bytes: &mut [u8]) -> Result<()> {
3098 use std::os::unix::fs::FileExt;
3099 while !bytes.is_empty() {
3100 let read = file.read_at(bytes, offset).map_err(io)?;
3101 if read == 0 {
3102 return Err(invalid("column page ends before its declared length"));
3103 }
3104 offset += read as u64;
3105 bytes = &mut bytes[read..];
3106 }
3107 Ok(())
3108}
3109
3110#[cfg(windows)]
3116fn read_at(file: &File, mut offset: u64, mut bytes: &mut [u8]) -> Result<()> {
3117 use std::os::windows::fs::FileExt;
3118 while !bytes.is_empty() {
3119 let read = file.seek_read(bytes, offset).map_err(io)?;
3120 if read == 0 {
3121 return Err(invalid("column page ends before its declared length"));
3122 }
3123 offset += read as u64;
3124 bytes = &mut bytes[read..];
3125 }
3126 Ok(())
3127}
3128
3129#[cfg(not(any(unix, windows)))]
3134fn read_at(file: &File, offset: u64, bytes: &mut [u8]) -> Result<()> {
3135 let mut file = file.try_clone().map_err(io)?;
3136 file.seek(SeekFrom::Start(offset)).map_err(io)?;
3137 file.read_exact(bytes).map_err(io)
3138}
3139
3140fn type_tag(ty: &LogicalType) -> Result<u8> {
3141 match ty {
3142 LogicalType::SmallInt => Ok(1),
3143 LogicalType::Integer => Ok(2),
3144 LogicalType::BigInt => Ok(3),
3145 LogicalType::Varchar => Ok(4),
3146 LogicalType::Date => Ok(5),
3147 LogicalType::Timestamp => Ok(6),
3148 LogicalType::Boolean => Ok(7),
3149 LogicalType::TinyInt => Ok(8),
3150 LogicalType::UTinyInt => Ok(9),
3151 LogicalType::USmallInt => Ok(10),
3152 LogicalType::UInteger => Ok(11),
3153 LogicalType::UBigInt => Ok(12),
3154 _ => Err(Error::not_implemented(format!("native storage for {ty}"))),
3155 }
3156}
3157
3158fn tag_type(tag: u8) -> Result<LogicalType> {
3159 match tag {
3160 1 => Ok(LogicalType::SmallInt),
3161 2 => Ok(LogicalType::Integer),
3162 3 => Ok(LogicalType::BigInt),
3163 4 => Ok(LogicalType::Varchar),
3164 5 => Ok(LogicalType::Date),
3165 6 => Ok(LogicalType::Timestamp),
3166 7 => Ok(LogicalType::Boolean),
3167 8 => Ok(LogicalType::TinyInt),
3168 9 => Ok(LogicalType::UTinyInt),
3169 10 => Ok(LogicalType::USmallInt),
3170 11 => Ok(LogicalType::UInteger),
3171 12 => Ok(LogicalType::UBigInt),
3172 _ => Err(invalid("column type tag is unknown")),
3173 }
3174}
3175
3176fn put_u16(out: &mut Vec<u8>, value: u16) {
3177 out.extend_from_slice(&value.to_le_bytes());
3178}
3179fn put_u32(out: &mut Vec<u8>, value: u32) {
3180 out.extend_from_slice(&value.to_le_bytes());
3181}
3182fn put_u64(out: &mut Vec<u8>, value: u64) {
3183 out.extend_from_slice(&value.to_le_bytes());
3184}
3185fn put_var_u64(out: &mut Vec<u8>, mut value: u64) {
3186 while value >= 0x80 {
3187 out.push((value as u8 & 0x7f) | 0x80);
3188 value >>= 7;
3189 }
3190 out.push(value as u8);
3191}
3192
3193fn frequency_order(left: FrequencyValue, right: FrequencyValue) -> Ordering {
3194 match (left, right) {
3195 (FrequencyValue::Null, FrequencyValue::Null) => Ordering::Equal,
3196 (FrequencyValue::Null, _) => Ordering::Less,
3197 (_, FrequencyValue::Null) => Ordering::Greater,
3198 (FrequencyValue::Integer(left), FrequencyValue::Integer(right)) => left.cmp(&right),
3199 (FrequencyValue::Code(left), FrequencyValue::Code(right)) => left.cmp(&right),
3200 (FrequencyValue::Integer(_), FrequencyValue::Code(_)) => Ordering::Less,
3201 (FrequencyValue::Code(_), FrequencyValue::Integer(_)) => Ordering::Greater,
3202 }
3203}
3204
3205fn code_frequency(dictionary: &GlobalDictionary) -> FrequencySummary {
3206 let mut entries = dictionary
3207 .counts
3208 .iter()
3209 .enumerate()
3210 .filter(|(_, count)| **count != 0)
3211 .map(|(code, &count)| FrequencyEntry { value: FrequencyValue::Code(code as u32), count })
3212 .collect::<Vec<_>>();
3213 if dictionary.nulls != 0 {
3214 entries.push(FrequencyEntry { value: FrequencyValue::Null, count: dictionary.nulls });
3215 }
3216 entries.sort_unstable_by(|left, right| {
3217 right.count.cmp(&left.count).then_with(|| frequency_order(left.value, right.value))
3218 });
3219 let omitted_max = entries.get(FREQUENCY_ENTRIES).map_or(0, |entry| entry.count);
3220 entries.truncate(FREQUENCY_ENTRIES);
3221 FrequencySummary { entries, omitted_max, ordinals: Vec::new() }
3222}
3223
3224fn encode_directory(table: &Table) -> Result<Vec<u8>> {
3225 let mut out = DIRECTORY.to_vec();
3226 let name = table.name.as_bytes();
3227 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("table name too long"))?);
3228 out.extend_from_slice(name);
3229 put_u16(&mut out, u16::try_from(table.fields.len()).map_err(|_| invalid("too many columns"))?);
3230 for field in &table.fields {
3231 let name = field.name.as_bytes();
3232 put_u16(&mut out, u16::try_from(name.len()).map_err(|_| invalid("column name too long"))?);
3233 out.extend_from_slice(name);
3234 out.push(type_tag(&field.ty)?);
3235 out.push(u8::from(field.not_null));
3236 }
3237 for dictionary in &table.dictionaries {
3238 match dictionary {
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 distinct in &table.distincts {
3249 match distinct {
3250 None => out.push(0),
3251 Some(count) => {
3252 out.push(1);
3253 put_u64(&mut out, *count);
3254 }
3255 }
3256 }
3257 put_u64(&mut out, u64::try_from(table.rows).map_err(|_| invalid("row count overflow"))?);
3258 put_u32(&mut out, u32::try_from(table.stripes.len()).map_err(|_| invalid("too many stripes"))?);
3259 for stripe in &table.stripes {
3260 put_u32(
3261 &mut out,
3262 u32::try_from(stripe.parts.len()).map_err(|_| invalid("too many parts in a stripe"))?,
3263 );
3264 for &rows in &stripe.parts {
3265 put_u32(&mut out, rows);
3266 }
3267 put_u64(&mut out, stripe.index.offset);
3268 put_u32(&mut out, stripe.index.length);
3269 for page in &stripe.pages {
3270 put_u64(&mut out, page.offset);
3271 put_u32(&mut out, page.length);
3272 }
3273 for (field, membership) in table.fields.iter().zip(&stripe.memberships) {
3274 if field.ty != LogicalType::Varchar {
3275 continue;
3276 }
3277 let page =
3278 membership.ok_or_else(|| invalid("string page has no code membership index"))?;
3279 put_u64(&mut out, page.offset);
3280 put_u32(&mut out, page.length);
3281 put_u64(&mut out, page.hash);
3282 }
3283 for sieve in &stripe.sieves {
3284 match sieve {
3285 None => out.push(0),
3286 Some(page) => {
3287 out.push(1);
3288 put_u64(&mut out, page.offset);
3289 put_u32(&mut out, page.length);
3290 put_u64(&mut out, page.hash);
3291 }
3292 }
3293 }
3294 for held in &stripe.part_ranges {
3295 match held {
3296 None => out.push(0),
3297 Some(page) => {
3298 out.push(1);
3299 put_u64(&mut out, page.offset);
3300 put_u32(&mut out, page.length);
3301 put_u64(&mut out, page.hash);
3302 }
3303 }
3304 }
3305 for range in stripe.zone.columns() {
3306 put_bound(&mut out, range.low.as_ref())?;
3307 put_bound(&mut out, range.high.as_ref())?;
3308 put_u32(
3309 &mut out,
3310 u32::try_from(range.nulls).map_err(|_| invalid("null count overflow"))?,
3311 );
3312 out.push(u8::from(range.exact));
3313 match range.sum {
3314 None => out.push(0),
3315 Some(total) => {
3316 out.push(1);
3317 out.extend_from_slice(&total.to_le_bytes());
3318 }
3319 }
3320 }
3321 }
3322 out.extend_from_slice(FREQUENCIES);
3323 put_u16(
3324 &mut out,
3325 u16::try_from(table.frequencies.len())
3326 .map_err(|_| invalid("too many frequency columns"))?,
3327 );
3328 for summary in &table.frequencies {
3329 let Some(summary) = summary else {
3330 out.push(0);
3331 continue;
3332 };
3333 out.push(1);
3334 put_u64(&mut out, summary.omitted_max);
3335 put_u32(
3336 &mut out,
3337 u32::try_from(summary.entries.len())
3338 .map_err(|_| invalid("too many frequency entries"))?,
3339 );
3340 for entry in &summary.entries {
3341 match entry.value {
3342 FrequencyValue::Null => out.push(0),
3343 FrequencyValue::Integer(value) => {
3344 out.push(1);
3345 out.extend_from_slice(&value.to_le_bytes());
3346 }
3347 FrequencyValue::Code(value) => {
3348 out.push(2);
3349 put_u32(&mut out, value);
3350 }
3351 }
3352 put_u64(&mut out, entry.count);
3353 }
3354 put_u32(
3355 &mut out,
3356 u32::try_from(summary.ordinals.len())
3357 .map_err(|_| invalid("too many frequency ordinals"))?,
3358 );
3359 let mut previous = 0_u64;
3360 for (at, &ordinal) in summary.ordinals.iter().enumerate() {
3361 let delta = if at == 0 {
3362 ordinal
3363 } else {
3364 ordinal
3365 .checked_sub(previous)
3366 .ok_or_else(|| invalid("frequency ordinals are not ordered"))?
3367 };
3368 if at != 0 && delta == 0 {
3369 return Err(invalid("frequency ordinals are not unique"));
3370 }
3371 put_var_u64(&mut out, delta);
3372 previous = ordinal;
3373 }
3374 }
3375 Ok(out)
3376}
3377
3378struct Cursor<'a> {
3379 bytes: &'a [u8],
3380 at: usize,
3381}
3382impl<'a> Cursor<'a> {
3383 fn take(&mut self, len: usize) -> Result<&'a [u8]> {
3384 let end = self.at.checked_add(len).ok_or_else(|| invalid("directory offset overflow"))?;
3385 let bytes =
3386 self.bytes.get(self.at..end).ok_or_else(|| invalid("directory is truncated"))?;
3387 self.at = end;
3388 Ok(bytes)
3389 }
3390 fn u8(&mut self) -> Result<u8> {
3391 Ok(self.take(1)?[0])
3392 }
3393 fn u16(&mut self) -> Result<u16> {
3394 Ok(u16::from_le_bytes(self.take(2)?.try_into().expect("two bytes")))
3395 }
3396 fn u32(&mut self) -> Result<u32> {
3397 Ok(u32::from_le_bytes(self.take(4)?.try_into().expect("four bytes")))
3398 }
3399 fn u64(&mut self) -> Result<u64> {
3400 Ok(u64::from_le_bytes(self.take(8)?.try_into().expect("eight bytes")))
3401 }
3402 fn var_u64(&mut self) -> Result<u64> {
3403 let mut value = 0_u64;
3404 for shift in (0..=63).step_by(7) {
3405 let byte = self.u8()?;
3406 let part = u64::from(byte & 0x7f);
3407 if shift == 63 && part > 1 {
3408 return Err(invalid("frequency ordinal varint overflows"));
3409 }
3410 value |= part << shift;
3411 if byte & 0x80 == 0 {
3412 return Ok(value);
3413 }
3414 }
3415 Err(invalid("frequency ordinal varint is too long"))
3416 }
3417 fn bound(&mut self) -> Result<Option<Bound>> {
3418 Ok(match self.u8()? {
3419 0 => None,
3420 1 => Some(Bound::Int(i128::from_le_bytes(
3421 self.take(16)?.try_into().expect("sixteen bytes"),
3422 ))),
3423 2 => Some(Bound::Real(f64::from_le_bytes(
3424 self.take(8)?.try_into().expect("eight bytes"),
3425 ))),
3426 3 => {
3427 let length = self.u32()? as usize;
3428 Some(Bound::Bytes(self.take(length)?.to_vec()))
3429 }
3430 4 => {
3431 let unscaled =
3432 i128::from_le_bytes(self.take(16)?.try_into().expect("sixteen bytes"));
3433 Some(Bound::Scaled { unscaled, scale: self.u8()? })
3434 }
3435 _ => return Err(invalid("bound tag differs")),
3436 })
3437 }
3438 fn text(&mut self) -> Result<String> {
3439 let len = self.u16()? as usize;
3440 String::from_utf8(self.take(len)?.to_vec()).map_err(|_| invalid("name is not UTF-8"))
3441 }
3442}
3443
3444fn decode_directory(bytes: &[u8], size: u64) -> Result<Table> {
3445 let mut cur = Cursor { bytes, at: 0 };
3446 if cur.take(8)? != DIRECTORY {
3447 return Err(invalid("directory magic differs"));
3448 }
3449 let name = cur.text()?;
3450 let width = cur.u16()? as usize;
3451 let mut fields = Vec::with_capacity(width);
3452 for _ in 0..width {
3453 let name = cur.text()?;
3454 let ty = tag_type(cur.u8()?)?;
3455 let not_null = match cur.u8()? {
3456 0 => false,
3457 1 => true,
3458 _ => return Err(invalid("nullability flag differs")),
3459 };
3460 fields.push(Field { name, ty, not_null });
3461 }
3462 let mut dictionaries = Vec::with_capacity(width);
3463 for _ in 0..width {
3464 dictionaries.push(match cur.u8()? {
3465 0 => None,
3466 1 => {
3467 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3468 let end = page
3469 .offset
3470 .checked_add(u64::from(page.length))
3471 .ok_or_else(|| invalid("dictionary page offset overflow"))?;
3472 if page.offset < HEADER || end > size {
3477 return Err(invalid("dictionary page range is outside the file"));
3478 }
3479 Some(page)
3480 }
3481 _ => return Err(invalid("dictionary page tag differs")),
3482 });
3483 }
3484 let mut distincts = Vec::with_capacity(width);
3485 for _ in 0..width {
3486 distincts.push(match cur.u8()? {
3487 0 => None,
3488 1 => Some(cur.u64()?),
3489 _ => return Err(invalid("distinct count tag differs")),
3490 });
3491 }
3492 let rows = usize::try_from(cur.u64()?).map_err(|_| invalid("row count does not fit"))?;
3493 let count = cur.u32()? as usize;
3494 let mut stripes = Vec::with_capacity(count);
3495 let mut total = 0_usize;
3496 for _ in 0..count {
3497 let count = cur.u32()? as usize;
3498 if count == 0 || count > STRIPE_PARTS {
3499 return Err(invalid("stripe part count is outside its bound"));
3500 }
3501 let mut parts = Vec::with_capacity(count);
3502 let mut stripe_rows = 0_usize;
3503 for _ in 0..count {
3504 let rows = cur.u32()?;
3505 if rows == 0 {
3506 return Err(invalid("empty part"));
3507 }
3508 parts.push(rows);
3509 stripe_rows = stripe_rows
3510 .checked_add(rows as usize)
3511 .ok_or_else(|| invalid("stripe row count overflow"))?;
3512 }
3513 total =
3514 total.checked_add(stripe_rows).ok_or_else(|| invalid("stripe row count overflow"))?;
3515 let index = Span { offset: cur.u64()?, length: cur.u32()? };
3516 let section = index_section(count)?;
3517 let wanted = section
3518 .checked_mul(width)
3519 .and_then(|bytes| u32::try_from(bytes).ok())
3520 .ok_or_else(|| invalid("index page length overflow"))?;
3521 let end = index
3522 .offset
3523 .checked_add(u64::from(index.length))
3524 .ok_or_else(|| invalid("index page offset overflow"))?;
3525 if index.offset < HEADER || end > size || index.length != wanted {
3526 return Err(invalid("index page range is outside the file"));
3527 }
3528 let mut pages = Vec::with_capacity(width);
3529 for _ in 0..width {
3530 let offset = cur.u64()?;
3531 let length = cur.u32()?;
3532 let end = offset
3533 .checked_add(u64::from(length))
3534 .ok_or_else(|| invalid("page offset overflow"))?;
3535 if offset < HEADER || end > size || length as usize > MAX_PAGE {
3536 return Err(invalid("page range is outside the file"));
3537 }
3538 pages.push(Span { offset, length });
3539 }
3540 let mut memberships = vec![None; width];
3541 for (column, field) in fields.iter().enumerate() {
3542 if field.ty != LogicalType::Varchar {
3543 continue;
3544 }
3545 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3546 let end = page
3547 .offset
3548 .checked_add(u64::from(page.length))
3549 .ok_or_else(|| invalid("membership page offset overflow"))?;
3550 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
3551 return Err(invalid("membership page range is outside the file"));
3552 }
3553 memberships[column] = Some(page);
3554 }
3555 let mut sieves = vec![None; width];
3556 for sieve in sieves.iter_mut().take(width) {
3557 match cur.u8()? {
3558 0 => continue,
3559 1 => {}
3560 _ => return Err(invalid("a sieve page has an unknown tag")),
3561 }
3562 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3563 let end = page
3564 .offset
3565 .checked_add(u64::from(page.length))
3566 .ok_or_else(|| invalid("sieve page offset overflow"))?;
3567 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
3568 return Err(invalid("sieve page range is outside the file"));
3569 }
3570 *sieve = Some(page);
3571 }
3572 let mut part_ranges = vec![None; width];
3573 for held in part_ranges.iter_mut().take(width) {
3574 match cur.u8()? {
3575 0 => continue,
3576 1 => {}
3577 _ => return Err(invalid("a part range page has an unknown tag")),
3578 }
3579 let page = Page { offset: cur.u64()?, length: cur.u32()?, hash: cur.u64()? };
3580 let end = page
3581 .offset
3582 .checked_add(u64::from(page.length))
3583 .ok_or_else(|| invalid("part range page offset overflow"))?;
3584 if page.offset < HEADER || end > size || page.length as usize > MAX_PAGE {
3585 return Err(invalid("part range page range is outside the file"));
3586 }
3587 *held = Some(page);
3588 }
3589 let mut ranges = Vec::with_capacity(width);
3590 for column in 0..width {
3591 let low = cur.bound()?;
3592 let high = cur.bound()?;
3593 let nulls = cur.u32()? as usize;
3594 if nulls > stripe_rows {
3595 return Err(invalid("null count exceeds stripe rows"));
3596 }
3597 let exact = cur.u8()? != 0;
3598 let sum = match cur.u8()? {
3599 0 => None,
3600 1 => Some(i128::from_le_bytes(
3601 cur.take(16)?.try_into().map_err(|_| invalid("a stripe sum is truncated"))?,
3602 )),
3603 _ => return Err(invalid("a stripe sum has an unknown tag")),
3604 };
3605 let ty = &fields.get(column).ok_or_else(|| invalid("a stripe range has no column"))?.ty;
3611 let low = low.map(|bound| scaled_as(bound, ty));
3612 let high = high.map(|bound| scaled_as(bound, ty));
3613 ranges.push(Range { low, high, nulls, exact, sum });
3614 }
3615 stripes.push(Stripe {
3616 rows: stripe_rows,
3617 parts,
3618 index,
3619 pages,
3620 memberships,
3621 sieves,
3622 part_ranges,
3623 zone: Zone::from_ranges(ranges),
3624 });
3625 }
3626 if total != rows {
3627 return Err(invalid("table row count differs from stripes"));
3628 }
3629 let frequencies = if cur.at == bytes.len() {
3630 vec![None; width]
3631 } else {
3632 if cur.take(8)? != FREQUENCIES {
3633 return Err(invalid("directory extension magic differs"));
3634 }
3635 if cur.u16()? as usize != width {
3636 return Err(invalid("frequency column count differs"));
3637 }
3638 let mut frequencies = Vec::with_capacity(width);
3639 for field in &fields {
3640 let summary = match cur.u8()? {
3641 0 => None,
3642 1 => {
3643 let omitted_max = cur.u64()?;
3644 let count = cur.u32()? as usize;
3645 if count > FREQUENCY_ENTRIES {
3646 return Err(invalid("frequency entry count exceeds its bound"));
3647 }
3648 let mut entries = Vec::with_capacity(count);
3649 for _ in 0..count {
3651 let value = match cur.u8()? {
3652 0 => FrequencyValue::Null,
3653 1 => FrequencyValue::Integer(i128::from_le_bytes(
3654 cur.take(16)?.try_into().expect("sixteen bytes"),
3655 )),
3656 2 => FrequencyValue::Code(cur.u32()?),
3657 _ => return Err(invalid("frequency value tag differs")),
3658 };
3659 let valid = matches!(
3660 (&field.ty, value),
3661 (_, FrequencyValue::Null)
3662 | (LogicalType::Varchar, FrequencyValue::Code(_))
3663 | (
3664 LogicalType::TinyInt
3665 | LogicalType::SmallInt
3666 | LogicalType::Integer
3667 | LogicalType::BigInt
3668 | LogicalType::UTinyInt
3669 | LogicalType::USmallInt
3670 | LogicalType::UInteger
3671 | LogicalType::UBigInt
3672 | LogicalType::Date
3673 | LogicalType::Timestamp,
3674 FrequencyValue::Integer(_),
3675 )
3676 );
3677 if !valid {
3678 return Err(invalid("frequency value does not match its column"));
3679 }
3680 let count = cur.u64()?;
3681 if count == 0 || count > rows as u64 {
3682 return Err(invalid("frequency count is outside the table"));
3683 }
3684 entries.push(FrequencyEntry { value, count });
3685 }
3686 if entries.windows(2).any(|pair| pair[0].count < pair[1].count) {
3687 return Err(invalid("frequency entries are not descending"));
3688 }
3689 let ordinals = {
3690 let ordinal_count = cur.u32()? as usize;
3691 if ordinal_count > FREQUENCY_ORDINALS || ordinal_count > rows {
3692 return Err(invalid("frequency ordinal count exceeds its bound"));
3693 }
3694 let mut ordinals = Vec::with_capacity(ordinal_count);
3695 let mut previous = 0_u64;
3696 for at in 0..ordinal_count {
3697 let delta = cur.var_u64()?;
3698 if at != 0 && delta == 0 {
3699 return Err(invalid("frequency ordinals are not increasing"));
3700 }
3701 let ordinal = if at == 0 {
3702 delta
3703 } else {
3704 previous
3705 .checked_add(delta)
3706 .ok_or_else(|| invalid("frequency ordinal overflows"))?
3707 };
3708 if ordinal >= rows as u64 {
3709 return Err(invalid("frequency ordinal is outside the table"));
3710 }
3711 ordinals.push(ordinal);
3712 previous = ordinal;
3713 }
3714 ordinals
3715 };
3716 Some(FrequencySummary { entries, omitted_max, ordinals })
3717 }
3718 _ => return Err(invalid("frequency summary tag differs")),
3719 };
3720 frequencies.push(summary);
3721 }
3722 frequencies
3723 };
3724 if cur.at != bytes.len() {
3725 return Err(invalid("directory has trailing bytes"));
3726 }
3727 Ok(Table { name, fields, stripes, rows, dictionaries, distincts, frequencies })
3728}
3729
3730fn put_bound(out: &mut Vec<u8>, bound: Option<&Bound>) -> Result<()> {
3731 match bound {
3732 None => out.push(0),
3733 Some(Bound::Int(value)) => {
3734 out.push(1);
3735 out.extend_from_slice(&value.to_le_bytes());
3736 }
3737 Some(Bound::Real(value)) => {
3738 out.push(2);
3739 out.extend_from_slice(&value.to_le_bytes());
3740 }
3741 Some(Bound::Bytes(value)) => {
3742 out.push(3);
3743 put_u32(out, u32::try_from(value.len()).map_err(|_| invalid("bound length overflow"))?);
3744 out.extend_from_slice(value);
3745 }
3746 Some(Bound::Scaled { unscaled, scale }) => {
3747 out.push(4);
3748 out.extend_from_slice(&unscaled.to_le_bytes());
3749 out.push(*scale);
3750 }
3751 }
3752 Ok(())
3753}
3754
3755#[derive(Debug)]
3772struct Codes;
3773
3774impl chooser::Chooser for Codes {
3775 fn name(&self) -> &'static str {
3776 "codes"
3777 }
3778
3779 fn narrow_strings(
3780 &self,
3781 _values: &[&[u8]],
3782 offered: &[string::Kind],
3783 _depth: u8,
3784 ) -> Vec<string::Kind> {
3785 offered.to_vec()
3788 }
3789
3790 fn narrow_integers(
3791 &self,
3792 _values: &[i64],
3793 offered: &[integer::Kind],
3794 depth: u8,
3795 ) -> Vec<integer::Kind> {
3796 let keep: &[integer::Kind] = if depth == 0 {
3797 &[integer::Kind::Constant, integer::Kind::Packed, integer::Kind::Rle]
3798 } else {
3799 &[integer::Kind::Constant, integer::Kind::Packed]
3800 };
3801 let narrowed: Vec<integer::Kind> =
3802 offered.iter().copied().filter(|kind| keep.contains(kind)).collect();
3803 if narrowed.is_empty() { offered.to_vec() } else { narrowed }
3806 }
3807}
3808
3809#[derive(Debug)]
3821struct Fixed;
3822
3823impl chooser::Chooser for Fixed {
3824 fn name(&self) -> &'static str {
3825 "fixed"
3826 }
3827
3828 fn narrow_strings(
3829 &self,
3830 _values: &[&[u8]],
3831 offered: &[string::Kind],
3832 _depth: u8,
3833 ) -> Vec<string::Kind> {
3834 offered.to_vec()
3835 }
3836
3837 fn narrow_integers(
3838 &self,
3839 _values: &[i64],
3840 offered: &[integer::Kind],
3841 depth: u8,
3842 ) -> Vec<integer::Kind> {
3843 let keep: &[integer::Kind] = if depth == 0 {
3844 &[
3845 integer::Kind::Constant,
3846 integer::Kind::Packed,
3847 integer::Kind::Delta,
3848 integer::Kind::Rle,
3849 integer::Kind::Sparse,
3850 integer::Kind::Strided,
3851 ]
3852 } else {
3853 &[integer::Kind::Constant, integer::Kind::Packed, integer::Kind::Delta]
3854 };
3855 let narrowed: Vec<integer::Kind> =
3856 offered.iter().copied().filter(|kind| keep.contains(kind)).collect();
3857 if narrowed.is_empty() { offered.to_vec() } else { narrowed }
3858 }
3859}
3860
3861fn widened(data: &Data) -> Option<Vec<i64>> {
3868 match data {
3869 Data::Int8(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
3870 Data::UInt8(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
3871 Data::Int16(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
3872 Data::UInt16(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
3873 Data::Int32(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
3874 Data::UInt32(values) => Some(values.iter().map(|value| i64::from(*value)).collect()),
3875 Data::Int64(values) => Some(values.to_vec()),
3876 _ => None,
3877 }
3878}
3879
3880trait Narrow: Copy {
3887 const BIASED: (u32, u64);
3892
3893 fn narrow(value: i64) -> Self;
3895}
3896
3897#[allow(clippy::cast_sign_loss, reason = "a residue is a bit pattern and not a number")]
3914fn residue<T: Narrow>(value: i64) -> u64 {
3915 let (bits, bias) = T::BIASED;
3916 (value as u64).wrapping_add(bias) >> bits
3917}
3918
3919macro_rules! narrows {
3924 ($($ty:ty => $bias:expr),* $(,)?) => {$(
3925 impl Narrow for $ty {
3926 const BIASED: (u32, u64) = (<$ty>::BITS, $bias);
3927
3928 #[allow(
3929 clippy::cast_possible_truncation,
3930 clippy::cast_sign_loss,
3931 reason = "the caller has checked the bits this truncates away"
3932 )]
3933 fn narrow(value: i64) -> Self {
3934 value as Self
3935 }
3936 }
3937 )*};
3938}
3939
3940narrows! {
3941 i8 => 1 << 7,
3942 u8 => 0,
3943 i16 => 1 << 15,
3944 u16 => 0,
3945 i32 => 1 << 31,
3946 u32 => 0,
3947}
3948
3949fn fit<T: Narrow>(values: &[i64]) -> Result<Vec<T>> {
3962 let mut spilled = 0u64;
3963 for value in values {
3964 spilled |= residue::<T>(*value);
3965 }
3966 if spilled != 0 {
3967 return Err(invalid("page value is not of its type"));
3968 }
3969 Ok(values.iter().map(|value| T::narrow(*value)).collect())
3970}
3971
3972fn narrowed(ty: &LogicalType, values: Vec<i64>) -> Result<Data> {
3977 Ok(match ty {
3978 LogicalType::TinyInt => Data::Int8(fit::<i8>(&values)?.into()),
3979 LogicalType::UTinyInt => Data::UInt8(fit::<u8>(&values)?.into()),
3980 LogicalType::SmallInt => Data::Int16(fit::<i16>(&values)?.into()),
3981 LogicalType::USmallInt => Data::UInt16(fit::<u16>(&values)?.into()),
3982 LogicalType::Integer | LogicalType::Date => Data::Int32(fit::<i32>(&values)?.into()),
3983 LogicalType::UInteger => Data::UInt32(fit::<u32>(&values)?.into()),
3984 LogicalType::BigInt | LogicalType::Timestamp => Data::Int64(values.into()),
3985 _ => return Err(invalid("cascade codec belongs to a page that is not integers")),
3986 })
3987}
3988
3989fn plain_width(ty: &LogicalType) -> Option<usize> {
3992 Some(match ty {
3993 LogicalType::TinyInt | LogicalType::UTinyInt => 1,
3994 LogicalType::SmallInt | LogicalType::USmallInt => 2,
3995 LogicalType::Integer | LogicalType::UInteger | LogicalType::Date => 4,
3996 LogicalType::BigInt | LogicalType::Timestamp => 8,
3997 _ => return None,
3998 })
3999}
4000
4001fn cascaded(
4007 flat: &Vector,
4008 ty: &LogicalType,
4009 packed: Option<&Packed<'_>>,
4010) -> Result<Option<Vec<u8>>> {
4011 let (Some(width), Some(data)) = (plain_width(ty), flat.data()) else { return Ok(None) };
4012 let Some(values) = widened(data) else { return Ok(None) };
4013 let plain = values.len().saturating_mul(width);
4014 let best = match packed {
4015 Some(packed) => plain.min(21 + size_of_val(packed.words())),
4017 None => plain,
4018 };
4019 let out = integer::encode_with(&values, &Fixed)?;
4020 Ok((out.len() < best).then_some(out))
4021}
4022
4023fn encoded_codes(codes: &[u32]) -> Result<Option<Vec<u8>>> {
4035 let wide: Vec<i64> = codes.iter().map(|code| i64::from(*code)).collect();
4036 let coded = integer::encode_with(&wide, &Codes)?;
4037 let plain = codes.len().saturating_mul(size_of::<u32>());
4038 Ok((coded.len() < plain).then_some(coded))
4039}
4040
4041fn encode(
4042 vector: &Vector,
4043 global: Option<&mut GlobalDictionary>,
4044) -> Result<(Vec<u8>, Option<Vec<u32>>)> {
4045 let ty = vector.logical_type();
4046 let flat = vector.flatten()?;
4048 let mut out = Vec::new();
4049 let mut global_codes = None;
4050 if let Some(global) = global {
4051 let mut codes = Vec::with_capacity(flat.len());
4052 for row in 0..flat.len() {
4053 let text = flat.text_at(row).unwrap_or("");
4054 let code = global.code(text)?;
4055 global.observe(code, flat.is_null_at(row))?;
4056 codes.push(code);
4057 }
4058 global_codes = Some(codes);
4059 }
4060 let membership = global_codes.as_deref().map(unique_codes);
4061 let dictionary = if global_codes.is_none() && ty == &LogicalType::Varchar {
4062 string_dictionary(&flat)?
4063 } else {
4064 None
4065 };
4066 let packed_vector = if dictionary.is_none() && global_codes.is_none() {
4067 Some(flat.bit_packed()?)
4068 } else {
4069 None
4070 };
4071 let packed = packed_vector.as_ref().and_then(Vector::packed_parts);
4072 let coded = match global_codes.as_deref() {
4073 Some(codes) => encoded_codes(codes)?,
4074 None => None,
4075 };
4076 let cascade = if dictionary.is_none() && global_codes.is_none() {
4080 cascaded(&flat, ty, packed.as_ref())?
4081 } else {
4082 None
4083 };
4084 out.push(if coded.is_some() {
4085 4
4086 } else if cascade.is_some() {
4087 5
4088 } else if global_codes.is_some() {
4089 3
4090 } else if dictionary.is_some() {
4091 1
4092 } else if packed.is_some() {
4093 2
4094 } else {
4095 0
4096 });
4097 let nulls = flat.validity();
4098 let flag = match nulls {
4099 Validity::AllValid => 0,
4100 Validity::AllInvalid => 1,
4101 Validity::Mask(_) => 2,
4102 };
4103 out.push(flag);
4104 if flag == 2 {
4105 for group in (0..vector.len()).step_by(8) {
4106 let mut bits = 0_u8;
4107 for bit in 0..8 {
4108 if group + bit < vector.len() && !flat.is_null_at(group + bit) {
4109 bits |= 1 << bit;
4110 }
4111 }
4112 out.push(bits);
4113 }
4114 }
4115 if let Some(coded) = coded {
4116 out.extend_from_slice(&coded);
4117 return Ok((out, membership));
4118 }
4119 if let Some(cascade) = cascade {
4120 out.extend_from_slice(&cascade);
4121 return Ok((out, membership));
4122 }
4123 if let Some(codes) = global_codes {
4124 for code in codes {
4125 put_u32(&mut out, code);
4126 }
4127 return Ok((out, membership));
4128 }
4129 if let Some(dictionary) = dictionary {
4130 out.extend_from_slice(&dictionary);
4131 return Ok((out, membership));
4132 }
4133 if let Some(packed) = packed {
4134 if packed.offset() != 0 {
4135 return Err(invalid("writer received a sliced packed vector"));
4136 }
4137 out.push(u8::try_from(packed.width()).map_err(|_| invalid("packed width overflow"))?);
4138 out.extend_from_slice(&packed.base().to_le_bytes());
4139 put_u32(
4140 &mut out,
4141 u32::try_from(packed.words().len()).map_err(|_| invalid("too many packed words"))?,
4142 );
4143 for word in packed.words() {
4144 put_u64(&mut out, *word);
4145 }
4146 return Ok((out, membership));
4147 }
4148 let data = flat.data().ok_or_else(|| invalid("scalar column did not flatten"))?;
4149 match (ty, data) {
4150 (LogicalType::TinyInt, Data::Int8(values)) => {
4151 for value in &**values {
4152 out.extend_from_slice(&value.to_le_bytes());
4153 }
4154 }
4155 (LogicalType::UTinyInt, Data::UInt8(values)) => {
4156 for value in &**values {
4157 out.extend_from_slice(&value.to_le_bytes());
4158 }
4159 }
4160 (LogicalType::SmallInt, Data::Int16(values)) => {
4161 for value in &**values {
4162 out.extend_from_slice(&value.to_le_bytes());
4163 }
4164 }
4165 (LogicalType::USmallInt, Data::UInt16(values)) => {
4166 for value in &**values {
4167 out.extend_from_slice(&value.to_le_bytes());
4168 }
4169 }
4170 (LogicalType::UInteger, Data::UInt32(values)) => {
4171 for value in &**values {
4172 out.extend_from_slice(&value.to_le_bytes());
4173 }
4174 }
4175 (LogicalType::UBigInt, Data::UInt64(values)) => {
4176 for value in &**values {
4177 out.extend_from_slice(&value.to_le_bytes());
4178 }
4179 }
4180 (LogicalType::Integer | LogicalType::Date, Data::Int32(values)) => {
4181 for value in &**values {
4182 out.extend_from_slice(&value.to_le_bytes());
4183 }
4184 }
4185 (LogicalType::BigInt | LogicalType::Timestamp, Data::Int64(values)) => {
4186 for value in &**values {
4187 out.extend_from_slice(&value.to_le_bytes());
4188 }
4189 }
4190 (LogicalType::Boolean, Data::Bool(values)) => {
4191 for value in &**values {
4192 out.push(u8::from(*value));
4193 }
4194 }
4195 (LogicalType::Varchar, Data::Varlen(values)) => {
4196 let mut bytes = Vec::new();
4197 put_u32(&mut out, 0);
4198 for row in 0..vector.len() {
4199 let value = values.bytes(row).ok_or_else(|| invalid("string view is invalid"))?;
4200 bytes.extend_from_slice(value);
4201 put_u32(
4202 &mut out,
4203 u32::try_from(bytes.len())
4204 .map_err(|_| invalid("string payload exceeds 4GiB"))?,
4205 );
4206 }
4207 out.extend_from_slice(&bytes);
4208 }
4209 _ => return Err(Error::not_implemented(format!("native page for {ty}"))),
4210 }
4211 Ok((out, membership))
4212}
4213
4214fn put_varint(out: &mut Vec<u8>, mut value: u32) {
4215 while value >= 0x80 {
4216 out.push((value as u8 & 0x7f) | 0x80);
4217 value >>= 7;
4218 }
4219 out.push(value as u8);
4220}
4221
4222fn unique_codes(codes: &[u32]) -> Vec<u32> {
4224 let mut unique = codes.to_vec();
4225 unique.sort_unstable();
4226 unique.dedup();
4227 unique
4228}
4229
4230fn merged_codes(lists: Vec<Vec<u32>>) -> Vec<u32> {
4236 let mut lists = lists;
4237 while lists.len() > 1 {
4238 let mut next = Vec::with_capacity(lists.len().div_ceil(2));
4239 for pair in lists.chunks(2) {
4240 match pair {
4241 [left, right] => next.push(merged_pair(left, right)),
4242 [only] => next.push(only.clone()),
4243 _ => {}
4244 }
4245 }
4246 lists = next;
4247 }
4248 lists.pop().unwrap_or_default()
4249}
4250
4251fn merged_pair(left: &[u32], right: &[u32]) -> Vec<u32> {
4252 let mut out = Vec::with_capacity(left.len().saturating_add(right.len()));
4253 let mut at = 0;
4254 let mut to = 0;
4255 while at < left.len() && to < right.len() {
4256 match left[at].cmp(&right[to]) {
4257 Ordering::Less => {
4258 out.push(left[at]);
4259 at += 1;
4260 }
4261 Ordering::Greater => {
4262 out.push(right[to]);
4263 to += 1;
4264 }
4265 Ordering::Equal => {
4266 out.push(left[at]);
4267 at += 1;
4268 to += 1;
4269 }
4270 }
4271 }
4272 out.extend_from_slice(&left[at..]);
4273 out.extend_from_slice(&right[to..]);
4274 out
4275}
4276
4277fn merged_range(ranges: impl Iterator<Item = Range>) -> Range {
4282 let mut merged = Range::default();
4283 let mut first = true;
4284 for range in ranges {
4285 merged.nulls = merged.nulls.saturating_add(range.nulls);
4286 merged.sum = match (merged.sum.take(), range.sum) {
4290 (Some(held), Some(next)) if !first => held.checked_add(next),
4291 (_, next) if first => next,
4292 _ => None,
4293 };
4294 merged.exact = if first { range.exact } else { merged.exact && range.exact };
4295 if first {
4296 merged.low = range.low;
4297 merged.high = range.high;
4298 first = false;
4299 continue;
4300 }
4301 merged.low = match (merged.low.take(), range.low) {
4302 (Some(held), Some(next)) => Some(held.smaller(next)),
4303 _ => None,
4304 };
4305 merged.high = match (merged.high.take(), range.high) {
4306 (Some(held), Some(next)) => Some(held.larger(next)),
4307 _ => None,
4308 };
4309 }
4310 merged
4311}
4312
4313fn shortened(bound: Option<Bound>, high: bool) -> Option<Bound> {
4326 match bound {
4327 Some(Bound::Bytes(mut value)) if value.len() > PART_BOUND_BYTES => {
4328 value.truncate(PART_BOUND_BYTES);
4329 if !high {
4330 return Some(Bound::Bytes(value));
4331 }
4332 while let Some(last) = value.pop() {
4333 if last < u8::MAX {
4334 value.push(last + 1);
4335 return Some(Bound::Bytes(value));
4336 }
4337 }
4338 None
4339 }
4340 other => other,
4341 }
4342}
4343
4344fn encode_part_ranges(ranges: &[Range]) -> Result<Vec<u8>> {
4352 let mut out = Vec::new();
4353 put_u32(
4354 &mut out,
4355 u32::try_from(ranges.len()).map_err(|_| invalid("too many parts in a stripe"))?,
4356 );
4357 for range in ranges {
4358 put_bound(&mut out, shortened(range.low.clone(), false).as_ref())?;
4359 put_bound(&mut out, shortened(range.high.clone(), true).as_ref())?;
4360 put_u32(&mut out, u32::try_from(range.nulls).map_err(|_| invalid("null count overflow"))?);
4361 }
4362 Ok(out)
4363}
4364
4365fn decode_part_ranges(bytes: &[u8]) -> Result<Vec<Range>> {
4367 let mut cur = Cursor { bytes, at: 0 };
4368 let parts = cur.u32()? as usize;
4369 let mut out = Vec::new();
4370 for _ in 0..parts {
4371 let low = cur.bound()?;
4372 let high = cur.bound()?;
4373 let nulls = cur.u32()? as usize;
4374 out.push(Range { low, high, nulls, exact: false, sum: None });
4375 }
4376 Ok(out)
4377}
4378
4379fn encode_sieves<'a>(sieves: impl Iterator<Item = &'a Option<Sieve>>) -> Result<Vec<u8>> {
4380 let held: Vec<&Option<Sieve>> = sieves.collect();
4381 let mut out = Vec::new();
4382 put_u32(
4383 &mut out,
4384 u32::try_from(held.len()).map_err(|_| invalid("too many parts in a stripe"))?,
4385 );
4386 for sieve in &held {
4387 let length = sieve.as_ref().map_or(0, Sieve::len);
4388 put_u32(&mut out, u32::try_from(length).map_err(|_| invalid("sieve length overflow"))?);
4389 }
4390 for sieve in held.into_iter().flatten() {
4392 out.extend_from_slice(&sieve.to_bytes());
4393 }
4394 Ok(out)
4395}
4396
4397fn decode_sieves(bytes: &[u8]) -> Result<Vec<Option<Sieve>>> {
4403 let parts = u32::from_le_bytes(
4404 bytes
4405 .get(..4)
4406 .ok_or_else(|| invalid("sieve page is truncated"))?
4407 .try_into()
4408 .map_err(|_| invalid("sieve page is truncated"))?,
4409 ) as usize;
4410 let mut lengths = Vec::with_capacity(parts);
4411 for part in 0..parts {
4412 let at = 4 + part * 4;
4413 let field = bytes.get(at..at + 4).ok_or_else(|| invalid("sieve page is truncated"))?;
4414 lengths.push(u32::from_le_bytes(
4415 field.try_into().map_err(|_| invalid("sieve page is truncated"))?,
4416 ) as usize);
4417 }
4418 let mut at = 4 + parts * 4;
4419 let mut out = Vec::with_capacity(parts);
4420 for length in lengths {
4421 if length == 0 {
4422 out.push(None);
4423 continue;
4424 }
4425 let end = at.checked_add(length).ok_or_else(|| invalid("sieve page is truncated"))?;
4426 let field = bytes.get(at..end).ok_or_else(|| invalid("sieve page is truncated"))?;
4427 out.push(Sieve::from_bytes(field));
4428 at = end;
4429 }
4430 if at != bytes.len() {
4431 return Err(invalid("sieve page has trailing bytes"));
4432 }
4433 Ok(out)
4434}
4435
4436fn encode_membership(unique: &[u32]) -> Vec<u8> {
4442 let mut out = Vec::with_capacity(unique.len().saturating_mul(2).saturating_add(5));
4443 put_varint(&mut out, u32::try_from(unique.len()).unwrap_or(u32::MAX));
4444 let mut previous = 0;
4445 for (at, &code) in unique.iter().enumerate() {
4446 put_varint(&mut out, if at == 0 { code } else { code - previous });
4447 previous = code;
4448 }
4449 out
4450}
4451
4452fn take_varint(bytes: &[u8], at: &mut usize) -> Result<u32> {
4453 let mut value = 0_u32;
4454 for shift in (0..35).step_by(7) {
4455 let byte = *bytes.get(*at).ok_or_else(|| invalid("membership varint is truncated"))?;
4456 *at += 1;
4457 let part = u32::from(byte & 0x7f);
4458 if shift == 28 && part > 0x0f {
4459 return Err(invalid("membership varint overflow"));
4460 }
4461 value = value
4462 .checked_add(
4463 part.checked_shl(shift).ok_or_else(|| invalid("membership varint overflow"))?,
4464 )
4465 .ok_or_else(|| invalid("membership varint overflow"))?;
4466 if byte & 0x80 == 0 {
4467 return Ok(value);
4468 }
4469 }
4470 Err(invalid("membership varint is too long"))
4471}
4472
4473fn decode_membership(bytes: &[u8]) -> Result<Vec<u32>> {
4474 let mut at = 0;
4475 let count = take_varint(bytes, &mut at)? as usize;
4476 let mut codes = Vec::with_capacity(count);
4477 let mut previous = 0_u32;
4478 for index in 0..count {
4479 let delta = take_varint(bytes, &mut at)?;
4480 let code = if index == 0 {
4481 delta
4482 } else {
4483 previous.checked_add(delta).ok_or_else(|| invalid("membership code overflow"))?
4484 };
4485 if index > 0 && code <= previous {
4486 return Err(invalid("membership codes are not increasing"));
4487 }
4488 codes.push(code);
4489 previous = code;
4490 }
4491 if at != bytes.len() {
4492 return Err(invalid("membership page has trailing bytes"));
4493 }
4494 Ok(codes)
4495}
4496
4497fn string_dictionary(vector: &Vector) -> Result<Option<Vec<u8>>> {
4498 let mut by_text = HashMap::new();
4499 let mut values = Vec::new();
4500 let mut codes = Vec::with_capacity(vector.len());
4501 let mut plain_bytes = 0_usize;
4502 for row in 0..vector.len() {
4503 let text = vector.text_at(row).unwrap_or("");
4504 plain_bytes = plain_bytes.saturating_add(text.len());
4505 let code = match by_text.get(text) {
4506 Some(&code) => code,
4507 None => {
4508 let code = u32::try_from(values.len())
4509 .map_err(|_| invalid("too many dictionary values"))?;
4510 by_text.insert(text, code);
4511 values.push(text);
4512 code
4513 }
4514 };
4515 codes.push(code);
4516 }
4517 let dictionary_bytes = values.iter().map(|value| value.len()).sum::<usize>();
4518 let encoded = 8_usize
4519 .saturating_add((values.len() + 1).saturating_mul(4))
4520 .saturating_add(dictionary_bytes)
4521 .saturating_add(codes.len().saturating_mul(4));
4522 let plain = (vector.len() + 1).saturating_mul(4).saturating_add(plain_bytes);
4523 if encoded >= plain {
4524 return Ok(None);
4525 }
4526 let mut out = Vec::with_capacity(encoded);
4527 put_u32(
4528 &mut out,
4529 u32::try_from(values.len()).map_err(|_| invalid("too many dictionary values"))?,
4530 );
4531 put_u32(
4532 &mut out,
4533 u32::try_from(dictionary_bytes).map_err(|_| invalid("dictionary payload exceeds 4GiB"))?,
4534 );
4535 let mut offset = 0_u32;
4536 put_u32(&mut out, offset);
4537 for value in &values {
4538 offset = offset
4539 .checked_add(
4540 u32::try_from(value.len()).map_err(|_| invalid("dictionary value is too long"))?,
4541 )
4542 .ok_or_else(|| invalid("dictionary payload exceeds 4GiB"))?;
4543 put_u32(&mut out, offset);
4544 }
4545 for value in values {
4546 out.extend_from_slice(value.as_bytes());
4547 }
4548 for code in codes {
4549 put_u32(&mut out, code);
4550 }
4551 Ok(Some(out))
4552}
4553
4554struct EncodedDictionary {
4555 index: Vec<u8>,
4556 ranks: Vec<u8>,
4557 payload: Vec<Vec<u8>>,
4560}
4561
4562fn head(bytes: &[u8]) -> u64 {
4564 let mut word = [0; 8];
4565 let take = bytes.len().min(8);
4566 word[..take].copy_from_slice(&bytes[..take]);
4567 u64::from_be_bytes(word)
4568}
4569
4570fn rankings(dictionaries: &[Option<GlobalDictionary>]) -> Result<Vec<Vec<(u64, u32)>>> {
4578 let present =
4579 dictionaries.iter().enumerate().filter(|(_, held)| held.is_some()).map(|(at, _)| at);
4580 let present = present.collect::<Vec<_>>();
4581 let mut orders = vec![Vec::new(); dictionaries.len()];
4582 let workers = std::thread::available_parallelism()
4583 .map_or(1, usize::from)
4584 .min(MAX_FREQUENCY_WORKERS)
4585 .min(present.len());
4586 if workers <= 1 {
4587 for at in present {
4588 if let Some(dictionary) = &dictionaries[at] {
4589 orders[at] = dictionary.ranked();
4590 }
4591 }
4592 return Ok(orders);
4593 }
4594 let width = present.len().div_ceil(workers);
4595 let pieces = std::thread::scope(|scope| {
4596 present
4597 .chunks(width)
4598 .map(|columns| {
4599 scope.spawn(|| {
4600 columns
4601 .iter()
4602 .filter_map(|&at| dictionaries[at].as_ref().map(|held| (at, held.ranked())))
4603 .collect::<Vec<_>>()
4604 })
4605 })
4606 .collect::<Vec<_>>()
4607 .into_iter()
4608 .map(|handle| {
4609 handle.join().map_err(|_| Error::internal("a dictionary sort worker panicked"))
4610 })
4611 .collect::<Result<Vec<_>>>()
4612 })?;
4613 for piece in pieces {
4614 for (at, order) in piece {
4615 orders[at] = order;
4616 }
4617 }
4618 Ok(orders)
4619}
4620
4621fn encode_global_dictionary(
4622 dictionary: GlobalDictionary,
4623 order: &[(u64, u32)],
4624) -> Result<EncodedDictionary> {
4625 let values = dictionary.offsets.len() - 1;
4626 if order.len() != values {
4627 return Err(invalid("global dictionary order does not cover its values"));
4628 }
4629 let blocks = values.div_ceil(TEXT_PAYLOAD_VALUES);
4630 let payload = encode_payload(&dictionary)?;
4631 if payload.len() != blocks {
4632 return Err(invalid("global dictionary payload is not the blocks it says it is"));
4633 }
4634 let (ranks, rank_ends) = encode_ranks(order, code_width(values))?;
4635 let rank_blocks = values.div_ceil(TEXT_RANK_BLOCK);
4636 let offset_bits = offset_width(&dictionary.offsets);
4637 let mut index = Vec::with_capacity(
4638 DICTIONARY_HEADER + offset_bytes(values, offset_bits) + (blocks + rank_blocks) * 16,
4639 );
4640 put_u32(
4641 &mut index,
4642 u32::try_from(values).map_err(|_| invalid("global dictionary has too many values"))?,
4643 );
4644 put_u32(&mut index, TEXT_PAYLOAD_VALUES as u32);
4645 put_u32(
4646 &mut index,
4647 u32::try_from(blocks).map_err(|_| invalid("global dictionary has too many blocks"))?,
4648 );
4649 put_u32(&mut index, offset_bits as u32);
4650 encode_offsets(&dictionary.offsets, offset_bits, &mut index)?;
4651 let mut at = 0_u64;
4655 for block in &payload {
4656 at = at
4657 .checked_add(block.len() as u64)
4658 .ok_or_else(|| invalid("global dictionary payload overflow"))?;
4659 put_u64(&mut index, at);
4660 }
4661 for block in &payload {
4662 put_u64(&mut index, checksum(block));
4663 }
4664 if rank_ends.len() != rank_blocks {
4667 return Err(invalid("global dictionary order is not the blocks it says it is"));
4668 }
4669 for end in &rank_ends {
4670 put_u64(&mut index, *end);
4671 }
4672 let mut at = 0_usize;
4673 for end in &rank_ends {
4674 let end = usize::try_from(*end).map_err(|_| invalid("global dictionary order overflow"))?;
4675 put_u64(&mut index, checksum(&ranks[at..end]));
4676 at = end;
4677 }
4678 Ok(EncodedDictionary { index, ranks, payload })
4679}
4680
4681const PAYLOAD_SAMPLE_BLOCKS: usize = 8;
4688
4689fn payload_shapes() -> Vec<chooser::Settled> {
4715 let integers = vec![integer::Kind::Packed];
4716 [
4717 vec![string::Kind::Front, string::Kind::Lz],
4718 vec![string::Kind::Lz, string::Kind::Fsst],
4719 vec![string::Kind::Lz, string::Kind::Plain],
4720 vec![string::Kind::Fsst],
4721 vec![string::Kind::Plain],
4722 ]
4723 .into_iter()
4724 .map(|strings| chooser::Settled::new(strings, integers.clone()))
4725 .collect()
4726}
4727
4728fn encode_payload(dictionary: &GlobalDictionary) -> Result<Vec<Vec<u8>>> {
4734 let values = dictionary.offsets.len() - 1;
4735 let blocks = values.div_ceil(TEXT_PAYLOAD_VALUES);
4736 let run = |block: usize| {
4737 let first = block * TEXT_PAYLOAD_VALUES;
4738 let last = (first + TEXT_PAYLOAD_VALUES).min(values);
4739 (first..last)
4740 .map(|value| {
4741 let from = dictionary.offsets[value] as usize;
4742 let to = dictionary.offsets[value + 1] as usize;
4743 &dictionary.payload[from..to]
4744 })
4745 .collect::<Vec<_>>()
4746 };
4747 let shape = (blocks > PAYLOAD_SAMPLE_BLOCKS).then(|| settle_shape(&run, blocks)).transpose()?;
4750 let one = |block: usize| match &shape {
4751 Some(shape) => string::encode_with(&run(block), shape),
4752 None => string::encode(&run(block)),
4753 };
4754 let workers = std::thread::available_parallelism()
4755 .map_or(1, usize::from)
4756 .min(MAX_FREQUENCY_WORKERS)
4757 .min(blocks);
4758 if workers <= 1 {
4759 return (0..blocks).map(one).collect();
4760 }
4761 let next = AtomicUsize::new(0);
4762 let pieces = std::thread::scope(|scope| {
4763 (0..workers)
4764 .map(|_| {
4765 scope.spawn(|| {
4766 let mut mine = Vec::new();
4767 loop {
4768 let block = next.fetch_add(1, Atomic::Relaxed);
4769 if block >= blocks {
4770 break;
4771 }
4772 mine.push((block, one(block)?));
4773 }
4774 Ok(mine)
4775 })
4776 })
4777 .collect::<Vec<_>>()
4778 .into_iter()
4779 .map(|handle| {
4780 handle.join().map_err(|_| Error::internal("a dictionary encode worker panicked"))?
4781 })
4782 .collect::<Result<Vec<_>>>()
4783 })?;
4784 let mut payload = vec![Vec::new(); blocks];
4785 for piece in pieces {
4786 for (block, bytes) in piece {
4787 payload[block] = bytes;
4788 }
4789 }
4790 Ok(payload)
4791}
4792
4793fn settle_shape<'a>(
4801 run: &dyn Fn(usize) -> Vec<&'a [u8]>,
4802 blocks: usize,
4803) -> Result<chooser::Settled> {
4804 let last = blocks - 1;
4805 let sample = (0..PAYLOAD_SAMPLE_BLOCKS)
4806 .map(|region| run(region * last / (PAYLOAD_SAMPLE_BLOCKS - 1)))
4807 .collect::<Vec<_>>();
4808 let mut best: Option<(chooser::Settled, usize)> = None;
4809 for shape in payload_shapes() {
4810 let mut size = 0;
4811 for block in &sample {
4812 size += string::encode_with(block, &shape)?.len();
4813 }
4814 if best.as_ref().is_none_or(|(_, smallest)| size < *smallest) {
4815 best = Some((shape, size));
4816 }
4817 }
4818 best.map(|(shape, _)| shape)
4819 .ok_or_else(|| invalid("no shape applies to a global dictionary payload"))
4820}
4821
4822fn encode_ranks(order: &[(u64, u32)], code_bits: usize) -> Result<(Vec<u8>, Vec<u64>)> {
4829 let mut out = Vec::with_capacity(order.len() * 4);
4830 let mut ends = Vec::with_capacity(order.len().div_ceil(TEXT_RANK_BLOCK));
4831 let mut heads = Vec::with_capacity(TEXT_RANK_BLOCK);
4832 let mut codes = Vec::with_capacity(TEXT_RANK_BLOCK);
4833 for block in order.chunks(TEXT_RANK_BLOCK) {
4834 let base = block.first().map_or(0, |&(head, _)| head);
4837 let span = block.last().map_or(0, |&(head, _)| head.wrapping_sub(base));
4838 let width = (u64::BITS - span.leading_zeros()) as usize;
4839 heads.clear();
4840 codes.clear();
4841 for &(head, code) in block {
4842 heads.push(head.wrapping_sub(base));
4843 codes.push(u64::from(code));
4844 }
4845 put_u64(&mut out, base);
4846 out.push(width as u8);
4847 bitpack::pack_tail(&heads, width, &mut out)
4848 .map_err(|_| invalid("global dictionary heads do not pack"))?;
4849 bitpack::pack_tail(&codes, code_bits, &mut out)
4850 .map_err(|_| invalid("global dictionary codes do not pack"))?;
4851 ends.push(out.len() as u64);
4852 }
4853 Ok((out, ends))
4854}
4855
4856fn open_global_dictionary(
4863 file: Arc<File>,
4864 page: Page,
4865 ty: &LogicalType,
4866 keep_budget: usize,
4867) -> Result<Vector> {
4868 if ty != &LogicalType::Varchar {
4869 return Err(invalid("global dictionary belongs to a non-string column"));
4870 }
4871 let mut header = [0; DICTIONARY_HEADER];
4872 read_at(&file, page.offset, &mut header)?;
4873 let count = u32::from_le_bytes(header[0..4].try_into().expect("four bytes")) as usize;
4874 let per_block = u32::from_le_bytes(header[4..8].try_into().expect("four bytes")) as usize;
4875 let blocks = u32::from_le_bytes(header[8..12].try_into().expect("four bytes")) as usize;
4876 let offset_bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
4877 if per_block != TEXT_PAYLOAD_VALUES {
4878 return Err(invalid("global dictionary block width differs"));
4879 }
4880 if blocks != count.div_ceil(TEXT_PAYLOAD_VALUES) {
4881 return Err(invalid("global dictionary block count differs from its value count"));
4882 }
4883 if offset_bits > u32::BITS as usize {
4884 return Err(invalid("global dictionary packs offsets past a payload"));
4885 }
4886 let offset_len = offset_bytes(count, offset_bits);
4887 let ranks = count;
4892 let rank_blocks = ranks.div_ceil(TEXT_RANK_BLOCK);
4893 let hash_len = blocks
4896 .checked_add(rank_blocks)
4897 .and_then(|words| words.checked_mul(16))
4898 .ok_or_else(|| invalid("global dictionary block count overflow"))?;
4899 let index_len = DICTIONARY_HEADER
4900 .checked_add(offset_len)
4901 .and_then(|len| len.checked_add(hash_len))
4902 .ok_or_else(|| invalid("global dictionary header overflow"))?;
4903 if index_len > page.length as usize {
4904 return Err(invalid("global dictionary offset index exceeds its page"));
4905 }
4906 let mut index = vec![0; index_len];
4907 index[..DICTIONARY_HEADER].copy_from_slice(&header);
4908 read_at(&file, page.offset + DICTIONARY_HEADER as u64, &mut index[DICTIONARY_HEADER..])?;
4909 if checksum(&index) != page.hash {
4910 return Err(invalid("global dictionary index checksum differs"));
4911 }
4912 let offsets = index[DICTIONARY_HEADER..DICTIONARY_HEADER + offset_len].to_vec();
4913 let mut words = index[DICTIONARY_HEADER + offset_len..]
4914 .chunks_exact(8)
4915 .map(|part| u64::from_le_bytes(part.try_into().expect("eight bytes")))
4916 .collect::<Vec<_>>();
4917 let mut hashes = words.split_off(blocks);
4918 let mut rank_ends = hashes.split_off(blocks);
4919 let rank_hashes = rank_ends.split_off(rank_blocks);
4920 let ends = words;
4921 if rank_ends.windows(2).any(|pair| pair[0] >= pair[1]) {
4924 return Err(invalid("global dictionary order blocks do not rise"));
4925 }
4926 let rank_len = usize::try_from(rank_ends.last().copied().unwrap_or_default())
4927 .map_err(|_| invalid("global dictionary rank overflow"))?;
4928 let body_len = index_len
4929 .checked_add(rank_len)
4930 .ok_or_else(|| invalid("global dictionary header overflow"))?;
4931 if body_len > page.length as usize {
4932 return Err(invalid("global dictionary order exceeds its page"));
4933 }
4934 let stored_len = page.length as usize - body_len;
4937 if ends.last().copied().unwrap_or_default() as usize != stored_len
4938 || ends.windows(2).any(|pair| pair[0] > pair[1])
4939 {
4940 return Err(invalid("global dictionary blocks do not bound the payload"));
4941 }
4942 Vector::external_text(
4943 LogicalType::Varchar,
4944 Arc::new(NativeText {
4945 file,
4946 values: count,
4947 offsets,
4948 offset_bits,
4949 ranks,
4950 rank_at: page.offset + index_len as u64,
4951 rank_ends,
4952 rank_hashes,
4953 rank_blocks: (0..rank_blocks).map(|_| OnceLock::new()).collect(),
4954 code_bits: code_width(count),
4955 code_ranks: OnceLock::new(),
4956 payload: page.offset + body_len as u64,
4957 ends,
4958 hashes,
4959 blocks: (0..blocks).map(|_| OnceLock::new()).collect(),
4960 keep_budget,
4961 payload_kept: AtomicUsize::new(0),
4962 }),
4963 )
4964}
4965
4966fn decode(
4967 ty: &LogicalType,
4968 rows: usize,
4969 bytes: &[u8],
4970 global: Option<Arc<Vector>>,
4971) -> Result<Vector> {
4972 let mut cur = Cursor { bytes, at: 0 };
4973 let codec = cur.u8()?;
4974 let flag = cur.u8()?;
4975 let validity = match flag {
4976 0 => Validity::AllValid,
4977 1 => Validity::AllInvalid,
4978 2 => {
4979 let mask = cur.take(rows.div_ceil(8))?;
4980 Validity::from_iter(rows, |row| mask[row / 8] >> (row % 8) & 1 == 1)
4981 }
4982 _ => return Err(invalid("page validity tag differs")),
4983 };
4984 if codec == 1 {
4985 if ty != &LogicalType::Varchar {
4986 return Err(invalid("dictionary codec belongs to a non-string page"));
4987 }
4988 let count = cur.u32()? as usize;
4989 let payload_len = cur.u32()? as usize;
4990 let offset_bytes = cur.take(
4991 (count + 1)
4992 .checked_mul(4)
4993 .ok_or_else(|| invalid("dictionary offset count overflow"))?,
4994 )?;
4995 let offsets = offset_bytes
4996 .chunks_exact(4)
4997 .map(|part| u32::from_le_bytes(part.try_into().expect("four bytes")))
4998 .collect::<Vec<_>>();
4999 let payload = cur.take(payload_len)?.to_vec();
5000 if offsets.first() != Some(&0)
5001 || offsets.last().copied().map(|last| last as usize) != Some(payload.len())
5002 || offsets.windows(2).any(|pair| pair[0] > pair[1])
5003 {
5004 return Err(invalid("dictionary offsets do not bound the payload"));
5005 }
5006 let mut strings = StringColumn::over(Buffer::from_vec(payload));
5007 for pair in offsets.windows(2) {
5008 strings.push_in_place(pair[0] as usize, (pair[1] - pair[0]) as usize)?;
5009 }
5010 let mut codes = Vec::with_capacity(rows);
5011 for _ in 0..rows {
5012 codes.push(cur.u32()?);
5013 }
5014 if codes.iter().any(|code| *code as usize >= count) {
5015 return Err(invalid("dictionary code is out of range"));
5016 }
5017 if cur.at != bytes.len() {
5018 return Err(invalid("dictionary page has trailing bytes"));
5019 }
5020 let dictionary = Vector::flat(LogicalType::Varchar, Data::Varlen(strings))?;
5021 return Ok(Vector::dictionary(codes, dictionary)?.with_validity(validity));
5022 }
5023 if codec == 3 || codec == 4 {
5024 let dictionary = global.ok_or_else(|| invalid("global code page has no dictionary"))?;
5025 let codes = if codec == 4 {
5026 let wide = integer::decode(&bytes[cur.at..])?;
5029 if wide.len() != rows {
5030 return Err(invalid("encoded code page holds the wrong number of rows"));
5031 }
5032 let mut codes = Vec::with_capacity(wide.len());
5039 let mut seen = 0_i64;
5040 for &code in &wide {
5041 seen |= code;
5042 codes.push(code as u32);
5043 }
5044 if seen < 0 || seen > i64::from(u32::MAX) {
5045 return Err(invalid("code is not a code"));
5046 }
5047 codes
5048 } else {
5049 let mut codes = Vec::with_capacity(rows);
5050 for _ in 0..rows {
5051 codes.push(cur.u32()?);
5052 }
5053 if cur.at != bytes.len() {
5054 return Err(invalid("global code page has trailing bytes"));
5055 }
5056 codes
5057 };
5058 let highest = codes.iter().copied().max();
5059 return Ok(Vector::stable_dictionary_validated(codes, dictionary, highest)?
5060 .with_validity(validity));
5061 }
5062 if codec == 5 {
5063 let values = integer::decode(&bytes[cur.at..])?;
5065 if values.len() != rows {
5066 return Err(invalid("cascade page holds the wrong number of rows"));
5067 }
5068 let data = narrowed(ty, values)?;
5069 return Ok(Vector::flat(ty.clone(), data)?.with_validity(validity));
5070 }
5071 if codec == 2 {
5072 let width = u32::from(cur.u8()?);
5073 let base = i128::from_le_bytes(cur.take(16)?.try_into().expect("sixteen bytes"));
5074 let count = cur.u32()? as usize;
5075 let mut words = Vec::with_capacity(count);
5076 for _ in 0..count {
5077 words.push(cur.u64()?);
5078 }
5079 if cur.at != bytes.len() {
5080 return Err(invalid("packed page has trailing bytes"));
5081 }
5082 return Ok(Vector::packed(ty.clone(), words, width, base, rows)?.with_validity(validity));
5083 }
5084 if codec != 0 {
5085 return Err(invalid("page codec is unknown"));
5086 }
5087 let data = match ty {
5088 LogicalType::TinyInt => {
5089 let values = cur.take(rows)?;
5090 Data::Int8(values.iter().map(|item| *item as i8).collect::<Vec<_>>().into())
5091 }
5092 LogicalType::UTinyInt => Data::UInt8(cur.take(rows)?.to_vec().into()),
5093 LogicalType::SmallInt => {
5094 let values =
5095 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5096 Data::Int16(
5097 values
5098 .chunks_exact(2)
5099 .map(|item| i16::from_le_bytes(item.try_into().expect("two bytes")))
5100 .collect::<Vec<_>>()
5101 .into(),
5102 )
5103 }
5104 LogicalType::USmallInt => {
5105 let values =
5106 cur.take(rows.checked_mul(2).ok_or_else(|| invalid("page size overflow"))?)?;
5107 Data::UInt16(
5108 values
5109 .chunks_exact(2)
5110 .map(|item| u16::from_le_bytes(item.try_into().expect("two bytes")))
5111 .collect::<Vec<_>>()
5112 .into(),
5113 )
5114 }
5115 LogicalType::UInteger => {
5116 let values =
5117 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5118 Data::UInt32(
5119 values
5120 .chunks_exact(4)
5121 .map(|item| u32::from_le_bytes(item.try_into().expect("four bytes")))
5122 .collect::<Vec<_>>()
5123 .into(),
5124 )
5125 }
5126 LogicalType::UBigInt => {
5127 let values =
5128 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5129 Data::UInt64(
5130 values
5131 .chunks_exact(8)
5132 .map(|item| u64::from_le_bytes(item.try_into().expect("eight bytes")))
5133 .collect::<Vec<_>>()
5134 .into(),
5135 )
5136 }
5137 LogicalType::Integer | LogicalType::Date => {
5138 let values =
5139 cur.take(rows.checked_mul(4).ok_or_else(|| invalid("page size overflow"))?)?;
5140 Data::Int32(
5141 values
5142 .chunks_exact(4)
5143 .map(|item| i32::from_le_bytes(item.try_into().expect("four bytes")))
5144 .collect::<Vec<_>>()
5145 .into(),
5146 )
5147 }
5148 LogicalType::BigInt | LogicalType::Timestamp => {
5149 let values =
5150 cur.take(rows.checked_mul(8).ok_or_else(|| invalid("page size overflow"))?)?;
5151 Data::Int64(
5152 values
5153 .chunks_exact(8)
5154 .map(|item| i64::from_le_bytes(item.try_into().expect("eight bytes")))
5155 .collect::<Vec<_>>()
5156 .into(),
5157 )
5158 }
5159 LogicalType::Boolean => {
5160 let values = cur.take(rows)?;
5161 if values.iter().any(|value| *value > 1) {
5162 return Err(invalid("boolean page has another value"));
5163 }
5164 Data::Bool(values.iter().map(|value| *value == 1).collect::<Vec<_>>().into())
5165 }
5166 LogicalType::Varchar => {
5167 let offset_bytes = cur
5168 .take((rows + 1).checked_mul(4).ok_or_else(|| invalid("offset count overflow"))?)?;
5169 let offsets = offset_bytes
5170 .chunks_exact(4)
5171 .map(|part| u32::from_le_bytes(part.try_into().expect("four bytes")))
5172 .collect::<Vec<_>>();
5173 let payload = cur.take(bytes.len() - cur.at)?.to_vec();
5174 if offsets.first() != Some(&0)
5175 || offsets.last().copied().map(|last| last as usize) != Some(payload.len())
5176 || offsets.windows(2).any(|pair| pair[0] > pair[1])
5177 {
5178 return Err(invalid("string offsets do not bound the payload"));
5179 }
5180 let mut values = StringColumn::over(Buffer::from_vec(payload));
5181 for pair in offsets.windows(2) {
5182 values.push_in_place(pair[0] as usize, (pair[1] - pair[0]) as usize)?;
5183 }
5184 Data::Varlen(values)
5185 }
5186 _ => return Err(Error::not_implemented(format!("native page for {ty}"))),
5187 };
5188 if cur.at != bytes.len() {
5189 return Err(invalid("page has trailing bytes"));
5190 }
5191 Ok(Vector::flat(ty.clone(), data)?.with_validity(validity))
5192}
5193
5194#[cfg(test)]
5195mod tests {
5196 use std::fs;
5197 use std::io::{Seek, SeekFrom, Write};
5198 use std::path::PathBuf;
5199 use std::time::{SystemTime, UNIX_EPOCH};
5200
5201 use rudb_common::Value;
5202 use rudb_common::bounds::Op;
5203
5204 use super::*;
5205
5206 #[test]
5207 fn checksum_matches_fixed_vectors() {
5208 assert_eq!(checksum(b""), 0xef46_db37_51d8_e999);
5209 assert_eq!(checksum(b"a"), 0xd24e_c4f1_a98c_6e5b);
5210 assert_eq!(checksum(b"abc"), 0x44bc_2cf5_ad77_0999);
5211 }
5212
5213 fn path(label: &str) -> PathBuf {
5214 let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
5215 std::env::temp_dir().join(format!("rudb-native-{label}-{}-{stamp}.rdb", std::process::id()))
5216 }
5217
5218 #[test]
5220 fn a_read_at_an_offset_ignores_where_another_thread_left_the_cursor() {
5221 const SPANS: usize = 64;
5222 const SPAN: usize = 512;
5223 let path = path("positional");
5224 let content: Vec<u8> =
5225 (0..SPANS).flat_map(|span| std::iter::repeat_n(span as u8, SPAN)).collect();
5226 fs::write(&path, &content).expect("the file is written");
5227 let file = Arc::new(File::open(&path).expect("the file opens"));
5228 std::thread::scope(|scope| {
5229 for _ in 0..8 {
5230 let file = Arc::clone(&file);
5231 scope.spawn(move || {
5232 for _ in 0..64 {
5233 for span in 0..SPANS {
5234 let mut bytes = [0_u8; SPAN];
5235 read_at(&file, (span * SPAN) as u64, &mut bytes)
5236 .expect("the span reads");
5237 assert!(
5238 bytes.iter().all(|byte| *byte == span as u8),
5239 "span {span} came back as {}",
5240 bytes[0],
5241 );
5242 }
5243 }
5244 });
5245 }
5246 });
5247 let mut past = [0_u8; SPAN];
5248 let end = (SPANS * SPAN) as u64;
5249 let error = read_at(&file, end, &mut past).expect_err("a read past the end is refused");
5250 assert!(error.message().contains("ends before its declared length"), "{error}");
5251 drop(file);
5252 let _ = fs::remove_file(&path);
5253 }
5254
5255 #[test]
5261 fn a_writer_puts_a_page_where_it_said_it_did_wherever_the_cursor_has_got_to() {
5262 let path = path("cursor");
5263 let mut writer = Writer::create(
5264 &path,
5265 "items",
5266 vec![
5267 Field::required("id", LogicalType::Integer),
5268 Field::new("text", LogicalType::Varchar),
5269 ],
5270 )
5271 .expect("new file");
5272 writer.append(&sample()).expect("first part");
5273 writer.file.seek(SeekFrom::Start(0)).expect("the cursor goes back to the header");
5274 writer.append(&sample()).expect("second part");
5275 writer.file.seek(SeekFrom::Start(1)).expect("and somewhere useless again");
5276 writer.finish().expect("commit");
5277 let reader = Reader::open(&path).expect("reopen from disk");
5278 assert_eq!(reader.table().rows(), 6);
5279 let ids = reader.read(0, &[0]).expect("the integer page reads back");
5280 assert_eq!(ids.value_at(0, 0), Value::Integer(4));
5281 assert_eq!(ids.value_at(2, 0), Value::Integer(-2));
5282 let text = reader.read(1, &[1]).expect("the text page reads back");
5283 assert_eq!(text.value_at(1, 0), Value::Null);
5284 assert_eq!(text.value_at(2, 0), Value::Varchar("long text after a slash".into()));
5285 let end = reader.table().stripes().iter().flat_map(|stripe| {
5288 stripe
5289 .pages
5290 .iter()
5291 .map(|page| page.offset + u64::from(page.length))
5292 .chain(std::iter::once(stripe.index.offset + u64::from(stripe.index.length)))
5293 });
5294 let last = end.fold(HEADER, u64::max);
5295 let directory = fs::metadata(&path).expect("the file is there").len();
5296 assert!(last <= directory, "a page runs to {last} in a file of {directory} bytes");
5297 fs::remove_file(path).expect("remove scratch file");
5298 }
5299
5300 fn dictionary_index_len(header: &[u8; DICTIONARY_HEADER]) -> u64 {
5306 let count = u64::from(u32::from_le_bytes(header[0..4].try_into().expect("four bytes")));
5307 let blocks = u64::from(u32::from_le_bytes(header[8..12].try_into().expect("four bytes")));
5308 let bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
5309 let rank_blocks = count.div_ceil(TEXT_RANK_BLOCK as u64);
5310 DICTIONARY_HEADER as u64
5311 + offset_bytes(count as usize, bits) as u64
5312 + (blocks + rank_blocks) * 16
5313 }
5314
5315 fn last_rank_end(file: &File, offset: u64, header: &[u8; DICTIONARY_HEADER]) -> u64 {
5317 let count = u64::from(u32::from_le_bytes(header[0..4].try_into().expect("four bytes")));
5318 let blocks = u64::from(u32::from_le_bytes(header[8..12].try_into().expect("four bytes")));
5319 let bits = u32::from_le_bytes(header[12..16].try_into().expect("four bytes")) as usize;
5320 let rank_blocks = count.div_ceil(TEXT_RANK_BLOCK as u64);
5321 let at = offset
5322 + DICTIONARY_HEADER as u64
5323 + offset_bytes(count as usize, bits) as u64
5324 + blocks * 16
5325 + (rank_blocks - 1) * 8;
5326 let mut end = [0; 8];
5327 read_at(file, at, &mut end).expect("the last rank block end");
5328 u64::from_le_bytes(end)
5329 }
5330
5331 fn sample() -> Chunk {
5332 Chunk::new(vec![
5333 Vector::from_values(
5334 LogicalType::Integer,
5335 &[Value::Integer(4), Value::Integer(9), Value::Integer(-2)],
5336 )
5337 .expect("integers"),
5338 Vector::from_values(
5339 LogicalType::Varchar,
5340 &[
5341 Value::Varchar("alpha".into()),
5342 Value::Null,
5343 Value::Varchar("long text after a slash".into()),
5344 ],
5345 )
5346 .expect("strings"),
5347 ])
5348 .expect("matching rows")
5349 }
5350
5351 fn sample_ids() -> Chunk {
5352 Chunk::new(vec![
5353 Vector::flat(LogicalType::Integer, Data::Int32(vec![7, 8, 9].into()))
5354 .expect("integers"),
5355 ])
5356 .expect("one column")
5357 }
5358
5359 #[test]
5360 fn committed_file_reopens_and_reads_only_requested_columns() {
5361 let path = path("reopen");
5362 let mut writer = Writer::create(
5363 &path,
5364 "items",
5365 vec![
5366 Field::required("id", LogicalType::Integer),
5367 Field::new("text", LogicalType::Varchar),
5368 ],
5369 )
5370 .expect("new file");
5371 writer.append(&sample()).expect("first part");
5372 writer.append(&sample()).expect("second part");
5373 writer.finish().expect("commit");
5374 let reader = Reader::open(&path).expect("reopen from disk");
5375 assert_eq!(reader.table().rows(), 6);
5376 assert_eq!(reader.table().stripes().len(), 1);
5379 assert_eq!(reader.parts(), 2);
5380 assert_eq!(reader.part_rows(0), 3);
5381 assert_eq!(reader.part_rows(1), 3);
5382 let text = reader.read(1, &[1]).expect("only text page");
5383 assert_eq!(text.width(), 1);
5384 assert_eq!(text.value_at(1, 0), Value::Null);
5385 assert_eq!(text.value_at(2, 0), Value::Varchar("long text after a slash".into()));
5386 let sparse = reader.read_sparse(1, &[1]).expect("one part without its whole page");
5387 assert_eq!(sparse.width(), 1);
5388 assert_eq!(sparse.value_at(1, 0), Value::Null);
5389 assert_eq!(sparse.value_at(2, 0), Value::Varchar("long text after a slash".into()));
5390 assert!(!reader.skips_codes(0, 1, &[0]).expect("alpha is in the stripe"));
5391 assert!(!reader.skips_codes(0, 1, &[2]).expect("long text is in the stripe"));
5392 assert!(reader.skips_codes(0, 1, &[3]).expect("unknown code is absent"));
5393 let count = reader.read(0, &[]).expect("no page is needed for count");
5394 assert_eq!(count.len(), 3);
5395 assert!(reader.skips(0, &[Probe { column: 0, op: Op::Greater, value: Bound::Int(100) }]));
5396 assert!(!reader.skips(0, &[Probe { column: 0, op: Op::Greater, value: Bound::Int(0) }]));
5397 let integers = reader.top_frequencies(0, 1).expect("valid integer synopsis").expect("kept");
5398 assert_eq!(
5399 integers,
5400 vec![(Value::Integer(-2), 2), (Value::Integer(4), 2), (Value::Integer(9), 2),]
5401 );
5402 let strings = reader.top_frequencies(1, 1).expect("valid string synopsis").expect("kept");
5403 assert_eq!(strings.len(), 3);
5404 assert!(strings.contains(&(Value::Null, 2)));
5405 assert!(strings.contains(&(Value::Varchar("alpha".into()), 2)));
5406 assert!(strings.contains(&(Value::Varchar("long text after a slash".into()), 2)));
5407 fs::remove_file(path).expect("remove scratch file");
5408 }
5409
5410 #[test]
5418 fn runs_handed_over_out_of_order_still_read_back_in_source_order() {
5419 let path = path("interleaved-runs");
5420 let mut writer =
5421 Writer::create(&path, "interleaved", vec![Field::new("v", LogicalType::BigInt)])
5422 .expect("new file");
5423 for morsel in [2_u64, 0, 3, 1] {
5424 let parts = (0..4_u64)
5425 .map(|chunk| {
5426 let first = i64::try_from(morsel * 32 + chunk * 8).expect("small");
5427 let values =
5428 (0..8_i64).map(|row| Value::BigInt(first + row)).collect::<Vec<_>>();
5429 let column =
5430 Vector::from_values(LogicalType::BigInt, &values).expect("a column");
5431 ((morsel, chunk), Chunk::new(vec![column]).expect("one column"))
5432 })
5433 .collect::<Vec<_>>();
5434 writer.append_stripe(parts).expect("a stripe");
5435 }
5436 writer.finish().expect("commit");
5437
5438 let reader = Reader::open(&path).expect("valid directory");
5439 assert_eq!(reader.table().stripes().len(), 4, "a run is a stripe of its own");
5440 assert_eq!(reader.table().rows(), 128);
5441 for part in 0..16_usize {
5442 let read = reader.read(part, &[0]).expect("a part back");
5443 for row in 0..8_usize {
5444 let want = i64::try_from(part * 8 + row).expect("small");
5445 assert_eq!(read.value_at(row, 0), Value::BigInt(want), "part {part} row {row}");
5446 }
5447 }
5448 fs::remove_file(path).expect("remove scratch file");
5449 }
5450
5451 #[test]
5454 fn runs_that_overlap_each_other_are_refused_at_commit() {
5455 let path = path("overlapping-runs");
5456 let mut writer =
5457 Writer::create(&path, "overlapping", vec![Field::new("v", LogicalType::BigInt)])
5458 .expect("new file");
5459 let one = |order: (u64, u64)| {
5460 let column =
5461 Vector::from_values(LogicalType::BigInt, &[Value::BigInt(1)]).expect("a column");
5462 (order, Chunk::new(vec![column]).expect("one column"))
5463 };
5464 writer.append_stripe(vec![one((0, 0)), one((0, 2))]).expect("a stripe");
5467 writer.append_stripe(vec![one((0, 1))]).expect("a stripe");
5468 let error = writer.finish().expect_err("the runs overlap");
5469 assert!(error.message().contains("source order"), "{error}");
5470 fs::remove_file(path).expect("remove scratch file");
5471 }
5472
5473 #[test]
5476 fn a_run_longer_than_a_stripe_is_refused() {
5477 let path = path("overlong-run");
5478 let mut writer =
5479 Writer::create(&path, "overlong", vec![Field::new("v", LogicalType::BigInt)])
5480 .expect("new file");
5481 let parts = (0..=STRIPE_PARTS)
5482 .map(|at| {
5483 let column = Vector::from_values(LogicalType::BigInt, &[Value::BigInt(1)])
5484 .expect("a column");
5485 let chunk = Chunk::new(vec![column]).expect("one column");
5486 ((0, u64::try_from(at).expect("small")), chunk)
5487 })
5488 .collect::<Vec<_>>();
5489 let error = writer.append_stripe(parts).expect_err("one part too many");
5490 assert!(error.message().contains("more parts than it holds"), "{error}");
5491 fs::remove_file(path).expect("remove scratch file");
5492 }
5493
5494 #[test]
5500 fn parts_past_the_stripe_bound_start_a_new_stripe() {
5501 let path = path("stripe-bound");
5502 let mut writer = Writer::create(
5503 &path,
5504 "items",
5505 vec![
5506 Field::required("id", LogicalType::Integer),
5507 Field::new("text", LogicalType::Varchar),
5508 ],
5509 )
5510 .expect("new file");
5511 let parts = STRIPE_PARTS * 2 + 3;
5512 for part in 0..parts {
5513 let id = part as i32;
5514 let chunk = Chunk::new(vec![
5515 Vector::from_values(
5516 LogicalType::Integer,
5517 &[Value::Integer(id), Value::Integer(-id)],
5518 )
5519 .expect("integers"),
5520 Vector::from_values(
5521 LogicalType::Varchar,
5522 &[Value::Varchar(format!("value {part}")), Value::Null],
5523 )
5524 .expect("strings"),
5525 ])
5526 .expect("matching rows");
5527 writer.append(&chunk).expect("one part");
5528 }
5529 writer.finish().expect("commit");
5530
5531 let reader = Reader::open(&path).expect("reopen from disk");
5532 assert_eq!(reader.parts(), parts);
5533 assert_eq!(reader.table().rows(), parts * 2);
5534 assert_eq!(reader.table().stripes().len(), parts.div_ceil(STRIPE_PARTS));
5535 assert_eq!(reader.table().stripes()[0].parts(), STRIPE_PARTS);
5536 assert_eq!(reader.table().stripes()[0].rows(), STRIPE_PARTS * 2);
5537 assert_eq!(reader.table().stripes()[2].parts(), 3);
5538 for part in (0..parts).rev() {
5541 let dense = reader.read(part, &[0, 1]).expect("a whole page read");
5542 let sparse = reader.read_sparse(part, &[0, 1]).expect("one part read");
5543 for chunk in [&dense, &sparse] {
5544 assert_eq!(chunk.len(), 2, "part {part} has its own row count");
5545 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
5546 assert_eq!(chunk.value_at(1, 0), Value::Integer(-(part as i32)));
5547 assert_eq!(chunk.value_at(0, 1), Value::Varchar(format!("value {part}")));
5548 assert_eq!(chunk.value_at(1, 1), Value::Null);
5549 }
5550 }
5551 let above = [Probe { column: 0, op: Op::Greater, value: Bound::Int(100) }];
5554 assert!(reader.skips(0, &above), "the first stripe stops at 63");
5555 assert!(!reader.skips(STRIPE_PARTS * 2, &above), "the third stripe reaches 130");
5556 fs::remove_file(path).expect("remove scratch file");
5557 }
5558
5559 fn scattered(n: i64) -> i64 {
5561 n.wrapping_mul(-7_046_029_254_386_353_131)
5562 }
5563
5564 #[test]
5570 fn a_part_is_skipped_when_its_sieve_does_not_hold_the_constant() {
5571 let path = path("sieve-skip");
5572 let mut writer =
5573 Writer::create(&path, "hits", vec![Field::required("id", LogicalType::BigInt)])
5574 .expect("new file");
5575 let parts = STRIPE_PARTS + 3;
5576 let per_part = 128;
5580 for part in 0..parts {
5581 let held: Vec<Value> = (0..per_part)
5582 .map(|row| Value::BigInt(scattered((part * per_part + row) as i64)))
5583 .collect();
5584 let chunk =
5585 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
5586 .expect("one column");
5587 writer.append(&chunk).expect("one part");
5588 }
5589 writer.finish().expect("commit");
5590
5591 let reader = Reader::open(&path).expect("reopen from disk");
5592 let probe = |value: i64| Probe {
5593 column: 0,
5594 op: Op::Equal,
5595 value: Bound::Int(i128::from(scattered(value))),
5596 };
5597 for wanted in [0_i64, (per_part + 1) as i64, (parts * per_part - 1) as i64] {
5598 let tests = [probe(wanted)];
5599 let kept: Vec<usize> = (0..parts).filter(|&part| !reader.skips(part, &tests)).collect();
5600 let home = wanted as usize / per_part;
5601 assert!(kept.contains(&home), "the part holding {wanted} is read");
5602 assert!(kept.len() <= 2, "{wanted} keeps {kept:?}, which is more than one stray part");
5606 }
5607 let absent = [probe((parts * per_part) as i64 + 1)];
5608 let kept = (0..parts).filter(|&part| !reader.skips(part, &absent)).count();
5609 assert!(kept <= 1, "{kept} parts of {parts} kept a value no part holds");
5610 let tests = [probe(0)];
5613 assert!(
5614 reader.table().stripes().iter().all(|stripe| !stripe.zone.skips(&tests)),
5615 "the bounds rule out no stripe at all"
5616 );
5617 fs::remove_file(path).expect("remove scratch file");
5618 }
5619
5620 #[test]
5626 fn a_part_is_skipped_when_its_own_bounds_rule_out_a_comparison_the_stripe_keeps() {
5627 let path = path("part-range-skip");
5628 let mut writer =
5629 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
5630 .expect("new file");
5631 let parts = STRIPE_PARTS + 3;
5632 let per_part = 128;
5633 for part in 0..parts {
5634 let held: Vec<Value> = (0..per_part)
5638 .map(|row| {
5639 Value::BigInt((part * 1_000) as i64 + (scattered(row as i64).rem_euclid(900)))
5640 })
5641 .collect();
5642 let chunk =
5643 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
5644 .expect("one column");
5645 writer.append(&chunk).expect("one part");
5646 }
5647 writer.finish().expect("commit");
5648
5649 let reader = Reader::open(&path).expect("reopen from disk");
5650 let under = [Probe { column: 0, op: Op::Less, value: Bound::Int(3_000) }];
5651 let kept: Vec<usize> = (0..parts).filter(|&part| !reader.skips(part, &under)).collect();
5652 assert_eq!(kept, vec![0, 1, 2], "only the three parts that start under three thousand");
5653 assert!(!reader.stripe_skips(0, &under), "the stripe reaches from zero and keeps itself");
5655 fs::remove_file(path).expect("remove scratch file");
5656 }
5657
5658 #[test]
5662 fn a_part_is_waved_through_when_its_own_bounds_pass_a_comparison_the_stripe_cannot() {
5663 let path = path("part-range-certain");
5664 let mut writer =
5665 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
5666 .expect("new file");
5667 let parts = STRIPE_PARTS + 3;
5668 let per_part = 128;
5669 for part in 0..parts {
5670 let held: Vec<Value> = (0..per_part)
5671 .map(|row| {
5672 Value::BigInt((part * 1_000) as i64 + (scattered(row as i64).rem_euclid(900)))
5673 })
5674 .collect();
5675 let chunk =
5676 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
5677 .expect("one column");
5678 writer.append(&chunk).expect("one part");
5679 }
5680 writer.finish().expect("commit");
5681
5682 let reader = Reader::open(&path).expect("reopen from disk");
5683 let under = [Probe { column: 0, op: Op::Less, value: Bound::Int(3_000) }];
5684 let waved: Vec<usize> = (0..parts).filter(|&part| reader.certain(part, &under)).collect();
5685 assert_eq!(waved, vec![0, 1, 2], "the three parts that end under three thousand");
5686 assert!(!reader.stripe_skips(0, &under), "the stripe straddles the comparison");
5689 fs::remove_file(path).expect("remove scratch file");
5690 }
5691
5692 #[test]
5695 fn a_stripe_of_one_part_writes_no_range_page_and_a_stripe_of_many_does() {
5696 for (parts, wanted) in [(1_usize, false), (STRIPE_PARTS, true)] {
5697 let path = path("part-range-page");
5698 let mut writer =
5699 Writer::create(&path, "hits", vec![Field::required("at", LogicalType::BigInt)])
5700 .expect("new file");
5701 for part in 0..parts {
5702 let held: Vec<Value> = (0..128)
5703 .map(|row| {
5704 Value::BigInt((part * 1_000) as i64 + scattered(row as i64).rem_euclid(900))
5705 })
5706 .collect();
5707 let chunk = Chunk::new(vec![
5708 Vector::from_values(LogicalType::BigInt, &held).expect("numbers"),
5709 ])
5710 .expect("one column");
5711 writer.append(&chunk).expect("one part");
5712 }
5713 writer.finish().expect("commit");
5714 let reader = Reader::open(&path).expect("reopen from disk");
5715 let bytes = reader.layout().columns[0].part_ranges;
5716 assert_eq!(bytes > 0, wanted, "{parts} parts wrote {bytes} bytes of ranges");
5717 fs::remove_file(path).expect("remove scratch file");
5718 }
5719 }
5720
5721 #[test]
5724 fn a_string_end_that_is_cut_down_still_covers_the_value_it_came_from() {
5725 let long = vec![b'a'; PART_BOUND_BYTES * 2];
5726 let low = shortened(Some(Bound::Bytes(long.clone())), false).expect("a low end");
5727 let high = shortened(Some(Bound::Bytes(long.clone())), true).expect("a high end");
5728 let Bound::Bytes(low) = low else { panic!("a string stays a string") };
5729 let Bound::Bytes(high) = high else { panic!("a string stays a string") };
5730 assert!(low.len() <= PART_BOUND_BYTES && high.len() <= PART_BOUND_BYTES);
5731 assert!(low.as_slice() <= long.as_slice(), "the low end is at or under the value");
5732 assert!(high.as_slice() >= long.as_slice(), "the high end is at or over the value");
5733 }
5734
5735 #[test]
5738 fn a_string_end_with_no_room_to_step_up_gives_up_the_bound() {
5739 let long = vec![u8::MAX; PART_BOUND_BYTES * 2];
5740 assert_eq!(shortened(Some(Bound::Bytes(long.clone())), true), None);
5741 let low = shortened(Some(Bound::Bytes(long)), false).expect("a low end is still a prefix");
5742 assert_eq!(low, Bound::Bytes(vec![u8::MAX; PART_BOUND_BYTES]));
5743 }
5744
5745 #[test]
5755 fn a_sieve_larger_than_the_part_it_indexes_is_not_written() {
5756 let path = path("sieve-pays");
5757 let fields = vec![
5758 Field::required("spread", LogicalType::BigInt),
5759 Field::required("repeated", LogicalType::BigInt),
5760 ];
5761 let mut writer = Writer::create(&path, "hits", fields).expect("new file");
5762 let parts = 3;
5763 let per_part = 1024;
5764 for part in 0..parts {
5765 let base = (part * per_part) as i64;
5766 let spread: Vec<Value> =
5767 (0..per_part).map(|row| Value::BigInt(scattered(base + row as i64))).collect();
5768 let repeated: Vec<Value> =
5769 (0..per_part).map(|row| Value::BigInt(scattered((row / 256) as i64))).collect();
5770 let chunk = Chunk::new(vec![
5771 Vector::from_values(LogicalType::BigInt, &spread).expect("numbers"),
5772 Vector::from_values(LogicalType::BigInt, &repeated).expect("numbers"),
5773 ])
5774 .expect("two columns");
5775 writer.append(&chunk).expect("one part");
5776 }
5777 writer.finish().expect("commit");
5778
5779 let reader = Reader::open(&path).expect("reopen from disk");
5780 let layout = reader.layout();
5781 let spread = &layout.columns[0];
5782 let repeated = &layout.columns[1];
5783 assert!(spread.sieves > 0, "a column whose parts are worth a filter keeps one");
5784 assert_eq!(
5785 repeated.sieves, 0,
5786 "a column whose filter costs more than its parts keeps none"
5787 );
5788 for column in &layout.columns {
5791 assert!(
5792 column.sieves < column.pages,
5793 "{} spends {} on sieves over {} of data",
5794 column.name,
5795 column.sieves,
5796 column.pages
5797 );
5798 }
5799 let absent = [Probe {
5801 column: 0,
5802 op: Op::Equal,
5803 value: Bound::Int(i128::from(scattered((parts * per_part) as i64 + 1))),
5804 }];
5805 assert!((0..parts).all(|part| reader.skips(part, &absent)), "no part holds it");
5806 fs::remove_file(path).expect("remove scratch file");
5807 }
5808
5809 #[test]
5815 fn a_damaged_sieve_page_is_read_through_rather_than_refused() {
5816 let path = path("sieve-damaged");
5817 let mut writer =
5818 Writer::create(&path, "hits", vec![Field::required("id", LogicalType::BigInt)])
5819 .expect("new file");
5820 let rows = 128;
5821 let held: Vec<Value> = (0..rows).map(|row| Value::BigInt(scattered(row))).collect();
5822 let chunk =
5823 Chunk::new(vec![Vector::from_values(LogicalType::BigInt, &held).expect("numbers")])
5824 .expect("one column");
5825 writer.append(&chunk).expect("one part");
5826 writer.finish().expect("commit");
5827
5828 let page =
5829 Reader::open(&path).expect("reopen").table.stripes[0].sieves[0].expect("a sieve page");
5830 let mut file = OpenOptions::new().write(true).open(&path).expect("open the sieve page");
5831 file.seek(SeekFrom::Start(page.offset + u64::from(page.length) - 1)).expect("seek");
5832 file.write_all(&[0xff]).expect("damage one byte");
5833 drop(file);
5834
5835 let reader = Reader::open(&path).expect("reopen the damaged file");
5836 let absent =
5837 [Probe { column: 0, op: Op::Equal, value: Bound::Int(i128::from(scattered(99))) }];
5838 assert!(!reader.skips(0, &absent), "a sieve that cannot be read skips nothing");
5839 assert_eq!(
5840 reader.read(0, &[0]).expect("the rows are untouched").len(),
5841 usize::try_from(rows).expect("a small count")
5842 );
5843 fs::remove_file(path).expect("remove scratch file");
5844 }
5845
5846 #[test]
5857 fn workers_that_want_the_same_stripe_read_it_once() {
5858 let path = path("single-flight");
5859 let mut writer =
5860 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
5861 .expect("new file");
5862 for part in 0..STRIPE_PARTS {
5863 let id = part as i32;
5864 let chunk = Chunk::new(vec![
5865 Vector::from_values(
5866 LogicalType::Integer,
5867 &[Value::Integer(id), Value::Integer(-id)],
5868 )
5869 .expect("integers"),
5870 ])
5871 .expect("matching rows");
5872 writer.append(&chunk).expect("one part");
5873 }
5874 writer.finish().expect("commit");
5875
5876 let reader = Reader::open(&path).expect("reopen from disk");
5877 assert_eq!(reader.table().stripes().len(), 1, "one stripe is the point of the test");
5878 let barrier = std::sync::Barrier::new(8);
5879 std::thread::scope(|scope| {
5880 for worker in 0..8 {
5881 let reader = &reader;
5882 let barrier = &barrier;
5883 scope.spawn(move || {
5884 barrier.wait();
5885 for part in (worker..STRIPE_PARTS).step_by(8) {
5886 let chunk = reader.read(part, &[0]).expect("a whole page read");
5887 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
5888 assert_eq!(chunk.value_at(1, 0), Value::Integer(-(part as i32)));
5889 }
5890 });
5891 }
5892 });
5893 assert_eq!(reader.pages.load(Atomic::Relaxed), 1, "one stripe, one page read, whoever won");
5894 fs::remove_file(path).expect("remove scratch file");
5895 }
5896
5897 #[test]
5910 fn opening_costs_the_same_over_a_thousand_times_the_rows() {
5911 let opened = |label: &str, rows_per_part: i32| {
5912 let path = path(label);
5913 let mut writer =
5914 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
5915 .expect("new file");
5916 for part in 0..STRIPE_PARTS * 3 {
5917 let values = (0..rows_per_part)
5921 .map(|row| {
5922 Value::Integer((part as i32 * rows_per_part + row).wrapping_mul(2_654_435))
5923 })
5924 .collect::<Vec<_>>();
5925 let chunk = Chunk::new(vec![
5926 Vector::from_values(LogicalType::Integer, &values).expect("integers"),
5927 ])
5928 .expect("matching rows");
5929 writer.append(&chunk).expect("one part");
5930 }
5931 writer.finish().expect("commit");
5932 let reader = Reader::open(&path).expect("reopen from disk");
5933 let size = fs::metadata(&path).expect("the file is there").len();
5934 let out = (reader.reads(), reader.table().stripes().len(), size);
5935 fs::remove_file(path).expect("remove scratch file");
5936 out
5937 };
5938
5939 let (thin, thin_stripes, thin_size) = opened("open-thin", 1);
5940 let (fat, fat_stripes, fat_size) = opened("open-fat", 1000);
5941 assert_eq!(
5942 thin_stripes, fat_stripes,
5943 "the same stripe count is what makes this a fair ask"
5944 );
5945 assert!(
5946 fat_size > thin_size * 50,
5947 "the fat file has to actually be larger, and it is {fat_size} against {thin_size}"
5948 );
5949
5950 assert_eq!(thin.opening.reads, fat.opening.reads, "the same reads either way");
5951 assert_eq!(thin.pages, 0, "opening read a page");
5952 assert_eq!(fat.pages, 0, "opening read a page");
5953 assert_eq!(thin.indexes, 0, "opening read an index");
5954 assert_eq!(fat.indexes, 0, "opening read an index");
5955 assert!(
5958 fat.opening.bytes < thin.opening.bytes * 2,
5959 "opening the thin file read {} bytes and the fat one read {}",
5960 thin.opening.bytes,
5961 fat.opening.bytes
5962 );
5963 }
5964
5965 #[test]
5973 fn two_opens_of_one_file_cost_the_same_and_the_second_is_not_cheaper() {
5974 let path = path("open-twice");
5975 let mut writer =
5976 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
5977 .expect("new file");
5978 for part in 0..STRIPE_PARTS * 3 {
5979 let chunk = Chunk::new(vec![
5980 Vector::from_values(LogicalType::Integer, &[Value::Integer(part as i32)])
5981 .expect("integers"),
5982 ])
5983 .expect("matching rows");
5984 writer.append(&chunk).expect("one part");
5985 }
5986 writer.finish().expect("commit");
5987
5988 let first = Reader::open(&path).expect("open");
5989 for part in 0..first.parts() {
5992 first.read(part, &[0]).expect("a part");
5993 }
5994 assert!(first.reads().pages > 0, "the scan has to have read something");
5995 let second = Reader::open(&path).expect("open again");
5996
5997 assert_eq!(first.reads().opening, second.reads().opening);
5998 assert_eq!(
5999 second.reads().pages,
6000 0,
6001 "the second open read a page off the back of the first"
6002 );
6003 assert_eq!(second.reads().indexes, 0, "the second open read an index it inherited");
6004 fs::remove_file(path).expect("remove scratch file");
6005 }
6006
6007 #[test]
6015 fn an_index_is_read_once_per_stripe_however_often_the_page_is_evicted() {
6016 let path = path("index-cache");
6017 let mut writer =
6018 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6019 .expect("new file");
6020 let parts = STRIPE_PARTS * (CACHED_STRIPES_PER_COLUMN + 2);
6021 for part in 0..parts {
6022 let id = part as i32;
6023 let chunk = Chunk::new(vec![
6024 Vector::from_values(LogicalType::Integer, &[Value::Integer(id)]).expect("integers"),
6025 ])
6026 .expect("matching rows");
6027 writer.append(&chunk).expect("one part");
6028 }
6029 writer.finish().expect("commit");
6030
6031 let reader = Reader::open(&path).expect("reopen from disk");
6032 let stripes = reader.table().stripes().len();
6033 assert!(stripes > CACHED_STRIPES_PER_COLUMN, "the page cache has to be too small for this");
6034 for _ in 0..2 {
6036 for part in 0..parts {
6037 let chunk = reader.read(part, &[0]).expect("a part");
6038 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6039 }
6040 }
6041 assert_eq!(reader.indexes.load(Atomic::Relaxed), stripes, "one index read per stripe");
6042 assert!(
6043 reader.pages.load(Atomic::Relaxed) > stripes,
6044 "the pages are the ones that get read again, which is what makes the index count mean \
6045 something"
6046 );
6047 fs::remove_file(path).expect("remove scratch file");
6048 }
6049
6050 #[test]
6059 fn a_worker_per_stripe_reads_its_page_once_when_the_cache_was_told_to_expect_it() {
6060 let workers = CACHED_STRIPES_PER_COLUMN + 4;
6061 let path = path("stripe-per-worker");
6062 let mut writer =
6063 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6064 .expect("new file");
6065 for part in 0..STRIPE_PARTS * workers {
6066 let chunk = Chunk::new(vec![
6067 Vector::from_values(LogicalType::Integer, &[Value::Integer(part as i32)])
6068 .expect("integers"),
6069 ])
6070 .expect("matching rows");
6071 writer.append(&chunk).expect("one part");
6072 }
6073 writer.finish().expect("commit");
6074
6075 let read = |told: bool| {
6076 let reader = Reader::open(&path).expect("reopen from disk");
6077 assert_eq!(reader.table().stripes().len(), workers, "a stripe per worker");
6078 if told {
6079 reader.keep_stripes(workers);
6080 }
6081 let barrier = std::sync::Barrier::new(workers);
6082 std::thread::scope(|scope| {
6083 for (worker, run) in reader.stripe_parts().into_iter().enumerate() {
6084 let reader = &reader;
6085 let barrier = &barrier;
6086 scope.spawn(move || {
6087 for part in run {
6088 barrier.wait();
6089 let chunk = reader.read(part, &[0]).expect("a part of my own stripe");
6090 assert_eq!(chunk.value_at(0, 0), Value::Integer(part as i32));
6091 }
6092 assert!(worker < workers);
6093 });
6094 }
6095 });
6096 reader.pages.load(Atomic::Relaxed)
6097 };
6098
6099 assert_eq!(read(true), workers, "one page read per stripe and no more");
6100 assert!(read(false) > workers, "a cache that small is read again on every part");
6101 fs::remove_file(path).expect("remove scratch file");
6102 }
6103
6104 #[test]
6109 fn a_damaged_index_page_is_an_error() {
6110 let path = path("damaged-index");
6111 let mut writer =
6112 Writer::create(&path, "items", vec![Field::required("id", LogicalType::Integer)])
6113 .expect("new file");
6114 writer.append(&sample_ids()).expect("first part");
6115 writer.append(&sample_ids()).expect("second part");
6116 writer.finish().expect("commit");
6117
6118 let reader = Reader::open(&path).expect("valid directory");
6119 let index = reader.table.stripes[0].index;
6120 let mut byte = [0; 1];
6121 read_at(&reader.file, index.offset, &mut byte).expect("the first part length");
6122 let mut file = OpenOptions::new().write(true).open(&path).expect("open index page");
6123 file.seek(SeekFrom::Start(index.offset)).expect("index start");
6124 file.write_all(&[!byte[0]]).expect("damage the first part length");
6125 let error = reader.read(1, &[0]).expect_err("a damaged index must not be used");
6126 assert!(error.message().contains("index page section checksum differs"), "{error}");
6127 fs::remove_file(path).expect("remove scratch file");
6128 }
6129
6130 #[test]
6137 fn every_integer_width_round_trips_through_a_page() {
6138 let path = path("integer-widths");
6139 let columns = [
6140 (LogicalType::TinyInt, vec![Value::TinyInt(i8::MIN), Value::TinyInt(i8::MAX)]),
6141 (LogicalType::UTinyInt, vec![Value::UTinyInt(0), Value::UTinyInt(u8::MAX)]),
6142 (LogicalType::SmallInt, vec![Value::SmallInt(i16::MIN), Value::SmallInt(i16::MAX)]),
6143 (LogicalType::USmallInt, vec![Value::USmallInt(0), Value::USmallInt(u16::MAX)]),
6144 (LogicalType::Integer, vec![Value::Integer(i32::MIN), Value::Integer(i32::MAX)]),
6145 (LogicalType::UInteger, vec![Value::UInteger(0), Value::UInteger(u32::MAX)]),
6146 (LogicalType::BigInt, vec![Value::BigInt(i64::MIN), Value::BigInt(i64::MAX)]),
6147 (LogicalType::UBigInt, vec![Value::UBigInt(0), Value::UBigInt(u64::MAX)]),
6148 ];
6149 let fields = columns
6150 .iter()
6151 .enumerate()
6152 .map(|(at, (ty, _))| Field::required(format!("c{at}"), ty.clone()))
6153 .collect::<Vec<_>>();
6154 let vectors = columns
6155 .iter()
6156 .map(|(ty, values)| Vector::from_values(ty.clone(), values).expect("a vector"))
6157 .collect::<Vec<_>>();
6158 let mut writer = Writer::create(&path, "widths", fields).expect("new file");
6159 writer.append(&Chunk::new(vectors).expect("matching rows")).expect("one stripe");
6160 writer.finish().expect("commit");
6161
6162 let reader = Reader::open(&path).expect("reopen from disk");
6163 let wanted = (0..columns.len()).collect::<Vec<_>>();
6164 let read = reader.read(0, &wanted).expect("every column");
6165 assert_eq!(read.len(), 2);
6166 for (at, (ty, values)) in columns.iter().enumerate() {
6168 assert_eq!(read.value_at(0, at), values[0], "the low end of {ty}");
6169 assert_eq!(read.value_at(1, at), values[1], "the high end of {ty}");
6170 }
6171 fs::remove_file(path).expect("remove scratch file");
6172 }
6173
6174 #[test]
6175 fn numeric_frequency_candidates_keep_bounded_row_ordinals() {
6176 let path = path("frequency-ordinals");
6177 let mut writer =
6178 Writer::create(&path, "items", vec![Field::required("id", LogicalType::BigInt)])
6179 .expect("new file");
6180 let mut values = Vec::new();
6181 for leader in 0..10_i64 {
6182 values.extend(std::iter::repeat_n(leader, 100));
6183 }
6184 values.extend(1_000_i64..41_000);
6185 for part in values.chunks(1_024) {
6186 let vector = Vector::flat(LogicalType::BigInt, Data::Int64(part.to_vec().into()))
6187 .expect("big integers");
6188 writer.append(&Chunk::new(vec![vector]).expect("one column")).expect("one stripe");
6189 }
6190 writer.finish().expect("commit");
6191
6192 let reader = Reader::open(&path).expect("reopen from disk");
6193 let occurrences =
6194 reader.frequency_occurrences(0).expect("valid metadata").expect("bounded ordinals");
6195 assert!(occurrences.omitted_max < 100);
6196 assert!(occurrences.ordinals.len() <= FREQUENCY_ORDINALS);
6197 assert!(occurrences.ordinals.windows(2).all(|pair| pair[0] < pair[1]));
6198 assert_eq!(&occurrences.ordinals[..1_000], &(0_u64..1_000).collect::<Vec<_>>());
6199 fs::remove_file(path).expect("remove scratch file");
6200 }
6201
6202 #[test]
6208 fn a_file_from_another_format_says_which_format_it_is() {
6209 let older = path("older-format");
6210 let mut writer =
6211 Writer::create(&older, "items", vec![Field::new("id", LogicalType::Integer)])
6212 .expect("new file");
6213 let chunk = Chunk::new(vec![
6214 Vector::flat(LogicalType::Integer, Data::Int32(vec![1, 2, 3].into()))
6215 .expect("integers"),
6216 ])
6217 .expect("chunk");
6218 writer.append(&chunk).expect("page written");
6219 writer.finish().expect("commit");
6220
6221 let mut file = OpenOptions::new().write(true).open(&older).expect("open for the header");
6222 file.seek(SeekFrom::Start(8)).expect("the version follows the magic");
6223 file.write_all(&(FORMAT - 1).to_le_bytes()).expect("write an older version");
6224 drop(file);
6225 let complaint = Reader::open(&older).expect_err("an older format is refused").to_string();
6226 assert!(complaint.contains(&format!("format {}", FORMAT - 1)), "{complaint}");
6227 assert!(complaint.contains(&format!("format {FORMAT}")), "{complaint}");
6228
6229 let mut file = OpenOptions::new().write(true).open(&older).expect("open for the header");
6230 file.seek(SeekFrom::Start(0)).expect("the magic is first");
6231 file.write_all(b"NOTRUDB!").expect("write another engine's magic");
6232 drop(file);
6233 let complaint = Reader::open(&older).expect_err("a foreign file is refused").to_string();
6234 assert!(complaint.contains("magic"), "{complaint}");
6235 assert!(!complaint.contains("format"), "a version has nothing to do with it: {complaint}");
6236 fs::remove_file(older).expect("remove scratch file");
6237 }
6238
6239 #[test]
6240 fn an_unfinished_or_damaged_file_does_not_answer_with_partial_rows() {
6241 let unfinished = path("unfinished");
6242 let mut writer =
6243 Writer::create(&unfinished, "items", vec![Field::new("id", LogicalType::Integer)])
6244 .expect("new file");
6245 let chunk = Chunk::new(vec![
6246 Vector::flat(LogicalType::Integer, Data::Int32(vec![1, 2, 3].into()))
6247 .expect("integers"),
6248 ])
6249 .expect("chunk");
6250 writer.append(&chunk).expect("page written");
6251 drop(writer);
6252 assert!(Reader::open(&unfinished).is_err(), "no directory was committed");
6253 fs::remove_file(unfinished).expect("remove scratch file");
6254
6255 let damaged = path("damaged");
6256 let mut writer =
6257 Writer::create(&damaged, "items", vec![Field::new("id", LogicalType::Integer)])
6258 .expect("new file");
6259 writer.append(&chunk).expect("page written");
6260 writer.finish().expect("commit");
6261 let reader = Reader::open(&damaged).expect("valid directory");
6262 let mut file =
6263 OpenOptions::new().write(true).open(&damaged).expect("open for a damaged page");
6264 file.seek(SeekFrom::Start(HEADER + 1)).expect("inside first page");
6265 file.write_all(&[255]).expect("damage one byte");
6266 assert!(reader.read(0, &[0]).is_err(), "page checksum rejects corruption");
6267 fs::remove_file(damaged).expect("remove scratch file");
6268 }
6269
6270 #[test]
6271 fn damaged_lazy_dictionary_payload_is_an_error() {
6272 let path = path("damaged-dictionary");
6273 let mut writer = Writer::create(
6274 &path,
6275 "items",
6276 vec![
6277 Field::required("id", LogicalType::Integer),
6278 Field::new("text", LogicalType::Varchar),
6279 ],
6280 )
6281 .expect("new file");
6282 writer.append(&sample()).expect("stripe written");
6283 writer.finish().expect("commit");
6284
6285 let reader = Reader::open(&path).expect("valid directory");
6286 let dictionary = reader.table.dictionaries[1].expect("string dictionary page");
6287 let mut header = [0; DICTIONARY_HEADER];
6290 read_at(&reader.file, dictionary.offset, &mut header).expect("dictionary header");
6291 let index_len = dictionary_index_len(&header);
6292 let rank_len = last_rank_end(&reader.file, dictionary.offset, &header);
6293 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
6294 file.seek(SeekFrom::Start(dictionary.offset + index_len + rank_len))
6295 .expect("inside dictionary payload");
6296 file.write_all(&[255]).expect("damage dictionary payload");
6297
6298 let chunk = reader.read(0, &[1]).expect("code page and dictionary index remain valid");
6299 let error =
6300 chunk.validate_external().expect_err("payload corruption must reach the caller");
6301 assert!(error.message().contains("payload checksum differs"), "{error}");
6302 fs::remove_file(path).expect("remove scratch file");
6303 }
6304
6305 #[test]
6312 fn a_dictionary_over_many_blocks_checks_every_block_of_it() {
6313 let path = path("dictionary-blocks");
6314 let value =
6315 |row: usize| format!("{row:07} a value long enough to be worth a payload block");
6316 let parts = 30;
6317 let per_part = 1000;
6318 let mut writer =
6319 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
6320 .expect("new file");
6321 for part in 0..parts {
6322 let values = (0..per_part)
6323 .map(|row| Value::Varchar(value(part * per_part + row)))
6324 .collect::<Vec<_>>();
6325 let chunk = Chunk::new(vec![
6326 Vector::from_values(LogicalType::Varchar, &values).expect("strings"),
6327 ])
6328 .expect("matching rows");
6329 writer.append(&chunk).expect("a part");
6330 }
6331 writer.finish().expect("commit");
6332
6333 let reader = Reader::open(&path).expect("reopen from disk");
6334 let dictionary = reader.table.dictionaries[0].expect("string dictionary page");
6335 assert!(
6336 parts * per_part > TEXT_PAYLOAD_VALUES * 4,
6337 "the dictionary has to be several blocks for this to be testing anything"
6338 );
6339 for part in [0, parts - 1] {
6340 let chunk = reader.read(part, &[0]).expect("a part");
6341 chunk.validate_external().expect("every payload block checks out");
6342 assert_eq!(chunk.value_at(0, 0), Value::Varchar(value(part * per_part)));
6343 }
6344
6345 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
6346 file.seek(SeekFrom::Start(dictionary.offset + u64::from(dictionary.length) - 4))
6347 .expect("the last bytes of the page are payload");
6348 file.write_all(&[255]).expect("damage the last payload block");
6349 let reader = Reader::open(&path).expect("the directory and the index are untouched");
6350 let chunk = reader.read(parts - 1, &[0]).expect("the code page remains valid");
6351 let error = chunk.validate_external().expect_err("the damage must reach the caller");
6352 assert!(error.message().contains("payload checksum differs"), "{error}");
6353 fs::remove_file(path).expect("remove scratch file");
6354 }
6355
6356 #[test]
6366 fn values_of_different_lengths_read_back_out_of_packed_offsets() {
6367 let path = path("dictionary-offsets");
6368 let value = |row: usize| {
6369 if row % 511 == 3 { String::new() } else { "x".repeat(row % 97) + &format!("{row:05}") }
6370 };
6371 let rows = 5_000;
6372 let mut writer =
6373 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
6374 .expect("new file");
6375 let values = (0..rows).map(|row| Value::Varchar(value(row))).collect::<Vec<_>>();
6376 for part in values.chunks(1_000) {
6377 let chunk =
6378 Chunk::new(vec![Vector::from_values(LogicalType::Varchar, part).expect("strings")])
6379 .expect("matching rows");
6380 writer.append(&chunk).expect("a part");
6381 }
6382 writer.finish().expect("commit");
6383
6384 let reader = Reader::open(&path).expect("reopen from disk");
6385 assert!(
6386 rows > TEXT_PAYLOAD_VALUES * 4,
6387 "the dictionary has to be several blocks for this to be testing anything"
6388 );
6389 for part in 0..rows / 1_000 {
6390 let chunk = reader.read(part, &[0]).expect("a part");
6391 for row in 0..1_000 {
6392 let row = part * 1_000 + row;
6393 assert_eq!(
6394 chunk.value_at(row % 1_000, 0),
6395 Value::Varchar(value(row)),
6396 "value {row}"
6397 );
6398 }
6399 }
6400 fs::remove_file(path).expect("remove scratch file");
6401 }
6402
6403 #[test]
6415 fn a_global_dictionary_is_opened_once_however_many_workers_ask_at_once() {
6416 let path = path("dictionary-once");
6417 let parts = 8;
6418 let per_part = 500;
6419 let value =
6420 |row: usize| format!("{row:07} a value long enough to be worth a payload block");
6421 let mut writer =
6422 Writer::create(&path, "items", vec![Field::required("text", LogicalType::Varchar)])
6423 .expect("new file");
6424 for part in 0..parts {
6425 let values = (0..per_part)
6426 .map(|row| Value::Varchar(value(part * per_part + row)))
6427 .collect::<Vec<_>>();
6428 let chunk = Chunk::new(vec![
6429 Vector::from_values(LogicalType::Varchar, &values).expect("strings"),
6430 ])
6431 .expect("matching rows");
6432 writer.append(&chunk).expect("a part");
6433 }
6434 writer.finish().expect("commit");
6435
6436 let reader = Reader::open(&path).expect("reopen from disk");
6437 assert!(reader.table.dictionaries[0].is_some(), "the column has to have one to share");
6438 assert_eq!(reader.reads().dictionaries, 0, "opening the file does not open a dictionary");
6439
6440 let workers = 16;
6441 let gate = std::sync::Barrier::new(workers);
6442 std::thread::scope(|scope| {
6443 for worker in 0..workers {
6444 let reader = reader.clone();
6445 let gate = &gate;
6446 scope.spawn(move || {
6447 gate.wait();
6448 let chunk = reader.read(worker % parts, &[0]).expect("a part");
6449 assert_eq!(
6450 chunk.value_at(0, 0),
6451 Value::Varchar(value((worker % parts) * per_part))
6452 );
6453 });
6454 }
6455 });
6456
6457 assert_eq!(reader.reads().dictionaries, 1, "sixteen workers, one dictionary, one open");
6458 fs::remove_file(path).expect("remove scratch file");
6459 }
6460
6461 #[test]
6466 fn a_damaged_sorted_order_is_an_error() {
6467 let path = path("damaged-order");
6468 let mut writer = Writer::create(
6469 &path,
6470 "items",
6471 vec![
6472 Field::required("id", LogicalType::Integer),
6473 Field::new("text", LogicalType::Varchar),
6474 ],
6475 )
6476 .expect("new file");
6477 writer.append(&sample()).expect("stripe written");
6478 writer.finish().expect("commit");
6479
6480 let reader = Reader::open(&path).expect("valid directory");
6481 let page = reader.table.dictionaries[1].expect("string dictionary page");
6482 let mut header = [0; DICTIONARY_HEADER];
6483 read_at(&reader.file, page.offset, &mut header).expect("dictionary header");
6484 let index_len = dictionary_index_len(&header);
6485 let mut file = OpenOptions::new().write(true).open(&path).expect("open dictionary page");
6486 file.seek(SeekFrom::Start(page.offset + index_len)).expect("the first head");
6487 file.write_all(&[255]).expect("damage the order");
6488
6489 let dictionary = reader.dictionary(1).expect("read").expect("a string column has one");
6490 let error = dictionary.compare_rank(0, b"anything").expect_err("a damaged order is caught");
6491 assert!(error.message().contains("rank checksum differs"), "{error}");
6492 fs::remove_file(path).expect("remove scratch file");
6493 }
6494
6495 #[test]
6499 fn a_global_dictionary_carries_the_sorted_order_of_its_values() {
6500 let spellings = ["overlong1z", "b", "", "overlong1a", "overlong", "ab", "a", "overlong1"];
6503 let path = path("dictionary-order");
6504 let mut writer =
6505 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
6506 .expect("new file");
6507 writer
6508 .append(
6509 &Chunk::new(vec![
6510 Vector::from_values(
6511 LogicalType::Varchar,
6512 &spellings.map(|text| Value::Varchar(text.into())),
6513 )
6514 .expect("strings"),
6515 ])
6516 .expect("one column"),
6517 )
6518 .expect("stripe written");
6519 writer.finish().expect("commit");
6520
6521 let reader = Reader::open(&path).expect("valid directory");
6522 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
6523 let count = dictionary.ranks().expect("a v10 file stores one");
6524 assert_eq!(count, spellings.len(), "every distinct value has a rank");
6525 let order = (0..count)
6526 .map(|rank| dictionary.code_at_rank(rank).expect("a code"))
6527 .collect::<Vec<_>>();
6528 let mut seen = order.clone();
6529 seen.sort_unstable();
6530 assert_eq!(seen, (0..spellings.len() as u32).collect::<Vec<_>>(), "a permutation of codes");
6531
6532 let ranked = order
6533 .iter()
6534 .map(|&code| {
6535 dictionary.try_bytes_at(code as usize).expect("read").expect("a value").to_vec()
6536 })
6537 .collect::<Vec<_>>();
6538 let mut expected = spellings.map(|text| text.as_bytes().to_vec()).to_vec();
6539 expected.sort();
6540 assert_eq!(ranked, expected, "rank order is value order");
6541
6542 for (rank, value) in expected.iter().enumerate() {
6545 assert_eq!(
6546 dictionary.compare_rank(rank, value).expect("compare"),
6547 Ordering::Equal,
6548 "rank {rank} is its own value"
6549 );
6550 if rank > 0 {
6551 assert_eq!(
6552 dictionary.compare_rank(rank - 1, value).expect("compare"),
6553 Ordering::Less,
6554 "rank {rank} follows the one before it"
6555 );
6556 }
6557 }
6558 fs::remove_file(path).expect("remove scratch file");
6559 }
6560
6561 #[test]
6569 fn a_dictionary_sweep_reads_every_value_and_keeps_it_under_the_budget() {
6570 let path = path("dictionary-sweep");
6571 let spellings = (0..2_500)
6574 .map(|index| Value::Varchar(format!("value {index:08} {}", "x".repeat(index % 40))))
6575 .collect::<Vec<_>>();
6576 let mut writer =
6577 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
6578 .expect("new file");
6579 for part in spellings.chunks(1_024) {
6582 writer
6583 .append(
6584 &Chunk::new(vec![
6585 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
6586 ])
6587 .expect("one column"),
6588 )
6589 .expect("stripe written");
6590 }
6591 writer.finish().expect("commit");
6592
6593 let reader = Reader::open(&path).expect("valid directory");
6594 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
6595 assert_eq!(dictionary.len(), spellings.len(), "every value is distinct");
6596
6597 let resting = dictionary.footprint();
6598 let mut swept: Vec<Vec<u8>> = Vec::new();
6599 let mut at = 0;
6600 let mut calls = 0;
6601 while at < dictionary.len() {
6602 let stopped = dictionary
6603 .sweep_text(at, dictionary.len(), &mut |index: usize, text: &[u8]| {
6604 assert_eq!(index, swept.len(), "a sweep hands its values over in order");
6605 swept.push(text.to_vec());
6606 Ok(())
6607 })
6608 .expect("a sweep reads");
6609 assert!(stopped > at, "a sweep moves");
6610 at = stopped;
6611 calls += 1;
6612 }
6613 assert_eq!(calls, 3, "a sweep hands over one block at a time");
6614 let after = dictionary.footprint();
6615 assert!(after > resting, "a sweep under the budget keeps what it decoded");
6616
6617 let read = (0..dictionary.len())
6618 .map(|code| dictionary.try_bytes_at(code).expect("read").expect("a value").to_vec())
6619 .collect::<Vec<_>>();
6620 assert_eq!(swept, read, "a sweep answers what a point read answers");
6621 assert_eq!(dictionary.footprint(), after, "a point read of a kept block decodes nothing");
6622 fs::remove_file(path).expect("remove scratch file");
6623 }
6624
6625 #[test]
6636 fn a_sweep_over_a_block_with_a_short_second_run_reads_what_a_point_read_reads() {
6637 let path = path("dictionary-sweep-short-run");
6638 let spellings = (0..2_800)
6639 .map(|index| Value::Varchar(format!("value {index:08} {}", "x".repeat(index % 40))))
6640 .collect::<Vec<_>>();
6641 let mut writer =
6642 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
6643 .expect("new file");
6644 for part in spellings.chunks(1_024) {
6645 writer
6646 .append(
6647 &Chunk::new(vec![
6648 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
6649 ])
6650 .expect("one column"),
6651 )
6652 .expect("stripe written");
6653 }
6654 writer.finish().expect("commit");
6655
6656 let reader = Reader::open(&path).expect("valid directory");
6657 let dictionary = reader.dictionary(0).expect("read").expect("a string column has one");
6658 assert_eq!(dictionary.len(), spellings.len(), "every value is distinct");
6659 let last = dictionary.len() % TEXT_PAYLOAD_VALUES;
6660 assert!(last > TEXT_OFFSET_RUN, "the last block has to reach into a second run of offsets");
6661 assert!(last < TEXT_PAYLOAD_VALUES, "and that second run has to be short of a whole one");
6662
6663 let mut swept: Vec<Vec<u8>> = Vec::new();
6664 let mut at = 0;
6665 while at < dictionary.len() {
6666 let stopped = dictionary
6667 .sweep_text(at, dictionary.len(), &mut |index: usize, text: &[u8]| {
6668 assert_eq!(index, swept.len(), "a sweep hands its values over in order");
6669 swept.push(text.to_vec());
6670 Ok(())
6671 })
6672 .expect("a sweep reads");
6673 assert!(stopped > at, "a sweep moves");
6674 at = stopped;
6675 }
6676 let read = (0..dictionary.len())
6677 .map(|code| dictionary.try_bytes_at(code).expect("read").expect("a value").to_vec())
6678 .collect::<Vec<_>>();
6679 assert_eq!(swept, read, "a sweep answers what a point read answers");
6680 fs::remove_file(path).expect("remove scratch file");
6681 }
6682
6683 #[test]
6693 fn narrowing_a_page_takes_what_fits_and_refuses_what_does_not() {
6694 assert_eq!(fit::<i8>(&[]).expect("an empty page fits anything"), Vec::<i8>::new());
6695 assert_eq!(fit::<i8>(&[-128, 0, 127]).expect("the edges fit"), vec![-128_i8, 0, 127]);
6696 fit::<i8>(&[128]).expect_err("one past the top does not fit");
6697 fit::<i8>(&[-129]).expect_err("one past the bottom does not fit");
6698 assert_eq!(fit::<u8>(&[0, 255]).expect("the edges fit"), vec![0_u8, 255]);
6699 fit::<u8>(&[256]).expect_err("one past the top does not fit");
6700 fit::<u8>(&[-1]).expect_err("a negative does not fit an unsigned page");
6701 assert_eq!(
6702 fit::<i16>(&[-32_768, 0, 32_767]).expect("the edges fit"),
6703 vec![-32_768_i16, 0, 32_767]
6704 );
6705 fit::<i16>(&[32_768]).expect_err("one past the top does not fit");
6706 fit::<i16>(&[-32_769]).expect_err("one past the bottom does not fit");
6707 assert_eq!(fit::<u16>(&[0, 65_535]).expect("the edges fit"), vec![0_u16, 65_535]);
6708 fit::<u16>(&[65_536]).expect_err("one past the top does not fit");
6709 fit::<u16>(&[-1]).expect_err("a negative does not fit an unsigned page");
6710 assert_eq!(
6711 fit::<i32>(&[i64::from(i32::MIN), 0, i64::from(i32::MAX)]).expect("the edges fit"),
6712 vec![i32::MIN, 0, i32::MAX]
6713 );
6714 fit::<i32>(&[i64::from(i32::MAX) + 1]).expect_err("one past the top does not fit");
6715 fit::<i32>(&[i64::from(i32::MIN) - 1]).expect_err("one past the bottom does not fit");
6716 assert_eq!(
6717 fit::<u32>(&[0, 4_294_967_295]).expect("the edges fit"),
6718 vec![0_u32, 4_294_967_295]
6719 );
6720 fit::<u32>(&[4_294_967_296]).expect_err("one past the top does not fit");
6721 fit::<u32>(&[-1]).expect_err("a negative does not fit an unsigned page");
6722
6723 fit::<i8>(&[0, 1, 2, 128, 3]).expect_err("one bad value spoils the page");
6726 }
6727
6728 #[test]
6735 fn the_residue_agrees_with_a_checked_conversion_everywhere() {
6736 for value in -70_000_i64..70_000 {
6737 assert_eq!(fit::<i8>(&[value]).is_ok(), i8::try_from(value).is_ok(), "{value} as i8");
6738 assert_eq!(fit::<u8>(&[value]).is_ok(), u8::try_from(value).is_ok(), "{value} as u8");
6739 assert_eq!(fit::<i16>(&[value]).is_ok(), i16::try_from(value).is_ok(), "{value} i16");
6740 assert_eq!(fit::<u16>(&[value]).is_ok(), u16::try_from(value).is_ok(), "{value} u16");
6741 }
6742 let wide = [i64::MIN, i64::MIN + 1, i64::from(i32::MIN), 0, i64::from(u32::MAX), i64::MAX];
6743 for edge in wide {
6744 for step in -2_i64..=2 {
6745 let value = edge.saturating_add(step);
6746 assert_eq!(
6747 fit::<i32>(&[value]).is_ok(),
6748 i32::try_from(value).is_ok(),
6749 "{value} as i32"
6750 );
6751 assert_eq!(
6752 fit::<u32>(&[value]).is_ok(),
6753 u32::try_from(value).is_ok(),
6754 "{value} as u32"
6755 );
6756 }
6757 }
6758 }
6759
6760 #[test]
6768 fn a_dictionary_at_its_budget_sweeps_without_keeping() {
6769 let path = path("dictionary-budget");
6770 let spellings = (0..2_500)
6771 .map(|index| Value::Varchar(format!("value {index:08} {}", "y".repeat(index % 40))))
6772 .collect::<Vec<_>>();
6773 let mut writer =
6774 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
6775 .expect("new file");
6776 for part in spellings.chunks(1_024) {
6777 writer
6778 .append(
6779 &Chunk::new(vec![
6780 Vector::from_values(LogicalType::Varchar, part).expect("strings"),
6781 ])
6782 .expect("one column"),
6783 )
6784 .expect("stripe written");
6785 }
6786 writer.finish().expect("commit");
6787
6788 let reader = Reader::open(&path).expect("valid directory");
6789 let page = reader.table.dictionaries[0].expect("a string column has one");
6790 let file = Arc::clone(&reader.file);
6791 let starved = open_global_dictionary(file, page, &LogicalType::Varchar, 0)
6792 .expect("a dictionary opens whatever it may keep");
6793
6794 let resting = starved.footprint();
6795 let mut swept: Vec<Vec<u8>> = Vec::new();
6796 let mut at = 0;
6797 while at < starved.len() {
6798 at = starved
6799 .sweep_text(at, starved.len(), &mut |_index: usize, text: &[u8]| {
6800 swept.push(text.to_vec());
6801 Ok(())
6802 })
6803 .expect("a sweep reads");
6804 }
6805 assert_eq!(swept.len(), spellings.len(), "a starved sweep still reads every value");
6806 assert_eq!(starved.footprint(), resting, "and keeps no block it decoded");
6807
6808 let generous = reader.dictionary(0).expect("read").expect("a string column has one");
6809 let read = (0..generous.len())
6810 .map(|code| generous.try_bytes_at(code).expect("read").expect("a value").to_vec())
6811 .collect::<Vec<_>>();
6812 assert_eq!(swept, read, "a starved sweep answers what a point read answers");
6813 fs::remove_file(path).expect("remove scratch file");
6814 }
6815
6816 #[test]
6817 fn damaged_membership_cannot_skip_a_string_page() {
6818 let path = path("damaged-membership");
6819 let mut writer = Writer::create(
6820 &path,
6821 "items",
6822 vec![
6823 Field::required("id", LogicalType::Integer),
6824 Field::new("text", LogicalType::Varchar),
6825 ],
6826 )
6827 .expect("new file");
6828 writer.append(&sample()).expect("stripe written");
6829 writer.finish().expect("commit");
6830
6831 let reader = Reader::open(&path).expect("valid directory");
6832 let membership = reader.table.stripes[0].memberships[1].expect("string membership");
6833 let mut file = OpenOptions::new().write(true).open(&path).expect("open membership page");
6834 file.seek(SeekFrom::Start(membership.offset)).expect("membership start");
6835 file.write_all(&[255]).expect("damage membership");
6836 let error = reader.skips_codes(0, 1, &[3]).expect_err("corruption must not skip rows");
6837 assert!(error.message().contains("membership page checksum differs"), "{error}");
6838 fs::remove_file(path).expect("remove scratch file");
6839 }
6840
6841 #[test]
6842 fn membership_delta_stream_is_sorted_exact_and_bounded() {
6843 let unique = unique_codes(&[900, 4, 4, 72, 9, u32::MAX]);
6844 assert_eq!(unique, [4, 9, 72, 900, u32::MAX]);
6845 let encoded = encode_membership(&unique);
6846 assert_eq!(
6847 decode_membership(&encoded).expect("valid membership"),
6848 [4, 9, 72, 900, u32::MAX]
6849 );
6850 let merged = merged_codes(vec![vec![4, 900], vec![9, 900, u32::MAX], vec![72]]);
6853 assert_eq!(merged, [4, 9, 72, 900, u32::MAX]);
6854 assert_eq!(
6855 decode_membership(&encode_membership(&merged)).expect("valid membership"),
6856 unique
6857 );
6858 assert!(decode_membership(&[1, 0x80]).is_err(), "a truncated varint is invalid");
6859 assert!(
6860 decode_membership(&[1, 0xff, 0xff, 0xff, 0xff, 0x10]).is_err(),
6861 "a value past u32 is invalid"
6862 );
6863 }
6864
6865 #[test]
6866 fn a_global_dictionary_may_be_larger_than_one_column_page() {
6867 let dictionary = Page {
6868 offset: HEADER,
6869 length: u32::try_from(MAX_PAGE + 1).expect("the page bound fits on disk"),
6870 hash: 0,
6871 };
6872 let table = Table {
6873 name: "items".to_owned(),
6874 fields: vec![Field::new("text", LogicalType::Varchar)],
6875 stripes: Vec::new(),
6876 rows: 0,
6877 dictionaries: vec![Some(dictionary)],
6878 distincts: vec![None],
6879 frequencies: vec![None],
6880 };
6881 let directory = encode_directory(&table).expect("directory");
6882 let file_size = dictionary.offset + u64::from(dictionary.length) + 1;
6883
6884 let decoded = decode_directory(&directory, file_size).expect("large lazy dictionary");
6885 assert_eq!(decoded.dictionaries[0].expect("dictionary").length, dictionary.length);
6886 }
6887
6888 #[test]
6889 fn a_column_with_one_value_everywhere_costs_almost_nothing_a_row() {
6890 let path = path("constant-codes");
6891 let mut writer =
6892 Writer::create(&path, "items", vec![Field::new("text", LogicalType::Varchar)])
6893 .expect("new file");
6894 let empty = vec![Value::Varchar(String::new()); 1024];
6895 for _ in 0..4 {
6896 let column = Vector::from_values(LogicalType::Varchar, &empty).expect("strings");
6897 writer.append(&Chunk::new(vec![column]).expect("one column")).expect("a part");
6898 }
6899 writer.finish().expect("commit");
6900
6901 let reader = Reader::open(&path).expect("valid directory");
6902 let pages = reader.layout().columns.first().expect("one column").pages;
6903 assert!(pages < 256, "{pages} bytes of pages for 4,096 rows of one value");
6907 let read = reader.read(3, &[0]).expect("the last part back");
6908 assert_eq!(read.value_at(0, 0), Value::Varchar(String::new()));
6909 assert_eq!(read.value_at(1023, 0), Value::Varchar(String::new()));
6910 fs::remove_file(path).expect("remove scratch file");
6911 }
6912
6913 #[test]
6914 fn a_cascade_value_too_wide_for_its_column_is_refused_rather_than_cut() {
6915 let over = vec![i64::from(i32::MAX) + 1];
6918 let error = narrowed(&LogicalType::Integer, over).expect_err("a page that disagrees");
6919 assert!(format!("{error}").contains("not of its type"), "{error}");
6920 assert!(narrowed(&LogicalType::BigInt, vec![i64::MIN]).is_ok(), "bigint holds all of i64");
6921 assert!(narrowed(&LogicalType::Varchar, vec![0]).is_err(), "strings are not integers");
6922 }
6923
6924 #[test]
6925 fn a_code_stream_the_cascade_cannot_shrink_is_left_alone() {
6926 let mut state: u32 = 0x9e37_79b9;
6930 let spread: Vec<u32> = (0..1024)
6931 .map(|_| {
6932 state ^= state << 13;
6933 state ^= state >> 17;
6934 state ^= state << 5;
6935 state
6936 })
6937 .collect();
6938 assert_eq!(encoded_codes(&spread).expect("no failure"), None);
6939 let near: Vec<u32> = (0..1024).collect();
6940 let coded = encoded_codes(&near).expect("no failure").expect("counting up is packable");
6941 assert!(coded.len() < near.len() * 4, "{} bytes for a run of 1,024", coded.len());
6942 }
6943
6944 #[test]
6950 fn two_writes_of_the_same_rows_give_the_same_bytes() {
6951 fn written(path: &PathBuf) {
6952 let fields = (0..40)
6953 .map(|column| {
6954 let ty =
6955 if column % 4 == 0 { LogicalType::Varchar } else { LogicalType::BigInt };
6956 Field::new(format!("c{column}"), ty)
6957 })
6958 .collect::<Vec<_>>();
6959 let mut writer = Writer::create(path, "wide", fields).expect("new file");
6960 for part in 0..70_u64 {
6961 let columns = (0..40)
6962 .map(|column| {
6963 let values = (0..64_u64)
6964 .map(|row| {
6965 let seed = part.wrapping_mul(31).wrapping_add(row);
6966 if column % 4 == 0 {
6967 Value::Varchar(format!("v{}", seed % 17))
6968 } else {
6969 Value::BigInt(i64::try_from(seed % 97).expect("small"))
6970 }
6971 })
6972 .collect::<Vec<_>>();
6973 let ty = if column % 4 == 0 {
6974 LogicalType::Varchar
6975 } else {
6976 LogicalType::BigInt
6977 };
6978 Vector::from_values(ty, &values).expect("a column")
6979 })
6980 .collect::<Vec<_>>();
6981 writer.append(&Chunk::new(columns).expect("forty columns")).expect("a part");
6982 }
6983 writer.finish().expect("commit");
6984 }
6985
6986 let first = path("repeatable-one");
6987 let second = path("repeatable-two");
6988 written(&first);
6989 written(&second);
6990 let left = fs::read(&first).expect("the first file");
6991 let right = fs::read(&second).expect("the second file");
6992 assert_eq!(left.len(), right.len(), "two writes of the same rows differ in length");
6993 assert!(left == right, "two writes of the same rows differ in their bytes");
6994
6995 let reader = Reader::open(&first).expect("valid directory");
6998 assert_eq!(reader.table().rows(), 70 * 64);
6999 let read = reader.read(0, &[0, 1]).expect("the first part back");
7000 assert_eq!(read.value_at(0, 0), Value::Varchar("v0".to_owned()));
7001 assert_eq!(read.value_at(0, 1), Value::BigInt(0));
7002 fs::remove_file(first).expect("remove scratch file");
7003 fs::remove_file(second).expect("remove scratch file");
7004 }
7005}