1use std::collections::HashMap;
45use std::ops::Deref;
46use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering as Atomic};
47use std::sync::{Arc, Mutex};
48
49use rudb_common::{Error, LogicalType, Result};
50use rudb_metrics::{LoadProfile, Stage};
51use rudb_storage::Range;
52use rudb_vector::{Bitmap, Chunk, Data, StringColumn, Validity, Vector};
53
54use super::{
55 ColumnStripe, DICTIONARY_CHECK_SEED, DICTIONARY_DECIDE_ROWS, DICTIONARY_DISTINCT_IN_TEN,
56 EncodedBlock, GlobalDictionary, MAX_ENCODE_WORKERS, MAX_PAGE, Part, PendingChunk, STRIPE_PARTS,
57 Settling, Spread, Unencoded, Writer, checksum, coded_page, invalid, push_validity,
58 seeded_checksum, stats, unique_codes, weight,
59};
60
61static BUSY: AtomicUsize = AtomicUsize::new(0);
68
69pub const DICTIONARY_CAP_BYTES: u64 = 512 * 1024 * 1024;
77
78#[derive(Debug)]
85pub(crate) struct Coding {
86 flags: Box<[AtomicBool]>,
87 growth: Box<[AtomicU64]>,
89 held: AtomicU64,
91 cap: AtomicU64,
92 ceilings: Box<[AtomicU64]>,
101}
102
103impl Coding {
104 pub(crate) fn new(flags: impl IntoIterator<Item = bool>) -> Self {
105 let flags = flags.into_iter().map(AtomicBool::new).collect::<Box<[_]>>();
106 let growth = flags.iter().map(|_| AtomicU64::new(0)).collect();
107 let ceilings = flags.iter().map(|_| AtomicU64::new(u64::MAX)).collect();
108 Self {
109 flags,
110 growth,
111 held: AtomicU64::new(0),
112 cap: AtomicU64::new(DICTIONARY_CAP_BYTES),
113 ceilings,
114 }
115 }
116
117 pub(crate) fn cap(&self, bytes: u64) {
119 self.cap.store(bytes, Atomic::Relaxed);
120 }
121
122 fn recount(&self, before: u64, now: u64) -> u64 {
125 if now >= before {
126 self.held.fetch_add(now - before, Atomic::Relaxed) + (now - before)
127 } else {
128 self.held.fetch_sub(before - now, Atomic::Relaxed).saturating_sub(before - now)
129 }
130 }
131
132 fn grew_most(&self, index: usize) -> bool {
139 let mine = self.growth[index].load(Atomic::Relaxed);
140 mine > 0 && self.growth.iter().all(|other| other.load(Atomic::Relaxed) <= mine)
141 }
142}
143
144impl Deref for Coding {
145 type Target = [AtomicBool];
146
147 fn deref(&self) -> &[AtomicBool] {
148 &self.flags
149 }
150}
151
152struct Share(usize);
154
155impl Share {
156 fn take(columns: usize, parts: usize) -> Self {
157 let busy = BUSY.fetch_add(1, Atomic::Relaxed) + 1;
158 let cores =
159 std::thread::available_parallelism().map_or(1, usize::from).min(MAX_ENCODE_WORKERS);
160 let workers = if parts <= 1 { 1 } else { (cores / busy).clamp(1, columns.max(1)) };
162 Self(workers)
163 }
164}
165
166impl Drop for Share {
167 fn drop(&mut self) {
168 BUSY.fetch_sub(1, Atomic::Relaxed);
169 }
170}
171
172#[derive(Debug, Clone)]
178pub struct Preparer {
179 types: Vec<LogicalType>,
180 coded: Arc<Coding>,
181 profile: Option<Arc<LoadProfile>>,
182}
183
184#[derive(Debug)]
186pub struct Prepared {
187 parts: Vec<Part>,
188 types: Vec<LogicalType>,
189 columns: Vec<Column>,
190 gathers: Vec<Option<stats::Gather>>,
191 profile: Option<Arc<LoadProfile>>,
192}
193
194#[derive(Debug)]
196pub struct Merged {
197 parts: Vec<Part>,
198 columns: Vec<Merge>,
199 blocks: Vec<Unencoded>,
200 profile: Option<Arc<LoadProfile>>,
201 counted: bool,
204}
205
206#[derive(Debug)]
208pub struct Paged {
209 parts: Vec<Part>,
210 columns: Vec<ColumnStripe>,
211 blocks: Vec<(usize, usize, EncodedBlock)>,
213 counted: bool,
214}
215
216enum Built {
218 Stripe(ColumnStripe),
219 Block(EncodedBlock),
220}
221
222#[derive(Debug)]
224enum Column {
225 Pages(ColumnStripe),
227 Coded(Local),
229}
230
231#[derive(Debug)]
233enum Merge {
234 Pages(ColumnStripe),
235 Codes {
237 parts: Vec<LocalPart>,
238 global: Vec<u32>,
239 },
240 Plain(Local),
244}
245
246const END: u32 = u32::MAX;
248
249#[derive(Debug, Default)]
255struct Local {
256 first: HashMap<u64, u32, Spread>,
258 next: Vec<u32>,
260 hashes: Vec<u64>,
261 checks: Vec<u64>,
262 bytes: Vec<u8>,
264 ends: Vec<usize>,
265 counts: Vec<u64>,
267 nulls: u64,
268 parts: Vec<LocalPart>,
269 blob: bool,
271}
272
273#[derive(Debug)]
275struct LocalPart {
276 codes: Vec<u32>,
277 validity: Vec<u8>,
279 range: Range,
280}
281
282impl Local {
283 #[cfg(test)]
285 fn code_column(index: usize, held: &[PendingChunk]) -> Result<Self> {
286 let mut local = Self::default();
287 let mut mapped = None;
288 for pending in held {
289 local.code_part(pending.chunk.column(index)?, &mut mapped)?;
290 }
291 local.done();
292 Ok(local)
293 }
294
295 fn code_part(
304 &mut self,
305 column: &Vector,
306 mapped: &mut Option<(Arc<Vector>, Vec<u32>)>,
307 ) -> Result<()> {
308 self.blob = column.logical_type() == &LogicalType::Blob;
309 if let Some(codes) = self.code_dictionary(column, mapped)? {
310 let mut validity = Vec::new();
311 push_validity(&mut validity, column);
312 self.parts.push(LocalPart { codes, validity, range: Range::of(column) });
313 return Ok(());
314 }
315 let flat = column.flatten()?;
317 let mut codes = Vec::with_capacity(flat.len());
318 let mut last = None;
319 for row in 0..flat.len() {
320 let text = flat.bytes_at(row).unwrap_or(b"");
323 let code = match last {
326 Some(code) if self.value(code) == text => code,
327 _ => self.code(text)?,
328 };
329 last = Some(code);
330 if flat.is_null_at(row) {
331 self.nulls += 1;
332 } else {
333 self.counts[code as usize] += 1;
334 }
335 codes.push(code);
336 }
337 let mut validity = Vec::new();
338 push_validity(&mut validity, &flat);
339 self.parts.push(LocalPart { codes, validity, range: Range::of(column) });
340 Ok(())
341 }
342
343 fn done(&mut self) {
345 self.first = HashMap::default();
348 self.next = Vec::new();
349 }
350
351 fn code_dictionary(
366 &mut self,
367 column: &Vector,
368 mapped: &mut Option<(Arc<Vector>, Vec<u32>)>,
369 ) -> Result<Option<Vec<u32>>> {
370 let Some((codes, values)) = column.shared_dictionary_parts() else { return Ok(None) };
371 if !matches!(values.validity(), Validity::AllValid) {
372 return Ok(None);
373 }
374 let Some(codes) = codes.get(..column.len()) else { return Ok(None) };
375 let fresh = !matches!(mapped, Some((held, _)) if Arc::ptr_eq(held, values));
376 if fresh {
377 *mapped = Some((Arc::clone(values), vec![END; values.len()]));
378 }
379 let Some((_, map)) = mapped.as_mut() else { return Ok(None) };
380 let every = matches!(column.validity(), Validity::AllValid);
381 let mut coded = Vec::with_capacity(codes.len());
382 for (row, &code) in codes.iter().enumerate() {
383 if !every && !column.validity().is_valid(row) {
384 let code = self.code(b"")?;
386 self.nulls += 1;
387 coded.push(code);
388 continue;
389 }
390 let slot = map
391 .get_mut(code as usize)
392 .ok_or_else(|| invalid("a dictionary code is out of range"))?;
393 if *slot == END {
394 *slot = self.code(values.bytes_at(code as usize).unwrap_or(b""))?;
395 }
396 self.counts[*slot as usize] += 1;
397 coded.push(*slot);
398 }
399 Ok(Some(coded))
400 }
401
402 fn rows(&self) -> Result<Vec<Vector>> {
408 self.parts
409 .iter()
410 .map(|part| {
411 let len = part.codes.len();
412 let mut column = StringColumn::with_capacity(len);
413 for &code in &part.codes {
414 column.push_bytes(self.value(code));
415 }
416 let validity = match part.validity.split_first() {
417 Some((0, _)) => Validity::AllValid,
418 Some((1, _)) => Validity::AllInvalid,
419 Some((2, bits)) => {
420 let mut mask = Bitmap::all_valid(len);
421 for row in (0..len).filter(|row| bits[row / 8] & (1 << (row % 8)) == 0) {
422 mask.set(row, false);
423 }
424 Validity::Mask(mask)
425 }
426 _ => return Err(Error::internal("a coded part has no validity")),
427 };
428 let ty = if self.blob { LogicalType::Blob } else { LogicalType::Varchar };
429 Ok(Vector::flat(ty, Data::Varlen(column))?.with_validity(validity))
430 })
431 .collect()
432 }
433
434 fn values(&self) -> usize {
435 self.ends.len()
436 }
437
438 fn value(&self, code: u32) -> &[u8] {
439 let code = code as usize;
440 let from = if code == 0 { 0 } else { self.ends[code - 1] };
441 &self.bytes[from..self.ends[code]]
442 }
443
444 fn code(&mut self, text: &[u8]) -> Result<u32> {
445 let hash = checksum(text);
446 let Some(&first) = self.first.get(&hash) else {
447 let code = self.push(text, hash)?;
448 self.first.insert(hash, code);
449 return Ok(code);
450 };
451 let mut at = first;
452 loop {
453 if self.value(at) == text {
454 return Ok(at);
455 }
456 match self.next[at as usize] {
457 END => break,
458 next => at = next,
459 }
460 }
461 let code = self.push(text, hash)?;
462 self.next[at as usize] = code;
463 Ok(code)
464 }
465
466 fn push(&mut self, text: &[u8], hash: u64) -> Result<u32> {
467 let code = u32::try_from(self.ends.len())
468 .ok()
469 .filter(|&code| code != END)
470 .ok_or_else(|| invalid("a stripe has too many values in one column"))?;
471 self.bytes.extend_from_slice(text);
472 self.ends.push(self.bytes.len());
473 self.next.push(END);
474 self.hashes.push(hash);
475 self.checks.push(seeded_checksum(text, DICTIONARY_CHECK_SEED));
476 self.counts.push(0);
477 Ok(code)
478 }
479
480 fn merge_into(&self, dictionary: &mut GlobalDictionary) -> Result<Vec<u32>> {
483 let mut global = Vec::with_capacity(self.values());
484 for (code, (&hash, &check)) in self.hashes.iter().zip(&self.checks).enumerate() {
485 let text = self.value(code as u32);
486 let at = dictionary.code_hashed(text, hash, check)?;
487 let count = dictionary
488 .counts
489 .get_mut(at as usize)
490 .ok_or_else(|| invalid("global dictionary count code is out of range"))?;
491 *count = count.saturating_add(self.counts[code]);
492 global.push(at);
493 }
494 dictionary.nulls = dictionary.nulls.saturating_add(self.nulls);
495 Ok(global)
496 }
497}
498
499fn drops_dictionary(rows: usize, distinct: usize) -> bool {
523 rows >= DICTIONARY_DECIDE_ROWS
524 && distinct.saturating_mul(10) > rows.saturating_mul(DICTIONARY_DISTINCT_IN_TEN)
525}
526
527fn demotes(rows: usize, new: usize, total: u64, coding: &Coding, index: usize) -> bool {
540 drops_dictionary(rows, new)
541 || (total > coding.cap.load(Atomic::Relaxed) && coding.grew_most(index))
542}
543
544fn fan_out<T: Send>(
554 jobs: Vec<usize>,
555 workers: usize,
556 profile: Option<&LoadProfile>,
557 work: impl Fn(usize) -> Result<T> + Sync,
558) -> Result<Vec<(usize, T)>> {
559 if workers <= 1 || jobs.len() <= 1 {
560 let _span = profile.map(|profile| profile.span(Stage::Pages));
561 return jobs.into_iter().map(|index| Ok((index, work(index)?))).collect();
562 }
563 let workers = workers.min(jobs.len());
564 let queue = Mutex::new(jobs);
565 let pieces = std::thread::scope(|scope| {
566 (0..workers)
567 .map(|_| {
568 scope.spawn(|| {
569 let _span = profile.map(|profile| profile.span(Stage::Pages));
570 let mut mine = Vec::new();
571 loop {
572 let taken = queue
573 .lock()
574 .map_err(|_| Error::internal("a native encode worker panicked"))?
575 .pop();
576 let Some(index) = taken else { break };
577 mine.push((index, work(index)?));
578 }
579 Ok(mine)
580 })
581 })
582 .collect::<Vec<_>>()
583 .into_iter()
584 .map(|handle| {
585 handle.join().map_err(|_| Error::internal("a native encode worker panicked"))?
586 })
587 .collect::<Result<Vec<Vec<_>>>>()
588 })?;
589 Ok(pieces.into_iter().flatten().collect())
590}
591
592impl Preparer {
593 pub fn prepare(&self, parts: Vec<((u64, u64), Chunk)>) -> Result<Prepared> {
604 if parts.len() > STRIPE_PARTS {
605 return Err(invalid("a stripe was handed more parts than it holds"));
606 }
607 let held = parts
608 .into_iter()
609 .filter(|(_, chunk)| !chunk.is_empty())
610 .map(|(order, chunk)| PendingChunk { order, chunk })
611 .collect::<Vec<_>>();
612 for pending in &held {
613 self.fits(&pending.chunk)?;
614 }
615 self.prepare_held(held)
616 }
617
618 fn fits(&self, chunk: &Chunk) -> Result<()> {
620 if chunk.width() != self.types.len() {
621 return Err(invalid("chunk width differs from table schema"));
622 }
623 for (index, ty) in self.types.iter().enumerate() {
624 if chunk.column(index)?.logical_type() != ty {
625 return Err(invalid("chunk type differs from table schema"));
626 }
627 }
628 Ok(())
629 }
630
631 pub(crate) fn prepare_held(&self, held: Vec<PendingChunk>) -> Result<Prepared> {
632 let mut building = self.start();
633 self.feed_held(&mut building, held)?;
634 self.finish(building)
635 }
636
637 #[must_use]
650 pub fn start(&self) -> Building {
651 let columns = (0..self.types.len())
652 .map(|index| {
653 let body = if self.coded[index].load(Atomic::Relaxed) {
654 Body::Coded(Local::default(), None)
655 } else {
656 Body::Pages(ColumnStripe::default(), Settling::default())
657 };
658 let mut gather = stats::Gather::new(&self.types[index], 0);
659 let ceiling = self.coded.ceilings[index].load(Atomic::Relaxed);
660 if let Some(gather) = gather.as_mut().filter(|_| ceiling < u64::MAX) {
661 gather.cap_at(ceiling);
662 }
663 Mutex::new(Growing { body, gather })
664 })
665 .collect();
666 Building { parts: Vec::new(), columns }
667 }
668
669 pub fn feed(&self, building: &mut Building, parts: Vec<((u64, u64), Chunk)>) -> Result<()> {
679 if building.parts.len().saturating_add(parts.len()) > STRIPE_PARTS {
680 return Err(invalid("a stripe was handed more parts than it holds"));
681 }
682 let held = parts
683 .into_iter()
684 .filter(|(_, chunk)| !chunk.is_empty())
685 .map(|(order, chunk)| PendingChunk { order, chunk })
686 .collect::<Vec<_>>();
687 for pending in &held {
688 self.fits(&pending.chunk)?;
689 }
690 self.feed_held(building, held)
691 }
692
693 fn feed_held(&self, building: &mut Building, held: Vec<PendingChunk>) -> Result<()> {
694 if held.is_empty() {
695 return Ok(());
696 }
697 let width = self.types.len();
698 let opening = building.parts.is_empty();
700 let key = held.first().map_or((0, 0), |pending| pending.order);
701 let share = Share::take(width, held.len());
702 let mut jobs = (0..width).collect::<Vec<_>>();
703 jobs.sort_by_key(|&index| weight(&self.types[index]));
704 let columns = &building.columns;
705 fan_out(jobs, share.0, self.profile.as_deref(), |index| {
706 let mut growing = columns[index]
707 .lock()
708 .map_err(|_| Error::internal("a native encode worker panicked"))?;
709 let Growing { body, gather } = &mut *growing;
710 if matches!(body, Body::Coded(..)) && !self.coded[index].load(Atomic::Relaxed) {
711 body.plain()?;
712 }
713 if let Some(gather) = gather.as_mut() {
714 if opening {
717 gather.open_stripe(key);
718 }
719 for pending in &held {
720 gather.part(pending.chunk.column(index)?);
721 }
722 }
723 match body {
724 Body::Coded(local, mapped) => {
725 for pending in &held {
726 local.code_part(pending.chunk.column(index)?, mapped)?;
727 }
728 }
729 Body::Pages(stripe, settling) => {
730 for pending in &held {
731 Writer::encode_page(stripe, settling, pending.chunk.column(index)?)?;
732 }
733 }
734 }
735 Ok(())
736 })?;
737 drop(share);
738 building.parts.extend(held.iter().map(Part::of));
739 Ok(())
740 }
741
742 pub fn finish(&self, building: Building) -> Result<Prepared> {
748 let empty = building.parts.is_empty();
749 let (columns, gathers) = building
750 .columns
751 .into_iter()
752 .enumerate()
753 .map(|(index, growing)| {
754 let Growing { body, gather } = growing
755 .into_inner()
756 .map_err(|_| Error::internal("a native encode worker panicked"))?;
757 let column = match body {
758 Body::Coded(mut local, _) => {
759 local.done();
760 Column::Coded(local)
761 }
762 Body::Pages(stripe, _) => Column::Pages(stripe),
763 };
764 let gather = gather.filter(|_| !empty).map(|mut gather| {
766 gather.close_stripe();
767 if let Some(ceiling) = gather.ceiling() {
768 self.coded.ceilings[index].fetch_min(ceiling, Atomic::Relaxed);
769 }
770 gather
771 });
772 Ok((column, gather))
773 })
774 .collect::<Result<Vec<_>>>()?
775 .into_iter()
776 .unzip();
777 Ok(Prepared {
778 parts: building.parts,
779 types: self.types.clone(),
780 columns,
781 gathers,
782 profile: self.profile.clone(),
783 })
784 }
785}
786
787#[derive(Debug)]
789pub struct Building {
790 parts: Vec<Part>,
791 columns: Vec<Mutex<Growing>>,
794}
795
796impl Building {
797 #[must_use]
799 pub fn parts(&self) -> usize {
800 self.parts.len()
801 }
802}
803
804#[derive(Debug)]
806struct Growing {
807 body: Body,
808 gather: Option<stats::Gather>,
809}
810
811#[derive(Debug)]
813enum Body {
814 Pages(ColumnStripe, Settling),
816 Coded(Local, Option<(Arc<Vector>, Vec<u32>)>),
819}
820
821impl Body {
822 fn plain(&mut self) -> Result<()> {
829 let Self::Coded(local, _) = self else { return Ok(()) };
830 let mut stripe = ColumnStripe::default();
831 let mut settling = Settling::default();
832 for rows in local.rows()? {
833 Writer::encode_page(&mut stripe, &mut settling, &rows)?;
834 }
835 *self = Self::Pages(stripe, settling);
836 Ok(())
837 }
838}
839
840enum Slot<'a> {
842 Owned(&'a mut Option<GlobalDictionary>, &'a mut Option<stats::Gather>),
844 Lent(&'a Mutex<LentColumn>, &'a Lent),
846}
847
848struct Step<'a> {
850 index: usize,
851 column: Column,
852 slot: Slot<'a>,
853 gather: Option<stats::Gather>,
855}
856
857impl Step<'_> {
858 fn cost(&self, coded: &Coding) -> usize {
862 match &self.column {
863 Column::Coded(local) if coded[self.index].load(Atomic::Relaxed) => {
864 local.values().saturating_add(1)
865 }
866 _ => 0,
867 }
868 }
869
870 fn run(
872 self,
873 rows: usize,
874 coded: &Coding,
875 profile: Option<&LoadProfile>,
876 ) -> Result<(usize, Merge, Vec<Unencoded>)> {
877 let Self { index, column, slot, gather } = self;
878 match slot {
879 Slot::Owned(dictionary, mine) => {
880 merge_column(index, column, gather, dictionary, mine, rows, coded, profile)
881 }
882 Slot::Lent(held, lent) => {
883 let mut held = held.lock().map_err(|_| Error::internal("a merge panicked"))?;
884 if lent.reclaimed.load(Atomic::Acquire) {
887 return Err(Error::internal("a stripe was merged after its table was closed"));
888 }
889 let LentColumn { dictionary, gather: mine } = &mut *held;
890 merge_column(index, column, gather, dictionary, mine, rows, coded, profile)
891 }
892 }
893 }
894}
895
896#[expect(clippy::too_many_arguments, reason = "one column's share of the stripe's merge state")]
898fn merge_column(
899 index: usize,
900 column: Column,
901 stripe: Option<stats::Gather>,
902 dictionary: &mut Option<GlobalDictionary>,
903 gather: &mut Option<stats::Gather>,
904 rows: usize,
905 coded: &Coding,
906 profile: Option<&LoadProfile>,
907) -> Result<(usize, Merge, Vec<Unencoded>)> {
908 if let (Some(mine), Some(stripe)) = (gather.as_mut(), stripe) {
909 mine.absorb(stripe);
910 }
911 let mut new = None;
912 let merge = match (column, dictionary.as_mut()) {
913 (Column::Pages(stripe), None) => Merge::Pages(stripe),
914 (Column::Pages(stripe), Some(global)) if global.demoted => Merge::Pages(stripe),
915 (Column::Pages(_), Some(_)) => {
916 return Err(Error::internal(
917 "a column with a global dictionary was prepared without one",
918 ));
919 }
920 (Column::Coded(local), None) => Merge::Plain(local),
921 (Column::Coded(local), Some(global)) if global.demoted => Merge::Plain(local),
923 (Column::Coded(local), Some(global)) => {
924 if global.values() == 0 && drops_dictionary(rows, local.values()) {
927 if let Some(profile) = profile {
928 profile.release(global.charged);
929 }
930 coded.recount(global.charged, 0);
931 *dictionary = None;
932 coded[index].store(false, Atomic::Relaxed);
933 Merge::Plain(local)
934 } else {
935 let before = global.values();
936 let codes = local.merge_into(global)?;
937 new = Some(global.values() - before);
938 Merge::Codes { parts: local.parts, global: codes }
939 }
940 }
941 };
942 if let (Some(new), Some(global)) = (new, dictionary.as_mut()) {
944 let now = global.held_bytes();
945 coded.growth[index].store(now.saturating_sub(global.charged), Atomic::Relaxed);
946 let total =
947 coded.held.load(Atomic::Relaxed).saturating_add(now).saturating_sub(global.charged);
948 if demotes(rows, new, total, coded, index) {
949 global.demote();
950 coded[index].store(false, Atomic::Relaxed);
951 coded.growth[index].store(0, Atomic::Relaxed);
952 }
953 }
954 let blocks = match dictionary {
959 Some(dictionary) => {
960 dictionary.settle()?;
961 let blocks = dictionary.hand_out(index);
962 let (before, now) = dictionary.recharge(profile);
963 coded.recount(before, now);
964 blocks
965 }
966 None => Vec::new(),
967 };
968 Ok((index, merge, blocks))
969}
970
971fn merge_columns(prepared: Prepared, slots: Vec<Slot<'_>>, coded: &Coding) -> Result<Merged> {
979 let Prepared { parts, columns, gathers, profile, .. } = prepared;
980 let timing = profile.as_deref().map(|profile| profile.span(Stage::Dictionary));
981 let rows: usize = parts.iter().map(|part| part.rows).sum();
982 let width = columns.len();
983 if slots.len() != width || gathers.len() != width {
984 return Err(Error::internal("a stripe was merged into a table of another width"));
985 }
986 let mut steps = columns
987 .into_iter()
988 .zip(gathers)
989 .zip(slots)
990 .enumerate()
991 .map(|(index, ((column, gather), slot))| Step { index, column, slot, gather })
992 .collect::<Vec<_>>();
993 steps.sort_by_key(|step| step.cost(coded));
996 let workers = std::thread::available_parallelism()
997 .map_or(1, usize::from)
998 .min(MAX_ENCODE_WORKERS)
999 .min(steps.iter().filter(|step| step.cost(coded) > 0).count())
1000 .max(1);
1001 let done = if workers <= 1 {
1002 steps
1003 .into_iter()
1004 .map(|step| step.run(rows, coded, profile.as_deref()))
1005 .collect::<Result<Vec<_>>>()?
1006 } else {
1007 let queue = Mutex::new(steps);
1008 let pieces = std::thread::scope(|scope| {
1009 (0..workers)
1010 .map(|_| {
1011 scope.spawn(|| {
1012 let mut mine = Vec::new();
1013 loop {
1014 let taken = queue
1015 .lock()
1016 .map_err(|_| Error::internal("a merge worker panicked"))?
1017 .pop();
1018 let Some(step) = taken else { break };
1019 mine.push(step.run(rows, coded, profile.as_deref())?);
1020 }
1021 Ok(mine)
1022 })
1023 })
1024 .collect::<Vec<_>>()
1025 .into_iter()
1026 .map(|handle| {
1027 handle.join().map_err(|_| Error::internal("a merge worker panicked"))?
1028 })
1029 .collect::<Result<Vec<Vec<_>>>>()
1030 })?;
1031 pieces.into_iter().flatten().collect()
1032 };
1033 let mut slots: Vec<Option<(Merge, Vec<Unencoded>)>> = (0..width).map(|_| None).collect();
1034 for (index, merge, blocks) in done {
1035 slots[index] = Some((merge, blocks));
1036 }
1037 let mut merged = Vec::with_capacity(width);
1038 let mut blocks = Vec::new();
1039 for slot in slots {
1040 let (merge, handed) = slot.ok_or_else(|| Error::internal("a column was never merged"))?;
1041 merged.push(merge);
1042 blocks.extend(handed);
1043 }
1044 drop(timing);
1045 Ok(Merged { parts, columns: merged, blocks, profile, counted: false })
1046}
1047
1048#[derive(Debug)]
1050pub(crate) struct Lent {
1051 columns: Box<[Mutex<LentColumn>]>,
1052 reclaimed: AtomicBool,
1054}
1055
1056#[derive(Debug)]
1058pub(crate) struct LentColumn {
1059 pub(crate) dictionary: Option<GlobalDictionary>,
1060 gather: Option<stats::Gather>,
1061}
1062
1063impl Lent {
1064 pub(crate) fn columns(&self) -> &[Mutex<LentColumn>] {
1065 &self.columns
1066 }
1067
1068 fn take_back(&self, blocks: Vec<(usize, usize, EncodedBlock)>) -> Result<()> {
1070 for (column, at, block) in blocks {
1071 self.columns
1072 .get(column)
1073 .ok_or_else(|| Error::internal("a dictionary block came back to no column"))?
1074 .lock()
1075 .map_err(|_| Error::internal("a merge panicked"))?
1076 .dictionary
1077 .as_mut()
1078 .ok_or_else(|| Error::internal("a dictionary block came back to no dictionary"))?
1079 .take_back(at, block)?;
1080 }
1081 Ok(())
1082 }
1083
1084 #[allow(clippy::type_complexity)]
1086 pub(crate) fn reclaim(
1087 &self,
1088 ) -> Result<(Vec<Option<GlobalDictionary>>, Vec<Option<stats::Gather>>)> {
1089 self.reclaimed.store(true, Atomic::Release);
1090 let mut dictionaries = Vec::with_capacity(self.columns.len());
1091 let mut gathers = Vec::with_capacity(self.columns.len());
1092 for column in &self.columns {
1093 let mut held = column.lock().map_err(|_| Error::internal("a merge panicked"))?;
1094 dictionaries.push(held.dictionary.take());
1095 gathers.push(held.gather.take());
1096 }
1097 Ok((dictionaries, gathers))
1098 }
1099}
1100
1101#[derive(Debug, Clone)]
1108pub struct Merger {
1109 lent: Arc<Lent>,
1110 types: Vec<LogicalType>,
1111 coded: Arc<Coding>,
1112}
1113
1114impl Merger {
1115 pub fn merge(&self, prepared: Prepared) -> Result<Merged> {
1121 if prepared.types != self.types {
1122 return Err(invalid("a stripe was prepared for a table of other columns"));
1123 }
1124 let slots = self.lent.columns.iter().map(|column| Slot::Lent(column, &self.lent)).collect();
1125 merge_columns(prepared, slots, &self.coded)
1126 }
1127
1128 pub fn give_back(&self, paged: &mut Paged) -> Result<()> {
1135 self.lent.take_back(std::mem::take(&mut paged.blocks))
1136 }
1137}
1138
1139impl Merged {
1140 pub fn pages(self) -> Result<Paged> {
1148 let Self { parts, columns, blocks, profile, counted } = self;
1149 let width = columns.len();
1150 let mut jobs = (width..width + blocks.len())
1154 .chain((0..width).filter(|&index| !matches!(columns[index], Merge::Pages(_))))
1155 .collect::<Vec<_>>();
1156 jobs.sort_by_key(|&index| index < width && matches!(columns[index], Merge::Plain(_)));
1158 let share = Share::take(jobs.len(), parts.len());
1159 let built = fan_out(jobs, share.0, profile.as_deref(), |index| {
1160 let Some(column) = columns.get(index) else {
1161 return Ok(Built::Block(blocks[index - width].encode()?));
1162 };
1163 Ok(Built::Stripe(match column {
1164 Merge::Codes { parts, global } => code_pages(parts, global)?,
1165 Merge::Plain(local) => {
1166 Writer::encode_pages(&local.rows()?.iter().collect::<Vec<_>>())?
1167 }
1168 Merge::Pages(_) => {
1169 return Err(Error::internal("a finished column was queued to be built"));
1170 }
1171 }))
1172 })?;
1173 drop(share);
1174 let mut slots: Vec<Option<ColumnStripe>> = (0..width).map(|_| None).collect();
1175 let mut encoded = Vec::with_capacity(blocks.len());
1176 for (index, one) in built {
1177 match one {
1178 Built::Stripe(stripe) => slots[index] = Some(stripe),
1179 Built::Block(block) => {
1180 let (column, at) = blocks[index - width].place();
1181 encoded.push((column, at, block));
1182 }
1183 }
1184 }
1185 let columns = columns
1186 .into_iter()
1187 .zip(slots)
1188 .map(|(column, slot)| match (column, slot) {
1189 (Merge::Pages(stripe), _) | (_, Some(stripe)) => Ok(stripe),
1190 _ => Err(Error::internal("a column was never encoded")),
1191 })
1192 .collect::<Result<Vec<_>>>()?;
1193 Ok(Paged { parts, columns, blocks: encoded, counted })
1194 }
1195}
1196
1197fn code_pages(parts: &[LocalPart], global: &[u32]) -> Result<ColumnStripe> {
1199 let mut stripe = ColumnStripe {
1200 pages: Vec::with_capacity(parts.len()),
1201 codes: Vec::with_capacity(parts.len()),
1202 sieves: Vec::with_capacity(parts.len()),
1203 ranges: Vec::with_capacity(parts.len()),
1204 };
1205 for part in parts {
1206 let codes = part
1207 .codes
1208 .iter()
1209 .map(|&code| global.get(code as usize).copied())
1210 .collect::<Option<Vec<_>>>()
1211 .ok_or_else(|| Error::internal("a stripe's code has no global code"))?;
1212 let bytes = coded_page(&codes, &part.validity)?;
1213 if bytes.len() > MAX_PAGE {
1214 return Err(invalid("column page exceeds the configured bound"));
1215 }
1216 stripe.pages.push(bytes);
1217 stripe.codes.push(Some(unique_codes(&codes)));
1218 stripe.sieves.push(None);
1222 stripe.ranges.push(part.range.clone());
1223 }
1224 Ok(stripe)
1225}
1226
1227impl Writer {
1228 #[must_use]
1233 pub fn preparer(&self) -> Preparer {
1234 Preparer {
1235 types: self.table.fields.iter().map(|field| field.ty.clone()).collect(),
1236 coded: Arc::clone(&self.coded),
1237 profile: self.profile.clone(),
1238 }
1239 }
1240
1241 pub fn merge(&mut self, prepared: Prepared) -> Result<Merged> {
1252 self.flush_pending()?;
1253 if prepared.columns.len() != self.table.fields.len()
1254 || prepared.types.iter().ne(self.table.fields.iter().map(|field| &field.ty))
1255 {
1256 return Err(invalid("a stripe was prepared for a table of other columns"));
1257 }
1258 self.table.rows = prepared
1259 .parts
1260 .iter()
1261 .try_fold(self.table.rows, |rows, part| rows.checked_add(part.rows))
1262 .ok_or_else(|| invalid("row count overflow"))?;
1263 self.merge_held(prepared)
1264 }
1265
1266 pub(crate) fn merge_held(&mut self, prepared: Prepared) -> Result<Merged> {
1275 let slots = match &self.lent {
1276 Some(lent) => lent.columns.iter().map(|column| Slot::Lent(column, lent)).collect(),
1277 None => self
1278 .dictionaries
1279 .iter_mut()
1280 .zip(self.gathers.iter_mut())
1281 .map(|(dictionary, gather)| Slot::Owned(dictionary, gather))
1282 .collect::<Vec<_>>(),
1283 };
1284 let mut merged = merge_columns(prepared, slots, &self.coded)?;
1285 merged.counted = true;
1286 Ok(merged)
1287 }
1288
1289 pub fn merger(&mut self) -> Result<Merger> {
1299 self.flush_pending()?;
1300 let lent = match &self.lent {
1301 Some(lent) => Arc::clone(lent),
1302 None => {
1303 let lent = Arc::new(Lent {
1304 columns: std::mem::take(&mut self.dictionaries)
1305 .into_iter()
1306 .zip(std::mem::take(&mut self.gathers))
1307 .map(|(dictionary, gather)| Mutex::new(LentColumn { dictionary, gather }))
1308 .collect(),
1309 reclaimed: AtomicBool::new(false),
1310 });
1311 self.lent = Some(Arc::clone(&lent));
1312 lent
1313 }
1314 };
1315 Ok(Merger {
1316 lent,
1317 types: self.table.fields.iter().map(|field| field.ty.clone()).collect(),
1318 coded: Arc::clone(&self.coded),
1319 })
1320 }
1321
1322 pub fn write(&mut self, paged: Paged) -> Result<()> {
1328 self.write_paged(paged)
1329 }
1330
1331 pub(crate) fn write_paged(&mut self, paged: Paged) -> Result<()> {
1332 let Paged { parts, columns, blocks, counted } = paged;
1333 if !counted {
1334 self.table.rows = parts
1335 .iter()
1336 .try_fold(self.table.rows, |rows, part| rows.checked_add(part.rows))
1337 .ok_or_else(|| invalid("row count overflow"))?;
1338 }
1339 if let Some(lent) = &self.lent {
1340 lent.take_back(blocks)?;
1341 } else {
1342 for (column, at, block) in blocks {
1343 self.dictionaries
1344 .get_mut(column)
1345 .and_then(Option::as_mut)
1346 .ok_or_else(|| {
1347 Error::internal("a dictionary block came back to no dictionary")
1348 })?
1349 .take_back(at, block)?;
1350 }
1351 }
1352 if parts.is_empty() {
1353 return self.place_blocks();
1354 }
1355 self.write_stripe(&parts, columns)
1356 }
1357
1358 pub fn append_prepared(&mut self, prepared: Prepared) -> Result<()> {
1364 let merged = self.merge(prepared)?;
1365 let paged = merged.pages()?;
1366 self.write(paged)
1367 }
1368}
1369
1370#[cfg(test)]
1371mod tests {
1372 use std::fs;
1373 use std::path::PathBuf;
1374 use std::time::{SystemTime, UNIX_EPOCH};
1375
1376 use rudb_common::{Field, Value};
1377 use rudb_vector::Vector;
1378
1379 use super::*;
1380 use crate::Reader;
1381
1382 const PART: usize = 1_000;
1383
1384 #[test]
1390 fn dictionary_parts_code_as_their_rows_would() {
1391 let texts = |values: &[&str]| {
1392 Arc::new(
1393 Vector::from_values(
1394 LogicalType::Varchar,
1395 &values
1396 .iter()
1397 .map(|text| Value::Varchar((*text).to_string()))
1398 .collect::<Vec<_>>(),
1399 )
1400 .expect("a dictionary"),
1401 )
1402 };
1403 let shared = texts(&["b", "a", "", "c", "unused"]);
1404 let other = texts(&["c", "d", "a"]);
1405 let mut nulls = Bitmap::all_valid(6);
1406 nulls.set(1, false);
1407 nulls.set(4, false);
1408 let parts = [
1409 Vector::dictionary_over(vec![3, 3, 1, 0, 2, 1], Arc::clone(&shared)).expect("codes"),
1410 Vector::dictionary_over(vec![0, 3, 1, 1, 3, 2], Arc::clone(&shared))
1411 .expect("codes")
1412 .with_validity(Validity::Mask(nulls)),
1413 Vector::dictionary_over(vec![1, 2, 0, 1], other).expect("codes"),
1414 ];
1415 let held = |flat: bool| {
1416 parts
1417 .iter()
1418 .enumerate()
1419 .map(|(at, part)| PendingChunk {
1420 order: (at as u64, 0),
1421 chunk: Chunk::new(vec![if flat {
1422 part.flatten().expect("flat")
1423 } else {
1424 part.clone()
1425 }])
1426 .expect("a chunk"),
1427 })
1428 .collect::<Vec<_>>()
1429 };
1430 let parquet = held(false);
1431 let before = rudb_common::slow::here();
1432 let coded = Local::code_column(0, &parquet).expect("coded");
1433 assert_eq!(
1434 rudb_common::slow::here().since(before).get(rudb_common::slow::Cause::Flatten),
1435 0,
1436 "a part that came in as codes was flattened",
1437 );
1438 let flat = Local::code_column(0, &held(true)).expect("coded");
1439 assert_eq!(coded.values(), flat.values());
1440 for code in 0..flat.values() as u32 {
1441 assert_eq!(coded.value(code), flat.value(code), "value {code}");
1442 }
1443 assert_eq!(coded.counts, flat.counts);
1444 assert_eq!(coded.nulls, flat.nulls);
1445 assert_eq!(coded.nulls, 2);
1446 assert_eq!(coded.hashes, flat.hashes);
1447 assert_eq!(coded.checks, flat.checks);
1448 assert_eq!(coded.parts.len(), flat.parts.len());
1449 for (coded, flat) in coded.parts.iter().zip(&flat.parts) {
1450 assert_eq!(coded.codes, flat.codes);
1451 assert_eq!(coded.validity, flat.validity);
1452 }
1453 }
1454
1455 fn path(label: &str) -> PathBuf {
1456 let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
1457 std::env::temp_dir()
1458 .join(format!("rudb-prepare-{label}-{}-{stamp}.rdb", std::process::id()))
1459 }
1460
1461 fn fields() -> Vec<Field> {
1462 vec![
1463 Field::required("id", LogicalType::BigInt),
1464 Field::new("city", LogicalType::Varchar),
1465 Field::new("note", LogicalType::Varchar),
1466 ]
1467 }
1468
1469 fn row(id: usize) -> [Value; 3] {
1474 let city = if id % 11 == 0 {
1475 Value::Null
1476 } else {
1477 Value::Varchar(format!("city {}", (id / 7) % 13))
1478 };
1479 let note = if id % 17 == 0 { Value::Null } else { Value::Varchar(format!("note {id}")) };
1480 [Value::BigInt(id as i64), city, note]
1481 }
1482
1483 fn stripe(first: usize, parts: usize) -> Vec<((u64, u64), Chunk)> {
1485 (first..first + parts)
1486 .map(|part| {
1487 let rows = (part * PART..(part + 1) * PART).map(row).collect::<Vec<_>>();
1488 let column = |at: usize| {
1489 let values = rows.iter().map(|row| row[at].clone()).collect::<Vec<_>>();
1490 Vector::from_values(fields()[at].ty.clone(), &values).expect("a column")
1491 };
1492 let chunk = Chunk::new(vec![column(0), column(1), column(2)]).expect("a chunk");
1493 ((part as u64, 0), chunk)
1494 })
1495 .collect()
1496 }
1497
1498 fn runs() -> Vec<Vec<((u64, u64), Chunk)>> {
1500 vec![stripe(5, 5), stripe(0, 5), stripe(10, 3)]
1501 }
1502
1503 fn check(path: &PathBuf) {
1504 let reader = Reader::open(path).expect("reopen");
1505 assert_eq!(reader.parts(), 13);
1506 for part in 0..13 {
1507 let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1508 for at in [0, 17, PART - 1] {
1509 let want = row(part * PART + at);
1510 for (column, value) in want.iter().enumerate() {
1511 assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1512 }
1513 }
1514 }
1515 }
1516
1517 #[test]
1525 fn stripes_prepared_before_any_is_merged_write_the_same_bytes_as_one_at_a_time() {
1526 let alone = path("alone");
1527 let mut writer = Writer::create(&alone, "t", fields()).expect("a file");
1528 for run in runs() {
1529 writer.append_stripe(run).expect("a stripe");
1530 }
1531 writer.finish().expect("commit");
1532
1533 let split = path("split");
1534 let mut writer = Writer::create(&split, "t", fields()).expect("a file");
1535 let preparer = writer.preparer();
1536 let prepared = runs()
1537 .into_iter()
1538 .map(|run| preparer.prepare(run).expect("prepared"))
1539 .collect::<Vec<_>>();
1540 for one in prepared {
1541 writer.append_prepared(one).expect("a stripe");
1542 }
1543 assert!(!preparer.coded[2].load(Atomic::Relaxed), "note lost its dictionary");
1544 assert!(preparer.coded[1].load(Atomic::Relaxed), "city kept its dictionary");
1545 writer.finish().expect("commit");
1546
1547 assert_eq!(fs::read(&alone).expect("read"), fs::read(&split).expect("read"));
1548 check(&split);
1549 fs::remove_file(alone).expect("remove");
1550 fs::remove_file(split).expect("remove");
1551 }
1552
1553 #[test]
1559 fn stripes_counted_under_an_earlier_ceiling_count_what_they_would_have() {
1560 let estimates = |writer: &Writer| {
1561 writer
1562 .gathers
1563 .iter()
1564 .map(|gather| gather.as_ref().and_then(stats::Gather::distinct))
1565 .collect::<Vec<_>>()
1566 };
1567 let want = fields()
1569 .iter()
1570 .enumerate()
1571 .map(|(column, field)| {
1572 let mut gather = stats::Gather::new(&field.ty, 0)?;
1573 for (_, chunk) in runs().into_iter().flatten() {
1574 gather.part(chunk.column(column).expect("a column"));
1575 }
1576 gather.distinct()
1577 })
1578 .collect::<Vec<_>>();
1579
1580 for reversed in [false, true] {
1581 let split = path("capped");
1582 let mut writer = Writer::create(&split, "t", fields()).expect("a file");
1583 let preparer = writer.preparer();
1584 let mut prepared = runs()
1585 .into_iter()
1586 .map(|run| preparer.prepare(run).expect("prepared"))
1587 .collect::<Vec<_>>();
1588 for column in [0, 2] {
1589 assert!(preparer.coded.ceilings[column].load(Atomic::Relaxed) < u64::MAX);
1590 }
1591 if reversed {
1592 prepared.reverse();
1593 }
1594 for one in prepared {
1595 writer.append_prepared(one).expect("a stripe");
1596 }
1597 assert_eq!(estimates(&writer), want, "reversed {reversed}");
1598 writer.finish().expect("commit");
1599 fs::remove_file(split).expect("remove");
1600 }
1601 }
1602
1603 #[test]
1607 fn stripes_written_in_another_order_than_they_were_merged_read_back() {
1608 let path = path("crossed");
1609 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1610 let preparer = writer.preparer();
1611 let mut merged = runs()
1612 .into_iter()
1613 .map(|run| writer.merge(preparer.prepare(run).expect("prepared")).expect("merged"))
1614 .map(|merged| merged.pages().expect("paged"))
1615 .collect::<Vec<_>>();
1616 merged.reverse();
1617 for paged in merged {
1618 writer.write(paged).expect("written");
1619 }
1620 writer.finish().expect("commit");
1621 check(&path);
1622 fs::remove_file(path).expect("remove");
1623 }
1624
1625 #[test]
1628 fn stripes_merged_through_a_merger_write_the_same_bytes_as_the_writer() {
1629 let alone = path("alone-merger");
1630 let mut writer = Writer::create(&alone, "t", fields()).expect("a file");
1631 for run in runs() {
1632 writer.append_stripe(run).expect("a stripe");
1633 }
1634 writer.finish().expect("commit");
1635
1636 let lent = path("lent");
1637 let mut writer = Writer::create(&lent, "t", fields()).expect("a file");
1638 let preparer = writer.preparer();
1639 let merger = writer.merger().expect("a merger");
1640 for run in runs() {
1641 let merged = merger.merge(preparer.prepare(run).expect("prepared")).expect("merged");
1642 let mut paged = merged.pages().expect("paged");
1643 merger.give_back(&mut paged).expect("given back");
1644 writer.write(paged).expect("written");
1645 }
1646 assert_eq!(writer.table.rows, 13 * PART);
1647 writer.finish().expect("commit");
1648
1649 assert_eq!(fs::read(&alone).expect("read"), fs::read(&lent).expect("read"));
1650 check(&lent);
1651 fs::remove_file(alone).expect("remove");
1652 fs::remove_file(lent).expect("remove");
1653 }
1654
1655 #[test]
1658 fn stripes_fed_in_batches_write_the_same_bytes_as_prepared_whole() {
1659 let whole = path("whole");
1660 let mut writer = Writer::create(&whole, "t", fields()).expect("a file");
1661 for run in runs() {
1662 writer.append_stripe(run).expect("a stripe");
1663 }
1664 writer.finish().expect("commit");
1665
1666 let fed = path("fed");
1667 let mut writer = Writer::create(&fed, "t", fields()).expect("a file");
1668 let preparer = writer.preparer();
1669 let merger = writer.merger().expect("a merger");
1670 let nothing = Chunk::new(
1671 fields()
1672 .iter()
1673 .map(|field| Vector::from_values(field.ty.clone(), &[]).expect("a column"))
1674 .collect(),
1675 )
1676 .expect("a chunk");
1677 for mut run in runs() {
1678 let mut building = preparer.start();
1679 preparer.feed(&mut building, Vec::new()).expect("fed nothing");
1680 while !run.is_empty() {
1681 let rest = run.split_off(2.min(run.len()));
1682 let mut batch = std::mem::replace(&mut run, rest);
1683 batch.push(((u64::MAX, 0), nothing.clone()));
1684 preparer.feed(&mut building, batch).expect("fed");
1685 }
1686 let merged =
1687 merger.merge(preparer.finish(building).expect("finished")).expect("merged");
1688 let mut paged = merged.pages().expect("paged");
1689 merger.give_back(&mut paged).expect("given back");
1690 writer.write(paged).expect("written");
1691 }
1692 writer.finish().expect("commit");
1693
1694 assert_eq!(fs::read(&whole).expect("read"), fs::read(&fed).expect("read"));
1695 check(&fed);
1696 fs::remove_file(whole).expect("remove");
1697 fs::remove_file(fed).expect("remove");
1698 }
1699
1700 #[test]
1703 fn a_column_dropped_while_its_stripe_is_built_turns_to_pages() {
1704 let path = path("dropped-while-built");
1705 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1706 let preparer = writer.preparer();
1707 let merger = writer.merger().expect("a merger");
1708 let write = |writer: &mut Writer, building: Building| {
1709 let merged =
1710 merger.merge(preparer.finish(building).expect("finished")).expect("merged");
1711 let mut paged = merged.pages().expect("paged");
1712 merger.give_back(&mut paged).expect("given back");
1713 writer.write(paged).expect("written");
1714 };
1715 let mut runs = runs().into_iter();
1716 let mut first = preparer.start();
1717 let mut second = preparer.start();
1718 let mut later = runs.next().expect("a run");
1719 preparer.feed(&mut second, later.drain(..2).collect()).expect("fed");
1720 preparer.feed(&mut first, runs.next().expect("a run")).expect("fed");
1721 write(&mut writer, first);
1722 let note = |building: &Building| {
1723 matches!(building.columns[2].lock().expect("unpoisoned").body, Body::Pages(..))
1724 };
1725 assert!(!note(&second), "still coded until it is fed again");
1726 preparer.feed(&mut second, later).expect("fed");
1727 assert!(note(&second), "turned to pages once fed after the drop");
1728 write(&mut writer, second);
1729 let mut last = preparer.start();
1730 preparer.feed(&mut last, runs.next().expect("a run")).expect("fed");
1731 write(&mut writer, last);
1732 writer.finish().expect("commit");
1733 check(&path);
1734 fs::remove_file(path).expect("remove");
1735 }
1736
1737 #[test]
1739 fn a_stripe_fed_more_parts_than_it_holds_is_refused() {
1740 let path = path("overfed");
1741 let writer = Writer::create(&path, "t", fields()).expect("a file");
1742 let preparer = writer.preparer();
1743 let mut building = preparer.start();
1744 preparer.feed(&mut building, stripe(0, STRIPE_PARTS - 1)).expect("fed");
1745 assert_eq!(building.parts(), STRIPE_PARTS - 1);
1746 assert!(preparer.feed(&mut building, stripe(STRIPE_PARTS, 2)).is_err());
1747 drop(writer);
1748 let _ = fs::remove_file(path);
1749 }
1750
1751 #[test]
1754 fn stripes_merged_on_several_threads_at_once_read_back() {
1755 let path = path("merged-at-once");
1756 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1757 let preparer = writer.preparer();
1758 let merger = writer.merger().expect("a merger");
1759 let writer = Mutex::new(writer);
1760 std::thread::scope(|scope| {
1761 for run in runs() {
1762 let (preparer, merger, writer) = (&preparer, &merger, &writer);
1763 scope.spawn(move || {
1764 let merged =
1765 merger.merge(preparer.prepare(run).expect("prepared")).expect("merged");
1766 let mut paged = merged.pages().expect("paged");
1767 merger.give_back(&mut paged).expect("given back");
1768 writer.lock().expect("the writer").write(paged).expect("written");
1769 });
1770 }
1771 });
1772 writer.into_inner().expect("the writer").finish().expect("commit");
1773 check(&path);
1774 fs::remove_file(path).expect("remove");
1775 }
1776
1777 #[test]
1780 fn a_merge_after_the_table_is_closed_is_refused() {
1781 let path = path("late");
1782 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1783 let preparer = writer.preparer();
1784 let merger = writer.merger().expect("a merger");
1785 writer.finish().expect("commit");
1786 let prepared = preparer.prepare(stripe(0, 2)).expect("prepared");
1787 assert!(merger.merge(prepared).is_err());
1788 fs::remove_file(path).expect("remove");
1789 }
1790
1791 #[test]
1794 fn a_stripe_of_another_table_is_refused_at_the_merge() {
1795 let path = path("refused");
1796 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1797 let other = Writer::create(path.with_extension("other"), "u", vec![fields().remove(0)])
1798 .expect("a file");
1799 let prepared = other.preparer().prepare(vec![]).expect("nothing to prepare");
1800 assert!(writer.merge(prepared).is_err());
1801 assert_eq!(writer.table.rows, 0);
1802 drop(other);
1803 fs::remove_file(path.with_extension("other")).expect("remove");
1804 fs::remove_file(path).expect("remove");
1805 }
1806
1807 fn turning(id: usize) -> [Value; 3] {
1811 let url = match id {
1812 _ if id % 13 == 0 => Value::Null,
1813 _ if id < 5 * PART => Value::Varchar(format!("https://example.com/{}", id % 20)),
1814 _ => Value::Varchar(format!("https://example.com/page/{id}")),
1815 };
1816 [Value::BigInt(id as i64), Value::Varchar(format!("city {}", id % 13)), url]
1817 }
1818
1819 fn turning_fields() -> Vec<Field> {
1820 vec![
1821 Field::required("id", LogicalType::BigInt),
1822 Field::new("city", LogicalType::Varchar),
1823 Field::new("url", LogicalType::Varchar),
1824 ]
1825 }
1826
1827 fn turning_stripe(
1828 rows: fn(usize) -> [Value; 3],
1829 first: usize,
1830 parts: usize,
1831 ) -> Vec<((u64, u64), Chunk)> {
1832 (first..first + parts)
1833 .map(|part| {
1834 let rows = (part * PART..(part + 1) * PART).map(rows).collect::<Vec<_>>();
1835 let column = |at: usize| {
1836 let values = rows.iter().map(|row| row[at].clone()).collect::<Vec<_>>();
1837 Vector::from_values(turning_fields()[at].ty.clone(), &values).expect("a column")
1838 };
1839 let chunk = Chunk::new(vec![column(0), column(1), column(2)]).expect("a chunk");
1840 ((part as u64, 0), chunk)
1841 })
1842 .collect()
1843 }
1844
1845 fn turning_runs() -> Vec<Vec<((u64, u64), Chunk)>> {
1846 vec![
1847 turning_stripe(turning, 0, 5),
1848 turning_stripe(turning, 5, 5),
1849 turning_stripe(turning, 10, 3),
1850 ]
1851 }
1852
1853 fn check_turned(path: &PathBuf) {
1855 let reader = Reader::open(path).expect("reopen");
1856 assert_eq!(reader.parts(), 13);
1857 assert_eq!(reader.table().demoted, [false, false, true], "only url is demoted");
1858 for part in 0..13 {
1859 let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1860 let url = chunk.column(2).expect("url");
1861 assert!(url.stable_dictionary_parts().is_none(), "part {part} hands out no codes");
1862 for at in 0..PART {
1863 let want = turning(part * PART + at);
1864 for (column, value) in want.iter().enumerate() {
1865 assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1866 }
1867 }
1868 }
1869 assert_eq!(reader.distinct_values(2).expect("asked"), None);
1871 assert_eq!(reader.text_extremes(2).expect("asked"), None);
1872 assert_eq!(reader.exact_frequencies(2).expect("asked"), None);
1873 assert_eq!(reader.top_frequencies(2, 5).expect("asked"), None);
1874 assert!(!reader.skips_codes(0, 2, &[0]).expect("asked"), "no code proves a value absent");
1875 assert_eq!(reader.distinct_values(1).expect("asked"), Some(13));
1877 assert!(reader.text_extremes(1).expect("asked").is_some());
1878 }
1879
1880 #[test]
1884 fn a_column_that_turns_unique_is_demoted_and_reads_back() {
1885 let alone = path("demoted-alone");
1886 let mut writer = Writer::create(&alone, "t", turning_fields()).expect("a file");
1887 let preparer = writer.preparer();
1888 let mut runs = turning_runs().into_iter();
1889 writer.append_stripe(runs.next().expect("a run")).expect("a stripe");
1890 assert!(preparer.coded[2].load(Atomic::Relaxed), "url repeats in its first stripe");
1891 for run in runs {
1892 writer.append_stripe(run).expect("a stripe");
1893 }
1894 assert!(!preparer.coded[2].load(Atomic::Relaxed), "url was demoted");
1895 assert!(preparer.coded[1].load(Atomic::Relaxed), "city kept its dictionary");
1896 writer.finish().expect("commit");
1897 check_turned(&alone);
1898
1899 let split = path("demoted-split");
1900 let mut writer = Writer::create(&split, "t", turning_fields()).expect("a file");
1901 let preparer = writer.preparer();
1902 let prepared = turning_runs()
1903 .into_iter()
1904 .map(|run| preparer.prepare(run).expect("prepared"))
1905 .collect::<Vec<_>>();
1906 for one in prepared {
1907 writer.append_prepared(one).expect("a stripe");
1908 }
1909 writer.finish().expect("commit");
1910 check_turned(&split);
1911
1912 fs::remove_file(alone).expect("remove");
1913 fs::remove_file(split).expect("remove");
1914 }
1915
1916 fn growing(id: usize) -> [Value; 3] {
1919 let url = Value::Varchar(format!("https://example.com/{}", id / 4));
1920 [Value::BigInt(id as i64), Value::Varchar(format!("city {}", id % 13)), url]
1921 }
1922
1923 #[test]
1926 fn the_dictionary_cap_demotes_the_column_that_grew_most() {
1927 let path = path("capped");
1928 let mut writer = Writer::create(&path, "t", turning_fields())
1929 .expect("a file")
1930 .with_dictionary_cap(1 << 30);
1931 let preparer = writer.preparer();
1932 writer.append_stripe(turning_stripe(growing, 0, 5)).expect("a stripe");
1933 assert!(preparer.coded[2].load(Atomic::Relaxed), "url is under the cap");
1934 assert!(preparer.coded[1].load(Atomic::Relaxed), "city is under the cap");
1935
1936 writer.coded.cap(1);
1937 writer.append_stripe(turning_stripe(growing, 5, 5)).expect("a stripe");
1938 assert!(!preparer.coded[2].load(Atomic::Relaxed), "url grew most and was demoted");
1939 assert!(preparer.coded[1].load(Atomic::Relaxed), "city grew nothing and keeps it");
1940 writer.append_stripe(turning_stripe(growing, 10, 3)).expect("a stripe");
1941 assert!(preparer.coded[1].load(Atomic::Relaxed), "city still grows nothing");
1942 writer.finish().expect("commit");
1943
1944 let reader = Reader::open(&path).expect("reopen");
1945 assert_eq!(reader.table().demoted, [false, false, true]);
1946 for part in 0..13 {
1947 let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1948 for at in 0..PART {
1949 let want = growing(part * PART + at);
1950 for (column, value) in want.iter().enumerate() {
1951 assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1952 }
1953 }
1954 }
1955 assert_eq!(reader.distinct_values(1).expect("asked"), Some(13));
1956 assert_eq!(reader.distinct_values(2).expect("asked"), None);
1957 fs::remove_file(path).expect("remove");
1958 }
1959}