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                // bytes_at rather than text_at: the rows were checked for UTF-8 when they came in,
212                // and checking every one again cost more than coding it.
213                let text = flat.bytes_at(row).unwrap_or(b"");
214                // A repeat of the row before is common enough on a sorted table to be worth a
215                // comparison before a hash, and the comparison fails on its first bytes when not.
216                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        // Only the coding needs to find a value by its bytes, and on a column of URLs the table
233        // that does it is as large as the codes.
234        local.first = HashMap::default();
235        local.next = Vec::new();
236        Ok(local)
237    }
238
239    /// The column's parts as the rows they were coded from, for a column that lost its global
240    /// dictionary after this stripe was coded against one.
241    ///
242    /// A null row comes back as a null over the empty string, which is what it was coded as, and
243    /// each part gets back the same form of validity it had, since the page records which it was.
244    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    /// Puts every value into `dictionary` in the order this stripe first held it, and says what
318    /// each one's code is there.
319    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
336/// Whether a column's first stripe says it should not have a global dictionary.
337///
338/// Every varchar column starts with one, because the writer cannot know what is in a column before
339/// it has seen some of it. A global dictionary is the right shape for a column of a few dozen
340/// values repeated down the table: the pages become small integers, a filter against a literal is
341/// one search of the sorted order rather than a comparison a row, and a group by is on the codes.
342/// It is the wrong shape for a column whose values are nearly all different. There the codes are
343/// as wide as row numbers, nothing is saved on the pages, and the membership index of a stripe is
344/// a list of very nearly every code in the column. On TPC-H the orders table written on its own
345/// 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`
346/// from 1.810 G instructions to 1.213 G, which is what the rudb parquet reader takes over the same
347/// values.
348///
349/// So the first stripe of a column is the sample and the decision is made once on it. Once, rather
350/// than per stripe, because the codes of one column have to mean the same thing in every page of
351/// it, and a column that changed its mind halfway would need its earlier stripes rewritten. The
352/// first stripe is encoded again when the answer comes out against the dictionary, which is the one
353/// stripe that pays for the decision, along with any stripe that was prepared before it was made.
354///
355/// The threshold is deliberately near the top. [`DICTIONARY_DISTINCT_IN_TEN`] of the sample has to
356/// be values never seen before, which is a column with essentially no repeats. Everything with real
357/// repetition keeps its dictionary and keeps every property that hangs off it, and nothing is
358/// claimed here about where between the two the crossover really sits.
359fn 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
364/// Runs `work` on every one of `jobs`, spread over `workers` threads, and hands back each job with
365/// what it came to, in no particular order.
366///
367/// The jobs are handed out through a queue rather than dealt in equal piles, because they are
368/// nothing like equal: `URL` on ClickBench is a string column of sixty one million distinct values
369/// and `IsMobile` is a byte. A pile that happened to hold the four large string columns would be
370/// the whole stripe and the other workers would be waiting on it. The caller hands the jobs over
371/// cheapest first and they are taken from the back, so the expensive ones go first, which is the
372/// classic answer to a last job that runs longer than everything before it.
373fn 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
412/// One column of every part of a stripe.
413fn column_of(held: &[PendingChunk], index: usize) -> Result<Vec<&Vector>> {
414    held.iter().map(|pending| pending.chunk.column(index)).collect()
415}
416
417/// Puts what [`fan_out`] handed back in column order.
418fn 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    /// Encodes a run of chunks as one stripe, as far as it can be without the writer.
431    ///
432    /// The run is what [`Writer::append_stripe`] takes, and the rules are the same: it is a stripe
433    /// of its own, and the orders have to come out in source order once the stripes are sorted. An
434    /// empty chunk is dropped.
435    ///
436    /// # Errors
437    ///
438    /// If the run is longer than [`STRIPE_PARTS`], a chunk's columns are not the table's, or one
439    /// cannot be encoded.
440    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    /// The check [`Writer::admit`] makes, here because the rows are gone by the merge.
456    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            // The statistics on the thread that is already walking the column, and in the same
476            // step, because the rows are in memory once and this is the moment they are.
477            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
507/// Where one column's dictionary and statistics are while a stripe is merged into them.
508enum Slot<'a> {
509    /// In the writer, which the caller holds.
510    Owned(&'a mut Option<GlobalDictionary>, &'a mut Option<stats::Gather>),
511    /// Lent to a [`Merger`], behind the column's own lock.
512    Lent(&'a Mutex<LentColumn>, &'a Lent),
513}
514
515/// One column of a stripe on its way through [`merge_columns`].
516struct Step<'a> {
517    index: usize,
518    column: Column,
519    slot: Slot<'a>,
520    /// The stripe's statistics for the column, when the column keeps them.
521    gather: Option<stats::Gather>,
522}
523
524impl Step<'_> {
525    /// Roughly what the merge costs: a hash a distinct value when there is a global dictionary to
526    /// merge into, and next to nothing otherwise. Read off `coded` rather than the dictionary, so a
527    /// lent column does not have to be locked to be sorted.
528    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    /// Merges the column, settles its dictionary's shape and hands out the blocks it filled.
538    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                // Checked with the column locked, so a merge either finishes before the writer
547                // takes this column back or is refused.
548                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
558/// One column of [`merge_columns`].
559fn 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            // Empty means nothing has been merged into it yet, so this is the column's first
581            // stripe and the only one the decision is allowed to be made on.
582            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    // Settled here rather than when the stripe is written, so that the blocks this merge filled go
593    // out with it already knowing their shape. A column still too small to settle one keeps its
594    // blocks until it can, which is at most `PAYLOAD_SAMPLE_BLOCKS` of them, because encoding them
595    // now would be encoding them without having looked at the column.
596    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
606/// Merges every column of a stripe into the dictionaries and statistics in `slots`.
607///
608/// Every column is merged on its own, because nothing one column's merge reads or writes belongs to
609/// another: its statistics, its global dictionary and its flag in `coded`. So the columns are
610/// spread over threads, and a stripe takes as long as its slowest column rather than all of them.
611/// The answer is the same in any order, because a column's merge only depends on the stripes
612/// merged into that column before it.
613fn 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    // Taken from the back, so the biggest merges start first and the last one to finish is
629    // small, the same reason `fan_out` hands its jobs over cheapest first.
630    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/// The dictionaries and statistics of a table while a [`Merger`] has them, one lock a column.
681#[derive(Debug)]
682pub(crate) struct Lent {
683    columns: Box<[Mutex<LentColumn>]>,
684    /// Set when the writer takes them back, after which a merge is refused.
685    reclaimed: AtomicBool,
686}
687
688/// One column of [`Lent`].
689#[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    /// Puts encoded blocks back into their dictionaries, each under its own column's lock.
701    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    /// Everything lent, handed back to the writer.
717    #[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/// Merges prepared stripes into a writer's dictionaries and statistics without the writer.
734///
735/// Handed out by [`Writer::merger`]. With it, a load that shares one writer between many threads
736/// holds the writer's lock only to write, and two stripes merge at once as long as they are on
737/// different columns. A stripe merged here is written with [`Writer::write`] as usual, and that is
738/// where its rows are counted in.
739#[derive(Debug, Clone)]
740pub struct Merger {
741    lent: Arc<Lent>,
742    types: Vec<LogicalType>,
743    coded: Arc<[AtomicBool]>,
744}
745
746impl Merger {
747    /// [`Writer::merge`], one column lock at a time instead of the writer.
748    ///
749    /// # Errors
750    ///
751    /// If the stripe was prepared for a table of other columns, or the table was closed.
752    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    /// Puts a stripe's encoded dictionary blocks back, so that [`Writer::write`] does not wait on a
761    /// column's lock while it holds its own.
762    ///
763    /// # Errors
764    ///
765    /// If a block comes back to a column without a dictionary, or comes back twice.
766    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    /// Builds the pages the merge left to build, which is every column coded against a global
773    /// dictionary and every column that lost one after the stripe was prepared, and encodes the
774    /// dictionary blocks the merge filled.
775    ///
776    /// # Errors
777    ///
778    /// If a column or a block cannot be encoded or a page comes out larger than a page may be.
779    pub fn pages(self) -> Result<Paged> {
780        let Self { parts, columns, blocks, profile, counted } = self;
781        let width = columns.len();
782        // The blocks go first so that they are taken last. One block is a thousand values, which is
783        // less than any column of a stripe, and small jobs at the end are what keeps the last
784        // worker from finishing long after the others.
785        let mut jobs = (width..width + blocks.len())
786            .chain((0..width).filter(|&index| !matches!(columns[index], Merge::Pages(_))))
787            .collect::<Vec<_>>();
788        // A column encoded again from its rows costs more than one whose codes only need building.
789        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
829/// One column's parts as pages of global codes.
830fn 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        // None, because the codes already give the stripe an exact membership index, and an
851        // approximate one beside it would cost a hash of every string to answer a question that
852        // is already answered.
853        stripe.sieves.push(None);
854        stripe.ranges.push(part.range.clone());
855    }
856    Ok(stripe)
857}
858
859impl Writer {
860    /// Something that encodes stripes for this writer without holding it. See [`Preparer::prepare`].
861    ///
862    /// It carries the profile the writer has when it is asked for, so a writer that is going to be
863    /// given one with [`Writer::with_profile`] should be given it first.
864    #[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    /// Takes a prepared stripe into the table's dictionaries and statistics and counts its rows in.
874    ///
875    /// This is the step that has to see the stripes one at a time, and it is a hash a distinct
876    /// value of each varchar column rather than two a row. Whatever [`Writer::append_at`] left
877    /// behind is written first as its own stripe, the same rule [`Writer::append_stripe`] has.
878    ///
879    /// # Errors
880    ///
881    /// If the stripe was prepared for a table of other columns, or the buffered stripe cannot be
882    /// written.
883    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    /// [`Writer::merge`] for a stripe whose rows are already counted in.
899    ///
900    /// Every column is merged on its own, because nothing one column's merge reads or writes
901    /// belongs to another: its statistics, its global dictionary and its flag in `coded`. So the
902    /// columns are spread over threads, and the lock is held for the slowest column rather than for
903    /// all of them. On ClickBench `hits` the lock was busy 98% of a load and the merge was four
904    /// fifths of that, while two thirds of the machine waited for it. The answer is the same in any
905    /// order, because a column's merge only depends on the stripes merged into it before.
906    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    /// Hands the dictionaries and the statistics to a [`Merger`], so that stripes can be merged
922    /// without this writer's lock.
923    ///
924    /// Whatever [`Writer::append_at`] left behind is written first, the same rule
925    /// [`Writer::merge`] has. The writer takes them back when the table is closed.
926    ///
927    /// # Errors
928    ///
929    /// If the buffered stripe cannot be written.
930    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    /// Writes a stripe whose pages are built.
955    ///
956    /// # Errors
957    ///
958    /// If the stripe was built for a table of another width or cannot be written.
959    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    /// All four steps one after the other, for a caller with nobody to share the writer with.
991    ///
992    /// # Errors
993    ///
994    /// The same as [`Writer::merge`], [`Merged::pages`] and [`Writer::write`].
995    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    /// The value every row holds, so a test can check a row it reads back without keeping the rows.
1031    ///
1032    /// `city` repeats a handful of values and has a null every so often, which keeps its dictionary.
1033    /// `note` is different on every row but its nulls, which loses it on the first stripe.
1034    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    /// A run of `parts` chunks starting at part `first`, as a caller hands them to the writer.
1045    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    /// The runs the tests hand over, out of source order so that the stripes are sorted at commit.
1060    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    /// Every stripe prepared before any of them is merged writes the file that handing the same
1079    /// runs to the writer one at a time writes, byte for byte.
1080    ///
1081    /// That is the claim the whole split rests on. The second and third stripes here are coded
1082    /// against a dictionary for `note`, which the first stripe to be merged then decides the column
1083    /// should not have, so they are encoded again without it. `city` keeps its dictionary and the
1084    /// later stripes' values go into it in the order the merges happen.
1085    #[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    /// Two stripes merged in one order and written in the other read back as the rows they held,
1115    /// which is what two instances sharing a writer do whenever the second one's pages are built
1116    /// first.
1117    #[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    /// Stripes merged through a [`Merger`] write the same bytes as the writer merging them itself,
1137    /// and their rows are counted in when they are written.
1138    #[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    /// Stripes merged on several threads at once through one [`Merger`] and written in whatever
1167    /// order they finish read back as the rows they held.
1168    #[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    /// A merge that comes after the table is closed is refused rather than merged into
1193    /// dictionaries nothing will write.
1194    #[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    /// A chunk that is not the table's is refused when it reaches the writer, and the writer is not
1207    /// left counting its rows.
1208    #[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}