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}
93
94impl Coding {
95 pub(crate) fn new(flags: impl IntoIterator<Item = bool>) -> Self {
96 let flags = flags.into_iter().map(AtomicBool::new).collect::<Box<[_]>>();
97 let growth = flags.iter().map(|_| AtomicU64::new(0)).collect();
98 Self { flags, growth, held: AtomicU64::new(0), cap: AtomicU64::new(DICTIONARY_CAP_BYTES) }
99 }
100
101 pub(crate) fn cap(&self, bytes: u64) {
103 self.cap.store(bytes, Atomic::Relaxed);
104 }
105
106 fn recount(&self, before: u64, now: u64) -> u64 {
109 if now >= before {
110 self.held.fetch_add(now - before, Atomic::Relaxed) + (now - before)
111 } else {
112 self.held.fetch_sub(before - now, Atomic::Relaxed).saturating_sub(before - now)
113 }
114 }
115
116 fn grew_most(&self, index: usize) -> bool {
123 let mine = self.growth[index].load(Atomic::Relaxed);
124 mine > 0 && self.growth.iter().all(|other| other.load(Atomic::Relaxed) <= mine)
125 }
126}
127
128impl Deref for Coding {
129 type Target = [AtomicBool];
130
131 fn deref(&self) -> &[AtomicBool] {
132 &self.flags
133 }
134}
135
136struct Share(usize);
138
139impl Share {
140 fn take(columns: usize, parts: usize) -> Self {
141 let busy = BUSY.fetch_add(1, Atomic::Relaxed) + 1;
142 let cores =
143 std::thread::available_parallelism().map_or(1, usize::from).min(MAX_ENCODE_WORKERS);
144 let workers = if parts <= 1 { 1 } else { (cores / busy).clamp(1, columns.max(1)) };
146 Self(workers)
147 }
148}
149
150impl Drop for Share {
151 fn drop(&mut self) {
152 BUSY.fetch_sub(1, Atomic::Relaxed);
153 }
154}
155
156#[derive(Debug, Clone)]
162pub struct Preparer {
163 types: Vec<LogicalType>,
164 coded: Arc<Coding>,
165 profile: Option<Arc<LoadProfile>>,
166}
167
168#[derive(Debug)]
170pub struct Prepared {
171 parts: Vec<Part>,
172 types: Vec<LogicalType>,
173 columns: Vec<Column>,
174 gathers: Vec<Option<stats::Gather>>,
175 profile: Option<Arc<LoadProfile>>,
176}
177
178#[derive(Debug)]
180pub struct Merged {
181 parts: Vec<Part>,
182 columns: Vec<Merge>,
183 blocks: Vec<Unencoded>,
184 profile: Option<Arc<LoadProfile>>,
185 counted: bool,
188}
189
190#[derive(Debug)]
192pub struct Paged {
193 parts: Vec<Part>,
194 columns: Vec<ColumnStripe>,
195 blocks: Vec<(usize, usize, EncodedBlock)>,
197 counted: bool,
198}
199
200enum Built {
202 Stripe(ColumnStripe),
203 Block(EncodedBlock),
204}
205
206#[derive(Debug)]
208enum Column {
209 Pages(ColumnStripe),
211 Coded(Local),
213}
214
215#[derive(Debug)]
217enum Merge {
218 Pages(ColumnStripe),
219 Codes {
221 parts: Vec<LocalPart>,
222 global: Vec<u32>,
223 },
224 Plain(Local),
228}
229
230const END: u32 = u32::MAX;
232
233#[derive(Debug, Default)]
239struct Local {
240 first: HashMap<u64, u32, Spread>,
242 next: Vec<u32>,
244 hashes: Vec<u64>,
245 checks: Vec<u64>,
246 bytes: Vec<u8>,
248 ends: Vec<usize>,
249 counts: Vec<u64>,
251 nulls: u64,
252 parts: Vec<LocalPart>,
253 blob: bool,
255}
256
257#[derive(Debug)]
259struct LocalPart {
260 codes: Vec<u32>,
261 validity: Vec<u8>,
263 range: Range,
264}
265
266impl Local {
267 #[cfg(test)]
269 fn code_column(index: usize, held: &[PendingChunk]) -> Result<Self> {
270 let mut local = Self::default();
271 let mut mapped = None;
272 for pending in held {
273 local.code_part(pending.chunk.column(index)?, &mut mapped)?;
274 }
275 local.done();
276 Ok(local)
277 }
278
279 fn code_part(
288 &mut self,
289 column: &Vector,
290 mapped: &mut Option<(Arc<Vector>, Vec<u32>)>,
291 ) -> Result<()> {
292 self.blob = column.logical_type() == &LogicalType::Blob;
293 if let Some(codes) = self.code_dictionary(column, mapped)? {
294 let mut validity = Vec::new();
295 push_validity(&mut validity, column);
296 self.parts.push(LocalPart { codes, validity, range: Range::of(column) });
297 return Ok(());
298 }
299 let flat = column.flatten()?;
301 let mut codes = Vec::with_capacity(flat.len());
302 let mut last = None;
303 for row in 0..flat.len() {
304 let text = flat.bytes_at(row).unwrap_or(b"");
307 let code = match last {
310 Some(code) if self.value(code) == text => code,
311 _ => self.code(text)?,
312 };
313 last = Some(code);
314 if flat.is_null_at(row) {
315 self.nulls += 1;
316 } else {
317 self.counts[code as usize] += 1;
318 }
319 codes.push(code);
320 }
321 let mut validity = Vec::new();
322 push_validity(&mut validity, &flat);
323 self.parts.push(LocalPart { codes, validity, range: Range::of(column) });
324 Ok(())
325 }
326
327 fn done(&mut self) {
329 self.first = HashMap::default();
332 self.next = Vec::new();
333 }
334
335 fn code_dictionary(
350 &mut self,
351 column: &Vector,
352 mapped: &mut Option<(Arc<Vector>, Vec<u32>)>,
353 ) -> Result<Option<Vec<u32>>> {
354 let Some((codes, values)) = column.shared_dictionary_parts() else { return Ok(None) };
355 if !matches!(values.validity(), Validity::AllValid) {
356 return Ok(None);
357 }
358 let Some(codes) = codes.get(..column.len()) else { return Ok(None) };
359 let fresh = !matches!(mapped, Some((held, _)) if Arc::ptr_eq(held, values));
360 if fresh {
361 *mapped = Some((Arc::clone(values), vec![END; values.len()]));
362 }
363 let Some((_, map)) = mapped.as_mut() else { return Ok(None) };
364 let every = matches!(column.validity(), Validity::AllValid);
365 let mut coded = Vec::with_capacity(codes.len());
366 for (row, &code) in codes.iter().enumerate() {
367 if !every && !column.validity().is_valid(row) {
368 let code = self.code(b"")?;
370 self.nulls += 1;
371 coded.push(code);
372 continue;
373 }
374 let slot = map
375 .get_mut(code as usize)
376 .ok_or_else(|| invalid("a dictionary code is out of range"))?;
377 if *slot == END {
378 *slot = self.code(values.bytes_at(code as usize).unwrap_or(b""))?;
379 }
380 self.counts[*slot as usize] += 1;
381 coded.push(*slot);
382 }
383 Ok(Some(coded))
384 }
385
386 fn rows(&self) -> Result<Vec<Vector>> {
392 self.parts
393 .iter()
394 .map(|part| {
395 let len = part.codes.len();
396 let mut column = StringColumn::with_capacity(len);
397 for &code in &part.codes {
398 column.push_bytes(self.value(code));
399 }
400 let validity = match part.validity.split_first() {
401 Some((0, _)) => Validity::AllValid,
402 Some((1, _)) => Validity::AllInvalid,
403 Some((2, bits)) => {
404 let mut mask = Bitmap::all_valid(len);
405 for row in (0..len).filter(|row| bits[row / 8] & (1 << (row % 8)) == 0) {
406 mask.set(row, false);
407 }
408 Validity::Mask(mask)
409 }
410 _ => return Err(Error::internal("a coded part has no validity")),
411 };
412 let ty = if self.blob { LogicalType::Blob } else { LogicalType::Varchar };
413 Ok(Vector::flat(ty, Data::Varlen(column))?.with_validity(validity))
414 })
415 .collect()
416 }
417
418 fn values(&self) -> usize {
419 self.ends.len()
420 }
421
422 fn value(&self, code: u32) -> &[u8] {
423 let code = code as usize;
424 let from = if code == 0 { 0 } else { self.ends[code - 1] };
425 &self.bytes[from..self.ends[code]]
426 }
427
428 fn code(&mut self, text: &[u8]) -> Result<u32> {
429 let hash = checksum(text);
430 let Some(&first) = self.first.get(&hash) else {
431 let code = self.push(text, hash)?;
432 self.first.insert(hash, code);
433 return Ok(code);
434 };
435 let mut at = first;
436 loop {
437 if self.value(at) == text {
438 return Ok(at);
439 }
440 match self.next[at as usize] {
441 END => break,
442 next => at = next,
443 }
444 }
445 let code = self.push(text, hash)?;
446 self.next[at as usize] = code;
447 Ok(code)
448 }
449
450 fn push(&mut self, text: &[u8], hash: u64) -> Result<u32> {
451 let code = u32::try_from(self.ends.len())
452 .ok()
453 .filter(|&code| code != END)
454 .ok_or_else(|| invalid("a stripe has too many values in one column"))?;
455 self.bytes.extend_from_slice(text);
456 self.ends.push(self.bytes.len());
457 self.next.push(END);
458 self.hashes.push(hash);
459 self.checks.push(seeded_checksum(text, DICTIONARY_CHECK_SEED));
460 self.counts.push(0);
461 Ok(code)
462 }
463
464 fn merge_into(&self, dictionary: &mut GlobalDictionary) -> Result<Vec<u32>> {
467 let mut global = Vec::with_capacity(self.values());
468 for (code, (&hash, &check)) in self.hashes.iter().zip(&self.checks).enumerate() {
469 let text = self.value(code as u32);
470 let at = dictionary.code_hashed(text, hash, check)?;
471 let count = dictionary
472 .counts
473 .get_mut(at as usize)
474 .ok_or_else(|| invalid("global dictionary count code is out of range"))?;
475 *count = count.saturating_add(self.counts[code]);
476 global.push(at);
477 }
478 dictionary.nulls = dictionary.nulls.saturating_add(self.nulls);
479 Ok(global)
480 }
481}
482
483fn drops_dictionary(rows: usize, distinct: usize) -> bool {
507 rows >= DICTIONARY_DECIDE_ROWS
508 && distinct.saturating_mul(10) > rows.saturating_mul(DICTIONARY_DISTINCT_IN_TEN)
509}
510
511fn demotes(rows: usize, new: usize, total: u64, coding: &Coding, index: usize) -> bool {
524 drops_dictionary(rows, new)
525 || (total > coding.cap.load(Atomic::Relaxed) && coding.grew_most(index))
526}
527
528fn fan_out<T: Send>(
538 jobs: Vec<usize>,
539 workers: usize,
540 profile: Option<&LoadProfile>,
541 work: impl Fn(usize) -> Result<T> + Sync,
542) -> Result<Vec<(usize, T)>> {
543 if workers <= 1 || jobs.len() <= 1 {
544 let _span = profile.map(|profile| profile.span(Stage::Pages));
545 return jobs.into_iter().map(|index| Ok((index, work(index)?))).collect();
546 }
547 let workers = workers.min(jobs.len());
548 let queue = Mutex::new(jobs);
549 let pieces = std::thread::scope(|scope| {
550 (0..workers)
551 .map(|_| {
552 scope.spawn(|| {
553 let _span = profile.map(|profile| profile.span(Stage::Pages));
554 let mut mine = Vec::new();
555 loop {
556 let taken = queue
557 .lock()
558 .map_err(|_| Error::internal("a native encode worker panicked"))?
559 .pop();
560 let Some(index) = taken else { break };
561 mine.push((index, work(index)?));
562 }
563 Ok(mine)
564 })
565 })
566 .collect::<Vec<_>>()
567 .into_iter()
568 .map(|handle| {
569 handle.join().map_err(|_| Error::internal("a native encode worker panicked"))?
570 })
571 .collect::<Result<Vec<Vec<_>>>>()
572 })?;
573 Ok(pieces.into_iter().flatten().collect())
574}
575
576impl Preparer {
577 pub fn prepare(&self, parts: Vec<((u64, u64), Chunk)>) -> Result<Prepared> {
588 if parts.len() > STRIPE_PARTS {
589 return Err(invalid("a stripe was handed more parts than it holds"));
590 }
591 let held = parts
592 .into_iter()
593 .filter(|(_, chunk)| !chunk.is_empty())
594 .map(|(order, chunk)| PendingChunk { order, chunk })
595 .collect::<Vec<_>>();
596 for pending in &held {
597 self.fits(&pending.chunk)?;
598 }
599 self.prepare_held(held)
600 }
601
602 fn fits(&self, chunk: &Chunk) -> Result<()> {
604 if chunk.width() != self.types.len() {
605 return Err(invalid("chunk width differs from table schema"));
606 }
607 for (index, ty) in self.types.iter().enumerate() {
608 if chunk.column(index)?.logical_type() != ty {
609 return Err(invalid("chunk type differs from table schema"));
610 }
611 }
612 Ok(())
613 }
614
615 pub(crate) fn prepare_held(&self, held: Vec<PendingChunk>) -> Result<Prepared> {
616 let mut building = self.start();
617 self.feed_held(&mut building, held)?;
618 self.finish(building)
619 }
620
621 #[must_use]
634 pub fn start(&self) -> Building {
635 let columns = (0..self.types.len())
636 .map(|index| {
637 let body = if self.coded[index].load(Atomic::Relaxed) {
638 Body::Coded(Local::default(), None)
639 } else {
640 Body::Pages(ColumnStripe::default(), Settling::default())
641 };
642 let gather = stats::Gather::new(&self.types[index], 0);
643 Mutex::new(Growing { body, gather })
644 })
645 .collect();
646 Building { parts: Vec::new(), columns }
647 }
648
649 pub fn feed(&self, building: &mut Building, parts: Vec<((u64, u64), Chunk)>) -> Result<()> {
659 if building.parts.len().saturating_add(parts.len()) > STRIPE_PARTS {
660 return Err(invalid("a stripe was handed more parts than it holds"));
661 }
662 let held = parts
663 .into_iter()
664 .filter(|(_, chunk)| !chunk.is_empty())
665 .map(|(order, chunk)| PendingChunk { order, chunk })
666 .collect::<Vec<_>>();
667 for pending in &held {
668 self.fits(&pending.chunk)?;
669 }
670 self.feed_held(building, held)
671 }
672
673 fn feed_held(&self, building: &mut Building, held: Vec<PendingChunk>) -> Result<()> {
674 if held.is_empty() {
675 return Ok(());
676 }
677 let width = self.types.len();
678 let opening = building.parts.is_empty();
680 let key = held.first().map_or((0, 0), |pending| pending.order);
681 let share = Share::take(width, held.len());
682 let mut jobs = (0..width).collect::<Vec<_>>();
683 jobs.sort_by_key(|&index| weight(&self.types[index]));
684 let columns = &building.columns;
685 fan_out(jobs, share.0, self.profile.as_deref(), |index| {
686 let mut growing = columns[index]
687 .lock()
688 .map_err(|_| Error::internal("a native encode worker panicked"))?;
689 let Growing { body, gather } = &mut *growing;
690 if matches!(body, Body::Coded(..)) && !self.coded[index].load(Atomic::Relaxed) {
691 body.plain()?;
692 }
693 if let Some(gather) = gather.as_mut() {
694 if opening {
697 gather.open_stripe(key);
698 }
699 for pending in &held {
700 gather.part(pending.chunk.column(index)?);
701 }
702 }
703 match body {
704 Body::Coded(local, mapped) => {
705 for pending in &held {
706 local.code_part(pending.chunk.column(index)?, mapped)?;
707 }
708 }
709 Body::Pages(stripe, settling) => {
710 for pending in &held {
711 Writer::encode_page(stripe, settling, pending.chunk.column(index)?)?;
712 }
713 }
714 }
715 Ok(())
716 })?;
717 drop(share);
718 building.parts.extend(held.iter().map(Part::of));
719 Ok(())
720 }
721
722 pub fn finish(&self, building: Building) -> Result<Prepared> {
728 let empty = building.parts.is_empty();
729 let (columns, gathers) = building
730 .columns
731 .into_iter()
732 .map(|growing| {
733 let Growing { body, gather } = growing
734 .into_inner()
735 .map_err(|_| Error::internal("a native encode worker panicked"))?;
736 let column = match body {
737 Body::Coded(mut local, _) => {
738 local.done();
739 Column::Coded(local)
740 }
741 Body::Pages(stripe, _) => Column::Pages(stripe),
742 };
743 let gather = gather.filter(|_| !empty).map(|mut gather| {
745 gather.close_stripe();
746 gather
747 });
748 Ok((column, gather))
749 })
750 .collect::<Result<Vec<_>>>()?
751 .into_iter()
752 .unzip();
753 Ok(Prepared {
754 parts: building.parts,
755 types: self.types.clone(),
756 columns,
757 gathers,
758 profile: self.profile.clone(),
759 })
760 }
761}
762
763#[derive(Debug)]
765pub struct Building {
766 parts: Vec<Part>,
767 columns: Vec<Mutex<Growing>>,
770}
771
772impl Building {
773 #[must_use]
775 pub fn parts(&self) -> usize {
776 self.parts.len()
777 }
778}
779
780#[derive(Debug)]
782struct Growing {
783 body: Body,
784 gather: Option<stats::Gather>,
785}
786
787#[derive(Debug)]
789enum Body {
790 Pages(ColumnStripe, Settling),
792 Coded(Local, Option<(Arc<Vector>, Vec<u32>)>),
795}
796
797impl Body {
798 fn plain(&mut self) -> Result<()> {
805 let Self::Coded(local, _) = self else { return Ok(()) };
806 let mut stripe = ColumnStripe::default();
807 let mut settling = Settling::default();
808 for rows in local.rows()? {
809 Writer::encode_page(&mut stripe, &mut settling, &rows)?;
810 }
811 *self = Self::Pages(stripe, settling);
812 Ok(())
813 }
814}
815
816enum Slot<'a> {
818 Owned(&'a mut Option<GlobalDictionary>, &'a mut Option<stats::Gather>),
820 Lent(&'a Mutex<LentColumn>, &'a Lent),
822}
823
824struct Step<'a> {
826 index: usize,
827 column: Column,
828 slot: Slot<'a>,
829 gather: Option<stats::Gather>,
831}
832
833impl Step<'_> {
834 fn cost(&self, coded: &Coding) -> usize {
838 match &self.column {
839 Column::Coded(local) if coded[self.index].load(Atomic::Relaxed) => {
840 local.values().saturating_add(1)
841 }
842 _ => 0,
843 }
844 }
845
846 fn run(
848 self,
849 rows: usize,
850 coded: &Coding,
851 profile: Option<&LoadProfile>,
852 ) -> Result<(usize, Merge, Vec<Unencoded>)> {
853 let Self { index, column, slot, gather } = self;
854 match slot {
855 Slot::Owned(dictionary, mine) => {
856 merge_column(index, column, gather, dictionary, mine, rows, coded, profile)
857 }
858 Slot::Lent(held, lent) => {
859 let mut held = held.lock().map_err(|_| Error::internal("a merge panicked"))?;
860 if lent.reclaimed.load(Atomic::Acquire) {
863 return Err(Error::internal("a stripe was merged after its table was closed"));
864 }
865 let LentColumn { dictionary, gather: mine } = &mut *held;
866 merge_column(index, column, gather, dictionary, mine, rows, coded, profile)
867 }
868 }
869 }
870}
871
872#[expect(clippy::too_many_arguments, reason = "one column's share of the stripe's merge state")]
874fn merge_column(
875 index: usize,
876 column: Column,
877 stripe: Option<stats::Gather>,
878 dictionary: &mut Option<GlobalDictionary>,
879 gather: &mut Option<stats::Gather>,
880 rows: usize,
881 coded: &Coding,
882 profile: Option<&LoadProfile>,
883) -> Result<(usize, Merge, Vec<Unencoded>)> {
884 if let (Some(mine), Some(stripe)) = (gather.as_mut(), stripe) {
885 mine.absorb(stripe);
886 }
887 let mut new = None;
888 let merge = match (column, dictionary.as_mut()) {
889 (Column::Pages(stripe), None) => Merge::Pages(stripe),
890 (Column::Pages(stripe), Some(global)) if global.demoted => Merge::Pages(stripe),
891 (Column::Pages(_), Some(_)) => {
892 return Err(Error::internal(
893 "a column with a global dictionary was prepared without one",
894 ));
895 }
896 (Column::Coded(local), None) => Merge::Plain(local),
897 (Column::Coded(local), Some(global)) if global.demoted => Merge::Plain(local),
899 (Column::Coded(local), Some(global)) => {
900 if global.values() == 0 && drops_dictionary(rows, local.values()) {
903 if let Some(profile) = profile {
904 profile.release(global.charged);
905 }
906 coded.recount(global.charged, 0);
907 *dictionary = None;
908 coded[index].store(false, Atomic::Relaxed);
909 Merge::Plain(local)
910 } else {
911 let before = global.values();
912 let codes = local.merge_into(global)?;
913 new = Some(global.values() - before);
914 Merge::Codes { parts: local.parts, global: codes }
915 }
916 }
917 };
918 if let (Some(new), Some(global)) = (new, dictionary.as_mut()) {
920 let now = global.held_bytes();
921 coded.growth[index].store(now.saturating_sub(global.charged), Atomic::Relaxed);
922 let total =
923 coded.held.load(Atomic::Relaxed).saturating_add(now).saturating_sub(global.charged);
924 if demotes(rows, new, total, coded, index) {
925 global.demote();
926 coded[index].store(false, Atomic::Relaxed);
927 coded.growth[index].store(0, Atomic::Relaxed);
928 }
929 }
930 let blocks = match dictionary {
935 Some(dictionary) => {
936 dictionary.settle()?;
937 let blocks = dictionary.hand_out(index);
938 let (before, now) = dictionary.recharge(profile);
939 coded.recount(before, now);
940 blocks
941 }
942 None => Vec::new(),
943 };
944 Ok((index, merge, blocks))
945}
946
947fn merge_columns(prepared: Prepared, slots: Vec<Slot<'_>>, coded: &Coding) -> Result<Merged> {
955 let Prepared { parts, columns, gathers, profile, .. } = prepared;
956 let timing = profile.as_deref().map(|profile| profile.span(Stage::Dictionary));
957 let rows: usize = parts.iter().map(|part| part.rows).sum();
958 let width = columns.len();
959 if slots.len() != width || gathers.len() != width {
960 return Err(Error::internal("a stripe was merged into a table of another width"));
961 }
962 let mut steps = columns
963 .into_iter()
964 .zip(gathers)
965 .zip(slots)
966 .enumerate()
967 .map(|(index, ((column, gather), slot))| Step { index, column, slot, gather })
968 .collect::<Vec<_>>();
969 steps.sort_by_key(|step| step.cost(coded));
972 let workers = std::thread::available_parallelism()
973 .map_or(1, usize::from)
974 .min(MAX_ENCODE_WORKERS)
975 .min(steps.iter().filter(|step| step.cost(coded) > 0).count())
976 .max(1);
977 let done = if workers <= 1 {
978 steps
979 .into_iter()
980 .map(|step| step.run(rows, coded, profile.as_deref()))
981 .collect::<Result<Vec<_>>>()?
982 } else {
983 let queue = Mutex::new(steps);
984 let pieces = std::thread::scope(|scope| {
985 (0..workers)
986 .map(|_| {
987 scope.spawn(|| {
988 let mut mine = Vec::new();
989 loop {
990 let taken = queue
991 .lock()
992 .map_err(|_| Error::internal("a merge worker panicked"))?
993 .pop();
994 let Some(step) = taken else { break };
995 mine.push(step.run(rows, coded, profile.as_deref())?);
996 }
997 Ok(mine)
998 })
999 })
1000 .collect::<Vec<_>>()
1001 .into_iter()
1002 .map(|handle| {
1003 handle.join().map_err(|_| Error::internal("a merge worker panicked"))?
1004 })
1005 .collect::<Result<Vec<Vec<_>>>>()
1006 })?;
1007 pieces.into_iter().flatten().collect()
1008 };
1009 let mut slots: Vec<Option<(Merge, Vec<Unencoded>)>> = (0..width).map(|_| None).collect();
1010 for (index, merge, blocks) in done {
1011 slots[index] = Some((merge, blocks));
1012 }
1013 let mut merged = Vec::with_capacity(width);
1014 let mut blocks = Vec::new();
1015 for slot in slots {
1016 let (merge, handed) = slot.ok_or_else(|| Error::internal("a column was never merged"))?;
1017 merged.push(merge);
1018 blocks.extend(handed);
1019 }
1020 drop(timing);
1021 Ok(Merged { parts, columns: merged, blocks, profile, counted: false })
1022}
1023
1024#[derive(Debug)]
1026pub(crate) struct Lent {
1027 columns: Box<[Mutex<LentColumn>]>,
1028 reclaimed: AtomicBool,
1030}
1031
1032#[derive(Debug)]
1034pub(crate) struct LentColumn {
1035 pub(crate) dictionary: Option<GlobalDictionary>,
1036 gather: Option<stats::Gather>,
1037}
1038
1039impl Lent {
1040 pub(crate) fn columns(&self) -> &[Mutex<LentColumn>] {
1041 &self.columns
1042 }
1043
1044 fn take_back(&self, blocks: Vec<(usize, usize, EncodedBlock)>) -> Result<()> {
1046 for (column, at, block) in blocks {
1047 self.columns
1048 .get(column)
1049 .ok_or_else(|| Error::internal("a dictionary block came back to no column"))?
1050 .lock()
1051 .map_err(|_| Error::internal("a merge panicked"))?
1052 .dictionary
1053 .as_mut()
1054 .ok_or_else(|| Error::internal("a dictionary block came back to no dictionary"))?
1055 .take_back(at, block)?;
1056 }
1057 Ok(())
1058 }
1059
1060 #[allow(clippy::type_complexity)]
1062 pub(crate) fn reclaim(
1063 &self,
1064 ) -> Result<(Vec<Option<GlobalDictionary>>, Vec<Option<stats::Gather>>)> {
1065 self.reclaimed.store(true, Atomic::Release);
1066 let mut dictionaries = Vec::with_capacity(self.columns.len());
1067 let mut gathers = Vec::with_capacity(self.columns.len());
1068 for column in &self.columns {
1069 let mut held = column.lock().map_err(|_| Error::internal("a merge panicked"))?;
1070 dictionaries.push(held.dictionary.take());
1071 gathers.push(held.gather.take());
1072 }
1073 Ok((dictionaries, gathers))
1074 }
1075}
1076
1077#[derive(Debug, Clone)]
1084pub struct Merger {
1085 lent: Arc<Lent>,
1086 types: Vec<LogicalType>,
1087 coded: Arc<Coding>,
1088}
1089
1090impl Merger {
1091 pub fn merge(&self, prepared: Prepared) -> Result<Merged> {
1097 if prepared.types != self.types {
1098 return Err(invalid("a stripe was prepared for a table of other columns"));
1099 }
1100 let slots = self.lent.columns.iter().map(|column| Slot::Lent(column, &self.lent)).collect();
1101 merge_columns(prepared, slots, &self.coded)
1102 }
1103
1104 pub fn give_back(&self, paged: &mut Paged) -> Result<()> {
1111 self.lent.take_back(std::mem::take(&mut paged.blocks))
1112 }
1113}
1114
1115impl Merged {
1116 pub fn pages(self) -> Result<Paged> {
1124 let Self { parts, columns, blocks, profile, counted } = self;
1125 let width = columns.len();
1126 let mut jobs = (width..width + blocks.len())
1130 .chain((0..width).filter(|&index| !matches!(columns[index], Merge::Pages(_))))
1131 .collect::<Vec<_>>();
1132 jobs.sort_by_key(|&index| index < width && matches!(columns[index], Merge::Plain(_)));
1134 let share = Share::take(jobs.len(), parts.len());
1135 let built = fan_out(jobs, share.0, profile.as_deref(), |index| {
1136 let Some(column) = columns.get(index) else {
1137 return Ok(Built::Block(blocks[index - width].encode()?));
1138 };
1139 Ok(Built::Stripe(match column {
1140 Merge::Codes { parts, global } => code_pages(parts, global)?,
1141 Merge::Plain(local) => {
1142 Writer::encode_pages(&local.rows()?.iter().collect::<Vec<_>>())?
1143 }
1144 Merge::Pages(_) => {
1145 return Err(Error::internal("a finished column was queued to be built"));
1146 }
1147 }))
1148 })?;
1149 drop(share);
1150 let mut slots: Vec<Option<ColumnStripe>> = (0..width).map(|_| None).collect();
1151 let mut encoded = Vec::with_capacity(blocks.len());
1152 for (index, one) in built {
1153 match one {
1154 Built::Stripe(stripe) => slots[index] = Some(stripe),
1155 Built::Block(block) => {
1156 let (column, at) = blocks[index - width].place();
1157 encoded.push((column, at, block));
1158 }
1159 }
1160 }
1161 let columns = columns
1162 .into_iter()
1163 .zip(slots)
1164 .map(|(column, slot)| match (column, slot) {
1165 (Merge::Pages(stripe), _) | (_, Some(stripe)) => Ok(stripe),
1166 _ => Err(Error::internal("a column was never encoded")),
1167 })
1168 .collect::<Result<Vec<_>>>()?;
1169 Ok(Paged { parts, columns, blocks: encoded, counted })
1170 }
1171}
1172
1173fn code_pages(parts: &[LocalPart], global: &[u32]) -> Result<ColumnStripe> {
1175 let mut stripe = ColumnStripe {
1176 pages: Vec::with_capacity(parts.len()),
1177 codes: Vec::with_capacity(parts.len()),
1178 sieves: Vec::with_capacity(parts.len()),
1179 ranges: Vec::with_capacity(parts.len()),
1180 };
1181 for part in parts {
1182 let codes = part
1183 .codes
1184 .iter()
1185 .map(|&code| global.get(code as usize).copied())
1186 .collect::<Option<Vec<_>>>()
1187 .ok_or_else(|| Error::internal("a stripe's code has no global code"))?;
1188 let bytes = coded_page(&codes, &part.validity)?;
1189 if bytes.len() > MAX_PAGE {
1190 return Err(invalid("column page exceeds the configured bound"));
1191 }
1192 stripe.pages.push(bytes);
1193 stripe.codes.push(Some(unique_codes(&codes)));
1194 stripe.sieves.push(None);
1198 stripe.ranges.push(part.range.clone());
1199 }
1200 Ok(stripe)
1201}
1202
1203impl Writer {
1204 #[must_use]
1209 pub fn preparer(&self) -> Preparer {
1210 Preparer {
1211 types: self.table.fields.iter().map(|field| field.ty.clone()).collect(),
1212 coded: Arc::clone(&self.coded),
1213 profile: self.profile.clone(),
1214 }
1215 }
1216
1217 pub fn merge(&mut self, prepared: Prepared) -> Result<Merged> {
1228 self.flush_pending()?;
1229 if prepared.columns.len() != self.table.fields.len()
1230 || prepared.types.iter().ne(self.table.fields.iter().map(|field| &field.ty))
1231 {
1232 return Err(invalid("a stripe was prepared for a table of other columns"));
1233 }
1234 self.table.rows = prepared
1235 .parts
1236 .iter()
1237 .try_fold(self.table.rows, |rows, part| rows.checked_add(part.rows))
1238 .ok_or_else(|| invalid("row count overflow"))?;
1239 self.merge_held(prepared)
1240 }
1241
1242 pub(crate) fn merge_held(&mut self, prepared: Prepared) -> Result<Merged> {
1251 let slots = match &self.lent {
1252 Some(lent) => lent.columns.iter().map(|column| Slot::Lent(column, lent)).collect(),
1253 None => self
1254 .dictionaries
1255 .iter_mut()
1256 .zip(self.gathers.iter_mut())
1257 .map(|(dictionary, gather)| Slot::Owned(dictionary, gather))
1258 .collect::<Vec<_>>(),
1259 };
1260 let mut merged = merge_columns(prepared, slots, &self.coded)?;
1261 merged.counted = true;
1262 Ok(merged)
1263 }
1264
1265 pub fn merger(&mut self) -> Result<Merger> {
1275 self.flush_pending()?;
1276 let lent = match &self.lent {
1277 Some(lent) => Arc::clone(lent),
1278 None => {
1279 let lent = Arc::new(Lent {
1280 columns: std::mem::take(&mut self.dictionaries)
1281 .into_iter()
1282 .zip(std::mem::take(&mut self.gathers))
1283 .map(|(dictionary, gather)| Mutex::new(LentColumn { dictionary, gather }))
1284 .collect(),
1285 reclaimed: AtomicBool::new(false),
1286 });
1287 self.lent = Some(Arc::clone(&lent));
1288 lent
1289 }
1290 };
1291 Ok(Merger {
1292 lent,
1293 types: self.table.fields.iter().map(|field| field.ty.clone()).collect(),
1294 coded: Arc::clone(&self.coded),
1295 })
1296 }
1297
1298 pub fn write(&mut self, paged: Paged) -> Result<()> {
1304 self.write_paged(paged)
1305 }
1306
1307 pub(crate) fn write_paged(&mut self, paged: Paged) -> Result<()> {
1308 let Paged { parts, columns, blocks, counted } = paged;
1309 if !counted {
1310 self.table.rows = parts
1311 .iter()
1312 .try_fold(self.table.rows, |rows, part| rows.checked_add(part.rows))
1313 .ok_or_else(|| invalid("row count overflow"))?;
1314 }
1315 if let Some(lent) = &self.lent {
1316 lent.take_back(blocks)?;
1317 } else {
1318 for (column, at, block) in blocks {
1319 self.dictionaries
1320 .get_mut(column)
1321 .and_then(Option::as_mut)
1322 .ok_or_else(|| {
1323 Error::internal("a dictionary block came back to no dictionary")
1324 })?
1325 .take_back(at, block)?;
1326 }
1327 }
1328 if parts.is_empty() {
1329 return self.place_blocks();
1330 }
1331 self.write_stripe(&parts, columns)
1332 }
1333
1334 pub fn append_prepared(&mut self, prepared: Prepared) -> Result<()> {
1340 let merged = self.merge(prepared)?;
1341 let paged = merged.pages()?;
1342 self.write(paged)
1343 }
1344}
1345
1346#[cfg(test)]
1347mod tests {
1348 use std::fs;
1349 use std::path::PathBuf;
1350 use std::time::{SystemTime, UNIX_EPOCH};
1351
1352 use rudb_common::{Field, Value};
1353 use rudb_vector::Vector;
1354
1355 use super::*;
1356 use crate::Reader;
1357
1358 const PART: usize = 1_000;
1359
1360 #[test]
1366 fn dictionary_parts_code_as_their_rows_would() {
1367 let texts = |values: &[&str]| {
1368 Arc::new(
1369 Vector::from_values(
1370 LogicalType::Varchar,
1371 &values
1372 .iter()
1373 .map(|text| Value::Varchar((*text).to_string()))
1374 .collect::<Vec<_>>(),
1375 )
1376 .expect("a dictionary"),
1377 )
1378 };
1379 let shared = texts(&["b", "a", "", "c", "unused"]);
1380 let other = texts(&["c", "d", "a"]);
1381 let mut nulls = Bitmap::all_valid(6);
1382 nulls.set(1, false);
1383 nulls.set(4, false);
1384 let parts = [
1385 Vector::dictionary_over(vec![3, 3, 1, 0, 2, 1], Arc::clone(&shared)).expect("codes"),
1386 Vector::dictionary_over(vec![0, 3, 1, 1, 3, 2], Arc::clone(&shared))
1387 .expect("codes")
1388 .with_validity(Validity::Mask(nulls)),
1389 Vector::dictionary_over(vec![1, 2, 0, 1], other).expect("codes"),
1390 ];
1391 let held = |flat: bool| {
1392 parts
1393 .iter()
1394 .enumerate()
1395 .map(|(at, part)| PendingChunk {
1396 order: (at as u64, 0),
1397 chunk: Chunk::new(vec![if flat {
1398 part.flatten().expect("flat")
1399 } else {
1400 part.clone()
1401 }])
1402 .expect("a chunk"),
1403 })
1404 .collect::<Vec<_>>()
1405 };
1406 let parquet = held(false);
1407 let before = rudb_common::slow::here();
1408 let coded = Local::code_column(0, &parquet).expect("coded");
1409 assert_eq!(
1410 rudb_common::slow::here().since(before).get(rudb_common::slow::Cause::Flatten),
1411 0,
1412 "a part that came in as codes was flattened",
1413 );
1414 let flat = Local::code_column(0, &held(true)).expect("coded");
1415 assert_eq!(coded.values(), flat.values());
1416 for code in 0..flat.values() as u32 {
1417 assert_eq!(coded.value(code), flat.value(code), "value {code}");
1418 }
1419 assert_eq!(coded.counts, flat.counts);
1420 assert_eq!(coded.nulls, flat.nulls);
1421 assert_eq!(coded.nulls, 2);
1422 assert_eq!(coded.hashes, flat.hashes);
1423 assert_eq!(coded.checks, flat.checks);
1424 assert_eq!(coded.parts.len(), flat.parts.len());
1425 for (coded, flat) in coded.parts.iter().zip(&flat.parts) {
1426 assert_eq!(coded.codes, flat.codes);
1427 assert_eq!(coded.validity, flat.validity);
1428 }
1429 }
1430
1431 fn path(label: &str) -> PathBuf {
1432 let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
1433 std::env::temp_dir()
1434 .join(format!("rudb-prepare-{label}-{}-{stamp}.rdb", std::process::id()))
1435 }
1436
1437 fn fields() -> Vec<Field> {
1438 vec![
1439 Field::required("id", LogicalType::BigInt),
1440 Field::new("city", LogicalType::Varchar),
1441 Field::new("note", LogicalType::Varchar),
1442 ]
1443 }
1444
1445 fn row(id: usize) -> [Value; 3] {
1450 let city = if id % 11 == 0 {
1451 Value::Null
1452 } else {
1453 Value::Varchar(format!("city {}", (id / 7) % 13))
1454 };
1455 let note = if id % 17 == 0 { Value::Null } else { Value::Varchar(format!("note {id}")) };
1456 [Value::BigInt(id as i64), city, note]
1457 }
1458
1459 fn stripe(first: usize, parts: usize) -> Vec<((u64, u64), Chunk)> {
1461 (first..first + parts)
1462 .map(|part| {
1463 let rows = (part * PART..(part + 1) * PART).map(row).collect::<Vec<_>>();
1464 let column = |at: usize| {
1465 let values = rows.iter().map(|row| row[at].clone()).collect::<Vec<_>>();
1466 Vector::from_values(fields()[at].ty.clone(), &values).expect("a column")
1467 };
1468 let chunk = Chunk::new(vec![column(0), column(1), column(2)]).expect("a chunk");
1469 ((part as u64, 0), chunk)
1470 })
1471 .collect()
1472 }
1473
1474 fn runs() -> Vec<Vec<((u64, u64), Chunk)>> {
1476 vec![stripe(5, 5), stripe(0, 5), stripe(10, 3)]
1477 }
1478
1479 fn check(path: &PathBuf) {
1480 let reader = Reader::open(path).expect("reopen");
1481 assert_eq!(reader.parts(), 13);
1482 for part in 0..13 {
1483 let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1484 for at in [0, 17, PART - 1] {
1485 let want = row(part * PART + at);
1486 for (column, value) in want.iter().enumerate() {
1487 assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1488 }
1489 }
1490 }
1491 }
1492
1493 #[test]
1501 fn stripes_prepared_before_any_is_merged_write_the_same_bytes_as_one_at_a_time() {
1502 let alone = path("alone");
1503 let mut writer = Writer::create(&alone, "t", fields()).expect("a file");
1504 for run in runs() {
1505 writer.append_stripe(run).expect("a stripe");
1506 }
1507 writer.finish().expect("commit");
1508
1509 let split = path("split");
1510 let mut writer = Writer::create(&split, "t", fields()).expect("a file");
1511 let preparer = writer.preparer();
1512 let prepared = runs()
1513 .into_iter()
1514 .map(|run| preparer.prepare(run).expect("prepared"))
1515 .collect::<Vec<_>>();
1516 for one in prepared {
1517 writer.append_prepared(one).expect("a stripe");
1518 }
1519 assert!(!preparer.coded[2].load(Atomic::Relaxed), "note lost its dictionary");
1520 assert!(preparer.coded[1].load(Atomic::Relaxed), "city kept its dictionary");
1521 writer.finish().expect("commit");
1522
1523 assert_eq!(fs::read(&alone).expect("read"), fs::read(&split).expect("read"));
1524 check(&split);
1525 fs::remove_file(alone).expect("remove");
1526 fs::remove_file(split).expect("remove");
1527 }
1528
1529 #[test]
1533 fn stripes_written_in_another_order_than_they_were_merged_read_back() {
1534 let path = path("crossed");
1535 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1536 let preparer = writer.preparer();
1537 let mut merged = runs()
1538 .into_iter()
1539 .map(|run| writer.merge(preparer.prepare(run).expect("prepared")).expect("merged"))
1540 .map(|merged| merged.pages().expect("paged"))
1541 .collect::<Vec<_>>();
1542 merged.reverse();
1543 for paged in merged {
1544 writer.write(paged).expect("written");
1545 }
1546 writer.finish().expect("commit");
1547 check(&path);
1548 fs::remove_file(path).expect("remove");
1549 }
1550
1551 #[test]
1554 fn stripes_merged_through_a_merger_write_the_same_bytes_as_the_writer() {
1555 let alone = path("alone-merger");
1556 let mut writer = Writer::create(&alone, "t", fields()).expect("a file");
1557 for run in runs() {
1558 writer.append_stripe(run).expect("a stripe");
1559 }
1560 writer.finish().expect("commit");
1561
1562 let lent = path("lent");
1563 let mut writer = Writer::create(&lent, "t", fields()).expect("a file");
1564 let preparer = writer.preparer();
1565 let merger = writer.merger().expect("a merger");
1566 for run in runs() {
1567 let merged = merger.merge(preparer.prepare(run).expect("prepared")).expect("merged");
1568 let mut paged = merged.pages().expect("paged");
1569 merger.give_back(&mut paged).expect("given back");
1570 writer.write(paged).expect("written");
1571 }
1572 assert_eq!(writer.table.rows, 13 * PART);
1573 writer.finish().expect("commit");
1574
1575 assert_eq!(fs::read(&alone).expect("read"), fs::read(&lent).expect("read"));
1576 check(&lent);
1577 fs::remove_file(alone).expect("remove");
1578 fs::remove_file(lent).expect("remove");
1579 }
1580
1581 #[test]
1584 fn stripes_fed_in_batches_write_the_same_bytes_as_prepared_whole() {
1585 let whole = path("whole");
1586 let mut writer = Writer::create(&whole, "t", fields()).expect("a file");
1587 for run in runs() {
1588 writer.append_stripe(run).expect("a stripe");
1589 }
1590 writer.finish().expect("commit");
1591
1592 let fed = path("fed");
1593 let mut writer = Writer::create(&fed, "t", fields()).expect("a file");
1594 let preparer = writer.preparer();
1595 let merger = writer.merger().expect("a merger");
1596 let nothing = Chunk::new(
1597 fields()
1598 .iter()
1599 .map(|field| Vector::from_values(field.ty.clone(), &[]).expect("a column"))
1600 .collect(),
1601 )
1602 .expect("a chunk");
1603 for mut run in runs() {
1604 let mut building = preparer.start();
1605 preparer.feed(&mut building, Vec::new()).expect("fed nothing");
1606 while !run.is_empty() {
1607 let rest = run.split_off(2.min(run.len()));
1608 let mut batch = std::mem::replace(&mut run, rest);
1609 batch.push(((u64::MAX, 0), nothing.clone()));
1610 preparer.feed(&mut building, batch).expect("fed");
1611 }
1612 let merged =
1613 merger.merge(preparer.finish(building).expect("finished")).expect("merged");
1614 let mut paged = merged.pages().expect("paged");
1615 merger.give_back(&mut paged).expect("given back");
1616 writer.write(paged).expect("written");
1617 }
1618 writer.finish().expect("commit");
1619
1620 assert_eq!(fs::read(&whole).expect("read"), fs::read(&fed).expect("read"));
1621 check(&fed);
1622 fs::remove_file(whole).expect("remove");
1623 fs::remove_file(fed).expect("remove");
1624 }
1625
1626 #[test]
1629 fn a_column_dropped_while_its_stripe_is_built_turns_to_pages() {
1630 let path = path("dropped-while-built");
1631 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1632 let preparer = writer.preparer();
1633 let merger = writer.merger().expect("a merger");
1634 let write = |writer: &mut Writer, building: Building| {
1635 let merged =
1636 merger.merge(preparer.finish(building).expect("finished")).expect("merged");
1637 let mut paged = merged.pages().expect("paged");
1638 merger.give_back(&mut paged).expect("given back");
1639 writer.write(paged).expect("written");
1640 };
1641 let mut runs = runs().into_iter();
1642 let mut first = preparer.start();
1643 let mut second = preparer.start();
1644 let mut later = runs.next().expect("a run");
1645 preparer.feed(&mut second, later.drain(..2).collect()).expect("fed");
1646 preparer.feed(&mut first, runs.next().expect("a run")).expect("fed");
1647 write(&mut writer, first);
1648 let note = |building: &Building| {
1649 matches!(building.columns[2].lock().expect("unpoisoned").body, Body::Pages(..))
1650 };
1651 assert!(!note(&second), "still coded until it is fed again");
1652 preparer.feed(&mut second, later).expect("fed");
1653 assert!(note(&second), "turned to pages once fed after the drop");
1654 write(&mut writer, second);
1655 let mut last = preparer.start();
1656 preparer.feed(&mut last, runs.next().expect("a run")).expect("fed");
1657 write(&mut writer, last);
1658 writer.finish().expect("commit");
1659 check(&path);
1660 fs::remove_file(path).expect("remove");
1661 }
1662
1663 #[test]
1665 fn a_stripe_fed_more_parts_than_it_holds_is_refused() {
1666 let path = path("overfed");
1667 let writer = Writer::create(&path, "t", fields()).expect("a file");
1668 let preparer = writer.preparer();
1669 let mut building = preparer.start();
1670 preparer.feed(&mut building, stripe(0, STRIPE_PARTS - 1)).expect("fed");
1671 assert_eq!(building.parts(), STRIPE_PARTS - 1);
1672 assert!(preparer.feed(&mut building, stripe(STRIPE_PARTS, 2)).is_err());
1673 drop(writer);
1674 let _ = fs::remove_file(path);
1675 }
1676
1677 #[test]
1680 fn stripes_merged_on_several_threads_at_once_read_back() {
1681 let path = path("merged-at-once");
1682 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1683 let preparer = writer.preparer();
1684 let merger = writer.merger().expect("a merger");
1685 let writer = Mutex::new(writer);
1686 std::thread::scope(|scope| {
1687 for run in runs() {
1688 let (preparer, merger, writer) = (&preparer, &merger, &writer);
1689 scope.spawn(move || {
1690 let merged =
1691 merger.merge(preparer.prepare(run).expect("prepared")).expect("merged");
1692 let mut paged = merged.pages().expect("paged");
1693 merger.give_back(&mut paged).expect("given back");
1694 writer.lock().expect("the writer").write(paged).expect("written");
1695 });
1696 }
1697 });
1698 writer.into_inner().expect("the writer").finish().expect("commit");
1699 check(&path);
1700 fs::remove_file(path).expect("remove");
1701 }
1702
1703 #[test]
1706 fn a_merge_after_the_table_is_closed_is_refused() {
1707 let path = path("late");
1708 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1709 let preparer = writer.preparer();
1710 let merger = writer.merger().expect("a merger");
1711 writer.finish().expect("commit");
1712 let prepared = preparer.prepare(stripe(0, 2)).expect("prepared");
1713 assert!(merger.merge(prepared).is_err());
1714 fs::remove_file(path).expect("remove");
1715 }
1716
1717 #[test]
1720 fn a_stripe_of_another_table_is_refused_at_the_merge() {
1721 let path = path("refused");
1722 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1723 let other = Writer::create(path.with_extension("other"), "u", vec![fields().remove(0)])
1724 .expect("a file");
1725 let prepared = other.preparer().prepare(vec![]).expect("nothing to prepare");
1726 assert!(writer.merge(prepared).is_err());
1727 assert_eq!(writer.table.rows, 0);
1728 drop(other);
1729 fs::remove_file(path.with_extension("other")).expect("remove");
1730 fs::remove_file(path).expect("remove");
1731 }
1732
1733 fn turning(id: usize) -> [Value; 3] {
1737 let url = match id {
1738 _ if id % 13 == 0 => Value::Null,
1739 _ if id < 5 * PART => Value::Varchar(format!("https://example.com/{}", id % 20)),
1740 _ => Value::Varchar(format!("https://example.com/page/{id}")),
1741 };
1742 [Value::BigInt(id as i64), Value::Varchar(format!("city {}", id % 13)), url]
1743 }
1744
1745 fn turning_fields() -> Vec<Field> {
1746 vec![
1747 Field::required("id", LogicalType::BigInt),
1748 Field::new("city", LogicalType::Varchar),
1749 Field::new("url", LogicalType::Varchar),
1750 ]
1751 }
1752
1753 fn turning_stripe(
1754 rows: fn(usize) -> [Value; 3],
1755 first: usize,
1756 parts: usize,
1757 ) -> Vec<((u64, u64), Chunk)> {
1758 (first..first + parts)
1759 .map(|part| {
1760 let rows = (part * PART..(part + 1) * PART).map(rows).collect::<Vec<_>>();
1761 let column = |at: usize| {
1762 let values = rows.iter().map(|row| row[at].clone()).collect::<Vec<_>>();
1763 Vector::from_values(turning_fields()[at].ty.clone(), &values).expect("a column")
1764 };
1765 let chunk = Chunk::new(vec![column(0), column(1), column(2)]).expect("a chunk");
1766 ((part as u64, 0), chunk)
1767 })
1768 .collect()
1769 }
1770
1771 fn turning_runs() -> Vec<Vec<((u64, u64), Chunk)>> {
1772 vec![
1773 turning_stripe(turning, 0, 5),
1774 turning_stripe(turning, 5, 5),
1775 turning_stripe(turning, 10, 3),
1776 ]
1777 }
1778
1779 fn check_turned(path: &PathBuf) {
1781 let reader = Reader::open(path).expect("reopen");
1782 assert_eq!(reader.parts(), 13);
1783 assert_eq!(reader.table().demoted, [false, false, true], "only url is demoted");
1784 for part in 0..13 {
1785 let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1786 let url = chunk.column(2).expect("url");
1787 assert!(url.stable_dictionary_parts().is_none(), "part {part} hands out no codes");
1788 for at in 0..PART {
1789 let want = turning(part * PART + at);
1790 for (column, value) in want.iter().enumerate() {
1791 assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1792 }
1793 }
1794 }
1795 assert_eq!(reader.distinct_values(2).expect("asked"), None);
1797 assert_eq!(reader.text_extremes(2).expect("asked"), None);
1798 assert_eq!(reader.exact_frequencies(2).expect("asked"), None);
1799 assert_eq!(reader.top_frequencies(2, 5).expect("asked"), None);
1800 assert!(!reader.skips_codes(0, 2, &[0]).expect("asked"), "no code proves a value absent");
1801 assert_eq!(reader.distinct_values(1).expect("asked"), Some(13));
1803 assert!(reader.text_extremes(1).expect("asked").is_some());
1804 }
1805
1806 #[test]
1810 fn a_column_that_turns_unique_is_demoted_and_reads_back() {
1811 let alone = path("demoted-alone");
1812 let mut writer = Writer::create(&alone, "t", turning_fields()).expect("a file");
1813 let preparer = writer.preparer();
1814 let mut runs = turning_runs().into_iter();
1815 writer.append_stripe(runs.next().expect("a run")).expect("a stripe");
1816 assert!(preparer.coded[2].load(Atomic::Relaxed), "url repeats in its first stripe");
1817 for run in runs {
1818 writer.append_stripe(run).expect("a stripe");
1819 }
1820 assert!(!preparer.coded[2].load(Atomic::Relaxed), "url was demoted");
1821 assert!(preparer.coded[1].load(Atomic::Relaxed), "city kept its dictionary");
1822 writer.finish().expect("commit");
1823 check_turned(&alone);
1824
1825 let split = path("demoted-split");
1826 let mut writer = Writer::create(&split, "t", turning_fields()).expect("a file");
1827 let preparer = writer.preparer();
1828 let prepared = turning_runs()
1829 .into_iter()
1830 .map(|run| preparer.prepare(run).expect("prepared"))
1831 .collect::<Vec<_>>();
1832 for one in prepared {
1833 writer.append_prepared(one).expect("a stripe");
1834 }
1835 writer.finish().expect("commit");
1836 check_turned(&split);
1837
1838 fs::remove_file(alone).expect("remove");
1839 fs::remove_file(split).expect("remove");
1840 }
1841
1842 fn growing(id: usize) -> [Value; 3] {
1845 let url = Value::Varchar(format!("https://example.com/{}", id / 4));
1846 [Value::BigInt(id as i64), Value::Varchar(format!("city {}", id % 13)), url]
1847 }
1848
1849 #[test]
1852 fn the_dictionary_cap_demotes_the_column_that_grew_most() {
1853 let path = path("capped");
1854 let mut writer = Writer::create(&path, "t", turning_fields())
1855 .expect("a file")
1856 .with_dictionary_cap(1 << 30);
1857 let preparer = writer.preparer();
1858 writer.append_stripe(turning_stripe(growing, 0, 5)).expect("a stripe");
1859 assert!(preparer.coded[2].load(Atomic::Relaxed), "url is under the cap");
1860 assert!(preparer.coded[1].load(Atomic::Relaxed), "city is under the cap");
1861
1862 writer.coded.cap(1);
1863 writer.append_stripe(turning_stripe(growing, 5, 5)).expect("a stripe");
1864 assert!(!preparer.coded[2].load(Atomic::Relaxed), "url grew most and was demoted");
1865 assert!(preparer.coded[1].load(Atomic::Relaxed), "city grew nothing and keeps it");
1866 writer.append_stripe(turning_stripe(growing, 10, 3)).expect("a stripe");
1867 assert!(preparer.coded[1].load(Atomic::Relaxed), "city still grows nothing");
1868 writer.finish().expect("commit");
1869
1870 let reader = Reader::open(&path).expect("reopen");
1871 assert_eq!(reader.table().demoted, [false, false, true]);
1872 for part in 0..13 {
1873 let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1874 for at in 0..PART {
1875 let want = growing(part * PART + at);
1876 for (column, value) in want.iter().enumerate() {
1877 assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1878 }
1879 }
1880 }
1881 assert_eq!(reader.distinct_values(1).expect("asked"), Some(13));
1882 assert_eq!(reader.distinct_values(2).expect("asked"), None);
1883 fs::remove_file(path).expect("remove");
1884 }
1885}