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.bytes_at(row).unwrap_or(b"");
214 let code = match last {
217 Some(code) if local.value(code) == text => code,
218 _ => local.code(text)?,
219 };
220 last = Some(code);
221 if flat.is_null_at(row) {
222 local.nulls += 1;
223 } else {
224 local.counts[code as usize] += 1;
225 }
226 codes.push(code);
227 }
228 let mut validity = Vec::new();
229 push_validity(&mut validity, &flat);
230 local.parts.push(LocalPart { codes, validity, range: Range::of(column) });
231 }
232 local.first = HashMap::default();
235 local.next = Vec::new();
236 Ok(local)
237 }
238
239 fn rows(&self) -> Result<Vec<Vector>> {
245 self.parts
246 .iter()
247 .map(|part| {
248 let len = part.codes.len();
249 let mut column = StringColumn::with_capacity(len);
250 for &code in &part.codes {
251 column.push_bytes(self.value(code));
252 }
253 let validity = match part.validity.split_first() {
254 Some((0, _)) => Validity::AllValid,
255 Some((1, _)) => Validity::AllInvalid,
256 Some((2, bits)) => {
257 let mut mask = Bitmap::all_valid(len);
258 for row in (0..len).filter(|row| bits[row / 8] & (1 << (row % 8)) == 0) {
259 mask.set(row, false);
260 }
261 Validity::Mask(mask)
262 }
263 _ => return Err(Error::internal("a coded part has no validity")),
264 };
265 Ok(Vector::flat(LogicalType::Varchar, Data::Varlen(column))?
266 .with_validity(validity))
267 })
268 .collect()
269 }
270
271 fn values(&self) -> usize {
272 self.ends.len()
273 }
274
275 fn value(&self, code: u32) -> &[u8] {
276 let code = code as usize;
277 let from = if code == 0 { 0 } else { self.ends[code - 1] };
278 &self.bytes[from..self.ends[code]]
279 }
280
281 fn code(&mut self, text: &[u8]) -> Result<u32> {
282 let hash = checksum(text);
283 let Some(&first) = self.first.get(&hash) else {
284 let code = self.push(text, hash)?;
285 self.first.insert(hash, code);
286 return Ok(code);
287 };
288 let mut at = first;
289 loop {
290 if self.value(at) == text {
291 return Ok(at);
292 }
293 match self.next[at as usize] {
294 END => break,
295 next => at = next,
296 }
297 }
298 let code = self.push(text, hash)?;
299 self.next[at as usize] = code;
300 Ok(code)
301 }
302
303 fn push(&mut self, text: &[u8], hash: u64) -> Result<u32> {
304 let code = u32::try_from(self.ends.len())
305 .ok()
306 .filter(|&code| code != END)
307 .ok_or_else(|| invalid("a stripe has too many values in one column"))?;
308 self.bytes.extend_from_slice(text);
309 self.ends.push(self.bytes.len());
310 self.next.push(END);
311 self.hashes.push(hash);
312 self.checks.push(seeded_checksum(text, DICTIONARY_CHECK_SEED));
313 self.counts.push(0);
314 Ok(code)
315 }
316
317 fn merge_into(&self, dictionary: &mut GlobalDictionary) -> Result<Vec<u32>> {
320 let mut global = Vec::with_capacity(self.values());
321 for (code, (&hash, &check)) in self.hashes.iter().zip(&self.checks).enumerate() {
322 let text = self.value(code as u32);
323 let at = dictionary.code_hashed(text, hash, check)?;
324 let count = dictionary
325 .counts
326 .get_mut(at as usize)
327 .ok_or_else(|| invalid("global dictionary count code is out of range"))?;
328 *count = count.saturating_add(self.counts[code]);
329 global.push(at);
330 }
331 dictionary.nulls = dictionary.nulls.saturating_add(self.nulls);
332 Ok(global)
333 }
334}
335
336fn drops_dictionary(rows: usize, distinct: usize) -> bool {
360 rows >= DICTIONARY_DECIDE_ROWS
361 && distinct.saturating_mul(10) > rows.saturating_mul(DICTIONARY_DISTINCT_IN_TEN)
362}
363
364fn fan_out<T: Send>(
374 jobs: Vec<usize>,
375 workers: usize,
376 profile: Option<&LoadProfile>,
377 work: impl Fn(usize) -> Result<T> + Sync,
378) -> Result<Vec<(usize, T)>> {
379 if workers <= 1 || jobs.len() <= 1 {
380 let _span = profile.map(|profile| profile.span(Stage::Pages));
381 return jobs.into_iter().map(|index| Ok((index, work(index)?))).collect();
382 }
383 let workers = workers.min(jobs.len());
384 let queue = Mutex::new(jobs);
385 let pieces = std::thread::scope(|scope| {
386 (0..workers)
387 .map(|_| {
388 scope.spawn(|| {
389 let _span = profile.map(|profile| profile.span(Stage::Pages));
390 let mut mine = Vec::new();
391 loop {
392 let taken = queue
393 .lock()
394 .map_err(|_| Error::internal("a native encode worker panicked"))?
395 .pop();
396 let Some(index) = taken else { break };
397 mine.push((index, work(index)?));
398 }
399 Ok(mine)
400 })
401 })
402 .collect::<Vec<_>>()
403 .into_iter()
404 .map(|handle| {
405 handle.join().map_err(|_| Error::internal("a native encode worker panicked"))?
406 })
407 .collect::<Result<Vec<Vec<_>>>>()
408 })?;
409 Ok(pieces.into_iter().flatten().collect())
410}
411
412fn column_of(held: &[PendingChunk], index: usize) -> Result<Vec<&Vector>> {
414 held.iter().map(|pending| pending.chunk.column(index)).collect()
415}
416
417fn in_order<T>(width: usize, done: Vec<(usize, T)>) -> Result<Vec<T>> {
419 let mut slots: Vec<Option<T>> = (0..width).map(|_| None).collect();
420 for (index, one) in done {
421 slots[index] = Some(one);
422 }
423 slots
424 .into_iter()
425 .map(|slot| slot.ok_or_else(|| Error::internal("a column was never encoded")))
426 .collect()
427}
428
429impl Preparer {
430 pub fn prepare(&self, parts: Vec<((u64, u64), Chunk)>) -> Result<Prepared> {
441 if parts.len() > STRIPE_PARTS {
442 return Err(invalid("a stripe was handed more parts than it holds"));
443 }
444 let held = parts
445 .into_iter()
446 .filter(|(_, chunk)| !chunk.is_empty())
447 .map(|(order, chunk)| PendingChunk { order, chunk })
448 .collect::<Vec<_>>();
449 for pending in &held {
450 self.fits(&pending.chunk)?;
451 }
452 self.prepare_held(held)
453 }
454
455 fn fits(&self, chunk: &Chunk) -> Result<()> {
457 if chunk.width() != self.types.len() {
458 return Err(invalid("chunk width differs from table schema"));
459 }
460 for (index, ty) in self.types.iter().enumerate() {
461 if chunk.column(index)?.logical_type() != ty {
462 return Err(invalid("chunk type differs from table schema"));
463 }
464 }
465 Ok(())
466 }
467
468 pub(crate) fn prepare_held(&self, held: Vec<PendingChunk>) -> Result<Prepared> {
469 let width = self.types.len();
470 let key = held.first().map_or((0, 0), |pending| pending.order);
471 let share = Share::take(width, held.len());
472 let mut jobs = (0..width).collect::<Vec<_>>();
473 jobs.sort_by_key(|&index| weight(&self.types[index]));
474 let done = fan_out(jobs, share.0, self.profile.as_deref(), |index| {
475 let gather = stats::Gather::new(&self.types[index], 0)
478 .filter(|_| !held.is_empty())
479 .map(|mut gather| {
480 gather.stripe(
481 key,
482 held.iter().filter_map(|pending| pending.chunk.column(index).ok()),
483 );
484 gather
485 });
486 let column = if self.coded[index].load(Atomic::Relaxed) {
487 Column::Coded(Local::code_column(index, &held)?)
488 } else {
489 Column::Pages(Writer::encode_pages(&column_of(&held, index)?)?)
490 };
491 Ok((column, gather))
492 })?;
493 drop(share);
494 let (columns, gathers) = in_order(width, done)?.into_iter().unzip();
495 let parts = held.iter().map(Part::of).collect();
496 drop(held);
497 Ok(Prepared {
498 parts,
499 types: self.types.clone(),
500 columns,
501 gathers,
502 profile: self.profile.clone(),
503 })
504 }
505}
506
507enum Slot<'a> {
509 Owned(&'a mut Option<GlobalDictionary>, &'a mut Option<stats::Gather>),
511 Lent(&'a Mutex<LentColumn>, &'a Lent),
513}
514
515struct Step<'a> {
517 index: usize,
518 column: Column,
519 slot: Slot<'a>,
520 gather: Option<stats::Gather>,
522}
523
524impl Step<'_> {
525 fn cost(&self, coded: &[AtomicBool]) -> usize {
529 match &self.column {
530 Column::Coded(local) if coded[self.index].load(Atomic::Relaxed) => {
531 local.values().saturating_add(1)
532 }
533 _ => 0,
534 }
535 }
536
537 fn run(self, rows: usize, coded: &[AtomicBool]) -> Result<(usize, Merge, Vec<Unencoded>)> {
539 let Self { index, column, slot, gather } = self;
540 match slot {
541 Slot::Owned(dictionary, mine) => {
542 merge_column(index, column, gather, dictionary, mine, rows, coded)
543 }
544 Slot::Lent(held, lent) => {
545 let mut held = held.lock().map_err(|_| Error::internal("a merge panicked"))?;
546 if lent.reclaimed.load(Atomic::Acquire) {
549 return Err(Error::internal("a stripe was merged after its table was closed"));
550 }
551 let LentColumn { dictionary, gather: mine } = &mut *held;
552 merge_column(index, column, gather, dictionary, mine, rows, coded)
553 }
554 }
555 }
556}
557
558fn merge_column(
560 index: usize,
561 column: Column,
562 stripe: Option<stats::Gather>,
563 dictionary: &mut Option<GlobalDictionary>,
564 gather: &mut Option<stats::Gather>,
565 rows: usize,
566 coded: &[AtomicBool],
567) -> Result<(usize, Merge, Vec<Unencoded>)> {
568 if let (Some(mine), Some(stripe)) = (gather.as_mut(), stripe) {
569 mine.absorb(stripe);
570 }
571 let merge = match (column, dictionary.as_mut()) {
572 (Column::Pages(stripe), None) => Merge::Pages(stripe),
573 (Column::Pages(_), Some(_)) => {
574 return Err(Error::internal(
575 "a column with a global dictionary was prepared without one",
576 ));
577 }
578 (Column::Coded(local), None) => Merge::Plain(local),
579 (Column::Coded(local), Some(global)) => {
580 if global.values() == 0 && drops_dictionary(rows, local.values()) {
583 *dictionary = None;
584 coded[index].store(false, Atomic::Relaxed);
585 Merge::Plain(local)
586 } else {
587 let global = local.merge_into(global)?;
588 Merge::Codes { parts: local.parts, global }
589 }
590 }
591 };
592 let blocks = match dictionary {
597 Some(dictionary) => {
598 dictionary.settle()?;
599 dictionary.hand_out(index)
600 }
601 None => Vec::new(),
602 };
603 Ok((index, merge, blocks))
604}
605
606fn merge_columns(prepared: Prepared, slots: Vec<Slot<'_>>, coded: &[AtomicBool]) -> Result<Merged> {
614 let Prepared { parts, columns, gathers, profile, .. } = prepared;
615 let timing = profile.as_deref().map(|profile| profile.span(Stage::Dictionary));
616 let rows: usize = parts.iter().map(|part| part.rows).sum();
617 let width = columns.len();
618 if slots.len() != width || gathers.len() != width {
619 return Err(Error::internal("a stripe was merged into a table of another width"));
620 }
621 let mut steps = columns
622 .into_iter()
623 .zip(gathers)
624 .zip(slots)
625 .enumerate()
626 .map(|(index, ((column, gather), slot))| Step { index, column, slot, gather })
627 .collect::<Vec<_>>();
628 steps.sort_by_key(|step| step.cost(coded));
631 let workers = std::thread::available_parallelism()
632 .map_or(1, usize::from)
633 .min(MAX_ENCODE_WORKERS)
634 .min(steps.iter().filter(|step| step.cost(coded) > 0).count())
635 .max(1);
636 let done = if workers <= 1 {
637 steps.into_iter().map(|step| step.run(rows, coded)).collect::<Result<Vec<_>>>()?
638 } else {
639 let queue = Mutex::new(steps);
640 let pieces = std::thread::scope(|scope| {
641 (0..workers)
642 .map(|_| {
643 scope.spawn(|| {
644 let mut mine = Vec::new();
645 loop {
646 let taken = queue
647 .lock()
648 .map_err(|_| Error::internal("a merge worker panicked"))?
649 .pop();
650 let Some(step) = taken else { break };
651 mine.push(step.run(rows, coded)?);
652 }
653 Ok(mine)
654 })
655 })
656 .collect::<Vec<_>>()
657 .into_iter()
658 .map(|handle| {
659 handle.join().map_err(|_| Error::internal("a merge worker panicked"))?
660 })
661 .collect::<Result<Vec<Vec<_>>>>()
662 })?;
663 pieces.into_iter().flatten().collect()
664 };
665 let mut slots: Vec<Option<(Merge, Vec<Unencoded>)>> = (0..width).map(|_| None).collect();
666 for (index, merge, blocks) in done {
667 slots[index] = Some((merge, blocks));
668 }
669 let mut merged = Vec::with_capacity(width);
670 let mut blocks = Vec::new();
671 for slot in slots {
672 let (merge, handed) = slot.ok_or_else(|| Error::internal("a column was never merged"))?;
673 merged.push(merge);
674 blocks.extend(handed);
675 }
676 drop(timing);
677 Ok(Merged { parts, columns: merged, blocks, profile, counted: false })
678}
679
680#[derive(Debug)]
682pub(crate) struct Lent {
683 columns: Box<[Mutex<LentColumn>]>,
684 reclaimed: AtomicBool,
686}
687
688#[derive(Debug)]
690pub(crate) struct LentColumn {
691 pub(crate) dictionary: Option<GlobalDictionary>,
692 gather: Option<stats::Gather>,
693}
694
695impl Lent {
696 pub(crate) fn columns(&self) -> &[Mutex<LentColumn>] {
697 &self.columns
698 }
699
700 fn take_back(&self, blocks: Vec<(usize, usize, EncodedBlock)>) -> Result<()> {
702 for (column, at, block) in blocks {
703 self.columns
704 .get(column)
705 .ok_or_else(|| Error::internal("a dictionary block came back to no column"))?
706 .lock()
707 .map_err(|_| Error::internal("a merge panicked"))?
708 .dictionary
709 .as_mut()
710 .ok_or_else(|| Error::internal("a dictionary block came back to no dictionary"))?
711 .take_back(at, block)?;
712 }
713 Ok(())
714 }
715
716 #[allow(clippy::type_complexity)]
718 pub(crate) fn reclaim(
719 &self,
720 ) -> Result<(Vec<Option<GlobalDictionary>>, Vec<Option<stats::Gather>>)> {
721 self.reclaimed.store(true, Atomic::Release);
722 let mut dictionaries = Vec::with_capacity(self.columns.len());
723 let mut gathers = Vec::with_capacity(self.columns.len());
724 for column in &self.columns {
725 let mut held = column.lock().map_err(|_| Error::internal("a merge panicked"))?;
726 dictionaries.push(held.dictionary.take());
727 gathers.push(held.gather.take());
728 }
729 Ok((dictionaries, gathers))
730 }
731}
732
733#[derive(Debug, Clone)]
740pub struct Merger {
741 lent: Arc<Lent>,
742 types: Vec<LogicalType>,
743 coded: Arc<[AtomicBool]>,
744}
745
746impl Merger {
747 pub fn merge(&self, prepared: Prepared) -> Result<Merged> {
753 if prepared.types != self.types {
754 return Err(invalid("a stripe was prepared for a table of other columns"));
755 }
756 let slots = self.lent.columns.iter().map(|column| Slot::Lent(column, &self.lent)).collect();
757 merge_columns(prepared, slots, &self.coded)
758 }
759
760 pub fn give_back(&self, paged: &mut Paged) -> Result<()> {
767 self.lent.take_back(std::mem::take(&mut paged.blocks))
768 }
769}
770
771impl Merged {
772 pub fn pages(self) -> Result<Paged> {
780 let Self { parts, columns, blocks, profile, counted } = self;
781 let width = columns.len();
782 let mut jobs = (width..width + blocks.len())
786 .chain((0..width).filter(|&index| !matches!(columns[index], Merge::Pages(_))))
787 .collect::<Vec<_>>();
788 jobs.sort_by_key(|&index| index < width && matches!(columns[index], Merge::Plain(_)));
790 let share = Share::take(jobs.len(), parts.len());
791 let built = fan_out(jobs, share.0, profile.as_deref(), |index| {
792 let Some(column) = columns.get(index) else {
793 return Ok(Built::Block(blocks[index - width].encode()?));
794 };
795 Ok(Built::Stripe(match column {
796 Merge::Codes { parts, global } => code_pages(parts, global)?,
797 Merge::Plain(local) => {
798 Writer::encode_pages(&local.rows()?.iter().collect::<Vec<_>>())?
799 }
800 Merge::Pages(_) => {
801 return Err(Error::internal("a finished column was queued to be built"));
802 }
803 }))
804 })?;
805 drop(share);
806 let mut slots: Vec<Option<ColumnStripe>> = (0..width).map(|_| None).collect();
807 let mut encoded = Vec::with_capacity(blocks.len());
808 for (index, one) in built {
809 match one {
810 Built::Stripe(stripe) => slots[index] = Some(stripe),
811 Built::Block(block) => {
812 let (column, at) = blocks[index - width].place();
813 encoded.push((column, at, block));
814 }
815 }
816 }
817 let columns = columns
818 .into_iter()
819 .zip(slots)
820 .map(|(column, slot)| match (column, slot) {
821 (Merge::Pages(stripe), _) | (_, Some(stripe)) => Ok(stripe),
822 _ => Err(Error::internal("a column was never encoded")),
823 })
824 .collect::<Result<Vec<_>>>()?;
825 Ok(Paged { parts, columns, blocks: encoded, counted })
826 }
827}
828
829fn code_pages(parts: &[LocalPart], global: &[u32]) -> Result<ColumnStripe> {
831 let mut stripe = ColumnStripe {
832 pages: Vec::with_capacity(parts.len()),
833 codes: Vec::with_capacity(parts.len()),
834 sieves: Vec::with_capacity(parts.len()),
835 ranges: Vec::with_capacity(parts.len()),
836 };
837 for part in parts {
838 let codes = part
839 .codes
840 .iter()
841 .map(|&code| global.get(code as usize).copied())
842 .collect::<Option<Vec<_>>>()
843 .ok_or_else(|| Error::internal("a stripe's code has no global code"))?;
844 let bytes = coded_page(&codes, &part.validity)?;
845 if bytes.len() > MAX_PAGE {
846 return Err(invalid("column page exceeds the configured bound"));
847 }
848 stripe.pages.push(bytes);
849 stripe.codes.push(Some(unique_codes(&codes)));
850 stripe.sieves.push(None);
854 stripe.ranges.push(part.range.clone());
855 }
856 Ok(stripe)
857}
858
859impl Writer {
860 #[must_use]
865 pub fn preparer(&self) -> Preparer {
866 Preparer {
867 types: self.table.fields.iter().map(|field| field.ty.clone()).collect(),
868 coded: Arc::clone(&self.coded),
869 profile: self.profile.clone(),
870 }
871 }
872
873 pub fn merge(&mut self, prepared: Prepared) -> Result<Merged> {
884 self.flush_pending()?;
885 if prepared.columns.len() != self.table.fields.len()
886 || prepared.types.iter().ne(self.table.fields.iter().map(|field| &field.ty))
887 {
888 return Err(invalid("a stripe was prepared for a table of other columns"));
889 }
890 self.table.rows = prepared
891 .parts
892 .iter()
893 .try_fold(self.table.rows, |rows, part| rows.checked_add(part.rows))
894 .ok_or_else(|| invalid("row count overflow"))?;
895 self.merge_held(prepared)
896 }
897
898 pub(crate) fn merge_held(&mut self, prepared: Prepared) -> Result<Merged> {
907 let slots = match &self.lent {
908 Some(lent) => lent.columns.iter().map(|column| Slot::Lent(column, lent)).collect(),
909 None => self
910 .dictionaries
911 .iter_mut()
912 .zip(self.gathers.iter_mut())
913 .map(|(dictionary, gather)| Slot::Owned(dictionary, gather))
914 .collect::<Vec<_>>(),
915 };
916 let mut merged = merge_columns(prepared, slots, &self.coded)?;
917 merged.counted = true;
918 Ok(merged)
919 }
920
921 pub fn merger(&mut self) -> Result<Merger> {
931 self.flush_pending()?;
932 let lent = match &self.lent {
933 Some(lent) => Arc::clone(lent),
934 None => {
935 let lent = Arc::new(Lent {
936 columns: std::mem::take(&mut self.dictionaries)
937 .into_iter()
938 .zip(std::mem::take(&mut self.gathers))
939 .map(|(dictionary, gather)| Mutex::new(LentColumn { dictionary, gather }))
940 .collect(),
941 reclaimed: AtomicBool::new(false),
942 });
943 self.lent = Some(Arc::clone(&lent));
944 lent
945 }
946 };
947 Ok(Merger {
948 lent,
949 types: self.table.fields.iter().map(|field| field.ty.clone()).collect(),
950 coded: Arc::clone(&self.coded),
951 })
952 }
953
954 pub fn write(&mut self, paged: Paged) -> Result<()> {
960 self.write_paged(paged)
961 }
962
963 pub(crate) fn write_paged(&mut self, paged: Paged) -> Result<()> {
964 let Paged { parts, columns, blocks, counted } = paged;
965 if !counted {
966 self.table.rows = parts
967 .iter()
968 .try_fold(self.table.rows, |rows, part| rows.checked_add(part.rows))
969 .ok_or_else(|| invalid("row count overflow"))?;
970 }
971 if let Some(lent) = &self.lent {
972 lent.take_back(blocks)?;
973 } else {
974 for (column, at, block) in blocks {
975 self.dictionaries
976 .get_mut(column)
977 .and_then(Option::as_mut)
978 .ok_or_else(|| {
979 Error::internal("a dictionary block came back to no dictionary")
980 })?
981 .take_back(at, block)?;
982 }
983 }
984 if parts.is_empty() {
985 return self.place_blocks();
986 }
987 self.write_stripe(&parts, columns)
988 }
989
990 pub fn append_prepared(&mut self, prepared: Prepared) -> Result<()> {
996 let merged = self.merge(prepared)?;
997 let paged = merged.pages()?;
998 self.write(paged)
999 }
1000}
1001
1002#[cfg(test)]
1003mod tests {
1004 use std::fs;
1005 use std::path::PathBuf;
1006 use std::time::{SystemTime, UNIX_EPOCH};
1007
1008 use rudb_common::{Field, Value};
1009 use rudb_vector::Vector;
1010
1011 use super::*;
1012 use crate::Reader;
1013
1014 const PART: usize = 1_000;
1015
1016 fn path(label: &str) -> PathBuf {
1017 let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
1018 std::env::temp_dir()
1019 .join(format!("rudb-prepare-{label}-{}-{stamp}.rdb", std::process::id()))
1020 }
1021
1022 fn fields() -> Vec<Field> {
1023 vec![
1024 Field::required("id", LogicalType::BigInt),
1025 Field::new("city", LogicalType::Varchar),
1026 Field::new("note", LogicalType::Varchar),
1027 ]
1028 }
1029
1030 fn row(id: usize) -> [Value; 3] {
1035 let city = if id % 11 == 0 {
1036 Value::Null
1037 } else {
1038 Value::Varchar(format!("city {}", (id / 7) % 13))
1039 };
1040 let note = if id % 17 == 0 { Value::Null } else { Value::Varchar(format!("note {id}")) };
1041 [Value::BigInt(id as i64), city, note]
1042 }
1043
1044 fn stripe(first: usize, parts: usize) -> Vec<((u64, u64), Chunk)> {
1046 (first..first + parts)
1047 .map(|part| {
1048 let rows = (part * PART..(part + 1) * PART).map(row).collect::<Vec<_>>();
1049 let column = |at: usize| {
1050 let values = rows.iter().map(|row| row[at].clone()).collect::<Vec<_>>();
1051 Vector::from_values(fields()[at].ty.clone(), &values).expect("a column")
1052 };
1053 let chunk = Chunk::new(vec![column(0), column(1), column(2)]).expect("a chunk");
1054 ((part as u64, 0), chunk)
1055 })
1056 .collect()
1057 }
1058
1059 fn runs() -> Vec<Vec<((u64, u64), Chunk)>> {
1061 vec![stripe(5, 5), stripe(0, 5), stripe(10, 3)]
1062 }
1063
1064 fn check(path: &PathBuf) {
1065 let reader = Reader::open(path).expect("reopen");
1066 assert_eq!(reader.parts(), 13);
1067 for part in 0..13 {
1068 let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1069 for at in [0, 17, PART - 1] {
1070 let want = row(part * PART + at);
1071 for (column, value) in want.iter().enumerate() {
1072 assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1073 }
1074 }
1075 }
1076 }
1077
1078 #[test]
1086 fn stripes_prepared_before_any_is_merged_write_the_same_bytes_as_one_at_a_time() {
1087 let alone = path("alone");
1088 let mut writer = Writer::create(&alone, "t", fields()).expect("a file");
1089 for run in runs() {
1090 writer.append_stripe(run).expect("a stripe");
1091 }
1092 writer.finish().expect("commit");
1093
1094 let split = path("split");
1095 let mut writer = Writer::create(&split, "t", fields()).expect("a file");
1096 let preparer = writer.preparer();
1097 let prepared = runs()
1098 .into_iter()
1099 .map(|run| preparer.prepare(run).expect("prepared"))
1100 .collect::<Vec<_>>();
1101 for one in prepared {
1102 writer.append_prepared(one).expect("a stripe");
1103 }
1104 assert!(!preparer.coded[2].load(Atomic::Relaxed), "note lost its dictionary");
1105 assert!(preparer.coded[1].load(Atomic::Relaxed), "city kept its dictionary");
1106 writer.finish().expect("commit");
1107
1108 assert_eq!(fs::read(&alone).expect("read"), fs::read(&split).expect("read"));
1109 check(&split);
1110 fs::remove_file(alone).expect("remove");
1111 fs::remove_file(split).expect("remove");
1112 }
1113
1114 #[test]
1118 fn stripes_written_in_another_order_than_they_were_merged_read_back() {
1119 let path = path("crossed");
1120 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1121 let preparer = writer.preparer();
1122 let mut merged = runs()
1123 .into_iter()
1124 .map(|run| writer.merge(preparer.prepare(run).expect("prepared")).expect("merged"))
1125 .map(|merged| merged.pages().expect("paged"))
1126 .collect::<Vec<_>>();
1127 merged.reverse();
1128 for paged in merged {
1129 writer.write(paged).expect("written");
1130 }
1131 writer.finish().expect("commit");
1132 check(&path);
1133 fs::remove_file(path).expect("remove");
1134 }
1135
1136 #[test]
1139 fn stripes_merged_through_a_merger_write_the_same_bytes_as_the_writer() {
1140 let alone = path("alone-merger");
1141 let mut writer = Writer::create(&alone, "t", fields()).expect("a file");
1142 for run in runs() {
1143 writer.append_stripe(run).expect("a stripe");
1144 }
1145 writer.finish().expect("commit");
1146
1147 let lent = path("lent");
1148 let mut writer = Writer::create(&lent, "t", fields()).expect("a file");
1149 let preparer = writer.preparer();
1150 let merger = writer.merger().expect("a merger");
1151 for run in runs() {
1152 let merged = merger.merge(preparer.prepare(run).expect("prepared")).expect("merged");
1153 let mut paged = merged.pages().expect("paged");
1154 merger.give_back(&mut paged).expect("given back");
1155 writer.write(paged).expect("written");
1156 }
1157 assert_eq!(writer.table.rows, 13 * PART);
1158 writer.finish().expect("commit");
1159
1160 assert_eq!(fs::read(&alone).expect("read"), fs::read(&lent).expect("read"));
1161 check(&lent);
1162 fs::remove_file(alone).expect("remove");
1163 fs::remove_file(lent).expect("remove");
1164 }
1165
1166 #[test]
1169 fn stripes_merged_on_several_threads_at_once_read_back() {
1170 let path = path("merged-at-once");
1171 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1172 let preparer = writer.preparer();
1173 let merger = writer.merger().expect("a merger");
1174 let writer = Mutex::new(writer);
1175 std::thread::scope(|scope| {
1176 for run in runs() {
1177 let (preparer, merger, writer) = (&preparer, &merger, &writer);
1178 scope.spawn(move || {
1179 let merged =
1180 merger.merge(preparer.prepare(run).expect("prepared")).expect("merged");
1181 let mut paged = merged.pages().expect("paged");
1182 merger.give_back(&mut paged).expect("given back");
1183 writer.lock().expect("the writer").write(paged).expect("written");
1184 });
1185 }
1186 });
1187 writer.into_inner().expect("the writer").finish().expect("commit");
1188 check(&path);
1189 fs::remove_file(path).expect("remove");
1190 }
1191
1192 #[test]
1195 fn a_merge_after_the_table_is_closed_is_refused() {
1196 let path = path("late");
1197 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1198 let preparer = writer.preparer();
1199 let merger = writer.merger().expect("a merger");
1200 writer.finish().expect("commit");
1201 let prepared = preparer.prepare(stripe(0, 2)).expect("prepared");
1202 assert!(merger.merge(prepared).is_err());
1203 fs::remove_file(path).expect("remove");
1204 }
1205
1206 #[test]
1209 fn a_stripe_of_another_table_is_refused_at_the_merge() {
1210 let path = path("refused");
1211 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1212 let other = Writer::create(path.with_extension("other"), "u", vec![fields().remove(0)])
1213 .expect("a file");
1214 let prepared = other.preparer().prepare(vec![]).expect("nothing to prepare");
1215 assert!(writer.merge(prepared).is_err());
1216 assert_eq!(writer.table.rows, 0);
1217 drop(other);
1218 fs::remove_file(path.with_extension("other")).expect("remove");
1219 fs::remove_file(path).expect("remove");
1220 }
1221}