Skip to main content

rudb_native/
prepare.rs

1//! A stripe encoded before the writer is asked for it.
2//!
3//! A load fed by many pipeline instances has one [`Writer`] behind one lock, and until this every
4//! instance encoded its stripe while holding that lock. On the 32 core box, loading the ClickBench
5//! 10m sample spent 68% of all its processor time in the stripe encode, all of it under the lock,
6//! and the instances waited 146 seconds between them for a load whose wall clock was 20.6 seconds.
7//! The machine was one encode at a time with thirty one readers queued behind it.
8//!
9//! Almost none of that work needs the writer. A plain column's pages, its sieves, its ranges and
10//! its statistics depend on the stripe's own rows and nothing else. The one thing a stripe shares
11//! with the rest of the table is a varchar column's global dictionary, because a code has to mean
12//! the same value in every page of the column. So a stripe is taken in four steps:
13//!
14//! 1. [`Preparer::prepare`], with no lock. Every column without a global dictionary is encoded to
15//!    its pages, and every column with one is coded against a dictionary of the stripe's own, which
16//!    holds each distinct value of the stripe once, in the order the rows first held it. The
17//!    statistics of every column are folded into a gather of the stripe's own.
18//! 2. [`Writer::merge`], under the lock. The stripe's dictionaries go into the global ones a
19//!    distinct value at a time, which gives back what each local code is globally, and the gathers
20//!    are absorbed. This is the only step that has to see the stripes one at a time. Every
21//!    dictionary block the merge filled is taken out with the stripe.
22//! 3. [`Merged::pages`], with no lock. The codes are turned into global ones and built into pages,
23//!    and the dictionary blocks the merge took out are encoded.
24//! 4. [`Writer::write`], under the lock. The dictionary blocks go back in order, and they and the
25//!    pages go into the file.
26//!
27//! The dictionary blocks were encoded in the fourth step, under the lock, until the 10m ClickBench
28//! load on the 32 core box was measured spending 2.7 of its 14 seconds there, on thirty two threads
29//! spawned for it every stripe, with every instance queued behind them.
30//!
31//! A value merged in the order the stripe first held it gets the code it would have got had the
32//! stripe been coded against the global dictionary row by row, because the rows before its first
33//! appearance hold only values that were already merged. So a writer taking the four steps one
34//! after the other writes the same bytes as one that coded every row against the global dictionary,
35//! which is what [`Writer::flush_pending`] does.
36//!
37//! A stripe lets go of its rows at the end of the first step. What it carries from there on is its
38//! pages, its stripe dictionaries and codes, and a few numbers a part, so the stripes queued for the
39//! lock are a fraction of the size of the rows they came from. A column that loses its dictionary
40//! after it was coded against one is rebuilt from the stripe dictionary, which holds every value
41//! the rows did. Keeping the rows until the write instead took the 10m ClickBench load on the 32
42//! core box from 3.5 GB resident to 9.9 GB, with thirty two stripes waiting at a time.
43
44use 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
60/// How many stripes are being prepared or paged right now, across every writer in the process.
61///
62/// A stripe's columns are spread over threads of their own, which is what a writer being fed by
63/// one caller needs, because that caller is the only one encoding. Thirty two callers each doing
64/// that at once would be a thousand threads on a machine with thirty two cores. So each one takes
65/// its share of the machine: the cores over however many stripes are being worked on right now.
66static BUSY: AtomicUsize = AtomicUsize::new(0);
67
68/// One stripe's share of the machine, held for as long as the stripe is being worked on.
69struct 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        // A stripe of one part is one small page a column, which is less than a thread is worth.
77        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/// Encodes stripes for one [`Writer`] without the writer.
89///
90/// Handed out by [`Writer::preparer`] and cheap to hold. It shares with its writer which varchar
91/// columns still have a global dictionary, so a stripe prepared after the first one decided a
92/// column should not have one is encoded plainly from the start.
93#[derive(Debug, Clone)]
94pub struct Preparer {
95    types: Vec<LogicalType>,
96    coded: Arc<[AtomicBool]>,
97    profile: Option<Arc<LoadProfile>>,
98}
99
100/// A stripe that has been through [`Preparer::prepare`] and is waiting for [`Writer::merge`].
101#[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/// A stripe that has been through [`Writer::merge`] and is waiting for [`Merged::pages`].
111#[derive(Debug)]
112pub struct Merged {
113    parts: Vec<Part>,
114    columns: Vec<Merge>,
115    blocks: Vec<Unencoded>,
116    profile: Option<Arc<LoadProfile>>,
117    /// Whether the rows are counted into the table yet. [`Writer::merge`] counts them, and a
118    /// [`Merger`] leaves them for [`Writer::write`], since it has no table to count them into.
119    counted: bool,
120}
121
122/// A stripe that has been through [`Merged::pages`] and is waiting for [`Writer::write`].
123#[derive(Debug)]
124pub struct Paged {
125    parts: Vec<Part>,
126    columns: Vec<ColumnStripe>,
127    /// Encoded dictionary blocks, each with its column and block number.
128    blocks: Vec<(usize, usize, EncodedBlock)>,
129    counted: bool,
130}
131
132/// What one job of [`Merged::pages`] built.
133enum Built {
134    Stripe(ColumnStripe),
135    Block(EncodedBlock),
136}
137
138/// One column of a prepared stripe.
139#[derive(Debug)]
140enum Column {
141    /// Finished, because the column has no global dictionary.
142    Pages(ColumnStripe),
143    /// Coded against the stripe's own dictionary, waiting to be merged into the global one.
144    Coded(Local),
145}
146
147/// One column of a merged stripe.
148#[derive(Debug)]
149enum Merge {
150    Pages(ColumnStripe),
151    /// The local codes of every part, and the global code of every local one.
152    Codes {
153        parts: Vec<LocalPart>,
154        global: Vec<u32>,
155    },
156    /// A column that was prepared against a dictionary it no longer has, which is every column
157    /// prepared before the first stripe decided it should not have one. Encoded again, plainly,
158    /// from the values its stripe dictionary holds.
159    Plain(Local),
160}
161
162/// No value after this one has its hash.
163const END: u32 = u32::MAX;
164
165/// A dictionary of one column of one stripe.
166///
167/// The values are compared by their bytes rather than by a second hash, because they are all here
168/// to compare. The global dictionary has two hashes to go on because its values are mostly in the
169/// file by now. Both hashes are taken here, once a distinct value, so that merging it takes none.
170#[derive(Debug, Default)]
171struct Local {
172    /// The first value holding each hash.
173    first: HashMap<u64, u32, Spread>,
174    /// The next value holding the same hash as this one, or [`END`].
175    next: Vec<u32>,
176    hashes: Vec<u64>,
177    checks: Vec<u64>,
178    /// The values back to back, and where each one ends.
179    bytes: Vec<u8>,
180    ends: Vec<usize>,
181    /// How many rows that are not null hold each value, and how many are null.
182    counts: Vec<u64>,
183    nulls: u64,
184    parts: Vec<LocalPart>,
185}
186
187/// One part of one column coded against its stripe's dictionary.
188#[derive(Debug)]
189struct LocalPart {
190    codes: Vec<u32>,
191    /// What [`push_validity`] wrote for the part, which is the page's second field onwards.
192    validity: Vec<u8>,
193    range: Range,
194}
195
196impl Local {
197    /// One column of a stripe, coded.
198    ///
199    /// A null row is coded as the empty string and counted as a null rather than against it, which
200    /// is what the writer has always done with one. The code is never read, since the page's
201    /// validity says the row is null, and giving it one keeps the page one code a row.
202    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            // flatten: the page is one code a row whatever form the rows came in.
207            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                // A repeat of the row before is common enough on a sorted table to be worth a
213                // comparison before a hash, and the comparison fails on its first bytes when not.
214                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        // Only the coding needs to find a value by its bytes, and on a column of URLs the table
231        // that does it is as large as the codes.
232        local.first = HashMap::default();
233        local.next = Vec::new();
234        Ok(local)
235    }
236
237    /// The column's parts as the rows they were coded from, for a column that lost its global
238    /// dictionary after this stripe was coded against one.
239    ///
240    /// A null row comes back as a null over the empty string, which is what it was coded as, and
241    /// each part gets back the same form of validity it had, since the page records which it was.
242    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    /// Puts every value into `dictionary` in the order this stripe first held it, and says what
316    /// each one's code is there.
317    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
334/// Whether a column's first stripe says it should not have a global dictionary.
335///
336/// Every varchar column starts with one, because the writer cannot know what is in a column before
337/// it has seen some of it. A global dictionary is the right shape for a column of a few dozen
338/// values repeated down the table: the pages become small integers, a filter against a literal is
339/// one search of the sorted order rather than a comparison a row, and a group by is on the codes.
340/// It is the wrong shape for a column whose values are nearly all different. There the codes are
341/// as wide as row numbers, nothing is saved on the pages, and the membership index of a stripe is
342/// a list of very nearly every code in the column. On TPC-H the orders table written on its own
343/// goes from 52.3 MB to 41.4 MB, the load from 6.9 s to 5.8 s, and `select o_comment from orders`
344/// from 1.810 G instructions to 1.213 G, which is what the rudb parquet reader takes over the same
345/// values.
346///
347/// So the first stripe of a column is the sample and the decision is made once on it. Once, rather
348/// than per stripe, because the codes of one column have to mean the same thing in every page of
349/// it, and a column that changed its mind halfway would need its earlier stripes rewritten. The
350/// first stripe is encoded again when the answer comes out against the dictionary, which is the one
351/// stripe that pays for the decision, along with any stripe that was prepared before it was made.
352///
353/// The threshold is deliberately near the top. [`DICTIONARY_DISTINCT_IN_TEN`] of the sample has to
354/// be values never seen before, which is a column with essentially no repeats. Everything with real
355/// repetition keeps its dictionary and keeps every property that hangs off it, and nothing is
356/// claimed here about where between the two the crossover really sits.
357fn 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
362/// Runs `work` on every one of `jobs`, spread over `workers` threads, and hands back each job with
363/// what it came to, in no particular order.
364///
365/// The jobs are handed out through a queue rather than dealt in equal piles, because they are
366/// nothing like equal: `URL` on ClickBench is a string column of sixty one million distinct values
367/// and `IsMobile` is a byte. A pile that happened to hold the four large string columns would be
368/// the whole stripe and the other workers would be waiting on it. The caller hands the jobs over
369/// cheapest first and they are taken from the back, so the expensive ones go first, which is the
370/// classic answer to a last job that runs longer than everything before it.
371fn 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
410/// One column of every part of a stripe.
411fn column_of(held: &[PendingChunk], index: usize) -> Result<Vec<&Vector>> {
412    held.iter().map(|pending| pending.chunk.column(index)).collect()
413}
414
415/// Puts what [`fan_out`] handed back in column order.
416fn 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    /// Encodes a run of chunks as one stripe, as far as it can be without the writer.
429    ///
430    /// The run is what [`Writer::append_stripe`] takes, and the rules are the same: it is a stripe
431    /// of its own, and the orders have to come out in source order once the stripes are sorted. An
432    /// empty chunk is dropped.
433    ///
434    /// # Errors
435    ///
436    /// If the run is longer than [`STRIPE_PARTS`], a chunk's columns are not the table's, or one
437    /// cannot be encoded.
438    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    /// The check [`Writer::admit`] makes, here because the rows are gone by the merge.
454    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            // The statistics on the thread that is already walking the column, and in the same
474            // step, because the rows are in memory once and this is the moment they are.
475            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
505/// Where one column's dictionary and statistics are while a stripe is merged into them.
506enum Slot<'a> {
507    /// In the writer, which the caller holds.
508    Owned(&'a mut Option<GlobalDictionary>, &'a mut Option<stats::Gather>),
509    /// Lent to a [`Merger`], behind the column's own lock.
510    Lent(&'a Mutex<LentColumn>, &'a Lent),
511}
512
513/// One column of a stripe on its way through [`merge_columns`].
514struct Step<'a> {
515    index: usize,
516    column: Column,
517    slot: Slot<'a>,
518    /// The stripe's statistics for the column, when the column keeps them.
519    gather: Option<stats::Gather>,
520}
521
522impl Step<'_> {
523    /// Roughly what the merge costs: a hash a distinct value when there is a global dictionary to
524    /// merge into, and next to nothing otherwise. Read off `coded` rather than the dictionary, so a
525    /// lent column does not have to be locked to be sorted.
526    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    /// Merges the column, settles its dictionary's shape and hands out the blocks it filled.
536    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                // Checked with the column locked, so a merge either finishes before the writer
545                // takes this column back or is refused.
546                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
556/// One column of [`merge_columns`].
557fn 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            // Empty means nothing has been merged into it yet, so this is the column's first
579            // stripe and the only one the decision is allowed to be made on.
580            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    // Settled here rather than when the stripe is written, so that the blocks this merge filled go
591    // out with it already knowing their shape. A column still too small to settle one keeps its
592    // blocks until it can, which is at most `PAYLOAD_SAMPLE_BLOCKS` of them, because encoding them
593    // now would be encoding them without having looked at the column.
594    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
604/// Merges every column of a stripe into the dictionaries and statistics in `slots`.
605///
606/// Every column is merged on its own, because nothing one column's merge reads or writes belongs to
607/// another: its statistics, its global dictionary and its flag in `coded`. So the columns are
608/// spread over threads, and a stripe takes as long as its slowest column rather than all of them.
609/// The answer is the same in any order, because a column's merge only depends on the stripes
610/// merged into that column before it.
611fn 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    // Taken from the back, so the biggest merges start first and the last one to finish is
627    // small, the same reason `fan_out` hands its jobs over cheapest first.
628    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/// The dictionaries and statistics of a table while a [`Merger`] has them, one lock a column.
679#[derive(Debug)]
680pub(crate) struct Lent {
681    columns: Box<[Mutex<LentColumn>]>,
682    /// Set when the writer takes them back, after which a merge is refused.
683    reclaimed: AtomicBool,
684}
685
686/// One column of [`Lent`].
687#[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    /// Puts encoded blocks back into their dictionaries, each under its own column's lock.
699    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    /// Everything lent, handed back to the writer.
715    #[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/// Merges prepared stripes into a writer's dictionaries and statistics without the writer.
732///
733/// Handed out by [`Writer::merger`]. With it, a load that shares one writer between many threads
734/// holds the writer's lock only to write, and two stripes merge at once as long as they are on
735/// different columns. A stripe merged here is written with [`Writer::write`] as usual, and that is
736/// where its rows are counted in.
737#[derive(Debug, Clone)]
738pub struct Merger {
739    lent: Arc<Lent>,
740    types: Vec<LogicalType>,
741    coded: Arc<[AtomicBool]>,
742}
743
744impl Merger {
745    /// [`Writer::merge`], one column lock at a time instead of the writer.
746    ///
747    /// # Errors
748    ///
749    /// If the stripe was prepared for a table of other columns, or the table was closed.
750    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    /// Puts a stripe's encoded dictionary blocks back, so that [`Writer::write`] does not wait on a
759    /// column's lock while it holds its own.
760    ///
761    /// # Errors
762    ///
763    /// If a block comes back to a column without a dictionary, or comes back twice.
764    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    /// Builds the pages the merge left to build, which is every column coded against a global
771    /// dictionary and every column that lost one after the stripe was prepared, and encodes the
772    /// dictionary blocks the merge filled.
773    ///
774    /// # Errors
775    ///
776    /// If a column or a block cannot be encoded or a page comes out larger than a page may be.
777    pub fn pages(self) -> Result<Paged> {
778        let Self { parts, columns, blocks, profile, counted } = self;
779        let width = columns.len();
780        // The blocks go first so that they are taken last. One block is a thousand values, which is
781        // less than any column of a stripe, and small jobs at the end are what keeps the last
782        // worker from finishing long after the others.
783        let mut jobs = (width..width + blocks.len())
784            .chain((0..width).filter(|&index| !matches!(columns[index], Merge::Pages(_))))
785            .collect::<Vec<_>>();
786        // A column encoded again from its rows costs more than one whose codes only need building.
787        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
827/// One column's parts as pages of global codes.
828fn 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        // None, because the codes already give the stripe an exact membership index, and an
849        // approximate one beside it would cost a hash of every string to answer a question that
850        // is already answered.
851        stripe.sieves.push(None);
852        stripe.ranges.push(part.range.clone());
853    }
854    Ok(stripe)
855}
856
857impl Writer {
858    /// Something that encodes stripes for this writer without holding it. See [`Preparer::prepare`].
859    ///
860    /// It carries the profile the writer has when it is asked for, so a writer that is going to be
861    /// given one with [`Writer::with_profile`] should be given it first.
862    #[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    /// Takes a prepared stripe into the table's dictionaries and statistics and counts its rows in.
872    ///
873    /// This is the step that has to see the stripes one at a time, and it is a hash a distinct
874    /// value of each varchar column rather than two a row. Whatever [`Writer::append_at`] left
875    /// behind is written first as its own stripe, the same rule [`Writer::append_stripe`] has.
876    ///
877    /// # Errors
878    ///
879    /// If the stripe was prepared for a table of other columns, or the buffered stripe cannot be
880    /// written.
881    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    /// [`Writer::merge`] for a stripe whose rows are already counted in.
897    ///
898    /// Every column is merged on its own, because nothing one column's merge reads or writes
899    /// belongs to another: its statistics, its global dictionary and its flag in `coded`. So the
900    /// columns are spread over threads, and the lock is held for the slowest column rather than for
901    /// all of them. On ClickBench `hits` the lock was busy 98% of a load and the merge was four
902    /// fifths of that, while two thirds of the machine waited for it. The answer is the same in any
903    /// order, because a column's merge only depends on the stripes merged into it before.
904    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    /// Hands the dictionaries and the statistics to a [`Merger`], so that stripes can be merged
920    /// without this writer's lock.
921    ///
922    /// Whatever [`Writer::append_at`] left behind is written first, the same rule
923    /// [`Writer::merge`] has. The writer takes them back when the table is closed.
924    ///
925    /// # Errors
926    ///
927    /// If the buffered stripe cannot be written.
928    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    /// Writes a stripe whose pages are built.
953    ///
954    /// # Errors
955    ///
956    /// If the stripe was built for a table of another width or cannot be written.
957    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    /// All four steps one after the other, for a caller with nobody to share the writer with.
989    ///
990    /// # Errors
991    ///
992    /// The same as [`Writer::merge`], [`Merged::pages`] and [`Writer::write`].
993    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    /// The value every row holds, so a test can check a row it reads back without keeping the rows.
1029    ///
1030    /// `city` repeats a handful of values and has a null every so often, which keeps its dictionary.
1031    /// `note` is different on every row but its nulls, which loses it on the first stripe.
1032    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    /// A run of `parts` chunks starting at part `first`, as a caller hands them to the writer.
1043    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    /// The runs the tests hand over, out of source order so that the stripes are sorted at commit.
1058    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    /// Every stripe prepared before any of them is merged writes the file that handing the same
1077    /// runs to the writer one at a time writes, byte for byte.
1078    ///
1079    /// That is the claim the whole split rests on. The second and third stripes here are coded
1080    /// against a dictionary for `note`, which the first stripe to be merged then decides the column
1081    /// should not have, so they are encoded again without it. `city` keeps its dictionary and the
1082    /// later stripes' values go into it in the order the merges happen.
1083    #[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    /// Two stripes merged in one order and written in the other read back as the rows they held,
1113    /// which is what two instances sharing a writer do whenever the second one's pages are built
1114    /// first.
1115    #[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    /// Stripes merged through a [`Merger`] write the same bytes as the writer merging them itself,
1135    /// and their rows are counted in when they are written.
1136    #[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    /// Stripes merged on several threads at once through one [`Merger`] and written in whatever
1165    /// order they finish read back as the rows they held.
1166    #[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    /// A merge that comes after the table is closed is refused rather than merged into
1191    /// dictionaries nothing will write.
1192    #[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    /// A chunk that is not the table's is refused when it reaches the writer, and the writer is not
1205    /// left counting its rows.
1206    #[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}