1use std::collections::HashMap;
45use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering as Atomic};
46use std::sync::{Arc, Mutex};
47
48use rudb_common::{Error, LogicalType, Result};
49use rudb_metrics::{LoadProfile, Stage};
50use rudb_storage::Range;
51use rudb_vector::{Bitmap, Chunk, Data, StringColumn, Validity, Vector};
52
53use super::{
54 ColumnStripe, DICTIONARY_CHECK_SEED, DICTIONARY_DECIDE_ROWS, DICTIONARY_DISTINCT_IN_TEN,
55 EncodedBlock, GlobalDictionary, MAX_ENCODE_WORKERS, MAX_PAGE, Part, PendingChunk, STRIPE_PARTS,
56 Spread, Unencoded, Writer, checksum, coded_page, invalid, push_validity, seeded_checksum,
57 stats, unique_codes, weight,
58};
59
60static BUSY: AtomicUsize = AtomicUsize::new(0);
67
68struct Share(usize);
70
71impl Share {
72 fn take(columns: usize, parts: usize) -> Self {
73 let busy = BUSY.fetch_add(1, Atomic::Relaxed) + 1;
74 let cores =
75 std::thread::available_parallelism().map_or(1, usize::from).min(MAX_ENCODE_WORKERS);
76 let workers = if parts <= 1 { 1 } else { (cores / busy).clamp(1, columns.max(1)) };
78 Self(workers)
79 }
80}
81
82impl Drop for Share {
83 fn drop(&mut self) {
84 BUSY.fetch_sub(1, Atomic::Relaxed);
85 }
86}
87
88#[derive(Debug, Clone)]
94pub struct Preparer {
95 types: Vec<LogicalType>,
96 coded: Arc<[AtomicBool]>,
97 profile: Option<Arc<LoadProfile>>,
98}
99
100#[derive(Debug)]
102pub struct Prepared {
103 parts: Vec<Part>,
104 types: Vec<LogicalType>,
105 columns: Vec<Column>,
106 gathers: Vec<Option<stats::Gather>>,
107 profile: Option<Arc<LoadProfile>>,
108}
109
110#[derive(Debug)]
112pub struct Merged {
113 parts: Vec<Part>,
114 columns: Vec<Merge>,
115 blocks: Vec<Unencoded>,
116 profile: Option<Arc<LoadProfile>>,
117 counted: bool,
120}
121
122#[derive(Debug)]
124pub struct Paged {
125 parts: Vec<Part>,
126 columns: Vec<ColumnStripe>,
127 blocks: Vec<(usize, usize, EncodedBlock)>,
129 counted: bool,
130}
131
132enum Built {
134 Stripe(ColumnStripe),
135 Block(EncodedBlock),
136}
137
138#[derive(Debug)]
140enum Column {
141 Pages(ColumnStripe),
143 Coded(Local),
145}
146
147#[derive(Debug)]
149enum Merge {
150 Pages(ColumnStripe),
151 Codes {
153 parts: Vec<LocalPart>,
154 global: Vec<u32>,
155 },
156 Plain(Local),
160}
161
162const END: u32 = u32::MAX;
164
165#[derive(Debug, Default)]
171struct Local {
172 first: HashMap<u64, u32, Spread>,
174 next: Vec<u32>,
176 hashes: Vec<u64>,
177 checks: Vec<u64>,
178 bytes: Vec<u8>,
180 ends: Vec<usize>,
181 counts: Vec<u64>,
183 nulls: u64,
184 parts: Vec<LocalPart>,
185}
186
187#[derive(Debug)]
189struct LocalPart {
190 codes: Vec<u32>,
191 validity: Vec<u8>,
193 range: Range,
194}
195
196impl Local {
197 fn code_column(index: usize, held: &[PendingChunk]) -> Result<Self> {
203 let mut local = Self::default();
204 for pending in held {
205 let column = pending.chunk.column(index)?;
206 let flat = column.flatten()?;
208 let mut codes = Vec::with_capacity(flat.len());
209 let mut last = None;
210 for row in 0..flat.len() {
211 let text = flat.text_at(row).unwrap_or("").as_bytes();
212 let code = match last {
215 Some(code) if local.value(code) == text => code,
216 _ => local.code(text)?,
217 };
218 last = Some(code);
219 if flat.is_null_at(row) {
220 local.nulls += 1;
221 } else {
222 local.counts[code as usize] += 1;
223 }
224 codes.push(code);
225 }
226 let mut validity = Vec::new();
227 push_validity(&mut validity, &flat);
228 local.parts.push(LocalPart { codes, validity, range: Range::of(column) });
229 }
230 local.first = HashMap::default();
233 local.next = Vec::new();
234 Ok(local)
235 }
236
237 fn rows(&self) -> Result<Vec<Vector>> {
243 self.parts
244 .iter()
245 .map(|part| {
246 let len = part.codes.len();
247 let mut column = StringColumn::with_capacity(len);
248 for &code in &part.codes {
249 column.push_bytes(self.value(code));
250 }
251 let validity = match part.validity.split_first() {
252 Some((0, _)) => Validity::AllValid,
253 Some((1, _)) => Validity::AllInvalid,
254 Some((2, bits)) => {
255 let mut mask = Bitmap::all_valid(len);
256 for row in (0..len).filter(|row| bits[row / 8] & (1 << (row % 8)) == 0) {
257 mask.set(row, false);
258 }
259 Validity::Mask(mask)
260 }
261 _ => return Err(Error::internal("a coded part has no validity")),
262 };
263 Ok(Vector::flat(LogicalType::Varchar, Data::Varlen(column))?
264 .with_validity(validity))
265 })
266 .collect()
267 }
268
269 fn values(&self) -> usize {
270 self.ends.len()
271 }
272
273 fn value(&self, code: u32) -> &[u8] {
274 let code = code as usize;
275 let from = if code == 0 { 0 } else { self.ends[code - 1] };
276 &self.bytes[from..self.ends[code]]
277 }
278
279 fn code(&mut self, text: &[u8]) -> Result<u32> {
280 let hash = checksum(text);
281 let Some(&first) = self.first.get(&hash) else {
282 let code = self.push(text, hash)?;
283 self.first.insert(hash, code);
284 return Ok(code);
285 };
286 let mut at = first;
287 loop {
288 if self.value(at) == text {
289 return Ok(at);
290 }
291 match self.next[at as usize] {
292 END => break,
293 next => at = next,
294 }
295 }
296 let code = self.push(text, hash)?;
297 self.next[at as usize] = code;
298 Ok(code)
299 }
300
301 fn push(&mut self, text: &[u8], hash: u64) -> Result<u32> {
302 let code = u32::try_from(self.ends.len())
303 .ok()
304 .filter(|&code| code != END)
305 .ok_or_else(|| invalid("a stripe has too many values in one column"))?;
306 self.bytes.extend_from_slice(text);
307 self.ends.push(self.bytes.len());
308 self.next.push(END);
309 self.hashes.push(hash);
310 self.checks.push(seeded_checksum(text, DICTIONARY_CHECK_SEED));
311 self.counts.push(0);
312 Ok(code)
313 }
314
315 fn merge_into(&self, dictionary: &mut GlobalDictionary) -> Result<Vec<u32>> {
318 let mut global = Vec::with_capacity(self.values());
319 for (code, (&hash, &check)) in self.hashes.iter().zip(&self.checks).enumerate() {
320 let text = self.value(code as u32);
321 let at = dictionary.code_hashed(text, hash, check)?;
322 let count = dictionary
323 .counts
324 .get_mut(at as usize)
325 .ok_or_else(|| invalid("global dictionary count code is out of range"))?;
326 *count = count.saturating_add(self.counts[code]);
327 global.push(at);
328 }
329 dictionary.nulls = dictionary.nulls.saturating_add(self.nulls);
330 Ok(global)
331 }
332}
333
334fn drops_dictionary(rows: usize, distinct: usize) -> bool {
358 rows >= DICTIONARY_DECIDE_ROWS
359 && distinct.saturating_mul(10) > rows.saturating_mul(DICTIONARY_DISTINCT_IN_TEN)
360}
361
362fn fan_out<T: Send>(
372 jobs: Vec<usize>,
373 workers: usize,
374 profile: Option<&LoadProfile>,
375 work: impl Fn(usize) -> Result<T> + Sync,
376) -> Result<Vec<(usize, T)>> {
377 if workers <= 1 || jobs.len() <= 1 {
378 let _span = profile.map(|profile| profile.span(Stage::Pages));
379 return jobs.into_iter().map(|index| Ok((index, work(index)?))).collect();
380 }
381 let workers = workers.min(jobs.len());
382 let queue = Mutex::new(jobs);
383 let pieces = std::thread::scope(|scope| {
384 (0..workers)
385 .map(|_| {
386 scope.spawn(|| {
387 let _span = profile.map(|profile| profile.span(Stage::Pages));
388 let mut mine = Vec::new();
389 loop {
390 let taken = queue
391 .lock()
392 .map_err(|_| Error::internal("a native encode worker panicked"))?
393 .pop();
394 let Some(index) = taken else { break };
395 mine.push((index, work(index)?));
396 }
397 Ok(mine)
398 })
399 })
400 .collect::<Vec<_>>()
401 .into_iter()
402 .map(|handle| {
403 handle.join().map_err(|_| Error::internal("a native encode worker panicked"))?
404 })
405 .collect::<Result<Vec<Vec<_>>>>()
406 })?;
407 Ok(pieces.into_iter().flatten().collect())
408}
409
410fn column_of(held: &[PendingChunk], index: usize) -> Result<Vec<&Vector>> {
412 held.iter().map(|pending| pending.chunk.column(index)).collect()
413}
414
415fn in_order<T>(width: usize, done: Vec<(usize, T)>) -> Result<Vec<T>> {
417 let mut slots: Vec<Option<T>> = (0..width).map(|_| None).collect();
418 for (index, one) in done {
419 slots[index] = Some(one);
420 }
421 slots
422 .into_iter()
423 .map(|slot| slot.ok_or_else(|| Error::internal("a column was never encoded")))
424 .collect()
425}
426
427impl Preparer {
428 pub fn prepare(&self, parts: Vec<((u64, u64), Chunk)>) -> Result<Prepared> {
439 if parts.len() > STRIPE_PARTS {
440 return Err(invalid("a stripe was handed more parts than it holds"));
441 }
442 let held = parts
443 .into_iter()
444 .filter(|(_, chunk)| !chunk.is_empty())
445 .map(|(order, chunk)| PendingChunk { order, chunk })
446 .collect::<Vec<_>>();
447 for pending in &held {
448 self.fits(&pending.chunk)?;
449 }
450 self.prepare_held(held)
451 }
452
453 fn fits(&self, chunk: &Chunk) -> Result<()> {
455 if chunk.width() != self.types.len() {
456 return Err(invalid("chunk width differs from table schema"));
457 }
458 for (index, ty) in self.types.iter().enumerate() {
459 if chunk.column(index)?.logical_type() != ty {
460 return Err(invalid("chunk type differs from table schema"));
461 }
462 }
463 Ok(())
464 }
465
466 pub(crate) fn prepare_held(&self, held: Vec<PendingChunk>) -> Result<Prepared> {
467 let width = self.types.len();
468 let key = held.first().map_or((0, 0), |pending| pending.order);
469 let share = Share::take(width, held.len());
470 let mut jobs = (0..width).collect::<Vec<_>>();
471 jobs.sort_by_key(|&index| weight(&self.types[index]));
472 let done = fan_out(jobs, share.0, self.profile.as_deref(), |index| {
473 let gather = stats::Gather::new(&self.types[index], 0)
476 .filter(|_| !held.is_empty())
477 .map(|mut gather| {
478 gather.stripe(
479 key,
480 held.iter().filter_map(|pending| pending.chunk.column(index).ok()),
481 );
482 gather
483 });
484 let column = if self.coded[index].load(Atomic::Relaxed) {
485 Column::Coded(Local::code_column(index, &held)?)
486 } else {
487 Column::Pages(Writer::encode_pages(&column_of(&held, index)?)?)
488 };
489 Ok((column, gather))
490 })?;
491 drop(share);
492 let (columns, gathers) = in_order(width, done)?.into_iter().unzip();
493 let parts = held.iter().map(Part::of).collect();
494 drop(held);
495 Ok(Prepared {
496 parts,
497 types: self.types.clone(),
498 columns,
499 gathers,
500 profile: self.profile.clone(),
501 })
502 }
503}
504
505enum Slot<'a> {
507 Owned(&'a mut Option<GlobalDictionary>, &'a mut Option<stats::Gather>),
509 Lent(&'a Mutex<LentColumn>, &'a Lent),
511}
512
513struct Step<'a> {
515 index: usize,
516 column: Column,
517 slot: Slot<'a>,
518 gather: Option<stats::Gather>,
520}
521
522impl Step<'_> {
523 fn cost(&self, coded: &[AtomicBool]) -> usize {
527 match &self.column {
528 Column::Coded(local) if coded[self.index].load(Atomic::Relaxed) => {
529 local.values().saturating_add(1)
530 }
531 _ => 0,
532 }
533 }
534
535 fn run(self, rows: usize, coded: &[AtomicBool]) -> Result<(usize, Merge, Vec<Unencoded>)> {
537 let Self { index, column, slot, gather } = self;
538 match slot {
539 Slot::Owned(dictionary, mine) => {
540 merge_column(index, column, gather, dictionary, mine, rows, coded)
541 }
542 Slot::Lent(held, lent) => {
543 let mut held = held.lock().map_err(|_| Error::internal("a merge panicked"))?;
544 if lent.reclaimed.load(Atomic::Acquire) {
547 return Err(Error::internal("a stripe was merged after its table was closed"));
548 }
549 let LentColumn { dictionary, gather: mine } = &mut *held;
550 merge_column(index, column, gather, dictionary, mine, rows, coded)
551 }
552 }
553 }
554}
555
556fn merge_column(
558 index: usize,
559 column: Column,
560 stripe: Option<stats::Gather>,
561 dictionary: &mut Option<GlobalDictionary>,
562 gather: &mut Option<stats::Gather>,
563 rows: usize,
564 coded: &[AtomicBool],
565) -> Result<(usize, Merge, Vec<Unencoded>)> {
566 if let (Some(mine), Some(stripe)) = (gather.as_mut(), stripe) {
567 mine.absorb(stripe);
568 }
569 let merge = match (column, dictionary.as_mut()) {
570 (Column::Pages(stripe), None) => Merge::Pages(stripe),
571 (Column::Pages(_), Some(_)) => {
572 return Err(Error::internal(
573 "a column with a global dictionary was prepared without one",
574 ));
575 }
576 (Column::Coded(local), None) => Merge::Plain(local),
577 (Column::Coded(local), Some(global)) => {
578 if global.values() == 0 && drops_dictionary(rows, local.values()) {
581 *dictionary = None;
582 coded[index].store(false, Atomic::Relaxed);
583 Merge::Plain(local)
584 } else {
585 let global = local.merge_into(global)?;
586 Merge::Codes { parts: local.parts, global }
587 }
588 }
589 };
590 let blocks = match dictionary {
595 Some(dictionary) => {
596 dictionary.settle()?;
597 dictionary.hand_out(index)
598 }
599 None => Vec::new(),
600 };
601 Ok((index, merge, blocks))
602}
603
604fn merge_columns(prepared: Prepared, slots: Vec<Slot<'_>>, coded: &[AtomicBool]) -> Result<Merged> {
612 let Prepared { parts, columns, gathers, profile, .. } = prepared;
613 let timing = profile.as_deref().map(|profile| profile.span(Stage::Dictionary));
614 let rows: usize = parts.iter().map(|part| part.rows).sum();
615 let width = columns.len();
616 if slots.len() != width || gathers.len() != width {
617 return Err(Error::internal("a stripe was merged into a table of another width"));
618 }
619 let mut steps = columns
620 .into_iter()
621 .zip(gathers)
622 .zip(slots)
623 .enumerate()
624 .map(|(index, ((column, gather), slot))| Step { index, column, slot, gather })
625 .collect::<Vec<_>>();
626 steps.sort_by_key(|step| step.cost(coded));
629 let workers = std::thread::available_parallelism()
630 .map_or(1, usize::from)
631 .min(MAX_ENCODE_WORKERS)
632 .min(steps.iter().filter(|step| step.cost(coded) > 0).count())
633 .max(1);
634 let done = if workers <= 1 {
635 steps.into_iter().map(|step| step.run(rows, coded)).collect::<Result<Vec<_>>>()?
636 } else {
637 let queue = Mutex::new(steps);
638 let pieces = std::thread::scope(|scope| {
639 (0..workers)
640 .map(|_| {
641 scope.spawn(|| {
642 let mut mine = Vec::new();
643 loop {
644 let taken = queue
645 .lock()
646 .map_err(|_| Error::internal("a merge worker panicked"))?
647 .pop();
648 let Some(step) = taken else { break };
649 mine.push(step.run(rows, coded)?);
650 }
651 Ok(mine)
652 })
653 })
654 .collect::<Vec<_>>()
655 .into_iter()
656 .map(|handle| {
657 handle.join().map_err(|_| Error::internal("a merge worker panicked"))?
658 })
659 .collect::<Result<Vec<Vec<_>>>>()
660 })?;
661 pieces.into_iter().flatten().collect()
662 };
663 let mut slots: Vec<Option<(Merge, Vec<Unencoded>)>> = (0..width).map(|_| None).collect();
664 for (index, merge, blocks) in done {
665 slots[index] = Some((merge, blocks));
666 }
667 let mut merged = Vec::with_capacity(width);
668 let mut blocks = Vec::new();
669 for slot in slots {
670 let (merge, handed) = slot.ok_or_else(|| Error::internal("a column was never merged"))?;
671 merged.push(merge);
672 blocks.extend(handed);
673 }
674 drop(timing);
675 Ok(Merged { parts, columns: merged, blocks, profile, counted: false })
676}
677
678#[derive(Debug)]
680pub(crate) struct Lent {
681 columns: Box<[Mutex<LentColumn>]>,
682 reclaimed: AtomicBool,
684}
685
686#[derive(Debug)]
688pub(crate) struct LentColumn {
689 pub(crate) dictionary: Option<GlobalDictionary>,
690 gather: Option<stats::Gather>,
691}
692
693impl Lent {
694 pub(crate) fn columns(&self) -> &[Mutex<LentColumn>] {
695 &self.columns
696 }
697
698 fn take_back(&self, blocks: Vec<(usize, usize, EncodedBlock)>) -> Result<()> {
700 for (column, at, block) in blocks {
701 self.columns
702 .get(column)
703 .ok_or_else(|| Error::internal("a dictionary block came back to no column"))?
704 .lock()
705 .map_err(|_| Error::internal("a merge panicked"))?
706 .dictionary
707 .as_mut()
708 .ok_or_else(|| Error::internal("a dictionary block came back to no dictionary"))?
709 .take_back(at, block)?;
710 }
711 Ok(())
712 }
713
714 #[allow(clippy::type_complexity)]
716 pub(crate) fn reclaim(
717 &self,
718 ) -> Result<(Vec<Option<GlobalDictionary>>, Vec<Option<stats::Gather>>)> {
719 self.reclaimed.store(true, Atomic::Release);
720 let mut dictionaries = Vec::with_capacity(self.columns.len());
721 let mut gathers = Vec::with_capacity(self.columns.len());
722 for column in &self.columns {
723 let mut held = column.lock().map_err(|_| Error::internal("a merge panicked"))?;
724 dictionaries.push(held.dictionary.take());
725 gathers.push(held.gather.take());
726 }
727 Ok((dictionaries, gathers))
728 }
729}
730
731#[derive(Debug, Clone)]
738pub struct Merger {
739 lent: Arc<Lent>,
740 types: Vec<LogicalType>,
741 coded: Arc<[AtomicBool]>,
742}
743
744impl Merger {
745 pub fn merge(&self, prepared: Prepared) -> Result<Merged> {
751 if prepared.types != self.types {
752 return Err(invalid("a stripe was prepared for a table of other columns"));
753 }
754 let slots = self.lent.columns.iter().map(|column| Slot::Lent(column, &self.lent)).collect();
755 merge_columns(prepared, slots, &self.coded)
756 }
757
758 pub fn give_back(&self, paged: &mut Paged) -> Result<()> {
765 self.lent.take_back(std::mem::take(&mut paged.blocks))
766 }
767}
768
769impl Merged {
770 pub fn pages(self) -> Result<Paged> {
778 let Self { parts, columns, blocks, profile, counted } = self;
779 let width = columns.len();
780 let mut jobs = (width..width + blocks.len())
784 .chain((0..width).filter(|&index| !matches!(columns[index], Merge::Pages(_))))
785 .collect::<Vec<_>>();
786 jobs.sort_by_key(|&index| index < width && matches!(columns[index], Merge::Plain(_)));
788 let share = Share::take(jobs.len(), parts.len());
789 let built = fan_out(jobs, share.0, profile.as_deref(), |index| {
790 let Some(column) = columns.get(index) else {
791 return Ok(Built::Block(blocks[index - width].encode()?));
792 };
793 Ok(Built::Stripe(match column {
794 Merge::Codes { parts, global } => code_pages(parts, global)?,
795 Merge::Plain(local) => {
796 Writer::encode_pages(&local.rows()?.iter().collect::<Vec<_>>())?
797 }
798 Merge::Pages(_) => {
799 return Err(Error::internal("a finished column was queued to be built"));
800 }
801 }))
802 })?;
803 drop(share);
804 let mut slots: Vec<Option<ColumnStripe>> = (0..width).map(|_| None).collect();
805 let mut encoded = Vec::with_capacity(blocks.len());
806 for (index, one) in built {
807 match one {
808 Built::Stripe(stripe) => slots[index] = Some(stripe),
809 Built::Block(block) => {
810 let (column, at) = blocks[index - width].place();
811 encoded.push((column, at, block));
812 }
813 }
814 }
815 let columns = columns
816 .into_iter()
817 .zip(slots)
818 .map(|(column, slot)| match (column, slot) {
819 (Merge::Pages(stripe), _) | (_, Some(stripe)) => Ok(stripe),
820 _ => Err(Error::internal("a column was never encoded")),
821 })
822 .collect::<Result<Vec<_>>>()?;
823 Ok(Paged { parts, columns, blocks: encoded, counted })
824 }
825}
826
827fn code_pages(parts: &[LocalPart], global: &[u32]) -> Result<ColumnStripe> {
829 let mut stripe = ColumnStripe {
830 pages: Vec::with_capacity(parts.len()),
831 codes: Vec::with_capacity(parts.len()),
832 sieves: Vec::with_capacity(parts.len()),
833 ranges: Vec::with_capacity(parts.len()),
834 };
835 for part in parts {
836 let codes = part
837 .codes
838 .iter()
839 .map(|&code| global.get(code as usize).copied())
840 .collect::<Option<Vec<_>>>()
841 .ok_or_else(|| Error::internal("a stripe's code has no global code"))?;
842 let bytes = coded_page(&codes, &part.validity)?;
843 if bytes.len() > MAX_PAGE {
844 return Err(invalid("column page exceeds the configured bound"));
845 }
846 stripe.pages.push(bytes);
847 stripe.codes.push(Some(unique_codes(&codes)));
848 stripe.sieves.push(None);
852 stripe.ranges.push(part.range.clone());
853 }
854 Ok(stripe)
855}
856
857impl Writer {
858 #[must_use]
863 pub fn preparer(&self) -> Preparer {
864 Preparer {
865 types: self.table.fields.iter().map(|field| field.ty.clone()).collect(),
866 coded: Arc::clone(&self.coded),
867 profile: self.profile.clone(),
868 }
869 }
870
871 pub fn merge(&mut self, prepared: Prepared) -> Result<Merged> {
882 self.flush_pending()?;
883 if prepared.columns.len() != self.table.fields.len()
884 || prepared.types.iter().ne(self.table.fields.iter().map(|field| &field.ty))
885 {
886 return Err(invalid("a stripe was prepared for a table of other columns"));
887 }
888 self.table.rows = prepared
889 .parts
890 .iter()
891 .try_fold(self.table.rows, |rows, part| rows.checked_add(part.rows))
892 .ok_or_else(|| invalid("row count overflow"))?;
893 self.merge_held(prepared)
894 }
895
896 pub(crate) fn merge_held(&mut self, prepared: Prepared) -> Result<Merged> {
905 let slots = match &self.lent {
906 Some(lent) => lent.columns.iter().map(|column| Slot::Lent(column, lent)).collect(),
907 None => self
908 .dictionaries
909 .iter_mut()
910 .zip(self.gathers.iter_mut())
911 .map(|(dictionary, gather)| Slot::Owned(dictionary, gather))
912 .collect::<Vec<_>>(),
913 };
914 let mut merged = merge_columns(prepared, slots, &self.coded)?;
915 merged.counted = true;
916 Ok(merged)
917 }
918
919 pub fn merger(&mut self) -> Result<Merger> {
929 self.flush_pending()?;
930 let lent = match &self.lent {
931 Some(lent) => Arc::clone(lent),
932 None => {
933 let lent = Arc::new(Lent {
934 columns: std::mem::take(&mut self.dictionaries)
935 .into_iter()
936 .zip(std::mem::take(&mut self.gathers))
937 .map(|(dictionary, gather)| Mutex::new(LentColumn { dictionary, gather }))
938 .collect(),
939 reclaimed: AtomicBool::new(false),
940 });
941 self.lent = Some(Arc::clone(&lent));
942 lent
943 }
944 };
945 Ok(Merger {
946 lent,
947 types: self.table.fields.iter().map(|field| field.ty.clone()).collect(),
948 coded: Arc::clone(&self.coded),
949 })
950 }
951
952 pub fn write(&mut self, paged: Paged) -> Result<()> {
958 self.write_paged(paged)
959 }
960
961 pub(crate) fn write_paged(&mut self, paged: Paged) -> Result<()> {
962 let Paged { parts, columns, blocks, counted } = paged;
963 if !counted {
964 self.table.rows = parts
965 .iter()
966 .try_fold(self.table.rows, |rows, part| rows.checked_add(part.rows))
967 .ok_or_else(|| invalid("row count overflow"))?;
968 }
969 if let Some(lent) = &self.lent {
970 lent.take_back(blocks)?;
971 } else {
972 for (column, at, block) in blocks {
973 self.dictionaries
974 .get_mut(column)
975 .and_then(Option::as_mut)
976 .ok_or_else(|| {
977 Error::internal("a dictionary block came back to no dictionary")
978 })?
979 .take_back(at, block)?;
980 }
981 }
982 if parts.is_empty() {
983 return self.place_blocks();
984 }
985 self.write_stripe(&parts, columns)
986 }
987
988 pub fn append_prepared(&mut self, prepared: Prepared) -> Result<()> {
994 let merged = self.merge(prepared)?;
995 let paged = merged.pages()?;
996 self.write(paged)
997 }
998}
999
1000#[cfg(test)]
1001mod tests {
1002 use std::fs;
1003 use std::path::PathBuf;
1004 use std::time::{SystemTime, UNIX_EPOCH};
1005
1006 use rudb_common::{Field, Value};
1007 use rudb_vector::Vector;
1008
1009 use super::*;
1010 use crate::Reader;
1011
1012 const PART: usize = 1_000;
1013
1014 fn path(label: &str) -> PathBuf {
1015 let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
1016 std::env::temp_dir()
1017 .join(format!("rudb-prepare-{label}-{}-{stamp}.rdb", std::process::id()))
1018 }
1019
1020 fn fields() -> Vec<Field> {
1021 vec![
1022 Field::required("id", LogicalType::BigInt),
1023 Field::new("city", LogicalType::Varchar),
1024 Field::new("note", LogicalType::Varchar),
1025 ]
1026 }
1027
1028 fn row(id: usize) -> [Value; 3] {
1033 let city = if id % 11 == 0 {
1034 Value::Null
1035 } else {
1036 Value::Varchar(format!("city {}", (id / 7) % 13))
1037 };
1038 let note = if id % 17 == 0 { Value::Null } else { Value::Varchar(format!("note {id}")) };
1039 [Value::BigInt(id as i64), city, note]
1040 }
1041
1042 fn stripe(first: usize, parts: usize) -> Vec<((u64, u64), Chunk)> {
1044 (first..first + parts)
1045 .map(|part| {
1046 let rows = (part * PART..(part + 1) * PART).map(row).collect::<Vec<_>>();
1047 let column = |at: usize| {
1048 let values = rows.iter().map(|row| row[at].clone()).collect::<Vec<_>>();
1049 Vector::from_values(fields()[at].ty.clone(), &values).expect("a column")
1050 };
1051 let chunk = Chunk::new(vec![column(0), column(1), column(2)]).expect("a chunk");
1052 ((part as u64, 0), chunk)
1053 })
1054 .collect()
1055 }
1056
1057 fn runs() -> Vec<Vec<((u64, u64), Chunk)>> {
1059 vec![stripe(5, 5), stripe(0, 5), stripe(10, 3)]
1060 }
1061
1062 fn check(path: &PathBuf) {
1063 let reader = Reader::open(path).expect("reopen");
1064 assert_eq!(reader.parts(), 13);
1065 for part in 0..13 {
1066 let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1067 for at in [0, 17, PART - 1] {
1068 let want = row(part * PART + at);
1069 for (column, value) in want.iter().enumerate() {
1070 assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1071 }
1072 }
1073 }
1074 }
1075
1076 #[test]
1084 fn stripes_prepared_before_any_is_merged_write_the_same_bytes_as_one_at_a_time() {
1085 let alone = path("alone");
1086 let mut writer = Writer::create(&alone, "t", fields()).expect("a file");
1087 for run in runs() {
1088 writer.append_stripe(run).expect("a stripe");
1089 }
1090 writer.finish().expect("commit");
1091
1092 let split = path("split");
1093 let mut writer = Writer::create(&split, "t", fields()).expect("a file");
1094 let preparer = writer.preparer();
1095 let prepared = runs()
1096 .into_iter()
1097 .map(|run| preparer.prepare(run).expect("prepared"))
1098 .collect::<Vec<_>>();
1099 for one in prepared {
1100 writer.append_prepared(one).expect("a stripe");
1101 }
1102 assert!(!preparer.coded[2].load(Atomic::Relaxed), "note lost its dictionary");
1103 assert!(preparer.coded[1].load(Atomic::Relaxed), "city kept its dictionary");
1104 writer.finish().expect("commit");
1105
1106 assert_eq!(fs::read(&alone).expect("read"), fs::read(&split).expect("read"));
1107 check(&split);
1108 fs::remove_file(alone).expect("remove");
1109 fs::remove_file(split).expect("remove");
1110 }
1111
1112 #[test]
1116 fn stripes_written_in_another_order_than_they_were_merged_read_back() {
1117 let path = path("crossed");
1118 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1119 let preparer = writer.preparer();
1120 let mut merged = runs()
1121 .into_iter()
1122 .map(|run| writer.merge(preparer.prepare(run).expect("prepared")).expect("merged"))
1123 .map(|merged| merged.pages().expect("paged"))
1124 .collect::<Vec<_>>();
1125 merged.reverse();
1126 for paged in merged {
1127 writer.write(paged).expect("written");
1128 }
1129 writer.finish().expect("commit");
1130 check(&path);
1131 fs::remove_file(path).expect("remove");
1132 }
1133
1134 #[test]
1137 fn stripes_merged_through_a_merger_write_the_same_bytes_as_the_writer() {
1138 let alone = path("alone-merger");
1139 let mut writer = Writer::create(&alone, "t", fields()).expect("a file");
1140 for run in runs() {
1141 writer.append_stripe(run).expect("a stripe");
1142 }
1143 writer.finish().expect("commit");
1144
1145 let lent = path("lent");
1146 let mut writer = Writer::create(&lent, "t", fields()).expect("a file");
1147 let preparer = writer.preparer();
1148 let merger = writer.merger().expect("a merger");
1149 for run in runs() {
1150 let merged = merger.merge(preparer.prepare(run).expect("prepared")).expect("merged");
1151 let mut paged = merged.pages().expect("paged");
1152 merger.give_back(&mut paged).expect("given back");
1153 writer.write(paged).expect("written");
1154 }
1155 assert_eq!(writer.table.rows, 13 * PART);
1156 writer.finish().expect("commit");
1157
1158 assert_eq!(fs::read(&alone).expect("read"), fs::read(&lent).expect("read"));
1159 check(&lent);
1160 fs::remove_file(alone).expect("remove");
1161 fs::remove_file(lent).expect("remove");
1162 }
1163
1164 #[test]
1167 fn stripes_merged_on_several_threads_at_once_read_back() {
1168 let path = path("merged-at-once");
1169 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1170 let preparer = writer.preparer();
1171 let merger = writer.merger().expect("a merger");
1172 let writer = Mutex::new(writer);
1173 std::thread::scope(|scope| {
1174 for run in runs() {
1175 let (preparer, merger, writer) = (&preparer, &merger, &writer);
1176 scope.spawn(move || {
1177 let merged =
1178 merger.merge(preparer.prepare(run).expect("prepared")).expect("merged");
1179 let mut paged = merged.pages().expect("paged");
1180 merger.give_back(&mut paged).expect("given back");
1181 writer.lock().expect("the writer").write(paged).expect("written");
1182 });
1183 }
1184 });
1185 writer.into_inner().expect("the writer").finish().expect("commit");
1186 check(&path);
1187 fs::remove_file(path).expect("remove");
1188 }
1189
1190 #[test]
1193 fn a_merge_after_the_table_is_closed_is_refused() {
1194 let path = path("late");
1195 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1196 let preparer = writer.preparer();
1197 let merger = writer.merger().expect("a merger");
1198 writer.finish().expect("commit");
1199 let prepared = preparer.prepare(stripe(0, 2)).expect("prepared");
1200 assert!(merger.merge(prepared).is_err());
1201 fs::remove_file(path).expect("remove");
1202 }
1203
1204 #[test]
1207 fn a_stripe_of_another_table_is_refused_at_the_merge() {
1208 let path = path("refused");
1209 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1210 let other = Writer::create(path.with_extension("other"), "u", vec![fields().remove(0)])
1211 .expect("a file");
1212 let prepared = other.preparer().prepare(vec![]).expect("nothing to prepare");
1213 assert!(writer.merge(prepared).is_err());
1214 assert_eq!(writer.table.rows, 0);
1215 drop(other);
1216 fs::remove_file(path.with_extension("other")).expect("remove");
1217 fs::remove_file(path).expect("remove");
1218 }
1219}