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 sums: Vec::with_capacity(parts.len()),
1202 codes: Vec::with_capacity(parts.len()),
1203 sieves: Vec::with_capacity(parts.len()),
1204 ranges: Vec::with_capacity(parts.len()),
1205 };
1206 for part in parts {
1207 let codes = part
1208 .codes
1209 .iter()
1210 .map(|&code| global.get(code as usize).copied())
1211 .collect::<Option<Vec<_>>>()
1212 .ok_or_else(|| Error::internal("a stripe's code has no global code"))?;
1213 let bytes = coded_page(&codes, &part.validity)?;
1214 if bytes.len() > MAX_PAGE {
1215 return Err(invalid("column page exceeds the configured bound"));
1216 }
1217 stripe.sums.push(checksum(&bytes));
1218 stripe.pages.push(bytes);
1219 stripe.codes.push(Some(unique_codes(&codes)));
1220 stripe.sieves.push(None);
1224 stripe.ranges.push(part.range.clone());
1225 }
1226 Ok(stripe)
1227}
1228
1229impl Writer {
1230 #[must_use]
1235 pub fn preparer(&self) -> Preparer {
1236 Preparer {
1237 types: self.table.fields.iter().map(|field| field.ty.clone()).collect(),
1238 coded: Arc::clone(&self.coded),
1239 profile: self.profile.clone(),
1240 }
1241 }
1242
1243 pub fn merge(&mut self, prepared: Prepared) -> Result<Merged> {
1254 self.flush_pending()?;
1255 if prepared.columns.len() != self.table.fields.len()
1256 || prepared.types.iter().ne(self.table.fields.iter().map(|field| &field.ty))
1257 {
1258 return Err(invalid("a stripe was prepared for a table of other columns"));
1259 }
1260 self.table.rows = prepared
1261 .parts
1262 .iter()
1263 .try_fold(self.table.rows, |rows, part| rows.checked_add(part.rows))
1264 .ok_or_else(|| invalid("row count overflow"))?;
1265 self.merge_held(prepared)
1266 }
1267
1268 pub(crate) fn merge_held(&mut self, prepared: Prepared) -> Result<Merged> {
1277 let slots = match &self.lent {
1278 Some(lent) => lent.columns.iter().map(|column| Slot::Lent(column, lent)).collect(),
1279 None => self
1280 .dictionaries
1281 .iter_mut()
1282 .zip(self.gathers.iter_mut())
1283 .map(|(dictionary, gather)| Slot::Owned(dictionary, gather))
1284 .collect::<Vec<_>>(),
1285 };
1286 let mut merged = merge_columns(prepared, slots, &self.coded)?;
1287 merged.counted = true;
1288 Ok(merged)
1289 }
1290
1291 pub fn merger(&mut self) -> Result<Merger> {
1301 self.flush_pending()?;
1302 let lent = match &self.lent {
1303 Some(lent) => Arc::clone(lent),
1304 None => {
1305 let lent = Arc::new(Lent {
1306 columns: std::mem::take(&mut self.dictionaries)
1307 .into_iter()
1308 .zip(std::mem::take(&mut self.gathers))
1309 .map(|(dictionary, gather)| Mutex::new(LentColumn { dictionary, gather }))
1310 .collect(),
1311 reclaimed: AtomicBool::new(false),
1312 });
1313 self.lent = Some(Arc::clone(&lent));
1314 lent
1315 }
1316 };
1317 Ok(Merger {
1318 lent,
1319 types: self.table.fields.iter().map(|field| field.ty.clone()).collect(),
1320 coded: Arc::clone(&self.coded),
1321 })
1322 }
1323
1324 pub fn write(&mut self, paged: Paged) -> Result<()> {
1330 self.write_paged(paged)
1331 }
1332
1333 pub(crate) fn write_paged(&mut self, paged: Paged) -> Result<()> {
1334 let Paged { parts, columns, blocks, counted } = paged;
1335 if !counted {
1336 self.table.rows = parts
1337 .iter()
1338 .try_fold(self.table.rows, |rows, part| rows.checked_add(part.rows))
1339 .ok_or_else(|| invalid("row count overflow"))?;
1340 }
1341 if let Some(lent) = &self.lent {
1342 lent.take_back(blocks)?;
1343 } else {
1344 for (column, at, block) in blocks {
1345 self.dictionaries
1346 .get_mut(column)
1347 .and_then(Option::as_mut)
1348 .ok_or_else(|| {
1349 Error::internal("a dictionary block came back to no dictionary")
1350 })?
1351 .take_back(at, block)?;
1352 }
1353 }
1354 if parts.is_empty() {
1355 return self.place_blocks();
1356 }
1357 self.write_stripe(&parts, columns)
1358 }
1359
1360 pub fn append_prepared(&mut self, prepared: Prepared) -> Result<()> {
1366 let merged = self.merge(prepared)?;
1367 let paged = merged.pages()?;
1368 self.write(paged)
1369 }
1370}
1371
1372#[cfg(test)]
1373mod tests {
1374 use std::fs;
1375 use std::path::PathBuf;
1376 use std::time::{SystemTime, UNIX_EPOCH};
1377
1378 use rudb_common::{Field, Value};
1379 use rudb_vector::Vector;
1380
1381 use super::*;
1382 use crate::Reader;
1383
1384 const PART: usize = 1_000;
1385
1386 #[test]
1392 fn dictionary_parts_code_as_their_rows_would() {
1393 let texts = |values: &[&str]| {
1394 Arc::new(
1395 Vector::from_values(
1396 LogicalType::Varchar,
1397 &values
1398 .iter()
1399 .map(|text| Value::Varchar((*text).to_string()))
1400 .collect::<Vec<_>>(),
1401 )
1402 .expect("a dictionary"),
1403 )
1404 };
1405 let shared = texts(&["b", "a", "", "c", "unused"]);
1406 let other = texts(&["c", "d", "a"]);
1407 let mut nulls = Bitmap::all_valid(6);
1408 nulls.set(1, false);
1409 nulls.set(4, false);
1410 let parts = [
1411 Vector::dictionary_over(vec![3, 3, 1, 0, 2, 1], Arc::clone(&shared)).expect("codes"),
1412 Vector::dictionary_over(vec![0, 3, 1, 1, 3, 2], Arc::clone(&shared))
1413 .expect("codes")
1414 .with_validity(Validity::Mask(nulls)),
1415 Vector::dictionary_over(vec![1, 2, 0, 1], other).expect("codes"),
1416 ];
1417 let held = |flat: bool| {
1418 parts
1419 .iter()
1420 .enumerate()
1421 .map(|(at, part)| PendingChunk {
1422 order: (at as u64, 0),
1423 chunk: Chunk::new(vec![if flat {
1424 part.flatten().expect("flat")
1425 } else {
1426 part.clone()
1427 }])
1428 .expect("a chunk"),
1429 })
1430 .collect::<Vec<_>>()
1431 };
1432 let parquet = held(false);
1433 let before = rudb_common::slow::here();
1434 let coded = Local::code_column(0, &parquet).expect("coded");
1435 assert_eq!(
1436 rudb_common::slow::here().since(before).get(rudb_common::slow::Cause::Flatten),
1437 0,
1438 "a part that came in as codes was flattened",
1439 );
1440 let flat = Local::code_column(0, &held(true)).expect("coded");
1441 assert_eq!(coded.values(), flat.values());
1442 for code in 0..flat.values() as u32 {
1443 assert_eq!(coded.value(code), flat.value(code), "value {code}");
1444 }
1445 assert_eq!(coded.counts, flat.counts);
1446 assert_eq!(coded.nulls, flat.nulls);
1447 assert_eq!(coded.nulls, 2);
1448 assert_eq!(coded.hashes, flat.hashes);
1449 assert_eq!(coded.checks, flat.checks);
1450 assert_eq!(coded.parts.len(), flat.parts.len());
1451 for (coded, flat) in coded.parts.iter().zip(&flat.parts) {
1452 assert_eq!(coded.codes, flat.codes);
1453 assert_eq!(coded.validity, flat.validity);
1454 }
1455 }
1456
1457 fn path(label: &str) -> PathBuf {
1458 let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
1459 std::env::temp_dir()
1460 .join(format!("rudb-prepare-{label}-{}-{stamp}.rdb", std::process::id()))
1461 }
1462
1463 fn fields() -> Vec<Field> {
1464 vec![
1465 Field::required("id", LogicalType::BigInt),
1466 Field::new("city", LogicalType::Varchar),
1467 Field::new("note", LogicalType::Varchar),
1468 ]
1469 }
1470
1471 fn row(id: usize) -> [Value; 3] {
1476 let city = if id % 11 == 0 {
1477 Value::Null
1478 } else {
1479 Value::Varchar(format!("city {}", (id / 7) % 13))
1480 };
1481 let note = if id % 17 == 0 { Value::Null } else { Value::Varchar(format!("note {id}")) };
1482 [Value::BigInt(id as i64), city, note]
1483 }
1484
1485 fn stripe(first: usize, parts: usize) -> Vec<((u64, u64), Chunk)> {
1487 (first..first + parts)
1488 .map(|part| {
1489 let rows = (part * PART..(part + 1) * PART).map(row).collect::<Vec<_>>();
1490 let column = |at: usize| {
1491 let values = rows.iter().map(|row| row[at].clone()).collect::<Vec<_>>();
1492 Vector::from_values(fields()[at].ty.clone(), &values).expect("a column")
1493 };
1494 let chunk = Chunk::new(vec![column(0), column(1), column(2)]).expect("a chunk");
1495 ((part as u64, 0), chunk)
1496 })
1497 .collect()
1498 }
1499
1500 fn runs() -> Vec<Vec<((u64, u64), Chunk)>> {
1502 vec![stripe(5, 5), stripe(0, 5), stripe(10, 3)]
1503 }
1504
1505 fn check(path: &PathBuf) {
1506 let reader = Reader::open(path).expect("reopen");
1507 assert_eq!(reader.parts(), 13);
1508 for part in 0..13 {
1509 let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1510 for at in [0, 17, PART - 1] {
1511 let want = row(part * PART + at);
1512 for (column, value) in want.iter().enumerate() {
1513 assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1514 }
1515 }
1516 }
1517 }
1518
1519 #[test]
1527 fn stripes_prepared_before_any_is_merged_write_the_same_bytes_as_one_at_a_time() {
1528 let alone = path("alone");
1529 let mut writer = Writer::create(&alone, "t", fields()).expect("a file");
1530 for run in runs() {
1531 writer.append_stripe(run).expect("a stripe");
1532 }
1533 writer.finish().expect("commit");
1534
1535 let split = path("split");
1536 let mut writer = Writer::create(&split, "t", fields()).expect("a file");
1537 let preparer = writer.preparer();
1538 let prepared = runs()
1539 .into_iter()
1540 .map(|run| preparer.prepare(run).expect("prepared"))
1541 .collect::<Vec<_>>();
1542 for one in prepared {
1543 writer.append_prepared(one).expect("a stripe");
1544 }
1545 assert!(!preparer.coded[2].load(Atomic::Relaxed), "note lost its dictionary");
1546 assert!(preparer.coded[1].load(Atomic::Relaxed), "city kept its dictionary");
1547 writer.finish().expect("commit");
1548
1549 assert_eq!(fs::read(&alone).expect("read"), fs::read(&split).expect("read"));
1550 check(&split);
1551 fs::remove_file(alone).expect("remove");
1552 fs::remove_file(split).expect("remove");
1553 }
1554
1555 #[test]
1561 fn stripes_counted_under_an_earlier_ceiling_count_what_they_would_have() {
1562 let estimates = |writer: &Writer| {
1563 writer
1564 .gathers
1565 .iter()
1566 .map(|gather| gather.as_ref().and_then(stats::Gather::distinct))
1567 .collect::<Vec<_>>()
1568 };
1569 let want = fields()
1571 .iter()
1572 .enumerate()
1573 .map(|(column, field)| {
1574 let mut gather = stats::Gather::new(&field.ty, 0)?;
1575 for (_, chunk) in runs().into_iter().flatten() {
1576 gather.part(chunk.column(column).expect("a column"));
1577 }
1578 gather.distinct()
1579 })
1580 .collect::<Vec<_>>();
1581
1582 for reversed in [false, true] {
1583 let split = path("capped");
1584 let mut writer = Writer::create(&split, "t", fields()).expect("a file");
1585 let preparer = writer.preparer();
1586 let mut prepared = runs()
1587 .into_iter()
1588 .map(|run| preparer.prepare(run).expect("prepared"))
1589 .collect::<Vec<_>>();
1590 for column in [0, 2] {
1591 assert!(preparer.coded.ceilings[column].load(Atomic::Relaxed) < u64::MAX);
1592 }
1593 if reversed {
1594 prepared.reverse();
1595 }
1596 for one in prepared {
1597 writer.append_prepared(one).expect("a stripe");
1598 }
1599 assert_eq!(estimates(&writer), want, "reversed {reversed}");
1600 writer.finish().expect("commit");
1601 fs::remove_file(split).expect("remove");
1602 }
1603 }
1604
1605 #[test]
1609 fn stripes_written_in_another_order_than_they_were_merged_read_back() {
1610 let path = path("crossed");
1611 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1612 let preparer = writer.preparer();
1613 let mut merged = runs()
1614 .into_iter()
1615 .map(|run| writer.merge(preparer.prepare(run).expect("prepared")).expect("merged"))
1616 .map(|merged| merged.pages().expect("paged"))
1617 .collect::<Vec<_>>();
1618 merged.reverse();
1619 for paged in merged {
1620 writer.write(paged).expect("written");
1621 }
1622 writer.finish().expect("commit");
1623 check(&path);
1624 fs::remove_file(path).expect("remove");
1625 }
1626
1627 #[test]
1630 fn stripes_merged_through_a_merger_write_the_same_bytes_as_the_writer() {
1631 let alone = path("alone-merger");
1632 let mut writer = Writer::create(&alone, "t", fields()).expect("a file");
1633 for run in runs() {
1634 writer.append_stripe(run).expect("a stripe");
1635 }
1636 writer.finish().expect("commit");
1637
1638 let lent = path("lent");
1639 let mut writer = Writer::create(&lent, "t", fields()).expect("a file");
1640 let preparer = writer.preparer();
1641 let merger = writer.merger().expect("a merger");
1642 for run in runs() {
1643 let merged = merger.merge(preparer.prepare(run).expect("prepared")).expect("merged");
1644 let mut paged = merged.pages().expect("paged");
1645 merger.give_back(&mut paged).expect("given back");
1646 writer.write(paged).expect("written");
1647 }
1648 assert_eq!(writer.table.rows, 13 * PART);
1649 writer.finish().expect("commit");
1650
1651 assert_eq!(fs::read(&alone).expect("read"), fs::read(&lent).expect("read"));
1652 check(&lent);
1653 fs::remove_file(alone).expect("remove");
1654 fs::remove_file(lent).expect("remove");
1655 }
1656
1657 #[test]
1660 fn stripes_fed_in_batches_write_the_same_bytes_as_prepared_whole() {
1661 let whole = path("whole");
1662 let mut writer = Writer::create(&whole, "t", fields()).expect("a file");
1663 for run in runs() {
1664 writer.append_stripe(run).expect("a stripe");
1665 }
1666 writer.finish().expect("commit");
1667
1668 let fed = path("fed");
1669 let mut writer = Writer::create(&fed, "t", fields()).expect("a file");
1670 let preparer = writer.preparer();
1671 let merger = writer.merger().expect("a merger");
1672 let nothing = Chunk::new(
1673 fields()
1674 .iter()
1675 .map(|field| Vector::from_values(field.ty.clone(), &[]).expect("a column"))
1676 .collect(),
1677 )
1678 .expect("a chunk");
1679 for mut run in runs() {
1680 let mut building = preparer.start();
1681 preparer.feed(&mut building, Vec::new()).expect("fed nothing");
1682 while !run.is_empty() {
1683 let rest = run.split_off(2.min(run.len()));
1684 let mut batch = std::mem::replace(&mut run, rest);
1685 batch.push(((u64::MAX, 0), nothing.clone()));
1686 preparer.feed(&mut building, batch).expect("fed");
1687 }
1688 let merged =
1689 merger.merge(preparer.finish(building).expect("finished")).expect("merged");
1690 let mut paged = merged.pages().expect("paged");
1691 merger.give_back(&mut paged).expect("given back");
1692 writer.write(paged).expect("written");
1693 }
1694 writer.finish().expect("commit");
1695
1696 assert_eq!(fs::read(&whole).expect("read"), fs::read(&fed).expect("read"));
1697 check(&fed);
1698 fs::remove_file(whole).expect("remove");
1699 fs::remove_file(fed).expect("remove");
1700 }
1701
1702 #[test]
1705 fn a_column_dropped_while_its_stripe_is_built_turns_to_pages() {
1706 let path = path("dropped-while-built");
1707 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1708 let preparer = writer.preparer();
1709 let merger = writer.merger().expect("a merger");
1710 let write = |writer: &mut Writer, building: Building| {
1711 let merged =
1712 merger.merge(preparer.finish(building).expect("finished")).expect("merged");
1713 let mut paged = merged.pages().expect("paged");
1714 merger.give_back(&mut paged).expect("given back");
1715 writer.write(paged).expect("written");
1716 };
1717 let mut runs = runs().into_iter();
1718 let mut first = preparer.start();
1719 let mut second = preparer.start();
1720 let mut later = runs.next().expect("a run");
1721 preparer.feed(&mut second, later.drain(..2).collect()).expect("fed");
1722 preparer.feed(&mut first, runs.next().expect("a run")).expect("fed");
1723 write(&mut writer, first);
1724 let note = |building: &Building| {
1725 matches!(building.columns[2].lock().expect("unpoisoned").body, Body::Pages(..))
1726 };
1727 assert!(!note(&second), "still coded until it is fed again");
1728 preparer.feed(&mut second, later).expect("fed");
1729 assert!(note(&second), "turned to pages once fed after the drop");
1730 write(&mut writer, second);
1731 let mut last = preparer.start();
1732 preparer.feed(&mut last, runs.next().expect("a run")).expect("fed");
1733 write(&mut writer, last);
1734 writer.finish().expect("commit");
1735 check(&path);
1736 fs::remove_file(path).expect("remove");
1737 }
1738
1739 #[test]
1741 fn a_stripe_fed_more_parts_than_it_holds_is_refused() {
1742 let path = path("overfed");
1743 let writer = Writer::create(&path, "t", fields()).expect("a file");
1744 let preparer = writer.preparer();
1745 let mut building = preparer.start();
1746 preparer.feed(&mut building, stripe(0, STRIPE_PARTS - 1)).expect("fed");
1747 assert_eq!(building.parts(), STRIPE_PARTS - 1);
1748 assert!(preparer.feed(&mut building, stripe(STRIPE_PARTS, 2)).is_err());
1749 drop(writer);
1750 let _ = fs::remove_file(path);
1751 }
1752
1753 #[test]
1756 fn stripes_merged_on_several_threads_at_once_read_back() {
1757 let path = path("merged-at-once");
1758 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1759 let preparer = writer.preparer();
1760 let merger = writer.merger().expect("a merger");
1761 let writer = Mutex::new(writer);
1762 std::thread::scope(|scope| {
1763 for run in runs() {
1764 let (preparer, merger, writer) = (&preparer, &merger, &writer);
1765 scope.spawn(move || {
1766 let merged =
1767 merger.merge(preparer.prepare(run).expect("prepared")).expect("merged");
1768 let mut paged = merged.pages().expect("paged");
1769 merger.give_back(&mut paged).expect("given back");
1770 writer.lock().expect("the writer").write(paged).expect("written");
1771 });
1772 }
1773 });
1774 writer.into_inner().expect("the writer").finish().expect("commit");
1775 check(&path);
1776 fs::remove_file(path).expect("remove");
1777 }
1778
1779 #[test]
1782 fn a_merge_after_the_table_is_closed_is_refused() {
1783 let path = path("late");
1784 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1785 let preparer = writer.preparer();
1786 let merger = writer.merger().expect("a merger");
1787 writer.finish().expect("commit");
1788 let prepared = preparer.prepare(stripe(0, 2)).expect("prepared");
1789 assert!(merger.merge(prepared).is_err());
1790 fs::remove_file(path).expect("remove");
1791 }
1792
1793 #[test]
1796 fn a_stripe_of_another_table_is_refused_at_the_merge() {
1797 let path = path("refused");
1798 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1799 let other = Writer::create(path.with_extension("other"), "u", vec![fields().remove(0)])
1800 .expect("a file");
1801 let prepared = other.preparer().prepare(vec![]).expect("nothing to prepare");
1802 assert!(writer.merge(prepared).is_err());
1803 assert_eq!(writer.table.rows, 0);
1804 drop(other);
1805 fs::remove_file(path.with_extension("other")).expect("remove");
1806 fs::remove_file(path).expect("remove");
1807 }
1808
1809 fn turning(id: usize) -> [Value; 3] {
1813 let url = match id {
1814 _ if id % 13 == 0 => Value::Null,
1815 _ if id < 5 * PART => Value::Varchar(format!("https://example.com/{}", id % 20)),
1816 _ => Value::Varchar(format!("https://example.com/page/{id}")),
1817 };
1818 [Value::BigInt(id as i64), Value::Varchar(format!("city {}", id % 13)), url]
1819 }
1820
1821 fn turning_fields() -> Vec<Field> {
1822 vec![
1823 Field::required("id", LogicalType::BigInt),
1824 Field::new("city", LogicalType::Varchar),
1825 Field::new("url", LogicalType::Varchar),
1826 ]
1827 }
1828
1829 fn turning_stripe(
1830 rows: fn(usize) -> [Value; 3],
1831 first: usize,
1832 parts: usize,
1833 ) -> Vec<((u64, u64), Chunk)> {
1834 (first..first + parts)
1835 .map(|part| {
1836 let rows = (part * PART..(part + 1) * PART).map(rows).collect::<Vec<_>>();
1837 let column = |at: usize| {
1838 let values = rows.iter().map(|row| row[at].clone()).collect::<Vec<_>>();
1839 Vector::from_values(turning_fields()[at].ty.clone(), &values).expect("a column")
1840 };
1841 let chunk = Chunk::new(vec![column(0), column(1), column(2)]).expect("a chunk");
1842 ((part as u64, 0), chunk)
1843 })
1844 .collect()
1845 }
1846
1847 fn turning_runs() -> Vec<Vec<((u64, u64), Chunk)>> {
1848 vec![
1849 turning_stripe(turning, 0, 5),
1850 turning_stripe(turning, 5, 5),
1851 turning_stripe(turning, 10, 3),
1852 ]
1853 }
1854
1855 fn check_turned(path: &PathBuf) {
1857 let reader = Reader::open(path).expect("reopen");
1858 assert_eq!(reader.parts(), 13);
1859 assert_eq!(reader.table().demoted, [false, false, true], "only url is demoted");
1860 for part in 0..13 {
1861 let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1862 let url = chunk.column(2).expect("url");
1863 assert!(url.stable_dictionary_parts().is_none(), "part {part} hands out no codes");
1864 for at in 0..PART {
1865 let want = turning(part * PART + at);
1866 for (column, value) in want.iter().enumerate() {
1867 assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1868 }
1869 }
1870 }
1871 assert_eq!(reader.distinct_values(2).expect("asked"), None);
1873 assert_eq!(reader.text_extremes(2).expect("asked"), None);
1874 assert_eq!(reader.exact_frequencies(2).expect("asked"), None);
1875 assert_eq!(reader.top_frequencies(2, 5).expect("asked"), None);
1876 assert!(!reader.skips_codes(0, 2, &[0]).expect("asked"), "no code proves a value absent");
1877 assert_eq!(reader.distinct_values(1).expect("asked"), Some(13));
1879 assert!(reader.text_extremes(1).expect("asked").is_some());
1880 }
1881
1882 #[test]
1886 fn a_column_that_turns_unique_is_demoted_and_reads_back() {
1887 let alone = path("demoted-alone");
1888 let mut writer = Writer::create(&alone, "t", turning_fields()).expect("a file");
1889 let preparer = writer.preparer();
1890 let mut runs = turning_runs().into_iter();
1891 writer.append_stripe(runs.next().expect("a run")).expect("a stripe");
1892 assert!(preparer.coded[2].load(Atomic::Relaxed), "url repeats in its first stripe");
1893 for run in runs {
1894 writer.append_stripe(run).expect("a stripe");
1895 }
1896 assert!(!preparer.coded[2].load(Atomic::Relaxed), "url was demoted");
1897 assert!(preparer.coded[1].load(Atomic::Relaxed), "city kept its dictionary");
1898 writer.finish().expect("commit");
1899 check_turned(&alone);
1900
1901 let split = path("demoted-split");
1902 let mut writer = Writer::create(&split, "t", turning_fields()).expect("a file");
1903 let preparer = writer.preparer();
1904 let prepared = turning_runs()
1905 .into_iter()
1906 .map(|run| preparer.prepare(run).expect("prepared"))
1907 .collect::<Vec<_>>();
1908 for one in prepared {
1909 writer.append_prepared(one).expect("a stripe");
1910 }
1911 writer.finish().expect("commit");
1912 check_turned(&split);
1913
1914 fs::remove_file(alone).expect("remove");
1915 fs::remove_file(split).expect("remove");
1916 }
1917
1918 fn growing(id: usize) -> [Value; 3] {
1921 let url = Value::Varchar(format!("https://example.com/{}", id / 4));
1922 [Value::BigInt(id as i64), Value::Varchar(format!("city {}", id % 13)), url]
1923 }
1924
1925 #[test]
1928 fn the_dictionary_cap_demotes_the_column_that_grew_most() {
1929 let path = path("capped");
1930 let mut writer = Writer::create(&path, "t", turning_fields())
1931 .expect("a file")
1932 .with_dictionary_cap(1 << 30);
1933 let preparer = writer.preparer();
1934 writer.append_stripe(turning_stripe(growing, 0, 5)).expect("a stripe");
1935 assert!(preparer.coded[2].load(Atomic::Relaxed), "url is under the cap");
1936 assert!(preparer.coded[1].load(Atomic::Relaxed), "city is under the cap");
1937
1938 writer.coded.cap(1);
1939 writer.append_stripe(turning_stripe(growing, 5, 5)).expect("a stripe");
1940 assert!(!preparer.coded[2].load(Atomic::Relaxed), "url grew most and was demoted");
1941 assert!(preparer.coded[1].load(Atomic::Relaxed), "city grew nothing and keeps it");
1942 writer.append_stripe(turning_stripe(growing, 10, 3)).expect("a stripe");
1943 assert!(preparer.coded[1].load(Atomic::Relaxed), "city still grows nothing");
1944 writer.finish().expect("commit");
1945
1946 let reader = Reader::open(&path).expect("reopen");
1947 assert_eq!(reader.table().demoted, [false, false, true]);
1948 for part in 0..13 {
1949 let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1950 for at in 0..PART {
1951 let want = growing(part * PART + at);
1952 for (column, value) in want.iter().enumerate() {
1953 assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1954 }
1955 }
1956 }
1957 assert_eq!(reader.distinct_values(1).expect("asked"), Some(13));
1958 assert_eq!(reader.distinct_values(2).expect("asked"), None);
1959 fs::remove_file(path).expect("remove");
1960 }
1961}